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 aoneshot::Senderchannel. The harness must execute the model call externally and invokeCallModel::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.respondexactly once perCallModelstep. Failure to do so results in aDriverError::ResponseDroppederror that terminates the stream. RuntimeModelsaccepts aHashMap<Category, Vec<ModelId>>defining which models are visible to the algorithm for this specific request.- Replace
switchyard_libsy_llm_client::runwith your own async function that returns aswitchyard_protocol::Responseto 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_streamemits aStepStreamthat yieldsCallModeloff-load requests and finalOutcomedecisions.- Rust harnesses poll the stream asynchronously, execute model calls through external services, and unblock the algorithm via
CallModel::respond. - Python integration uses
PyAlgorithm.run_streamwith identical semantics, exposing async generators for compatibility withasyncioloops. - Mandatory response protocol requires exactly one
respondcall perCallModelstep to avoidDriverError::ResponseDropped. - Source authority resides in
crates/libsy/src/core/algorithm.rs, with Python bindings incrates/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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →