diff --git a/crates/libsy-llm-client/src/observation.rs b/crates/libsy-llm-client/src/observation.rs index da448eaf2..308ad4b28 100644 --- a/crates/libsy-llm-client/src/observation.rs +++ b/crates/libsy-llm-client/src/observation.rs @@ -6,6 +6,7 @@ use std::sync::Arc; use std::time::Duration; +use switchyard_libsy::OutcomeMetadata; use switchyard_protocol::{ModelId, Usage}; /// One completed model call observed while serving an algorithm run. @@ -21,9 +22,11 @@ pub struct LlmCallObservation { pub usage: Option, } -/// One request-scoped observation emitted by the algorithm runner. +/// Events emitted inline while [`crate::run`] serves a routing request. #[derive(Clone, Debug)] pub enum RunObservation { + /// Metadata attached to the completed routing outcome. + Outcome(OutcomeMetadata), /// A completed model call requested by the algorithm for routing work. LlmCall(LlmCallObservation), /// A completed terminal model call made from the routing outcome. diff --git a/crates/libsy-llm-client/src/run.rs b/crates/libsy-llm-client/src/run.rs index 380e9f602..94296ddcf 100644 --- a/crates/libsy-llm-client/src/run.rs +++ b/crates/libsy-llm-client/src/run.rs @@ -33,7 +33,8 @@ use crate::{metrics, observability}; /// Run one request to completion, serving every offloaded model call with `client`. /// /// Returns the model selected by the algorithm and the final [`Response`]. `observer`, when -/// present, receives each completed routing or answer call and the routing overhead. +/// present, receives metadata for the routing outcome, each completed routing or answer call, +/// and the routing overhead. /// /// `clients` resolves each offloaded call to the client for the target the algorithm /// selected — an algorithm may route among targets served by different providers, so this is @@ -71,6 +72,11 @@ pub async fn run( metrics::record_routing_overhead(&algorithm_name, overhead); let selected_model_id = outcome.selected_model_id()?.clone(); + if let Some(observer) = &observer + && let Some(metadata) = outcome.metadata + { + observer(RunObservation::Outcome(metadata)); + } let (result, answer_duration) = if let Some(response) = outcome.response { (Ok(response), None) } else { @@ -661,9 +667,13 @@ mod tests { assert!(matches!(observations[0], RunObservation::AnswerCall(_))); assert!(matches!( observations[1], + RunObservation::Outcome(ref metadata) if metadata.algorithm == "answered_test" + )); + assert!(matches!( + observations[2], RunObservation::RoutingOverhead(_) )); - assert_eq!(observations.len(), 2); + assert_eq!(observations.len(), 3); Ok(()) } diff --git a/crates/libsy-llm-client/tests/observability.rs b/crates/libsy-llm-client/tests/observability.rs index 473658281..71a35d5d1 100644 --- a/crates/libsy-llm-client/tests/observability.rs +++ b/crates/libsy-llm-client/tests/observability.rs @@ -1218,15 +1218,20 @@ async fn observed_run_reports_one_successful_routed_call() -> switchyard_libsy:: Some(Some(MODEL)) ); let observations = observations.lock(); - assert_eq!(observations.len(), 2); - let RunObservation::AnswerCall(observation) = &observations[0] else { + // Outcome metadata precedes the answer call and final overhead observation. + assert_eq!(observations.len(), 3); + let RunObservation::Outcome(metadata) = &observations[0] else { + return Err(test_error("expected an outcome observation")); + }; + assert_eq!(metadata.algorithm, ALGO); + let RunObservation::AnswerCall(observation) = &observations[1] else { return Err(test_error("expected an answer-call observation")); }; assert_eq!(observation.selected_model, MODEL); assert!(observation.is_success); assert!(observation.usage.is_some()); assert!(matches!( - observations[1], + observations[2], RunObservation::RoutingOverhead(_) )); Ok(()) diff --git a/crates/switchyard-nemo-relay-plugin/README.md b/crates/switchyard-nemo-relay-plugin/README.md index b1a4b48ef..fa1aaedd9 100644 --- a/crates/switchyard-nemo-relay-plugin/README.md +++ b/crates/switchyard-nemo-relay-plugin/README.md @@ -219,10 +219,14 @@ meaning requires a new schema version. | `switchyard.routing.requested` | `algorithm` | | `switchyard.routing.llm_call` | `call_index`, `selected_model`, `call_role`, `outcome`, `latency_ms` | | `switchyard.routing.overhead` | `latency_ms` | -| `switchyard.routing.decision` | `algorithm`, `selected_model`, nullable `served_model`, nullable `fallback_used` | -| `switchyard.routing.error` | `failure_kind`; optional `category`, `phase`, `upstream_status`, and `target` | +| `switchyard.routing.decision` | `algorithm`, optional `outcome_id`, `selected_model`, nullable `served_model`, nullable `fallback_used`, and optional `evidence` | +| `switchyard.routing.error` | `failure_kind`; route-execution failures also include `category`, `phase`, nullable `upstream_status` and `target`, and may include `outcome_id` and `evidence` | `served_model` and `fallback_used` are `null` when serving metadata is unavailable. +`outcome_id` is present when the algorithm runner supplies outcome metadata. +`evidence` is an object containing the supported string fields `source`, `verdict`, +`trigger`, and `reason_code`, and numeric fields `score`, `confidence`, and `threshold`. +String values longer than 64 bytes are omitted and should be stable, non-sensitive labels. ## Failure policy diff --git a/crates/switchyard-nemo-relay-plugin/src/runtime.rs b/crates/switchyard-nemo-relay-plugin/src/runtime.rs index 0781ab5e8..c29935170 100644 --- a/crates/switchyard-nemo-relay-plugin/src/runtime.rs +++ b/crates/switchyard-nemo-relay-plugin/src/runtime.rs @@ -23,6 +23,7 @@ use crate::config::SwitchyardConfig; use crate::translation; const ROUTING_MARK_SCHEMA_VERSION: &str = "1"; +const MAX_EVIDENCE_STRING_BYTES: usize = 64; #[derive(Debug)] pub(crate) struct RoutingMark { @@ -70,6 +71,11 @@ pub(crate) struct Execution { pub(crate) events: Vec, } +struct RoutedResponse { + response: Response, + outcome_fields: Map, +} + pub(crate) struct SwitchyardRuntime { runner: Runner, translation: TranslationEngine, @@ -122,11 +128,11 @@ impl SwitchyardRuntime { let request_extensions = request.llm_request.extensions.clone(); let Execution { result, mut events } = self.execute(inbound, request).await; let (result, finalization_failed) = match result { - Ok(response) => { + Ok(routed) => { let result = finalize_buffered_response( &self.translation, inbound, - response, + routed.response, &request_extensions, ); let failed = result.is_err(); @@ -149,7 +155,7 @@ impl SwitchyardRuntime { let request_extensions = request.llm_request.extensions.clone(); let Execution { result, mut events } = self.execute(inbound, request).await; let (result, finalization_failed) = match result { - Ok(response) => { + Ok(routed) => { let metadata = events .iter() .find_map(|event| match event { @@ -157,8 +163,14 @@ impl SwitchyardRuntime { RoutingEvent::Metric(_) => None, }) .unwrap_or_else(|| Json::Object(Map::new())); - let result = - returned_events(response, inbound, &request_extensions, metadata, emit_event); + let result = returned_events( + routed.response, + inbound, + &request_extensions, + metadata, + routed.outcome_fields, + emit_event, + ); let failed = result.is_err(); (result, failed) } @@ -170,7 +182,7 @@ impl SwitchyardRuntime { Execution { result, events } } - async fn execute(&self, inbound: WireFormat, request: Request) -> Execution { + async fn execute(&self, inbound: WireFormat, request: Request) -> Execution { let Some(route) = self.route(&request) else { return Execution { result: Err("Switchyard has no route for this request model".into()), @@ -202,31 +214,48 @@ impl SwitchyardRuntime { }); match route.execute(request, Some(observer)).await { Ok(output) => { - self.emit_observations(&mut events, take_observations(&observations), &metadata); + let outcome_fields = self.emit_observations( + &mut events, + take_observations(&observations), + &metadata, + ); let served_model = output.response.served_model().map(|model| model.as_str()); + let mut decision_outcome_fields = outcome_fields.clone(); + let mut data = json!({ + "algorithm": route.algorithm_name(), + "selected_model": output.selected_model.as_str(), + "served_model": served_model, + "fallback_used": served_model + .map(|model| model != output.selected_model.as_str()), + }); + if let Json::Object(data) = &mut data { + data.append(&mut decision_outcome_fields); + } events.push(RoutingEvent::Mark(RoutingMark { name: "switchyard.routing.decision".into(), - data: json!({ - "algorithm": route.algorithm_name(), - "selected_model": output.selected_model.as_str(), - "served_model": served_model, - "fallback_used": served_model - .map(|model| model != output.selected_model.as_str()), - }), + data, metadata, severity: Some(LogSeverity::Info), })); Execution { - result: Ok(output.response), + result: Ok(RoutedResponse { + response: output.response, + outcome_fields, + }), events, } } Err(error) => { - self.emit_observations(&mut events, take_observations(&observations), &metadata); + let outcome_fields = self.emit_observations( + &mut events, + take_observations(&observations), + &metadata, + ); self.route_execution_error_mark( &mut events, &error.execution_error_summary(), None, + outcome_fields, ); Execution { result: Err("Switchyard route execution failed".into()), @@ -249,10 +278,20 @@ impl SwitchyardRuntime { events: &mut Vec, observations: Vec, metadata: &Json, - ) { + ) -> Map { let mut call_index = 0; + let mut outcome_fields = Map::new(); for observation in observations { match observation { + RunObservation::Outcome(outcome) => { + outcome_fields.insert( + "outcome_id".into(), + Json::String(outcome.outcome_id().into()), + ); + if let Some(evidence) = evidence_for_mark(outcome.evidence) { + outcome_fields.insert("evidence".into(), evidence); + } + } RunObservation::LlmCall(call) => { call_index += 1; self.routing_call_events(events, call, call_index, metadata); @@ -287,6 +326,7 @@ impl SwitchyardRuntime { } } } + outcome_fields } fn routing_call_events( @@ -338,14 +378,34 @@ impl SwitchyardRuntime { events: &mut Vec, summary: &RouteErrorSummary, metadata: Option<&Json>, + outcome_fields: Map, ) { let metadata = metadata .cloned() .unwrap_or_else(|| event_metadata(events).unwrap_or_else(|| Json::Object(Map::new()))); - events.extend(route_execution_error_events(summary, metadata)); + events.extend(route_execution_error_events( + summary, + metadata, + outcome_fields, + )); } } +/// Projects documented, bounded evidence into Relay mark data. +fn evidence_for_mark(evidence: Option) -> Option { + let Some(Json::Object(mut evidence)) = evidence else { + return None; + }; + evidence.retain(|name, value| match name.as_str() { + "source" | "verdict" | "trigger" | "reason_code" => value + .as_str() + .is_some_and(|value| value.len() <= MAX_EVIDENCE_STRING_BYTES), + "score" | "confidence" | "threshold" => value.is_number(), + _ => false, + }); + (!evidence.is_empty()).then_some(Json::Object(evidence)) +} + pub(crate) fn emit_events(runtime: &PluginRuntime, events: Vec) { for event in events { emit_event(runtime, event); @@ -400,6 +460,7 @@ fn returned_events( inbound: WireFormat, request_extensions: &ProviderExtensions, metadata: Json, + mut outcome_fields: Map, emit_event: RoutingEventEmitter, ) -> Result { let served_model = response.served_model().cloned(); @@ -425,6 +486,7 @@ fn returned_events( for event in route_execution_error_events( &stream_error_summary(error, served_model.as_ref()), metadata.clone(), + std::mem::take(&mut outcome_fields), ) { emit_event(event); } @@ -491,24 +553,40 @@ fn relay_stream_error(error: LlmStreamError) -> String { } } -fn route_execution_error_mark(summary: &RouteErrorSummary, metadata: Json) -> RoutingMark { +fn route_execution_error_mark( + summary: &RouteErrorSummary, + metadata: Json, + mut outcome_fields: Map, +) -> RoutingMark { + let mut data = json!({ + "failure_kind": "route_execution", + "category": summary.kind.as_str(), + "phase": summary.phase.as_str(), + "upstream_status": summary.upstream_status, + "target": summary.target.as_ref().map(|target| target.as_str()), + }); + if let Json::Object(data) = &mut data { + data.append(&mut outcome_fields); + } RoutingMark { name: "switchyard.routing.error".into(), - data: json!({ - "failure_kind": "route_execution", - "category": summary.kind.as_str(), - "phase": summary.phase.as_str(), - "upstream_status": summary.upstream_status, - "target": summary.target.as_ref().map(|target| target.as_str()), - }), + data, metadata, severity: Some(LogSeverity::Error), } } -fn route_execution_error_events(summary: &RouteErrorSummary, metadata: Json) -> Vec { +fn route_execution_error_events( + summary: &RouteErrorSummary, + metadata: Json, + outcome_fields: Map, +) -> Vec { vec![ - RoutingEvent::Mark(route_execution_error_mark(summary, metadata.clone())), + RoutingEvent::Mark(route_execution_error_mark( + summary, + metadata.clone(), + outcome_fields, + )), failure_metric( "route_execution", Some(summary.kind.as_str()), @@ -894,12 +972,91 @@ mod tests { assert_eq!(decision["selected_model"], "target/model"); assert_eq!(decision["served_model"], "target/model"); assert_eq!(decision["fallback_used"], false); + assert!(decision["outcome_id"].is_string()); let response = execution.result.expect("target call should succeed"); assert_eq!(response["object"], "response"); assert_eq!(response["model"], "target/model"); } } + // Routing succeeds before answer candidates exhaust, so the error keeps its metadata. + #[tokio::test] + async fn failed_answer_keeps_routing_outcome_in_error_mark() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(ResponseTemplate::new(503).set_body_string("unavailable")) + .expect(2) + .mount(&server) + .await; + let deployment = json!({ + "schema_version": 1, + "llm_clients": { + "local": { + "format": "openai_chat", + "base_url": format!("{}/v1", server.uri()), + "max_retries": 0, + } + }, + "targets": { + "strong": {"id": "poc/strong", "llm_client": "local"}, + "weak": {"id": "poc/weak", "llm_client": "local"}, + }, + "routes": { + "stage": { + "id": "switchyard/stage", + "type": "stage_router", + "capable_target": "strong", + "efficient_target": "weak", + "picker": "efficient_first", + "confidence_threshold": 0.5, + } + }, + }); + let runtime = SwitchyardRuntime::new(crate::config::SwitchyardConfig { + priority: 0, + switchyard_config_path: None, + switchyard_config: Some( + deployment + .as_object() + .expect("deployment should be an object") + .clone(), + ), + }) + .expect("stage runtime should load"); + let request = runtime + .decode_request( + WireFormat::OpenAiChat, + RelayRequest { + headers: Map::new(), + content: json!({ + "model": "switchyard/stage", + "messages": [{"role": "user", "content": "hello"}], + }), + }, + false, + ) + .expect("request should decode"); + + let execution = runtime + .execute_buffered(WireFormat::OpenAiChat, request) + .await; + + assert!(execution.result.is_err()); + let error = execution + .events + .iter() + .find_map(|event| match event { + RoutingEvent::Mark(mark) if mark.name == "switchyard.routing.error" => { + Some(&mark.data) + } + RoutingEvent::Mark(_) | RoutingEvent::Metric(_) => None, + }) + .expect("error mark should be emitted"); + assert!(error["outcome_id"].is_string()); + assert_eq!(error["evidence"], json!({"source": "fall_open"})); + } + #[tokio::test] async fn buffered_responses_restore_codex_tool_namespaces() { let server = MockServer::start().await; @@ -1019,6 +1176,7 @@ mod tests { let mark = route_execution_error_mark( &error.execution_error_summary(), json!({"session_id": "session"}), + Map::new(), ); assert_eq!(mark.name, "switchyard.routing.error"); @@ -1033,6 +1191,33 @@ mod tests { assert!(!mark.data.to_string().contains(secret)); } + // Custom algorithms can supply evidence, so the Relay boundary filters it. + #[test] + fn decision_evidence_keeps_only_documented_bounded_fields() { + let evidence = evidence_for_mark(Some(json!({ + "source": "llm-classifier", + "score": 0.9, + "confidence": "wrong type", + "threshold": 0.5, + "verdict": "continue", + "trigger": "turn", + "reason_code": "x".repeat(MAX_EVIDENCE_STRING_BYTES + 1), + "unknown": "ignored", + "prompt": {"text": "ignored"}, + }))); + + assert_eq!( + evidence, + Some(json!({ + "source": "llm-classifier", + "score": 0.9, + "threshold": 0.5, + "verdict": "continue", + "trigger": "turn", + })) + ); + } + #[test] fn routing_observations_emit_debug_marks_and_metrics() { let runtime = runtime_for("switchyard"); @@ -1271,6 +1456,7 @@ mod tests { WireFormat::OpenAiChat, &ProviderExtensions::default(), json!({"session_id": "session"}), + Map::new(), Arc::new(move |event| emitted.lock().unwrap().push(event)), ) .expect("stream setup should succeed"); @@ -1324,6 +1510,7 @@ mod tests { WireFormat::OpenAiChat, &ProviderExtensions::default(), json!({}), + Map::new(), Arc::new(move |event| emitted.lock().unwrap().push(event)), ) .expect("stream setup should succeed"); @@ -1354,11 +1541,16 @@ mod tests { }; let captured = Arc::new(Mutex::new(Vec::new())); let emitted = Arc::clone(&captured); + // Routing has completed before this response stream reports its failure. let stream = returned_events( response, WireFormat::OpenAiChat, &ProviderExtensions::default(), json!({"session_id": "session"}), + Map::from_iter([ + ("outcome_id".into(), json!("outcome-test")), + ("evidence".into(), json!({"source": "fall_open"})), + ]), Arc::new(move |mark| emitted.lock().unwrap().push(mark)), ) .expect("stream setup should succeed"); @@ -1379,6 +1571,8 @@ mod tests { assert_eq!(mark.data["category"], "context_window_exceeded"); assert_eq!(mark.data["phase"], "during_stream"); assert_eq!(mark.data["target"], "strong"); + assert_eq!(mark.data["outcome_id"], "outcome-test"); + assert_eq!(mark.data["evidence"], json!({"source": "fall_open"})); assert_eq!(mark.severity, Some(LogSeverity::Error)); assert!(!mark.data.to_string().contains(secret)); let RoutingEvent::Metric(metric) = &events[1] else { @@ -1418,6 +1612,7 @@ mod tests { WireFormat::OpenAiChat, &ProviderExtensions::default(), json!({"session_id": "session"}), + Map::new(), Arc::new(move |event| emitted.lock().unwrap().push(event)), ) .expect("stream setup should succeed"); diff --git a/crates/switchyard-server/src/lib.rs b/crates/switchyard-server/src/lib.rs index 235329b76..b951a1cf0 100644 --- a/crates/switchyard-server/src/lib.rs +++ b/crates/switchyard-server/src/lib.rs @@ -430,6 +430,7 @@ fn stats_observer( classifier_log: Option<(SharedRoutingLog, routing_log::RoutingLogContext)>, ) -> RunObserver { Arc::new(move |observation| match observation { + RunObservation::Outcome(_) => {} RunObservation::AnswerCall(call) => { let latency_ms = call.duration.as_secs_f64() * 1_000.0; if call.is_success { diff --git a/docs/integrations/nemo_relay.md b/docs/integrations/nemo_relay.md index 83e17828e..be96a9da5 100644 --- a/docs/integrations/nemo_relay.md +++ b/docs/integrations/nemo_relay.md @@ -316,8 +316,8 @@ describes the surrounding event envelope. | `switchyard.routing.requested` | Info | Routing `algorithm` for a managed request. | | `switchyard.routing.llm_call` | Debug | `call_index`, `selected_model`, `call_role` (`routing` or `answer`), `outcome`, and `latency_ms` for each observed model call. | | `switchyard.routing.overhead` | Info | `latency_ms` spent producing the routing outcome, including routing-model calls. This is not the end-to-end request duration. | -| `switchyard.routing.decision` | Info | `algorithm`, initial `selected_model`, nullable final `served_model`, and nullable `fallback_used`. | -| `switchyard.routing.error` | Error | Generic failures contain `failure_kind`. Route-execution failures also contain `category` and `phase`, plus nullable `upstream_status` and `target`. | +| `switchyard.routing.decision` | Info | `algorithm`, optional `outcome_id`, initial `selected_model`, nullable final `served_model`, nullable `fallback_used`, and optional `evidence`. | +| `switchyard.routing.error` | Error | Generic failures contain `failure_kind`. Route-execution failures also contain `category`, `phase`, nullable `upstream_status` and `target`, and may contain `outcome_id` and `evidence`. | Call marks describe Switchyard observations, not every HTTP retry made inside a client. `call_role` records whether Switchyard classified the call as routing @@ -336,6 +336,12 @@ telemetry can report the model that answered. selection and `false` when they match. It and `served_model` are `null` when the response does not provide serving metadata. If route execution fails before a response is available, the error mark describes the terminal failure instead. +`outcome_id` is present when the algorithm runner supplies outcome metadata. +When the algorithm supplies evidence, the plugin includes an object containing +supported string fields (`source`, `verdict`, `trigger`, and `reason_code`) and +numeric fields (`score`, `confidence`, and `threshold`). Other fields and values +of the wrong type are omitted. String values longer than 64 bytes are also omitted +and should be stable, non-sensitive labels. ### Metrics