Skip to content
Merged
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
43 changes: 42 additions & 1 deletion crates/tinyagents-harness/src/observability/langfuse/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority medium security confident

Scope first-delta timestamps to each call instance

The first-delta map is collected once for the entire observation slice and is keyed only by call_id. Call IDs are reused across turns, so a later generation can inherit an earlier invocation's first-delta timestamp (and potentially a timestamp after its own completion). Track first deltas per concrete call instance, using the invocation's run/turn scope together with the call ID, and look up that scoped key when emitting each generation.

[RULE] invocation-scoped-correlation ·


let mut batch = Vec::with_capacity(observations.len() + 2);
batch.push(json!({
Expand Down Expand Up @@ -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 }))
Expand Down Expand Up @@ -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> {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority medium critique confident

Scope first-delta timestamps to each call instance

collect_first_deltas keys timestamps only by call_id, even though call IDs are reused across runs/turns. If two observations contain ModelDelta { call_id: "model-1", ... } in different runs, the first run's timestamp is reused for the second run's generation whenever the second call has no earlier matching delta (and the map's entry retains the earliest timestamp). This can produce a completionStartTime outside the generation's own start/end window. Include the run identity (or another invocation identity) in the key and use the same scope when looking it up.


Additional security observation

priority medium confident

Correlate first deltas with the specific model invocation

[RULE] invocation-scoped-correlation

collect_first_deltas keys timestamps only by call_id, which is not globally unique and can be reused across turns or runs. When a batch contains multiple invocations with the same call ID, a delta from one invocation can populate completionStartTime for another generation, producing incorrect TTFT telemetry. Key the map by the same invocation scope used to identify the corresponding ModelStarted/ModelCompleted event (at minimum the run/trace identity plus call ID), and use that scoped key when reading it in observation_event.

[RULE] unscoped-call-correlation ·

let mut first = BTreeMap::new();
for obs in observations {
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);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority medium security confident

Scope first-delta timestamps to each call instance

The map is collected across the entire observation slice and keyed only by call_id. Call IDs are reused across turns, so a later generation can inherit an earlier invocation's first-delta timestamp, potentially reporting a completion start before its own start or after its completion. Track first deltas by the concrete call instance, using the invocation's run/turn scope together with the call ID, and use that scoped key when emitting each generation.

[RULE] invocation-scoping ·

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority medium security confident

Correlate first deltas with the specific model invocation

A call_id alone does not uniquely identify a model invocation. When a batch contains multiple invocations with the same call ID, a delta from one invocation can populate completionStartTime for another generation, producing incorrect TTFT telemetry. Key the map by the same invocation scope used to identify the corresponding ModelStarted/ModelCompleted event, at minimum the run or trace identity plus call ID, and use that scoped key during lookup.

[RULE] invocation-correlation ·

}
}
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
Expand Down Expand Up @@ -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)),

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority medium security confident

Correlate first deltas with the specific model invocation

This lookup uses only call_id, which does not uniquely identify a model invocation in a batch containing multiple runs or reused call IDs. A delta from one invocation can populate completionStartTime for another generation, producing incorrect TTFT telemetry. Look up the first delta using the same invocation scope used to identify the corresponding ModelStarted and ModelCompleted events.


Additional critique observation

priority medium confident

Correlate first deltas with the specific model invocation

[RULE] incorrect-call-correlation

The lookup associates a first delta with a generation using only call_id. A call ID can be reused for sequential model invocations within the same run (for example, a retry or repeated turn), so the first invocation's delta timestamp can be attached to the later invocation's ModelCompleted event. The resulting TTFT is incorrect and may precede that generation's startTime. Track deltas per invocation, using run/call identity plus an occurrence or lifecycle pairing, rather than a bare call ID.

[RULE] incorrect-correlation-key ·

"usage": usage.map(langfuse_usage),
"usageDetails": usage.map(langfuse_usage_details),
"costDetails": usage.and_then(langfuse_cost_details),
"input": input,
"output": output,
"metadata": metadata,
Expand Down Expand Up @@ -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<Value> {
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
Expand Down
71 changes: 71 additions & 0 deletions crates/tinyagents-harness/src/observability/langfuse/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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 =
Expand Down
Loading