diff --git a/crates/libsy/src/core/algorithm.rs b/crates/libsy/src/core/algorithm.rs index 0d329601..5704aab5 100644 --- a/crates/libsy/src/core/algorithm.rs +++ b/crates/libsy/src/core/algorithm.rs @@ -5,10 +5,15 @@ //! routing/optimization algorithm implements, and the offload channel it makes model //! calls and publishes [`Decision`]s over. See the crate root for the narrative model. -use std::{pin::Pin, sync::Arc, time::Instant}; +use std::{ + pin::Pin, + sync::Arc, + time::{Duration, Instant}, +}; use async_trait::async_trait; use futures::{Stream, StreamExt}; +use parking_lot::Mutex; use tracing::Instrument; /// The request/response protocol types, re-exported from [`switchyard_protocol`]. @@ -111,6 +116,8 @@ impl CallLlmRequest { #[derive(Clone)] pub struct Driver { driver: TypeErasedDriver, + // How long the call that served this run took. We need this to calculate routing overhead. + routed_call: Arc>>, } impl Driver { @@ -119,9 +126,15 @@ impl Driver { pub(crate) fn new() -> Self { Self { driver: TypeErasedDriver::new(), + routed_call: Arc::new(Mutex::new(None)), } } + /// How long the call that served this run took, if one has succeeded. + pub(crate) fn routed_call_duration(&self) -> Option { + *self.routed_call.lock() + } + /// Offload a model call: publish `routed` as a [`Step::CallLlm`] and await the /// consumer's [`Response`]. The call's context travels inside /// [`routed.ctx`](RoutedRequest::ctx). Errors if the stream is closed or the call failed. @@ -155,15 +168,21 @@ impl Driver { .driver .fulfill_request::(routed.ctx.clone(), routed) .await; + let elapsed = started.elapsed(); observability::record_llm_call( &algorithm, &selected_model, tier.as_deref(), is_routed, - started.elapsed(), + elapsed, &result, &tracing::Span::current(), ); + // Classifier and judge calls are routing overhead. + // And don't record time for failed calls. + if is_routed && result.is_ok() { + *self.routed_call.lock() = Some(elapsed); + } result } @@ -407,10 +426,12 @@ pub trait Algorithm: Send + Sync + 'static { // `libsy.llm_call` spans and decision logs nest inside it via `tracing`'s // contextual parenting. let span = observability::run_span(self.name(), request.metadata.as_ref()); + let observed_driver = task_driver.clone(); let handle = tokio::spawn( async move { observability::observe_run( task_ctx.clone(), + observed_driver, self.create_run_task(task_ctx, task_driver, request), ) .await diff --git a/crates/libsy/src/observability.rs b/crates/libsy/src/observability.rs index 2fbf79c7..8c5de7c1 100644 --- a/crates/libsy/src/observability.rs +++ b/crates/libsy/src/observability.rs @@ -39,7 +39,7 @@ use opentelemetry::metrics::{Meter, ObservableGauge}; use opentelemetry::{global, KeyValue}; use tracing::Span; -use crate::{Context, Decision, Metadata, Response, Result}; +use crate::{Context, Decision, Driver, Metadata, Response, Result}; const METRICS_SCOPE: &str = "switchyard"; const TRACING_TARGET: &str = "libsy"; @@ -133,20 +133,22 @@ pub(crate) fn run_span(algorithm: &str, metadata: Option<&Metadata>) -> Span { } /// Runs one algorithm task to completion, recording the run counter, duration -/// histogram, span outcome, and failure log when it resolves. Executes inside -/// the `libsy.run` span its caller instruments the task with. +/// histogram, routing overhead, span outcome, and failure log when it resolves. +/// Executes inside the `libsy.run` span its caller instruments the task with. +/// `driver` is the run's own, holding the duration of the call that served it. pub(crate) async fn observe_run( ctx: Context, + driver: Driver, run: impl Future>, ) -> Result { let started = Instant::now(); let result = run.await; - record_run( - algorithm_label(&ctx), - started.elapsed(), - &result, - &Span::current(), - ); + let duration = started.elapsed(); + let algorithm = algorithm_label(&ctx); + record_run(algorithm, duration, &result, &Span::current()); + if result.is_ok() { + record_routing_overhead(algorithm, duration, driver.routed_call_duration()); + } result } @@ -192,6 +194,25 @@ fn record_run(algorithm: &str, duration: Duration, result: &Result, sp .record(duration.as_secs_f64() * 1000.0, &attributes); } +/// Records what routing cost on top of the call that served the run: classifier +/// calls, target resolution, decision publishing. A run with no routed call has +/// nothing to subtract, so it records nothing. +fn record_routing_overhead(algorithm: &str, run: Duration, routed_call: Option) { + let Some(routed_call) = routed_call else { + return; + }; + // Saturating: the two clocks start a moment apart, so a run that is all + // routed call can come out fractionally negative. + let overhead_ms = run.saturating_sub(routed_call).as_secs_f64() * 1000.0; + meter() + .f64_histogram("switchyard.routing_overhead_ms") + .build() + .record( + overhead_ms, + &[KeyValue::new("algorithm", algorithm.to_string())], + ); +} + /// Records the resolution of one offloaded model call: the call counter and /// latency histogram, the `outcome`/`error`/token fields on `span`, and a warn /// log when the call failed. diff --git a/crates/libsy/tests/observability.rs b/crates/libsy/tests/observability.rs index 3d540827..813c9453 100644 --- a/crates/libsy/tests/observability.rs +++ b/crates/libsy/tests/observability.rs @@ -13,6 +13,7 @@ use std::collections::BTreeMap; use std::fmt; use std::sync::{Arc, OnceLock}; +use std::time::Duration; use async_trait::async_trait; use futures::StreamExt; @@ -255,6 +256,22 @@ fn f64_histogram_count( }) } +/// Latest cumulative sample sum of an `f64` histogram, in whole milliseconds. +fn f64_histogram_sum_ms( + snapshots: &[ResourceMetrics], + name: &str, + wanted: &[(&str, &str)], +) -> Option { + latest_metric_value(snapshots, name, |data| match data { + AggregatedMetrics::F64(MetricData::Histogram(histogram)) => histogram + .data_points() + .filter(|point| attributes_match(point.attributes(), wanted)) + .map(|point| point.sum() as u64) + .collect(), + _ => Vec::new(), + }) +} + /// Latest value of a `u64` observable gauge. fn u64_gauge_value(snapshots: &[ResourceMetrics], name: &str) -> Option { latest_metric_value(snapshots, name, |data| match data { @@ -288,8 +305,12 @@ struct UsageClient { usage: Usage, } -/// Client that returns a weak classifier verdict. -struct ClassifierClient; +/// Client that returns a weak classifier verdict. The delays let a test tell +/// classifier time apart from routed-call time. +struct ClassifierClient { + classifier_delay: Duration, + routed_delay: Duration, +} #[async_trait] impl RoutedLlmClient for ClassifierClient { @@ -301,8 +322,10 @@ impl RoutedLlmClient for ClassifierClient { ) -> Result { let model = decision.selected_model().to_string(); let completion = if decision.is_routed_call() { + tokio::time::sleep(self.routed_delay).await; "routed response" } else { + tokio::time::sleep(self.classifier_delay).await; r#"{"recommended_route":"weak","p_solve":0.9,"confidence":0.9,"abstain":false,"capability_boundary":"supported","primary_rule":"SUP-1","crux":"bounded task"}"# }; Ok(Response { @@ -483,6 +506,15 @@ async fn successful_run_records_metrics_spans_and_decision_log() -> switchyard_l u64_gauge_value(&snapshots, "switchyard.total_errors"), Some(total_errors_before) ); + // One overhead observation per run, keyed by algorithm alone. + assert_eq!( + f64_histogram_count( + &snapshots, + "switchyard.routing_overhead_ms", + &[("algorithm", ALGO)] + ), + Some(1) + ); // Spans: one run span carrying the correlation ids and outcome, one child // llm_call span carrying the selection, outcome, and token counts. @@ -637,6 +669,15 @@ async fn failed_call_records_error_outcome_and_warn_logs() -> switchyard_libsy:: u64_gauge_value(&snapshots, "switchyard.total_errors"), Some(total_errors_before + 1) ); + // Nothing was served, so there is nothing to measure routing against. + assert_eq!( + f64_histogram_count( + &snapshots, + "switchyard.routing_overhead_ms", + &[("algorithm", ALGO)] + ), + None + ); // Spans: both spans carry outcome=error and the propagated error text. let spans = store.spans(); @@ -692,7 +733,10 @@ async fn classifier_metrics_count_only_the_final_routed_call() -> switchyard_lib let total_requests_before = u64_gauge_value(&before, "switchyard.total_requests").unwrap_or_default(); - let client = Arc::new(ClassifierClient); + let client = Arc::new(ClassifierClient { + classifier_delay: Duration::from_millis(60), + routed_delay: Duration::from_millis(200), + }); let target = |name: &str| LlmTarget { semantic_name: name.to_string(), llm_client: Some(client.clone()), @@ -770,5 +814,17 @@ async fn classifier_metrics_count_only_the_final_routed_call() -> switchyard_lib u64_gauge_value(&snapshots, "switchyard.total_requests"), Some(total_requests_before + 1) ); + // The classifier call is the router's own work but the routed call is not, + // so overhead lands near the classifier's 60ms, not their 260ms sum. + let overhead = f64_histogram_sum_ms( + &snapshots, + "switchyard.routing_overhead_ms", + &[("algorithm", "llm_task_classifier")], + ) + .unwrap_or_default(); + assert!( + (60..200).contains(&overhead), + "expected roughly the classifier's 60ms, got {overhead}ms" + ); Ok(()) } diff --git a/crates/switchyard-server/README.md b/crates/switchyard-server/README.md index 476de372..a7060d9c 100644 --- a/crates/switchyard-server/README.md +++ b/crates/switchyard-server/README.md @@ -103,6 +103,7 @@ Routed-call compatibility metrics are: | `switchyard_cache_creation_tokens_total` | counter | `model`, optional `tier` | Cache-creation input tokens | | `switchyard_reasoning_tokens_total` | counter | `model`, optional `tier` | Reasoning output tokens | | `switchyard_total_latency_ms` | histogram | `model`, optional `tier` | Full-turn latency for successful routed responses | +| `switchyard_routing_overhead_ms` | histogram | `algorithm` | Algorithm run time minus the call that served it | | `switchyard_client_responses_total` | counter | `outcome` | Final LLM-route responses | | `switchyard_upstream_attempts_total` | counter | `outcome`, `code` | Actual upstream HTTP attempts | | `switchyard_router_retry_recovered_total` | counter | none | Retry recoveries (currently always zero) | @@ -118,4 +119,16 @@ the server sees the request. The Rust server exports this metric as a histogram, server exports its counterpart as a summary; this matches the existing histogram/summary difference for model-call latency. +`switchyard_routing_overhead_ms` is what routing cost on top of the model call: the algorithm's run +time minus the call that served the request. Classifier calls are not subtracted, so an +LLM-classifier route reports its classification time here while `passthrough` and `random` report +the sub-millisecond cost of picking a target. It carries only `algorithm`, since the number +describes the router and not the target it chose, and a run that served nothing records nothing. Its +buckets start at 0.1 ms via a view in the server; the SDK defaults start at 5 ms. + +Both clocks stop when the routed call resolves, which for a streamed response is when the stream +handle arrives rather than when the stream ends, so SSE relay time is in neither term. The Python +summary of the same name measures its total through stream completion, making its streaming values +mostly generation time. + See [CONFIGURATION.md](CONFIGURATION.md) to add an LLM client, target, or algorithm. diff --git a/crates/switchyard-server/src/metrics.rs b/crates/switchyard-server/src/metrics.rs index d27b95b9..8cd3df14 100644 --- a/crates/switchyard-server/src/metrics.rs +++ b/crates/switchyard-server/src/metrics.rs @@ -6,12 +6,19 @@ use std::sync::OnceLock; use opentelemetry::{global, KeyValue}; -use opentelemetry_sdk::metrics::SdkMeterProvider; +use opentelemetry_sdk::metrics::{Aggregation, Instrument, SdkMeterProvider, Stream}; use prometheus::{Encoder, Registry, TextEncoder}; use switchyard_llm_client::metrics::{http_outcome_label, http_status_code_label}; pub(crate) const CONTENT_TYPE: &str = "text/plain; version=0.0.4; charset=utf-8"; +/// Bucket boundaries for `switchyard.routing_overhead_ms`. +/// Need a broad range because some algos call an LLM (classifier), and some +/// do very little (passthrough). +const ROUTING_OVERHEAD_BUCKETS_MS: &[f64] = &[ + 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 25.0, 50.0, 100.0, 250.0, 500.0, 1000.0, 2500.0, 5000.0, +]; + struct Metrics { registry: Registry, _provider: SdkMeterProvider, @@ -33,7 +40,10 @@ fn initialize() -> Result { .with_registry(registry.clone()) .build() .map_err(|error| format!("failed to initialize Prometheus metrics: {error}"))?; - let provider = SdkMeterProvider::builder().with_reader(exporter).build(); + let provider = SdkMeterProvider::builder() + .with_reader(exporter) + .with_view(routing_overhead_buckets) + .build(); global::set_meter_provider(provider.clone()); libsy::initialize_metrics(); global::meter("switchyard") @@ -47,6 +57,20 @@ fn initialize() -> Result { }) } +fn routing_overhead_buckets(instrument: &Instrument) -> Option { + if instrument.name() != "switchyard.routing_overhead_ms" { + return None; + } + Stream::builder() + .with_aggregation(Aggregation::ExplicitBucketHistogram { + boundaries: ROUTING_OVERHEAD_BUCKETS_MS.to_vec(), + // Cumulative min/max cover the whole process, so they aren't useful. + record_min_max: false, + }) + .build() + .ok() +} + /// Make the metrics exist before they get a hit. Nicer for dashboards but not really necessary. /// The HTTP status codes we seed are somewhat arbitrary. fn seed_outcome_metrics() { diff --git a/crates/switchyard-server/tests/server.rs b/crates/switchyard-server/tests/server.rs index 9735b3cf..72fc3608 100644 --- a/crates/switchyard-server/tests/server.rs +++ b/crates/switchyard-server/tests/server.rs @@ -254,6 +254,7 @@ async fn metrics_exposes_switchyard_otel_instruments() -> TestResult { "# TYPE switchyard_completion_tokens_total counter", "# TYPE switchyard_cached_tokens_total counter", "# TYPE switchyard_total_latency_ms histogram", + "# TYPE switchyard_routing_overhead_ms histogram", "algorithm=\"random\"", &format!("selected_model=\"{MODEL}\""), ] { @@ -274,6 +275,13 @@ async fn metrics_exposes_switchyard_otel_instruments() -> TestResult { "unexpected delta for {name}" ); } + // A sub-millisecond boundary exists only because of the server's bucket view. + assert!(metric_line( + metrics, + "switchyard_routing_overhead_ms_bucket", + &[("algorithm", "random"), ("le", "0.1")] + ) + .is_some()); assert!(metric_line( metrics, "switchyard_cache_creation_tokens_total",