How to Integrate NVIDIA NeMo Switchyard Algorithm::run_stream with a Custom Harness

Algorithm::run_stream returns a boxed StepStream that yields CallModel requests and final Outcome decisions, enabling custom harnesses to interleave routing logic with external LLM execution.

NVIDIA NeMo Switchyard provides a Rust-native routing engine for LLM request orchestration. The Algorithm::run_stream method in crates/libsy/src/core/algorithm.rs serves as the primary entry point for driving routing algorithms, emitting a stream of steps that a custom harness must consume to execute model calls and finalize responses.

Understanding the Algorithm::run_stream Architecture

The run_stream function defined at lines 526-540 in crates/libsy/src/core/algorithm.rs transforms a routing algorithm into an asynchronous stream of execution steps. When invoked, it spawns the algorithm's internal route future and returns a StepStream that the host application drives to completion.

The StepStream Contract

StepStream is a type alias for Pin<Box<dyn Stream<Item = Result<Step, DriverError>> + Send>>. Each Step variant represents a distinct phase of the routing lifecycle:

  • CallModel – Contains the LLM request, candidate model IDs, and a oneshot::Sender channel. The harness must execute the model call externally and invoke CallModel::respond (defined at lines 95-125) to unblock the algorithm.
  • Outcome – Signals routing completion with the selected model IDs, optional pre-generated response, and metadata for telemetry.

Internal Driver Mechanics

Inside run_stream, the method instantiates a Driver struct (lines 182-210) that owns an mpsc::Sender channel. The driver exposes call_model, which algorithms use to request off-loaded inference. The driver wraps the receiver into a ReceiverStream, boxes it, and returns it as the StepStream. This architecture decouples the routing policy from the inference implementation, allowing harnesses to substitute custom HTTP clients, batching logic, or hardware-specific executors.

Implementing a Rust Custom Harness

A Rust harness must create the RuntimeModels mapping, invoke run_stream, and poll the stream until completion. The harness shown below uses the built-in fall-through algorithm from crates/libsy/src/algorithms/fall_through.rs and delegates actual HTTP calls to switchyard-llm-client.

use std::sync::Arc;
use switchyard_libsy::{
    Algorithm, Driver, RuntimeModels, Request, Step, StepStream,
};
use switchyard_libsy_llm_client::run;

#[tokio::main]
async fn main() -> anyhow::Result<()> {
    // 1. Initialize the algorithm implementation
    let algo = Arc::new(switchyard_libsy::algorithms::fall_through::FallThrough::default());

    // 2. Construct the request object
    let request = Request::from_json(r#"{
        "category": "chat",
        "llm_request": {"prompt":"Hello world!", "max_new_tokens": 10}
    }"#)?;

    // 3. Define available models by category
    let mut by_cat = std::collections::HashMap::new();
    by_cat.insert(switchyard_protocol::Category::Chat,
                  vec!["gpt-4".into(), "gpt-3.5".into()]);
    let models = Arc::new(RuntimeModels::new(by_cat));

    // 4. Obtain the step stream from Algorithm::run_stream
    let mut stream: StepStream = algo.run_stream(request, models);

    // 5. Drive the stream to completion
    while let Some(step_res) = stream.next().await {
        let step = step_res?;
        match step {
            Step::CallModel(call) => {
                // Execute the model call via your inference service
                let response = run(call.request.clone(), call.models.clone()).await?;
                // Mandatory: respond to unblock the algorithm
                call.respond(Ok(response))?;
            }
            Step::Outcome(outcome) => {
                println!("Selected model(s): {:?}", outcome.selected_model_ids);
                if let Some(resp) = outcome.response {
                    println!("Pre-generated answer: {}", resp);
                }
                break;
            }
        }
    }
    Ok(())
}

Critical integration points:

  • The harness must call call.respond exactly once per CallModel step. Failure to do so results in a DriverError::ResponseDropped error that terminates the stream.
  • RuntimeModels accepts a HashMap<Category, Vec<ModelId>> defining which models are visible to the algorithm for this specific request.
  • Replace switchyard_libsy_llm_client::run with your own async function that returns a switchyard_protocol::Response to integrate proprietary inference stacks.

Building a Python Harness with PyO3 Bindings

Switchyard exposes Algorithm::run_stream to Python through PyAlgorithm in crates/switchyard-py/src/libsy_bindings.rs (lines 681-701). The Python API mirrors the Rust interface, allowing data science pipelines to embed routing logic without reimplementing the runtime.

import asyncio
from switchyard_py import PyAlgorithm, Request, RuntimeModels, Category

async def custom_harness():
    # 1. Load the fall-through algorithm

    algo = PyAlgorithm.fall_through()

    # 2. Build the request

    request = Request(
        category=Category.Chat,
        llm_request={"prompt": "Hello world!", "max_new_tokens": 10}
    )

    # 3. Declare available models

    models = RuntimeModels(
        by_category={Category.Chat: ["gpt-4", "gpt-3.5"]},
    )

    # 4. Drive the async generator

    step_gen = algo.run_stream(request, models)

    async for step in step_gen:
        if step.is_call_model():
            call = step.call_model()
            # Insert custom inference logic here

            dummy_response = {
                "choices": [{"text": "Hello from the custom harness!"}]
            }
            await call.respond(dummy_response)
        elif step.is_outcome():
            outcome = step.outcome()
            print("Routing completed")
            print("Selected models:", outcome.selected_model_ids)
            if outcome.response:
                print("Pre-generated answer:", outcome.response)
            break

asyncio.run(custom_harness())

The Python Step objects provide is_call_model() and is_outcome() predicates for type discrimination, implemented as thin wrappers around the Rust enums via PyO3.

Key Source Files and Method Signatures

Component File Path Line Range Purpose
Algorithm::run_stream crates/libsy/src/core/algorithm.rs 526-540 Core stream generator exposing the routing lifecycle
Driver::call_model crates/libsy/src/core/algorithm.rs 182-210 Channel-based off-loading mechanism for LLM requests
CallModel::respond crates/libsy/src/core/algorithm.rs 95-125 Unblocks the algorithm after harness execution
FallThrough algorithm crates/libsy/src/algorithms/fall_through.rs - Reference implementation for stateless routing
Python bindings crates/switchyard-py/src/libsy_bindings.rs 681-701 PyO3 wrappers for PyAlgorithm.run_stream
LLM client consumer crates/libsy-llm-client/src/run.rs - Ready-made HTTP client for CallModel execution

The method signature for run_stream requires self: Arc<dyn Algorithm>, ensuring the algorithm can be shared across async boundaries while maintaining thread safety through the Send bound on the returned stream.

Summary

  • Algorithm::run_stream emits a StepStream that yields CallModel off-load requests and final Outcome decisions.
  • Rust harnesses poll the stream asynchronously, execute model calls through external services, and unblock the algorithm via CallModel::respond.
  • Python integration uses PyAlgorithm.run_stream with identical semantics, exposing async generators for compatibility with asyncio loops.
  • Mandatory response protocol requires exactly one respond call per CallModel step to avoid DriverError::ResponseDropped.
  • Source authority resides in crates/libsy/src/core/algorithm.rs, with Python bindings in crates/switchyard-py/src/libsy_bindings.rs.

Frequently Asked Questions

What is the difference between a CallModel step and an Outcome step?

A CallModel step represents an intermediate request for inference that the algorithm cannot complete internally; it contains the prompt, candidate models, and a response channel. An Outcome step represents the terminal routing decision, containing the final model selection and optional pre-generated response. The harness must handle CallModel by executing the model call and responding, whereas Outcome signals that routing has concluded.

How does the harness communicate results back to the routing algorithm?

The harness calls CallModel::respond with a Result<Response, Error> to unblock the algorithm's internal future. This method sends the result through the oneshot::Sender channel created by the Driver. If the harness drops the CallModel without calling respond, the algorithm receives a DriverError::ResponseDropped and may emit a fallback outcome or terminate the stream.

Can I use Algorithm::run_stream with a custom algorithm implementation rather than the built-in fall-through?

Yes. Any type implementing the Algorithm trait can be wrapped in an Arc and passed to run_stream. Implement the route method to define your policy, using Driver::call_model to request inference when needed. The harness consuming the stream remains agnostic to the specific algorithm logic, interacting only with the Step variants.

What happens if my harness fails to respond to a CallModel step?

If respond is not called exactly once, the oneshot::Sender is dropped, causing the algorithm's internal await on the corresponding receiver to resolve with a RecvError. The Driver converts this into a DriverError::ResponseDropped, which propagates through the StepStream as an error variant. Robust harnesses should wrap model execution in match blocks to ensure respond is invoked even when the external inference service fails.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →