From e2b6ee41ddac18c8e85e4de532e17d5e139adac0 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 25 Sep 2026 22:14:27 +0530 Subject: [PATCH 1/2] feat: export native Langfuse generation metrics --- .../src/observability/langfuse/mod.rs | 43 ++++++++++- .../src/observability/langfuse/test.rs | 71 +++++++++++++++++++ 2 files changed, 113 insertions(+), 1 deletion(-) diff --git a/crates/tinyagents-harness/src/observability/langfuse/mod.rs b/crates/tinyagents-harness/src/observability/langfuse/mod.rs index 8207fd71a..c5568cc13 100644 --- a/crates/tinyagents-harness/src/observability/langfuse/mod.rs +++ b/crates/tinyagents-harness/src/observability/langfuse/mod.rs @@ -119,6 +119,7 @@ impl LangfuseClient { // `body["model"]`; without it Langfuse can't map pricing and every // generation's cost is $0. let call_models = collect_call_models(observations); + let first_deltas = collect_first_deltas(observations); let mut batch = Vec::with_capacity(observations.len() + 2); batch.push(json!({ @@ -156,7 +157,12 @@ impl LangfuseClient { if is_run_lifecycle(&obs.event) { continue; } - batch.push(observation_event(&trace_id, obs, &call_models)); + batch.push(observation_event( + &trace_id, + obs, + &call_models, + &first_deltas, + )); } Ok(json!({ "batch": batch })) @@ -465,10 +471,25 @@ fn collect_call_models(observations: &[AgentObservation]) -> BTreeMap<&str, &str models } +/// First streamed output for each model call. This is the observable TTFT; +/// non-streaming calls have no first-delta timestamp to report. +fn collect_first_deltas(observations: &[AgentObservation]) -> BTreeMap<&str, u64> { + let mut first = BTreeMap::new(); + for obs in observations { + if let AgentEvent::ModelDelta { call_id, delta, .. } = &obs.event { + if !delta.text.is_empty() || !delta.reasoning.is_empty() || delta.tool_call.is_some() { + first.entry(call_id.as_str()).or_insert(obs.ts_ms); + } + } + } + first +} + fn observation_event( trace_id: &str, obs: &AgentObservation, call_models: &BTreeMap<&str, &str>, + first_deltas: &BTreeMap<&str, u64>, ) -> Value { let timestamp = iso_ms(obs.ts_ms); // Every per-call observation nests under its run's span so the trace renders @@ -527,7 +548,10 @@ fn observation_event( // existed. "startTime": started_at_ms.map(iso_ms).unwrap_or_else(|| timestamp.clone()), "endTime": timestamp, + "completionStartTime": first_deltas.get(call_id.as_str()).map(|ms| iso_ms(*ms)), "usage": usage.map(langfuse_usage), + "usageDetails": usage.map(langfuse_usage_details), + "costDetails": usage.and_then(langfuse_cost_details), "input": input, "output": output, "metadata": metadata, @@ -646,6 +670,23 @@ fn langfuse_usage(usage: Usage) -> Value { }) } +fn langfuse_usage_details(usage: Usage) -> Value { + json!({ + "input": usage.input_tokens.saturating_sub(usage.cache_read_tokens), + "output": usage.output_tokens, + "total": usage.total_tokens, + "cache_read_input_tokens": usage.cache_read_tokens, + "cache_creation_input_tokens": usage.cache_creation_tokens, + "reasoning_output_tokens": usage.reasoning_tokens, + }) +} + +fn langfuse_cost_details(usage: Usage) -> Option { + usage + .charged_amount + .map(|amount| json!({ "total": amount.micros as f64 / 1_000_000.0 })) +} + /// Drops every `null`-valued key from a top-level JSON object, in place at /// one level (not recursive). Used before every ingestion body so Langfuse /// never sees an explicit `null` for a field this exporter chose not to diff --git a/crates/tinyagents-harness/src/observability/langfuse/test.rs b/crates/tinyagents-harness/src/observability/langfuse/test.rs index 3289ee901..ceb42316d 100644 --- a/crates/tinyagents-harness/src/observability/langfuse/test.rs +++ b/crates/tinyagents-harness/src/observability/langfuse/test.rs @@ -3,6 +3,8 @@ use super::*; use crate::events::AgentEvent; use crate::ids::{CallId, EventId, RunId}; +use tinyinference_llm::message::MessageDelta; +use tinyinference_llm::usage::ChargedAmount; fn obs(offset: u64, event: AgentEvent) -> AgentObservation { AgentObservation { @@ -141,6 +143,75 @@ fn generation_carries_model_from_model_started() { assert_eq!(generation["body"]["metadata"]["call_id"], "model-call"); } +#[test] +fn generation_reports_native_usage_cost_and_first_token_time() { + let client = LangfuseClient::proxy("https://backend.test", "t").unwrap(); + let call_id = CallId::new("model-call"); + let batch = client + .build_ingestion_batch( + LangfuseTraceConfig::default(), + &[ + obs( + 0, + AgentEvent::ModelStarted { + call_id: call_id.clone(), + model: "chat-v1".into(), + }, + ), + obs( + 25, + AgentEvent::ModelDelta { + run_id: RunId::new("run-1"), + call_id: call_id.clone(), + delta: MessageDelta::text("first token"), + }, + ), + obs( + 30, + AgentEvent::ModelDelta { + run_id: RunId::new("run-1"), + call_id: call_id.clone(), + delta: MessageDelta::text("more"), + }, + ), + obs( + 100, + AgentEvent::ModelCompleted { + call_id, + started_at_ms: Some(1_704_067_200_000), + usage: Some(Usage { + input_tokens: 100, + output_tokens: 20, + total_tokens: 120, + cache_read_tokens: 40, + reasoning_tokens: 5, + charged_amount: Some(ChargedAmount::usd_micros(125)), + ..Default::default() + }), + input: None, + output: None, + }, + ), + ], + ) + .unwrap(); + let generation = batch["batch"] + .as_array() + .unwrap() + .iter() + .find(|event| event["type"] == "generation-create") + .unwrap(); + let body = &generation["body"]; + assert_eq!(body["model"], "chat-v1"); + assert_eq!(body["completionStartTime"], iso_ms(1_704_067_200_025)); + assert_eq!(body["startTime"], iso_ms(1_704_067_200_000)); + assert_eq!(body["endTime"], iso_ms(1_704_067_200_100)); + assert_eq!(body["usageDetails"]["input"], 60); + assert_eq!(body["usageDetails"]["cache_read_input_tokens"], 40); + assert_eq!(body["usageDetails"]["reasoning_output_tokens"], 5); + assert_eq!(body["costDetails"]["total"], 0.000125); +} + #[test] fn populates_generation_and_tool_io_when_captured() { let client = From a525123aa189aa78148c9b318b8e1b3a5ef4c53d Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 25 Sep 2026 22:33:45 +0530 Subject: [PATCH 2/2] fix: satisfy Clippy in Langfuse delta scan --- .../tinyagents-harness/src/observability/langfuse/mod.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/crates/tinyagents-harness/src/observability/langfuse/mod.rs b/crates/tinyagents-harness/src/observability/langfuse/mod.rs index c5568cc13..f71434c95 100644 --- a/crates/tinyagents-harness/src/observability/langfuse/mod.rs +++ b/crates/tinyagents-harness/src/observability/langfuse/mod.rs @@ -476,10 +476,10 @@ fn collect_call_models(observations: &[AgentObservation]) -> BTreeMap<&str, &str fn collect_first_deltas(observations: &[AgentObservation]) -> BTreeMap<&str, u64> { let mut first = BTreeMap::new(); for obs in observations { - if let AgentEvent::ModelDelta { call_id, delta, .. } = &obs.event { - if !delta.text.is_empty() || !delta.reasoning.is_empty() || delta.tool_call.is_some() { - first.entry(call_id.as_str()).or_insert(obs.ts_ms); - } + if let AgentEvent::ModelDelta { call_id, delta, .. } = &obs.event + && (!delta.text.is_empty() || !delta.reasoning.is_empty() || delta.tool_call.is_some()) + { + first.entry(call_id.as_str()).or_insert(obs.ts_ms); } } first