Skip to content

fix(libsy): cancel abandoned offloaded calls - #948

Open
nachiketb-nvidia wants to merge 1 commit into
mainfrom
nachiketb/fix-call-cancellation
Open

nachiketb-nvidia wants to merge 1 commit into
mainfrom
nachiketb/fix-call-cancellation

Conversation

@nachiketb-nvidia

@nachiketb-nvidia nachiketb-nvidia commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

What

Let an algorithm cancel one model or decision call without cancelling its other calls or ending the run.

  • CallModel::respond and CallDecision::respond now accept the work future and are async. They stop that work when the algorithm drops its waiting call future.
  • The HTTP client passes its request future into respond, rather than finishing the request first. Cancelled calls return normally instead of failing the run with ResponseDropped.
  • Python ModelCall.respond and DecisionCall.respond accept an awaitable. The binding cancels the Python task and awaits its cleanup when the algorithm stops waiting. Python fail is now awaitable too.
  • Update every repository caller, Python type declaration, the example, and the relevant docs. Add focused Rust and Python cancellation regressions.

Algorithm, Driver, run_stream, and drive keep their existing signatures. No new dependency or cancellation-token API is needed.

Why

Parallel calls already work, but abandoning one call did not stop its HTTP work while the algorithm continued. The host awaited the request before replying, so the reply channel could only detect cancellation after the work finished. A late reply could then fail the entire run.

Passing the future to respond lets the existing reply channel control that call's lifetime. This supports racing models, abandoning a slow call, and continuing with another call.

API Migration

This changes the low-level host API. Pass the work itself, not its awaited result.

Rust:

// Before
let result = client_call().await;
call.respond(result)?;

// After
call.respond(client_call()).await?;

// For a result that is already available
call.respond(std::future::ready(result)).await?;

Python:

# Before
call.respond(await client.call(call.request, call.models[0]))

# After
await call.respond(client.call(call.request, call.models[0]))

# For concurrent run_stream consumers
handler = asyncio.create_task(
    call.respond(client.call(call.request, call.models[0]))
)

# Explicit Python failures are now awaitable
await call.fail(error)

Custom stream consumers must keep polling handlers concurrently, then cancel and await outstanding tasks when the run ends. The Python example shows this cleanup. Dropping a spawned task handle alone does not cancel it.

Model errors follow recover_errors inside respond: recoverable errors go back to the algorithm; other errors return to the host. The standard HTTP client's existing policy stays the same. Decision errors still go back to the algorithm for fallback.

Cancellation from the Algorithm

The algorithm cancels a call by dropping its waiting future. There is no separate cancel() method. This works for both Driver::call_model and Driver::call_decision.

Race two calls

tokio::select! polls both calls concurrently. When one finishes, it drops the other call future and signals the host to stop that work:

let (model, response) = tokio::select! {
    result = driver.call_model(request.clone(), vec!["fast".into()]) =>
        ("fast", result?),

    result = driver.call_model(request.clone(), vec!["strong".into()]) =>
        ("strong", result?),
};

This races the first completion, not the first successful answer. These examples use ? to propagate errors.

Cancel a decision after a timeout

A timeout drops the decision-call future if its deadline expires. The algorithm can then apply its own fallback:

let decision = tokio::time::timeout(
    std::time::Duration::from_millis(100),
    driver.call_decision(decision_request, "judge".into()),
).await;

Cancel based on another model's output

The algorithm can inspect the fast response before deciding whether to keep the strong call. Here, is_acceptable is the algorithm's own check, not a libsy API:

let mut fast = Box::pin(
    driver.call_model(request.clone(), vec!["fast".into()])
);
let mut strong = Box::pin(
    driver.call_model(request.clone(), vec!["strong".into()])
);

let (model, response) = tokio::select! {
    result = &mut fast => {
        let response = result?;
        if is_acceptable(&response) {
            drop(strong);
            ("fast", response)
        } else {
            ("strong", strong.await?)
        }
    }

    result = &mut strong => {
        drop(fast);
        ("strong", result?)
    }
};

Both calls run concurrently. If the fast answer passes the check, the strong call is cancelled. If it does not, the strong call keeps running. If strong finishes first, fast is cancelled. The selected model and response can then form the normal RoutingOutcome.

Use a decision to cancel a model call

A decision call can also run alongside a speculative model call. In this example, the algorithm starts strong while the judge decides whether fast is sufficient. choose_fast is the algorithm's interpretation of the decision response:

let mut decision = Box::pin(
    driver.call_decision(decision_request, "judge".into())
);
let mut strong = Box::pin(
    driver.call_model(request.clone(), vec!["strong".into()])
);

let (model, response) = tokio::select! {
    result = &mut decision => {
        let decision = result?;
        if choose_fast(&decision) {
            drop(strong);
            let response = driver
                .call_model(request.clone(), vec!["fast".into()])
                .await?;
            ("fast", response)
        } else {
            ("strong", strong.await?)
        }
    }

    result = &mut strong => {
        drop(decision);
        ("strong", result?)
    }
};

If the judge chooses fast first, the algorithm cancels strong and calls fast. If it chooses strong, the existing strong call continues. If strong finishes first, its answer is used and the judge is cancelled.

When selecting on &mut future, the future remains owned by the algorithm. Explicitly drop the owned future or leave its scope to cancel it. The Box::pin examples above make that ownership clear; dropping only a borrowed Pin<&mut _> does not cancel the underlying future.

Notes for reviewers

  • Start with crates/libsy/src/core/algorithm.rs: reply-channel closure drops the work future; a late reply after cancellation no longer aborts unrelated work.
  • Check crates/libsy-llm-client/src/run.rs next: the cancellable future owns the actual provider call and its completed-call observations.
  • Check crates/switchyard-py/src/libsy_bindings.rs and switchyard_rust/_libsy_async.py together. Dropping the Rust adapter alone does not cancel an asyncio task, so the small Python helper owns cancellation and cleanup.
  • The Rust regression abandons a model or decision call, waits for the host work to drop, then completes another call. Python tests cover stream abandonment, explicit handler cancellation, awaited cleanup, decision replies/errors, and duplicate replies.
  • A delivered response stream still belongs to its consumer. Cancellation stops local work; it cannot guarantee that a provider stops generating or billing after disconnect.
  • Existing call metrics remain in place. Cancelled calls use the existing error/drop accounting.

LoC Breakdown

Diff line counts across 20 files. These include comments and blank lines. Inline Rust test changes are counted under tests, not runtime implementation.

Area Added Removed Net
Rust runtime (libsy and HTTP client) 65 42 +23
Python bindings, async helper, and types 138 40 +98
New cancellation tests 250 0 +250
Existing tests and test helpers 94 65 +29
Documentation 65 1 +64
Python example 28 17 +11
Total 640 165 +475

Runtime implementation totals 203 added / 82 removed (+121 net) across Rust and Python. The two new cancellation test files account for 250 added lines: 132 Rust and 118 Python. Existing test changes mainly migrate callers to the async response API and preserve the tested error behavior.

Validation

Check Result
cargo test --workspace --locked Passed
cargo clippy --workspace --all-targets --locked -- -D warnings Passed
Server Clippy with prefill-router enabled Passed
Runner tests with prefill-router enabled Passed
cargo fmt --all --check Passed
Non-integration Python suite, Python 3.12 193 passed; 2 integration tests deselected
Non-integration Python suite, Python 3.14 193 passed; 2 integration tests deselected
Installed-wheel binding checks, Python 3.14 26 passed
Ruff and mypy Passed
mkdocs build --strict Passed

Live checks against NVIDIA Inference with GPT-OSS 20B and 120B also passed: concurrent completion, first-response race, whole-run cancellation, and cancellation of one call while the algorithm continues. In the continuation check, the abandoned HTTP future dropped within about 1 ms and the next call completed successfully. These checks verify local future cleanup, not upstream compute or billing cancellation.

Summary by CodeRabbit

  • New Features

    • Rust and Python model and decision-call responses now support asynchronous work, allowing handlers to run concurrently with algorithm execution.
    • If an algorithm stops waiting for a response, its unfinished handler work is cancelled and cleaned up. Ending a run also cancels outstanding handlers.
    • Cancellation is local: disconnecting after a response stream is delivered does not guarantee that the provider stops generating.
  • Documentation

    • Updated Rust and Python guides with asynchronous response handling, concurrency, and cancellation behavior.

Signed-off-by: nachiketb <nachiketb@nvidia.com>
@nachiketb-nvidia
nachiketb-nvidia requested a review from a team as a code owner October 7, 2026 20:21
@github-actions

github-actions Bot commented Oct 7, 2026

