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

> Learn how NVIDIA NeMo Switchyard's Algorithm::run_stream works and integrates with custom harnesses. Interleave routing logic with external LLM execution effectively.

- Repository: [NVIDIA-NeMo/Switchyard](https://github.com/NVIDIA-NeMo/Switchyard)
- Tags: how-to-guide
- Published: 2026-09-13

---

**`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`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/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`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/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`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/crates/libsy/src/algorithms/fall_through.rs) and delegates actual HTTP calls to `switchyard-llm-client`.

```rust
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`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/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.

```python
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`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/crates/libsy/src/core/algorithm.rs) | 526-540 | Core stream generator exposing the routing lifecycle |
| `Driver::call_model` | [`crates/libsy/src/core/algorithm.rs`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/crates/libsy/src/core/algorithm.rs) | 182-210 | Channel-based off-loading mechanism for LLM requests |
| `CallModel::respond` | [`crates/libsy/src/core/algorithm.rs`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/crates/libsy/src/core/algorithm.rs) | 95-125 | Unblocks the algorithm after harness execution |
| `FallThrough` algorithm | [`crates/libsy/src/algorithms/fall_through.rs`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/crates/libsy/src/algorithms/fall_through.rs) | - | Reference implementation for stateless routing |
| Python bindings | [`crates/switchyard-py/src/libsy_bindings.rs`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/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`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/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`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/crates/libsy/src/core/algorithm.rs), with Python bindings in [`crates/switchyard-py/src/libsy_bindings.rs`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/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.