# How NVIDIA Switchyard’s Algorithm::run_stream Handles Internal Routing Decisions

> Learn how NVIDIA Switchyard's Algorithm::run_stream manages internal routing decisions by spawning Drivers for efficient, non-blocking I/O and model execution.

- Repository: [NVIDIA-NeMo/Switchyard](https://github.com/NVIDIA-NeMo/Switchyard)
- Tags: internals
- Published: 2026-09-12

---

**`Algorithm::run_stream` converts routing logic into an asynchronous stream of `Step` events by spawning an isolated `Driver` per request, allowing the algorithm to off-load model calls to the host and assemble a final `RoutingOutcome` without blocking on I/O.**

The NVIDIA-NeMo/Switchyard framework decouples routing intelligence from execution infrastructure, enabling algorithms to remain pure decision-making logic while the host manages network requests. At the center of this architecture lies `run_stream`, which orchestrates the lifecycle of a routing request from initiation to final model selection.

## Architecture of the Routing Stream

The `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 bridge between algorithmic intent and host execution. Instead of directly awaiting model responses, the algorithm emits **steps**—discrete events that the host consumes and fulfills. This design ensures that the core library remains I/O-free, testable, and agnostic to the underlying transport mechanism.

Each invocation creates a completely isolated execution context. This isolation guarantees that parallel routing requests cannot interfere with each other, even when processing high-concurrency workloads across multipleTokio tasks.

## Step-by-Step Execution Flow

Understanding how `run_stream` manages routing requires examining its internal pipeline from driver instantiation to outcome delivery.

### Isolated Driver Initialization

Every call to `run_stream` constructs a fresh **Driver** via `Driver::new`, which initializes its own step channel and maintains a reference to the shared `RuntimeModels` (lines 1,526–1,532 in [`algorithm.rs`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/algorithm.rs)). This driver acts as the algorithm's sole interface to the outside world, ensuring that state pollution between concurrent runs is impossible.

The driver owns the **mpsc channel** through which all routing steps flow, including model call requests and the final completion signal.

### Async Task Spawning and Panic Isolation

Once the driver is ready, `run_stream` spawns the algorithm’s `route` method as a Tokio task. To prevent a panicking algorithm from crashing the entire process, the future is wrapped in `AssertUnwindSafe(...).catch_unwind()` (lines 1,33–1,35). If the algorithm panics, the error is captured and converted into a stream error rather than propagating to the runtime.

An **AbortOnDrop** guard ensures that if the consumer drops the stream handle, the background task terminates immediately, preventing resource leaks (lines 1,50–1,56).

### Off-Loading Model Calls via Steps

Inside the `route` implementation, algorithms request model inferences through `driver.call_model(request, models)`. This method performs several critical actions:

1. Stamps the request with the first candidate model
2. Constructs a `CallModel` promise representing the pending operation
3. Emits a `Step::CallModel` event onto the driver’s step channel (lines 1,60–1,78)

The algorithm then awaits the promise, effectively yielding control back to the host. The host fulfills the promise by executing the HTTP request and signaling completion, at which point the algorithm receives the response and continues its decision logic.

### Stream Consumption and Outcome Assembly

`run_stream` returns a stream that the host consumes through the `drive` helper or manual iteration. The host collects `Step::CallModel` events into a **FuturesUnordered** set of in-flight requests, enabling concurrent model evaluation. When the algorithm concludes, the driver calls `driver.finish(result)`, which attaches outcome metadata and emits the terminal `Step::Done` event containing the **RoutingOutcome** (lines 1,28–1,49).

The outcome itself is constructed using helpers such as `RoutingOutcome::route_to` or `RoutingOutcome::answered`, which embed the selected model identifier and optional generated responses directly into the request structure (lines 1,37–1,78).

## Observability and Error Handling

Switchyard provides rich instrumentation throughout the routing lifecycle. A `libsy.run` tracing span encapsulates the entire execution, while each model call generates a nested `libsy.llm_call` span for granular performance analysis.

Error handling operates at multiple levels:

- **Dropped promises**: If the host abandons a `CallModel` promise, `Driver::call_model` returns `DriverError::ResponseDropped`
- **Task cancellation**: Stream abandonment triggers immediate task abortion via `AbortOnDrop`
- **Panic recovery**: Unwinding panics are caught and streamed as error events rather than crashing the service

## Example: Implementing a Custom Router

The following implementation demonstrates how to build a routing algorithm that selects the first available completion model and returns the generated response:

```rust
use std::sync::Arc;
use libsy::prelude::*;

struct FirstModelAlgo;

#[async_trait::async_trait]
impl Algorithm for FirstModelAlgo {
    fn name(&self) -> &str { "first-model" }

    async fn route(
        self: Arc<Self>,
        driver: Driver,
        request: Request,
    ) -> Result<RoutingOutcome> {
        // Select models from the Completion category
        let models = driver.models_for(&Category::Completion).to_vec();
        
        // Off-load the call to the host via the step stream
        let response = driver
            .call_model(request.clone(), models)
            .await?;
        
        // Attach routing evidence for observability
        driver.set_evidence(json!({"source": "first-model"}));
        
        // Return an outcome embedding the response
        Ok(RoutingOutcome::answered(
            driver.first_model_for(&Category::Completion)?.clone(),
            request,
            response,
        ))
    }
}

// Host-side consumption of the stream
async fn execute_routing(
    alg: Arc<dyn Algorithm>,
    request: Request,
    models: Arc<RuntimeModels>,
) -> Result<()> {
    let mut stream = alg.run_stream(request, models);
    
    while let Some(step) = stream.next().await {
        match step? {
            Step::CallModel(call) => {
                // Host fulfills the model call via HTTP client
                let resp = http_client
                    .execute(call.request, call.models)
                    .await?;
                call.respond(Ok(resp))?;
            }
            Step::Done(outcome) => {
                println!("Selected model: {:?}", outcome.selected_model_ids);
                return Ok(());
            }
        }
    }
    Ok(())
}

```

## Summary

NVIDIA Switchyard’s `Algorithm::run_stream` manages internal routing decisions through a carefully orchestrated sequence:

- **Driver isolation** creates independent execution contexts per request in [`crates/libsy/src/core/algorithm.rs`](https://github.com/NVIDIA-NeMo/Switchyard/blob/main/crates/libsy/src/core/algorithm.rs)
- **Async step emission** allows algorithms to request model calls without blocking on I/O
- **Panic safety** via `catch_unwind` prevents algorithmic errors from crashing the host
- **Stream-based completion** delivers `RoutingOutcome` through a terminal `Step::Done` event
- **Lifecycle management** ensures proper cleanup via `AbortOnDrop` and response promise handling

## Frequently Asked Questions

### What is the role of the Driver in Switchyard routing?

The **Driver** acts as the algorithm's runtime interface, managing the step channel and mediating all external interactions. It provides `call_model` for requesting inferences, `set_evidence` for attaching metadata, and `finish` for signaling completion. Each `run_stream` invocation receives its own isolated driver instance to prevent cross-request contamination.

### How does run_stream handle panics in routing algorithms?

The method wraps the algorithm’s `route` future in `AssertUnwindSafe(...).catch_unwind()` (lines 1,33–1,35), which intercepts unwinding panics and converts them into stream errors. This ensures that a faulty algorithm implementation cannot terminate the host process or disrupt other concurrent routing operations.

### What happens if the host drops a model call promise?

When the host drops a `CallModel` promise without calling `respond()`, the awaiting algorithm receives `DriverError::ResponseDropped`. This signals that the host is no longer interested in the result, allowing the algorithm to abort, retry with alternative models, or return a degraded routing outcome.

### How are routing decisions finalized in the stream?

The algorithm concludes by returning a `RoutingOutcome`, which the driver processes through `driver.finish(result)`. This method attaches metadata and emits `Step::Done`, the terminal stream event containing the final model selection and any generated responses. The host receives this event and terminates the stream consumption loop.