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

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 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). 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:

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
  • 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.

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 →