diff --git a/crates/libsy-llm-client/tests/observability.rs b/crates/libsy-llm-client/tests/observability.rs index 8f5071ccd..5c3207a49 100644 --- a/crates/libsy-llm-client/tests/observability.rs +++ b/crates/libsy-llm-client/tests/observability.rs @@ -1603,6 +1603,140 @@ async fn upstream_body_is_redacted_from_the_client_call_span() -> switchyard_lib Ok(()) } +#[tokio::test] +async fn decision_calls_record_each_terminal_path_once() -> switchyard_libsy::Result<()> { + use switchyard_protocol::{DecisionRequest, DecisionResponse}; + + let _guard = serialize_test().lock().await; + let (store, exporter, provider, _, _) = telemetry(); + struct DecisionAlgo(String); + + #[async_trait] + impl Algorithm for DecisionAlgo { + fn name(&self) -> &str { + &self.0 + } + + async fn route( + self: Arc, + driver: Driver, + request: Request, + ) -> switchyard_libsy::Result { + driver + .call_decision( + DecisionRequest { + model: None, + context: json!(LEAKED_CONTENT), + questions: Default::default(), + }, + self.0.clone().into(), + ) + .await?; + Ok(RoutingOutcome::route_to("answer".into(), vec![], request)) + } + } + + for mode in ["reply", "error", "fail", "drop", "cancel"] { + let algorithm = format!("obs-decision-{mode}"); + let name = algorithm.as_str(); + let mut stream = Box::pin(Arc::new(DecisionAlgo(algorithm.clone())).run_stream( + request_with_metadata("decision-session", "decision-correlation"), + Arc::new(RuntimeModels::default()), + )); + let Some(Ok(Step::CallDecision(call))) = stream.next().await else { + return Err(test_error("expected a decision call")); + }; + assert_eq!( + u64_counter_value( + &flushed_metrics(exporter, provider), + "switchyard.decision_calls", + &[("algorithm", name)], + ), + None, + ); + tokio::time::sleep(Duration::from_millis(10)).await; + if mode == "cancel" { + drop(stream); + drop(call); + } else { + match mode { + "reply" => call.respond(Ok(DecisionResponse { + id: Some("decision-response".into()), + model: Some("provider-model".into()), + answers: Default::default(), + usage: Usage { + input_tokens: Some(42), + output_tokens: Some(5), + ..Usage::default() + }, + }))?, + "error" => call.respond(Err(test_error(LEAKED_CONTENT)))?, + "fail" => assert!(call.fail(test_error(LEAKED_CONTENT)).is_err()), + "drop" => drop(call), + _ => unreachable!(), + } + while stream.next().await.is_some() {} + } + let outcome = if mode == "reply" { "ok" } else { "error" }; + let attrs = [ + ("algorithm", name), + ("selected_model", name), + ("outcome", outcome), + ]; + let snapshots = flushed_metrics(exporter, provider); + assert_eq!( + u64_counter_value(&snapshots, "switchyard.decision_calls", &attrs), + Some(1) + ); + assert_eq!( + u64_counter_value( + &snapshots, + "switchyard.decision_calls", + &[ + ("algorithm", name), + ("outcome", if mode == "reply" { "error" } else { "ok" }) + ], + ), + None, + ); + assert_eq!( + f64_histogram_count(&snapshots, "switchyard.decision_call_duration_ms", &attrs), + Some(1) + ); + assert!( + f64_histogram_sum_ms(&snapshots, "switchyard.decision_call_duration_ms", &attrs) + .unwrap_or(0) + >= 10 + ); + assert_eq!( + u64_counter_value(&snapshots, "switchyard.llm_calls", &[("algorithm", name)]), + None + ); + + let span = find_span(&store.spans(), "libsy.decision_call", "algorithm", name); + assert_eq!(span.parent.as_deref(), Some("libsy.run")); + assert_eq!( + span.fields.get("outcome").map(String::as_str), + Some(outcome) + ); + for (field, value) in [ + ("input_tokens", "42"), + ("output_tokens", "5"), + ("gen_ai.response.id", "decision-response"), + ("gen_ai.response.model", "provider-model"), + ] { + assert_eq!( + span.fields.get(field).map(String::as_str), + (mode == "reply").then_some(value) + ); + } + assert!(!span.fields.contains_key("total_tokens")); + assert!(!span.fields.contains_key("error")); + assert!(!format!("{span:?}").contains(LEAKED_CONTENT)); + } + Ok(()) +} + #[tokio::test] async fn failed_call_records_metrics_without_error_details() -> switchyard_libsy::Result<()> { let _guard = serialize_test().lock().await; diff --git a/crates/libsy/src/core/algorithm.rs b/crates/libsy/src/core/algorithm.rs index aabab68a3..249805a01 100644 --- a/crates/libsy/src/core/algorithm.rs +++ b/crates/libsy/src/core/algorithm.rs @@ -161,6 +161,7 @@ impl Drop for CallModel { } /// A host-owned decision call. Dropping it without replying yields [`DriverError::ResponseDropped`]. +/// Completion or drop records one call and its duration, including time waiting for the host. pub struct CallDecision { /// Algorithm name for attributing host telemetry. pub algorithm: String, @@ -168,21 +169,51 @@ pub struct CallDecision { pub request: DecisionRequest, /// Target ID for client lookup. pub model: ModelId, - reply: oneshot::Sender>, + reply: Option>>, + started: Instant, + // Retain the originating span even if the waiting algorithm is cancelled first. + span: tracing::Span, } impl CallDecision { /// Return a response or provider error to the algorithm so it can continue or fall back. - pub fn respond(self, result: Result) -> Result<()> { + pub fn respond(mut self, result: Result) -> Result<()> { + self.record(result.is_ok()); + if let Ok(response) = &result { + observability::record_decision_response(response, &self.span); + } self.reply + .take() + .ok_or(DriverError::ResponseDropped)? .send(result) .map_err(|_| DriverError::ResponseDropped.into()) } /// Returning this error from the host handler aborts [`drive`]. - pub fn fail(self, error: LibsyError) -> Result<()> { + pub fn fail(mut self, error: LibsyError) -> Result<()> { + self.reply = None; + self.record(false); Err(error) } + + fn record(&self, is_ok: bool) { + observability::record_decision_call( + &self.algorithm, + &self.model, + self.started.elapsed(), + is_ok, + ); + self.span + .record("outcome", if is_ok { "ok" } else { "error" }); + } +} + +impl Drop for CallDecision { + fn drop(&mut self) { + if self.reply.is_some() { + self.record(false); + } + } } /// The terminal result of routing. @@ -364,7 +395,18 @@ impl Driver { target = "libsy", name = "libsy.decision_call", skip_all, - fields(algorithm = self.algorithm, selected_model = %model), + fields( + algorithm = self.algorithm, + selected_model = %model, + openinference.span.kind = "CHAIN", + outcome = tracing::field::Empty, + input_tokens = tracing::field::Empty, + output_tokens = tracing::field::Empty, + total_tokens = tracing::field::Empty, + reasoning_tokens = tracing::field::Empty, + gen_ai.response.id = tracing::field::Empty, + gen_ai.response.model = tracing::field::Empty, + ), )] pub async fn call_decision( &self, @@ -372,12 +414,15 @@ impl Driver { model: ModelId, ) -> Result { request.model = Some(model.clone()); + let started = Instant::now(); let (reply, response) = oneshot::channel(); let call = CallDecision { algorithm: self.algorithm.clone(), request, model, - reply, + reply: Some(reply), + started, + span: tracing::Span::current(), }; self.step_tx .send(Ok(Step::CallDecision(Box::new(call)))) @@ -598,7 +643,9 @@ impl RoutingIdentity { /// # Observability /// /// [`run_stream`](Self::run_stream) creates a `libsy.run` span, and each offloaded model -/// call creates a nested `libsy.llm_call` span. Successful outcomes record their +/// call creates a nested span for its call kind, with outcome and available token usage. +/// Call metrics include host queueing and count unfulfilled drops as errors. +/// Successful outcomes record their /// [`OutcomeMetadata::outcome_id`](crate::OutcomeMetadata::outcome_id) on `libsy.run`, /// alongside `selected_model_ids` (an ordered OpenTelemetry string array). /// `algorithm` and `switchyard.algorithm` retain the run's [`Algorithm::name`]. diff --git a/crates/libsy/src/observability.rs b/crates/libsy/src/observability.rs index 4779540da..4775256c3 100644 --- a/crates/libsy/src/observability.rs +++ b/crates/libsy/src/observability.rs @@ -41,7 +41,7 @@ use tracing::Span; use tracing_opentelemetry::OpenTelemetrySpanExt; use crate::{OutcomeMetadata, Result}; -use switchyard_protocol::{ModelId, Request, Response}; +use switchyard_protocol::{DecisionResponse, ModelId, Request, Response, Usage}; const METRICS_SCOPE: &str = "switchyard"; const TRACING_TARGET: &str = "libsy"; @@ -265,15 +265,51 @@ pub(crate) fn record_llm_call_span(result: &Result, span: &Span) { let Some(usage) = response.llm_response.as_agg().map(|agg| &agg.usage) else { return; }; - for (field, value) in [ - ("input_tokens", usage.input_tokens), - ("output_tokens", usage.output_tokens), - ("total_tokens", usage.total_tokens), - ("reasoning_tokens", usage.reasoning_tokens), - ] { - if let Some(value) = value { - span.record(field, value); - } + record_call_usage(usage, span); + } +} + +pub(crate) fn record_decision_call( + algorithm: &str, + selected_model: &str, + duration: Duration, + is_ok: bool, +) { + let attributes = [ + KeyValue::new("algorithm", algorithm.to_string()), + KeyValue::new("selected_model", selected_model.to_string()), + KeyValue::new("outcome", if is_ok { "ok" } else { "error" }), + ]; + let meter = meter(); + meter + .u64_counter("switchyard.decision_calls") + .build() + .add(1, &attributes); + meter + .f64_histogram("switchyard.decision_call_duration_ms") + .build() + .record(duration.as_secs_f64() * 1000.0, &attributes); +} + +pub(crate) fn record_decision_response(response: &DecisionResponse, span: &Span) { + if let Some(id) = &response.id { + span.record("gen_ai.response.id", id.as_str()); + } + if let Some(model) = &response.model { + span.record("gen_ai.response.model", model.as_str()); + } + record_call_usage(&response.usage, span); +} + +fn record_call_usage(usage: &Usage, span: &Span) { + for (field, value) in [ + ("input_tokens", usage.input_tokens), + ("output_tokens", usage.output_tokens), + ("total_tokens", usage.total_tokens), + ("reasoning_tokens", usage.reasoning_tokens), + ] { + if let Some(value) = value { + span.record(field, value); } } } diff --git a/docs/reference/opentelemetry.md b/docs/reference/opentelemetry.md index 3e4c4657e..2d798b41d 100644 --- a/docs/reference/opentelemetry.md +++ b/docs/reference/opentelemetry.md @@ -16,10 +16,16 @@ and [OTLP configuration](https://opentelemetry.io/docs/specs/otel/protocol/expor |---|---|---| | `libsy.run` | Libsy | One algorithm run, including routing-time work. OpenInference kind `CHAIN`. | | `libsy.llm_call` | Libsy driver | Waiting for the host to fulfill an offloaded call. Includes host queueing. OpenInference kind `CHAIN`. | +| `libsy.decision_call` | Libsy driver | An offloaded Decision Model call, including host queueing. OpenInference kind `CHAIN`. | | `libsy.client_call`, exported as `chat ` | LLM client driver | One candidate model call, including that candidate's retries. OTel kind `CLIENT`; OpenInference kind `LLM`. | Hosts driving `run_stream` without the LLM client driver instrument their own model I/O. +`libsy.decision_call` records `algorithm`, `selected_model`, and terminal `outcome`. +Successful replies add available `input_tokens`, `output_tokens`, `total_tokens`, +`reasoning_tokens`, `gen_ai.response.id`, and `gen_ai.response.model` fields. +Unknown values are omitted; request content, answers, and error details are not recorded. + ### Routing outcome fields On `libsy.run`: @@ -83,8 +89,10 @@ Metrics use the `switchyard` meter scope. The tables use OTel instrument names. | `switchyard.run_duration_ms` | Histogram | `algorithm`, `outcome` | Algorithm-task duration in milliseconds. | | `switchyard.algorithms_in_flight` | UpDownCounter | `algorithm` | Active algorithm tasks; exported as a Prometheus gauge. | | `switchyard.decisions` | Counter | `algorithm`, `selected_model` | Published routing decisions. | -| `switchyard.llm_calls` | Counter | `algorithm`, `selected_model`, `outcome` | Routing-time offloaded calls and each terminal answer candidate, including failures. | -| `switchyard.llm_call_duration_ms` | Histogram | `algorithm`, `selected_model`, `outcome` | One duration sample in milliseconds per offloaded call or terminal answer candidate, including failures; see streaming limits below. | +| `switchyard.llm_calls` | Counter | `algorithm`, `selected_model`, `outcome` | Routing-time LLM calls and each terminal answer candidate, including failures. | +| `switchyard.llm_call_duration_ms` | Histogram | `algorithm`, `selected_model`, `outcome` | One duration sample in milliseconds per offloaded LLM call or terminal answer candidate, including failures; see streaming limits below. | +| `switchyard.decision_calls` | Counter | `algorithm`, `selected_model`, `outcome` | Offloaded Decision Model calls, recorded once on reply, failure, or unfulfilled drop. | +| `switchyard.decision_call_duration_ms` | Histogram | `algorithm`, `selected_model`, `outcome` | One duration sample in milliseconds per Decision Model call, including host queueing. | | `switchyard.total_requests` | ObservableGauge | none | Process-wide total of successful and failed answer candidates after routing, including reused routing responses. | | `switchyard.total_errors` | ObservableGauge | none | Process-wide total of failures after routing, including failed answer candidates and reused routing responses. | | `switchyard.requests` | Counter | `model` | Successful answer candidates or reused routing responses, by model ID. | @@ -165,6 +173,7 @@ Advisor Gate instruments use the prefix `switchyard.advisor_gate.`: ### Timing and streaming - `libsy.run` may finish before the answer call. Nested algorithms have separate run spans. +- Decision calls record `error` when failed or dropped without a reply, including cancellation. A returned response records its outcome even if the waiting algorithm has already gone away. These metrics work with any host and do not require a client observer. - `libsy.llm_call` and routing-time call metrics end when the host fulfills the offloaded call. They include host queueing. The LLM client driver buffers routing streams before fulfilling the call. The span's `input_tokens`, `output_tokens`, `total_tokens`, and `reasoning_tokens` fields are buffered-response only. - For terminal answer candidates, `switchyard.llm_calls` and `switchyard.llm_call_duration_ms` record when the response stream ends or is dropped. Duration includes that candidate's retries and stream consumption. Stream errors and unfinished drops record `outcome=error`; a terminal message permits a successful drop. - `libsy.client_call` remains open while its stream is consumed. IDs, usage, and finish reasons update from normalized events. An unfinished stream dropped by its consumer records `cancelled`; a stream error records `error`.