Copy link
Copy Markdown
PR Preview Action v1.8.1

🚀 View preview at
https://NVIDIA-NeMo.github.io/Switchyard/pr-preview/pr-948/

Built to branch gh-pages at 2026-10-07 20:22 UTC.
Preview will be ready when the GitHub Pages deployment is complete.

@coderabbitai

coderabbitai Bot commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

Review in Change Stack →

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration
  • Configuration used: Repository: NVIDIA-NeMo/Switchyard/.coderabbit.yaml
  • Review profile: CHILL
  • Plan: Enterprise
  • Run ID: 12d9a090-24a3-435f-8657-db1f4da22cba
📥 Commits

Reviewing files that changed from the base of the PR and between 3a918cb and 7566b05.

📒 Files selected for processing (20)
  • crates/libsy-llm-client/README.md
  • crates/libsy-llm-client/src/run.rs
  • crates/libsy-llm-client/tests/observability.rs
  • crates/libsy/README.md
  • crates/libsy/src/algorithms/advisor_gate/tests.rs
  • crates/libsy/src/algorithms/escalation.rs
  • crates/libsy/src/algorithms/llm_class.rs
  • crates/libsy/src/algorithms/util/llm_judge.rs
  • crates/libsy/src/core/algorithm.rs
  • crates/libsy/src/core/testing.rs
  • crates/libsy/tests/cancellation.rs
  • crates/prefill-router/tests/unit/algorithm.rs
  • crates/switchyard-py/src/libsy_bindings.rs
  • docs/getting_started.md
  • docs/reference/opentelemetry.md
  • examples/libsy.py
  • switchyard_rust/_libsy_async.py
  • switchyard_rust/libsy.py
  • tests/test_libsy_cancellation.py
  • tests/test_libsy_minimal_bindings.py

Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 11 remain after this review.


Walkthrough

Rust and Python model and decision response methods now accept asynchronous work. The implementation cancels pending work when the algorithm stops waiting. Client integrations, tests, examples, and documentation now use or describe the asynchronous response contract.

Changes

Asynchronous Responses

Layer / File(s) Summary
Rust response contract and cancellation
crates/libsy/src/core/algorithm.rs, crates/libsy/src/core/testing.rs, crates/libsy/src/algorithms/*, crates/libsy/tests/cancellation.rs, crates/prefill-router/tests/unit/algorithm.rs, crates/libsy/README.md, docs/getting_started.md, docs/reference/opentelemetry.md
Rust CallModel::respond and CallDecision::respond now accept futures. They stop polling work when the algorithm stops waiting. Helpers, tests, call sites, and documentation use or describe the updated contract.
Python async response bridge
switchyard_rust/libsy.py, switchyard_rust/_libsy_async.py, crates/switchyard-py/src/libsy_bindings.rs, tests/test_libsy_cancellation.py, tests/test_libsy_minimal_bindings.py, examples/libsy.py
Python response methods now accept awaitables. The async helper schedules and cleans up work, and the Rust bindings await and convert model and decision responses. Tests and the example use asynchronous handlers.
LLM client routing and observability
crates/libsy-llm-client/src/run.rs, crates/libsy-llm-client/tests/observability.rs, crates/libsy-llm-client/README.md
Decision and routing work now runs inside the future passed to respond. Observability tests and routing documentation reflect the async response and cancellation behavior.

Priority: ➖ Normal

Estimated code review effort: 4 (Complex) | ~45 minutes

Merge Risk: ⚪ Minimal · up to 7566b

The change makes abandoned model and decision calls cancel their host-side work, and the callers, tests, and docs are updated to match. No concrete merge-blocking risk was identified.

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 42.86% which is insufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 56 functions across 16 files. (4 skipped:… Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly and concisely describes the main change: cancelling abandoned offloaded calls in libsy.
Full details: Docstring Coverage

Explanation

Docstring coverage is 42.86% which is insufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 56 functions across 16 files. (4 skipped: 4 unsupported.)

  • Fix all pre-merge checks with AI
  • Autopilot · Keep fixing CodeRabbit findings and required CI, and resolving merge conflicts

A rabbit taps a future into flight.
The waiting call lets go of work.
Rust and Python tidy tasks away.
Streams hop on while readers stay.
Tests trace each turn and cancellation.
The rabbit rests beside the docs.

Comment @coderabbitai help to get the list of available commands.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant