Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
134 changes: 134 additions & 0 deletions crates/libsy-llm-client/tests/observability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Self>,
driver: Driver,
request: Request,
) -> switchyard_libsy::Result<RoutingOutcome> {
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;
Expand Down
59 changes: 53 additions & 6 deletions crates/libsy/src/core/algorithm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,28 +161,59 @@ 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,
/// Its model is set to [`Self::model`] before dispatch.
pub request: DecisionRequest,
/// Target ID for client lookup.
pub model: ModelId,
reply: oneshot::Sender<Result<DecisionResponse>>,
reply: Option<oneshot::Sender<Result<DecisionResponse>>>,
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<DecisionResponse>) -> Result<()> {
pub fn respond(mut self, result: Result<DecisionResponse>) -> 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.
Expand Down Expand Up @@ -364,20 +395,34 @@ 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,
mut request: DecisionRequest,
model: ModelId,
) -> Result<DecisionResponse> {
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))))
Expand Down Expand Up @@ -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`].
Expand Down
56 changes: 46 additions & 10 deletions crates/libsy/src/observability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -265,15 +265,51 @@ pub(crate) fn record_llm_call_span(result: &Result<Response>, 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);
}
}
}
Expand Down
13 changes: 11 additions & 2 deletions docs/reference/opentelemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <model_id>` | 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`:
Expand Down Expand Up @@ -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. |
Expand Down Expand Up @@ -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`.
Expand Down
Loading