From eccbafac656322c273bc7e1dd34a5e09b8cc6824 Mon Sep 17 00:00:00 2001 From: adhaile Date: Tue, 29 Sep 2026 10:56:54 -0700 Subject: [PATCH] feat(vgr): hand off escalated agentic runs Behind agentic_handoff, compact_handoff and local_turn_budget_seconds (all off by default): tell the capable tier once per user turn that it inherits unverified tool-using work, optionally condense the local history into a cache-stable digest, and escalate a user turn once its local wall-clock budget is spent. --- crates/libsy/src/algorithms/vgr.rs | 10 +- crates/libsy/src/algorithms/vgr/compaction.rs | 379 ++++++++++++++++++ crates/libsy/src/algorithms/vgr/config.rs | 12 + crates/libsy/src/algorithms/vgr/runtime.rs | 125 +++++- .../libsy/src/algorithms/vgr/runtime/tests.rs | 59 +++ crates/switchyard-runner/src/algorithm.rs | 21 + docs/routing_algorithms/vgr_routing.md | 18 + 7 files changed, 612 insertions(+), 12 deletions(-) create mode 100644 crates/libsy/src/algorithms/vgr/compaction.rs diff --git a/crates/libsy/src/algorithms/vgr.rs b/crates/libsy/src/algorithms/vgr.rs index 65f39035d..8a03fbb47 100644 --- a/crates/libsy/src/algorithms/vgr.rs +++ b/crates/libsy/src/algorithms/vgr.rs @@ -22,6 +22,7 @@ use crate::algorithms::util::affinity::AffinityRouter; use crate::core::algorithm::{Algorithm, Driver, RoutingOutcome}; use crate::core::state::State; +mod compaction; mod config; mod decide; mod readout; @@ -57,12 +58,17 @@ impl Vgr { .with_release_on_user_turn() .with_latch_only([cloud.clone()]), ); + let mut route = FallThrough::new_with_state().with_name("vgr"); + if config.compact_handoff { + // First, so every later classifier and the served call see the + // condensed history. + route = route.with_classifier(Arc::new(compaction::Compactor)); + } let classifier = Arc::new(runtime::VgrClassifier { breaker: safety::CircuitBreaker::new(config.breaker), config, }); - let route = FallThrough::new_with_state() - .with_name("vgr") + let route = route .with_processor(turn_affinity.clone()) .with_classifier(turn_affinity) .with_classifier(classifier); diff --git a/crates/libsy/src/algorithms/vgr/compaction.rs b/crates/libsy/src/algorithms/vgr/compaction.rs new file mode 100644 index 000000000..d5378536b --- /dev/null +++ b/crates/libsy/src/algorithms/vgr/compaction.rs @@ -0,0 +1,379 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Condensing the local tier's history when the capable tier takes over. +//! +//! An escalated session otherwise hands the capable tier every token the local +//! tier produced, uncached, and then re-sends that history on every later call. +//! Compaction replaces the local work with a digest appended to the task +//! message, once, at handoff. +//! +//! The client keeps sending its full history, so the same rewrite must be +//! reapplied to every later request of the user turn. It is reapplied only when +//! the client's history still starts with exactly the messages that were +//! condensed; anything else is sent untouched, which is always safe. The digest +//! is stored rather than recomputed so the rewritten prefix is byte-identical +//! across calls and stays in the capable tier's prompt cache. + +use std::collections::hash_map::DefaultHasher; +use std::hash::{Hash, Hasher}; + +use async_trait::async_trait; +use switchyard_protocol::{ContentBlock, Message, Request, Response, Role}; + +use super::text::clip_mid; +use crate::Result; +use crate::algorithms::util::affinity::has_new_user_turn; +use crate::algorithms::util::prompts::drop_exact_replay; +use crate::core::algorithm::Driver; +use crate::core::classifier::{Classification, Classifier}; +use crate::core::state::{State, StateValue}; + +const HEAD_KEY: &str = "vgr.compaction.head"; +const CUT_KEY: &str = "vgr.compaction.cut"; +const ANCHOR_KEY: &str = "vgr.compaction.anchor"; +const SUMMARY_KEY: &str = "vgr.compaction.summary"; + +/// Steps rendered in full at the end of the digest; earlier ones get one line. +/// +/// The digest is re-sent on every later capable-tier call, so its size is paid +/// once per step, not once per handoff. About 2k tokens at most. +const RECENT_STEPS: usize = 3; +const RECENT_ARGUMENT_BUDGET: usize = 600; +const RECENT_RESULT_BUDGET: usize = 800; +const EARLY_ARGUMENT_BUDGET: usize = 100; +const EARLY_RESULT_BUDGET: usize = 80; +const EARLY_SECTION_BUDGET: usize = 2_500; +const NARRATION_BUDGET: usize = 200; +const RECENT_USER_TEXTS: usize = 2; + +/// Condenses everything after the task message and records how, so later +/// requests of this user turn can be rewritten identically. +/// +/// Returns `false`, leaving the request untouched, when there is no task +/// message with local work after it. +pub(super) fn compact(request: &mut Request, state: &mut State, note: &str) -> bool { + let messages = &request.llm_request.messages; + let Some(task) = messages.iter().position(is_task_message) else { + return false; + }; + let head = task + 1; + let cut = messages.len(); + if cut <= head { + return false; + } + let mut summary = digest(&messages[head..cut]); + summary.push_str(note); + let anchor = fingerprint(&messages[..cut]); + let (Ok(head_count), Ok(cut_count)) = (u32::try_from(head), u32::try_from(cut)) else { + return false; + }; + state + .extra + .insert(HEAD_KEY.into(), StateValue::Count(head_count)); + state + .extra + .insert(CUT_KEY.into(), StateValue::Count(cut_count)); + state + .extra + .insert(ANCHOR_KEY.into(), StateValue::String(anchor)); + rewrite(request, head, cut, &summary); + state + .extra + .insert(SUMMARY_KEY.into(), StateValue::String(summary)); + true +} + +/// Reapplies a recorded compaction to a later request of the same user turn. +/// +/// Registered ahead of every other classifier and never votes. +pub(super) struct Compactor; + +#[async_trait] +impl Classifier for Compactor { + async fn score( + &self, + state: &mut State, + request: &mut Request, + _driver: &Driver, + ) -> Result<(Classification, Option)> { + if has_new_user_turn(&request.llm_request.messages) { + for key in [HEAD_KEY, CUT_KEY, ANCHOR_KEY, SUMMARY_KEY] { + state.extra.remove(key); + } + } else if let Some((head, cut, summary)) = recorded(state, request) { + rewrite(request, head, cut, &summary); + } + Ok((Classification::Scores(Vec::new()), None)) + } +} + +/// The recorded compaction, when the request still begins with the condensed history. +fn recorded(state: &State, request: &Request) -> Option<(usize, usize, String)> { + let count = |key| match state.extra.get(key) { + Some(StateValue::Count(value)) => usize::try_from(*value).ok(), + _ => None, + }; + let text = |key| match state.extra.get(key) { + Some(StateValue::String(value)) => Some(value.clone()), + _ => None, + }; + let (head, cut) = (count(HEAD_KEY)?, count(CUT_KEY)?); + let messages = &request.llm_request.messages; + (head < cut && cut <= messages.len() && fingerprint(&messages[..cut]) == text(ANCHOR_KEY)?) + .then(|| text(SUMMARY_KEY)) + .flatten() + .map(|summary| (head, cut, summary)) +} + +fn rewrite(request: &mut Request, head: usize, cut: usize, summary: &str) { + let messages = &mut request.llm_request.messages; + messages[head - 1].content.push(ContentBlock::Text { + text: summary.to_string(), + }); + messages.drain(head..cut); + drop_exact_replay(request); +} + +fn is_task_message(message: &Message) -> bool { + message.role == Role::User + && message + .content + .iter() + .any(|block| matches!(block, ContentBlock::Text { .. })) +} + +fn fingerprint(messages: &[Message]) -> String { + let mut hasher = DefaultHasher::new(); + serde_json::to_string(messages) + .unwrap_or_default() + .hash(&mut hasher); + format!("{:016x}", hasher.finish()) +} + +/// One tool call and what came back, in conversation order. +struct Step { + narration: String, + name: String, + arguments: String, + result: Option<(bool, String)>, +} + +fn digest(messages: &[Message]) -> String { + let mut steps: Vec = Vec::new(); + let mut narration = String::new(); + let mut user_text = Vec::new(); + for message in messages { + for block in &message.content { + match block { + ContentBlock::Text { text } if message.role == Role::Assistant => { + narration.push_str(text); + } + ContentBlock::Text { text } if message.role == Role::User => { + user_text.push(clip_mid(text, RECENT_RESULT_BUDGET, 0.5)); + } + ContentBlock::ToolCall(call) => steps.push(Step { + narration: std::mem::take(&mut narration), + name: call.name.clone(), + arguments: call.arguments.to_string(), + result: None, + }), + ContentBlock::ToolResult(result) => { + let body = result + .content + .iter() + .map(|part| match part { + ContentBlock::Text { text } => text.as_str(), + _ => "[non-text content]", + }) + .collect::>() + .join("\n"); + let failed = result.is_error == Some(true); + if let Some(step) = steps.iter_mut().rev().find(|step| step.result.is_none()) { + step.result = Some((failed, body)); + } + } + _ => {} + } + } + } + + let recent_from = steps.len().saturating_sub(RECENT_STEPS); + let mut early = String::new(); + for (index, step) in steps[..recent_from].iter().enumerate() { + let (status, first_line) = match &step.result { + Some((failed, body)) => ( + if *failed { "error" } else { "ok" }, + body.lines() + .find(|line| !line.trim().is_empty()) + .unwrap_or(""), + ), + None => ("no result", ""), + }; + early.push_str(&format!( + "[{}] {} {} -> {status}: {}\n", + index + 1, + step.name, + clip_mid(&step.arguments, EARLY_ARGUMENT_BUDGET, 0.7), + clip_mid(first_line, EARLY_RESULT_BUDGET, 0.7), + )); + } + let mut out = String::from( + "\n\n[Context note from the serving infrastructure, not from the user] The tool-using \ + work that followed this message was condensed into the digest below. Earlier steps \ + are one line each; the most recent steps are shown in detail.\n\n", + ); + out.push_str(&clip_mid(&early, EARLY_SECTION_BUDGET, 0.3)); + for (offset, step) in steps[recent_from..].iter().enumerate() { + if !step.narration.trim().is_empty() { + out.push_str(&format!( + "assistant: {}\n", + clip_mid(step.narration.trim(), NARRATION_BUDGET, 0.5) + )); + } + out.push_str(&format!( + "[{}] {} {}\n", + recent_from + offset + 1, + step.name, + clip_mid(&step.arguments, RECENT_ARGUMENT_BUDGET, 0.5) + )); + match &step.result { + Some((failed, body)) => out.push_str(&format!( + "result{}:\n{}\n", + if *failed { " (error)" } else { "" }, + clip_mid(body, RECENT_RESULT_BUDGET, 0.3) + )), + None => out.push_str("result: none\n"), + } + } + if !narration.trim().is_empty() { + out.push_str(&format!( + "assistant: {}\n", + clip_mid(narration.trim(), NARRATION_BUDGET, 0.5) + )); + } + for text in &user_text[user_text.len().saturating_sub(RECENT_USER_TEXTS)..] { + out.push_str(&format!("user: {text}\n")); + } + out.push_str(""); + out +} + +#[cfg(test)] +mod tests { + use super::*; + use switchyard_protocol::{LlmRequest, ToolCall, ToolResult}; + + fn call(id: &str, command: &str) -> Message { + Message { + role: Role::Assistant, + content: vec![ContentBlock::ToolCall(ToolCall { + id: id.into(), + name: "bash".into(), + arguments: serde_json::json!({ "command": command }), + })], + } + } + + fn result(id: &str, text: &str) -> Message { + Message { + role: Role::Tool, + content: vec![ContentBlock::ToolResult(ToolResult { + tool_call_id: id.into(), + content: vec![ContentBlock::Text { text: text.into() }], + is_error: None, + })], + } + } + + fn session(steps: usize) -> Request { + let mut messages = vec![Message::text(Role::User, "build the thing")]; + for index in 0..steps { + let id = format!("call-{index}"); + messages.push(call(&id, &format!("step {index}"))); + messages.push(result(&id, &format!("output {index}"))); + } + Request { + llm_request: LlmRequest { + messages, + ..LlmRequest::default() + }, + raw_request: None, + metadata: None, + } + } + + async fn reapply(state: &mut State, request: &mut Request) { + let driver = Driver::new( + "test", + std::sync::Arc::new(crate::core::algorithm::RuntimeModels::new( + std::collections::HashMap::new(), + )), + ) + .0; + Compactor + .score(state, request, &driver) + .await + .expect("the compactor never fails"); + } + + #[tokio::test] + async fn later_requests_see_the_same_condensed_prefix() { + let mut state = State::default(); + let mut handoff = session(12); + assert!(compact(&mut handoff, &mut state, "NOTE")); + assert_eq!(handoff.llm_request.messages.len(), 1); + let condensed = handoff.llm_request.messages[0].clone(); + let digest = condensed.text_content("").expect("task text"); + assert!(digest.contains("step 11") && digest.ends_with("NOTE")); + + let mut follow_up = session(12); + follow_up.llm_request.messages.push(call("cloud-1", "ls")); + follow_up + .llm_request + .messages + .push(result("cloud-1", "a b")); + reapply(&mut state, &mut follow_up).await; + + let messages = &follow_up.llm_request.messages; + assert_eq!( + messages.len(), + 3, + "task, then only the capable tier's own work" + ); + assert_eq!( + messages[0], condensed, + "the cached prefix is byte-identical" + ); + } + + #[tokio::test] + async fn a_rewritten_history_is_sent_untouched() { + let mut state = State::default(); + let mut handoff = session(4); + assert!(compact(&mut handoff, &mut state, "NOTE")); + + let mut edited = session(4); + edited.llm_request.messages[2] = result("call-0", "the client rewrote this"); + let before = edited.llm_request.messages.clone(); + reapply(&mut state, &mut edited).await; + assert_eq!(edited.llm_request.messages, before); + } + + #[tokio::test] + async fn a_new_user_turn_forgets_the_compaction() { + let mut state = State::default(); + let mut handoff = session(4); + assert!(compact(&mut handoff, &mut state, "NOTE")); + + let mut next_turn = session(4); + next_turn + .llm_request + .messages + .push(Message::text(Role::User, "now do another thing")); + let before = next_turn.llm_request.messages.clone(); + reapply(&mut state, &mut next_turn).await; + assert_eq!(next_turn.llm_request.messages, before); + assert!(!state.extra.contains_key(SUMMARY_KEY)); + } +} diff --git a/crates/libsy/src/algorithms/vgr/config.rs b/crates/libsy/src/algorithms/vgr/config.rs index d171a3913..085dadb16 100644 --- a/crates/libsy/src/algorithms/vgr/config.rs +++ b/crates/libsy/src/algorithms/vgr/config.rs @@ -61,6 +61,15 @@ pub struct VgrConfig { /// Lets a long agentic run with earlier tool errors commit on a capable-tier /// confirmation once it ends in this many consecutive clean tool results. pub confirmed_recovery_min_clean_tail: Option, + /// Tells the capable tier, once per user turn, that it inherits unverified + /// tool-using work. + pub agentic_handoff: bool, + /// Condenses the local tier's work into a digest at handoff. Requires + /// `agentic_handoff`. + pub compact_handoff: bool, + /// Wall-clock budget for the local tier within one user turn, including the + /// client's tool execution. + pub local_turn_budget: Option, } impl VgrConfig { @@ -80,6 +89,9 @@ impl VgrConfig { task_typing: true, local_supports_images: false, confirmed_recovery_min_clean_tail: None, + agentic_handoff: false, + compact_handoff: false, + local_turn_budget: None, } } diff --git a/crates/libsy/src/algorithms/vgr/runtime.rs b/crates/libsy/src/algorithms/vgr/runtime.rs index a9ace1439..59d06eca2 100644 --- a/crates/libsy/src/algorithms/vgr/runtime.rs +++ b/crates/libsy/src/algorithms/vgr/runtime.rs @@ -9,7 +9,7 @@ //! by what is left of the turn's deadline, and a call that fails or times out //! is evidence never gathered, which does not commit. -use std::time::{Duration, Instant}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use async_trait::async_trait; use switchyard_protocol::{ @@ -21,8 +21,10 @@ use super::decide::{self, AgenticRun, Route, Signals, Tri}; use super::rungs::{self, Question}; use super::safety::{CircuitBreaker, endpoint_failure, fallback_eligible}; use super::{Branch, Capabilities, TaskType, derive_capabilities, readout, text}; +use crate::algorithms::util::affinity::has_new_user_turn; use crate::algorithms::util::buffered_response::buffer_response; use crate::algorithms::util::decisive; +use crate::algorithms::util::prompts::append_note; use crate::core::algorithm::Driver; use crate::core::classifier::{Classification, Classifier}; use crate::core::state::{State, StateValue}; @@ -30,6 +32,8 @@ use crate::{LibsyError, Result}; const TURN_STREAK_KEY: &str = "vgr.turn_verification.streak"; const TURN_LATCHED_KEY: &str = "vgr.turn_verification.latched"; +const TURN_STARTED_KEY: &str = "vgr.user_turn.started_unix"; +const HANDOFF_SENT_KEY: &str = "vgr.user_turn.handoff_sent"; /// Consecutive escalation votes that latch a session to the capable tier. const TURN_CONFIRMATIONS: u32 = 2; const TURN_ESCALATE_AT: f64 = 0.5; @@ -51,6 +55,11 @@ impl Classifier for VgrClassifier { ) -> Result { let started = Instant::now(); let local = self.config.targets.local.clone(); + if has_new_user_turn(&request.llm_request.messages) { + state.extra.remove(TURN_STARTED_KEY); + state.extra.remove(HANDOFF_SENT_KEY); + } + let turn_budget_spent = self.turn_budget_spent(state); let short_circuit = if self.config.mode == ServingMode::Off { Some("mode_off") } else if self @@ -64,17 +73,19 @@ impl Classifier for VgrClassifier { Some("breaker_open") } else if turn_latched(state) { Some("turn_verification_latched") + } else if turn_budget_spent { + Some("local_turn_budget") } else if !self.config.local_supports_images && request_has_image(request) { Some("local_image_unsupported") } else { None }; if let Some(reason) = short_circuit { - return Ok(self.escalate(reason, None)); + return Ok(self.escalate(reason, None, request, state, None)); } let Some(budget) = self.remaining(started) else { - return Ok(self.escalate("local_timed_out", None)); + return Ok(self.escalate("local_timed_out", None, request, state, None)); }; let attempted = tokio::time::timeout(budget, async { let response = driver @@ -95,26 +106,32 @@ impl Classifier for VgrClassifier { if !fallback_eligible(&error) { return Err(error); } - return Ok(self.escalate("local_unavailable", None)); + return Ok(self.escalate("local_unavailable", None, request, state, None)); } Err(_) => { self.breaker.failure(); - return Ok(self.escalate("local_timed_out", None)); + return Ok(self.escalate("local_timed_out", None, request, state, None)); } }; if rungs::has_tool_call(&buffered.agg) { if !tool_use_is_complete(&buffered.agg) { - return Ok(self.escalate("malformed_tool_call", None)); + return Ok(self.escalate("malformed_tool_call", None, request, state, None)); } if self .verify_turn(driver, request, &buffered.agg, state, started) .await { - return Ok(self.escalate("turn_verification_escalated", None)); + return Ok(self.escalate( + "turn_verification_escalated", + None, + request, + state, + None, + )); } if self.config.mode == ServingMode::Shadow { - return Ok(self.escalate("turn_verification_complete", None)); + return Ok(self.escalate("turn_verification_complete", None, request, state, None)); } log("turn_verification_complete", None, Route::Local); return Ok((decisive(&local), Some(buffered.into_response()))); @@ -126,7 +143,7 @@ impl Classifier for VgrClassifier { let signals = self.gather(driver, &caps, request, started).await; let route = decide::decide(&caps, &signals, self.config.confirmed_min_clean_tail()); if route == Route::Cloud || self.config.mode == ServingMode::Shadow { - return Ok(self.escalate("decided", Some(caps.branch))); + return Ok(self.escalate("decided", Some(caps.branch), request, state, Some(&attempt))); } log("decided", Some(caps.branch), Route::Local); Ok((decisive(&local), Some(buffered.into_response()))) @@ -134,11 +151,62 @@ impl Classifier for VgrClassifier { } impl VgrClassifier { - fn escalate(&self, reason: &'static str, branch: Option) -> Scored { + fn escalate( + &self, + reason: &'static str, + branch: Option, + request: &mut Request, + state: &mut State, + unverified_final: Option<&str>, + ) -> Scored { + if reason != "mode_off" { + self.hand_off(request, state, unverified_final); + } log(reason, branch, Route::Cloud); (decisive(&self.config.targets.cloud), None) } + /// Whether the current user turn has used up its local wall-clock budget. + /// + /// The clock starts at the first request of the user turn and includes the + /// client's tool execution, since that is the time the client's own limit + /// is spending. + fn turn_budget_spent(&self, state: &mut State) -> bool { + let Some(budget) = self.config.local_turn_budget else { + return false; + }; + let now = unix_seconds(); + let started = match state.extra.get(TURN_STARTED_KEY) { + Some(StateValue::Count(started)) => *started, + _ => { + state + .extra + .insert(TURN_STARTED_KEY.into(), StateValue::Count(now)); + now + } + }; + u64::from(now.saturating_sub(started)) >= budget.as_secs() + } + + /// Tells the capable tier it is inheriting unverified tool-using work. + /// + /// Sent at most once per user turn, on its first escalation. + fn hand_off(&self, request: &mut Request, state: &mut State, unverified_final: Option<&str>) { + if !self.config.agentic_handoff + || state.extra.contains_key(HANDOFF_SENT_KEY) + || !text::has_tool_trajectory(request) + { + return; + } + state + .extra + .insert(HANDOFF_SENT_KEY.into(), StateValue::Count(1)); + let note = handoff_note(unverified_final); + if !(self.config.compact_handoff && super::compaction::compact(request, state, ¬e)) { + append_note(request, ¬e); + } + } + /// Judges one proposed tool call; `true` escalates the rest of the session. async fn verify_turn( &self, @@ -373,6 +441,43 @@ impl VgrClassifier { } } +fn unix_seconds() -> u32 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(0, |elapsed| { + u32::try_from(elapsed.as_secs()).unwrap_or(u32::MAX) + }) +} + +/// The note a capable tier receives when it takes over a tool-using session. +fn handoff_note(unverified_final: Option<&str>) -> String { + let mut note = String::from( + "\n\n[Routing notice from the serving infrastructure, not from the user] Until now the \ + work in this session was carried out by a smaller, faster model. That work has NOT been \ + verified: it may be incomplete or incorrect, and it may have left the environment in a \ + broken state. You are taking over from here. Do not assume that any earlier step \ + succeeded or that the task is finished. Inspect the actual current state of the \ + environment against every requirement of the original task, fix whatever is missing or \ + wrong, and give your final answer only after you have confirmed that the requirements \ + are met.", + ); + if let Some(claim) = unverified_final.filter(|claim| !claim.trim().is_empty()) { + const CLIPPING_MARKER_RESERVE: usize = 64; + let claim = text::clip_mid( + claim, + super::render::ATTEMPT_BUDGET.saturating_sub(CLIPPING_MARKER_RESERVE), + 1.0 / 3.0, + ); + note.push_str( + "\n\nThe smaller model was about to finish with the following unchecked message:\n\ + \n", + ); + note.push_str(&claim); + note.push_str("\n"); + } + note +} + fn log(reason: &'static str, branch: Option, served: Route) { tracing::info!(target: "libsy", reason, branch = ?branch, served = ?served, "vgr decision"); } diff --git a/crates/libsy/src/algorithms/vgr/runtime/tests.rs b/crates/libsy/src/algorithms/vgr/runtime/tests.rs index 7944d1d5d..4ffdf8791 100644 --- a/crates/libsy/src/algorithms/vgr/runtime/tests.rs +++ b/crates/libsy/src/algorithms/vgr/runtime/tests.rs @@ -250,3 +250,62 @@ async fn an_unavailable_local_tier_escalates_and_other_failures_surface() -> Res } Ok(()) } + +#[tokio::test] +async fn an_escalated_agentic_run_hands_off_once_with_the_unverified_claim() -> Result<()> { + let route = route(|config| config.agentic_handoff = true)?; + let sent = Arc::new(std::sync::Mutex::new(Vec::new())); + let task = || { + session(vec![ + Message::text(Role::User, "fix the build"), + call("c1"), + result("c1", "ok"), + ]) + }; + for _ in 0..2 { + let sent = Arc::clone(&sent); + test_drive_with_models( + Arc::clone(&route), + task(), + models(), + move |target: ModelId, request: Request| { + let sent = Arc::clone(&sent); + async move { + Ok(match target.as_str() { + "local" => reply("the build passes"), + "judge" => scored(0.0), + "cloud" => { + let text = serde_json::to_string(&request.llm_request.messages) + .unwrap_or_default(); + sent.lock().map(|mut sent| sent.push(text)).ok(); + reply("cloud answer") + } + _ => reply("no"), + }) + } + }, + ) + .await?; + } + let sent = sent.lock().map(|sent| sent.clone()).unwrap_or_default(); + assert_eq!(sent.len(), 2); + assert!(sent[0].contains("Routing notice") && sent[0].contains("the build passes")); + assert!( + !sent[1].contains("Routing notice"), + "sent once per user turn" + ); + Ok(()) +} + +#[tokio::test] +async fn a_spent_local_turn_budget_skips_the_local_attempt() -> Result<()> { + let route = route(|config| config.local_turn_budget = Some(std::time::Duration::ZERO))?; + let calls = Arc::new(Calls::default()); + let task = session(vec![Message::text(Role::User, "fix the build")]); + assert_eq!( + drive(&route, task, calls.clone(), complete_call, 0.0).await?, + "cloud" + ); + assert_eq!(calls.0.load(Ordering::Relaxed), 0); + Ok(()) +} diff --git a/crates/switchyard-runner/src/algorithm.rs b/crates/switchyard-runner/src/algorithm.rs index 2a981a0fc..cbaff1db2 100644 --- a/crates/switchyard-runner/src/algorithm.rs +++ b/crates/switchyard-runner/src/algorithm.rs @@ -476,6 +476,16 @@ pub struct VgrRouteConfig { /// Clean trailing tool results that let a long run recover on the cloud judge. #[serde(default)] pub confirmed_recovery_min_clean_tail: Option, + /// Tells the capable tier it inherits unverified tool-using work. + #[serde(default)] + pub agentic_handoff: bool, + /// Condenses the local tier's work into a digest at handoff. Requires + /// `agentic_handoff`. + #[serde(default)] + pub compact_handoff: bool, + /// Wall-clock budget for the local tier within one user turn. + #[serde(default)] + pub local_turn_budget_seconds: Option, } #[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)] @@ -1564,6 +1574,17 @@ fn build_vgr( runtime.deadline = duration(route_name, "deadline_seconds", config.deadline_seconds)?; runtime.task_typing = config.task_typing; runtime.confirmed_recovery_min_clean_tail = config.confirmed_recovery_min_clean_tail; + if config.compact_handoff && !config.agentic_handoff { + return Err(AlgorithmConfigError::new(format!( + "vgr route {route_name}: compact_handoff requires agentic_handoff" + ))); + } + runtime.agentic_handoff = config.agentic_handoff; + runtime.compact_handoff = config.compact_handoff; + runtime.local_turn_budget = config + .local_turn_budget_seconds + .map(|seconds| duration(route_name, "local_turn_budget_seconds", seconds)) + .transpose()?; runtime.breaker = BreakerConfig { threshold: config.breaker_threshold, cooldown: duration( diff --git a/docs/routing_algorithms/vgr_routing.md b/docs/routing_algorithms/vgr_routing.md index ee8729e3f..0933e46ae 100644 --- a/docs/routing_algorithms/vgr_routing.md +++ b/docs/routing_algorithms/vgr_routing.md @@ -36,6 +36,24 @@ backend reports its live context capacity, VGR republishes it through run that recovered from tool errors commit locally once that many trailing tool results are clean and the cloud judge confirms the evidence. +## Agentic handoff + +These settings are off by default: + +```toml +agentic_handoff = true +compact_handoff = true +local_turn_budget_seconds = 600 +``` + +`agentic_handoff` tells the cloud tier, once per user turn, that it is taking +over unverified tool-using work, including the local tier's unchecked final +message. `compact_handoff` also condenses the local tier's history into a +digest and reapplies that same digest to every later request of the user turn, +which keeps the cloud tier's prompt cache warm. `local_turn_budget_seconds` +escalates a user turn once the local tier has spent that much wall-clock time +on it. + ## Serving modes - `off` skips candidate generation and serves cloud. This is the default.