mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-07 02:59:05 +00:00
fix(tracing): retain native logs and headless tool results (#44894)
* fix(tracing): retain native logs and headless tool results * fix(tracing): bound uncorrelated log annotations after decoding * fix(tracing): canonicalize absent native log context
This commit is contained in:
parent
8f041a8f06
commit
5086fb3038
4 changed files with 260 additions and 31 deletions
|
|
@ -159,10 +159,23 @@ fn exported_tool_results(attributes: &BTreeMap<String, String>) -> String {
|
|||
let Ok(body) = serde_json::from_str::<Value>(attr(attributes, "body")) else {
|
||||
return json!({"warning": "Claude's API body export is missing or truncated. Some tool results may be unavailable."}).to_string();
|
||||
};
|
||||
let results: Vec<Value> = body.get("messages").and_then(Value::as_array)
|
||||
.and_then(|messages| messages.last())
|
||||
.filter(|message| message.get("role").and_then(Value::as_str) == Some("user"))
|
||||
.and_then(|message| message.get("content").and_then(Value::as_array))
|
||||
let message = body
|
||||
.get("messages")
|
||||
.and_then(Value::as_array)
|
||||
.and_then(|messages| {
|
||||
messages
|
||||
.iter()
|
||||
.rev()
|
||||
.find(|message| message.get("role").and_then(Value::as_str) != Some("system"))
|
||||
})
|
||||
.filter(|message| message.get("role").and_then(Value::as_str) == Some("user"));
|
||||
let Some(content) = message.and_then(|message| message.get("content")) else {
|
||||
return json!({"warning": "Claude's API body export has an unexpected message shape. Some tool results may be unavailable."}).to_string();
|
||||
};
|
||||
if !content.is_array() && !content.is_string() {
|
||||
return json!({"warning": "Claude's API body export has an unexpected content shape. Some tool results may be unavailable."}).to_string();
|
||||
}
|
||||
let results: Vec<Value> = content.as_array()
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_result"))
|
||||
|
|
|
|||
|
|
@ -28,16 +28,31 @@ fn text<'a>(attributes: &'a [KeyValue], key: &str) -> &'a str {
|
|||
}
|
||||
}
|
||||
|
||||
fn message(record: LogRecord) -> Span {
|
||||
fn absent_id(id: &[u8], length: usize) -> bool {
|
||||
id.is_empty() || (id.len() == length && id.iter().all(|byte| *byte == 0))
|
||||
}
|
||||
|
||||
struct LogContext {
|
||||
missing_trace: bool,
|
||||
missing_parent: bool,
|
||||
}
|
||||
|
||||
fn message(record: LogRecord) -> (Span, LogContext) {
|
||||
let timestamp = if record.time_unix_nano == 0 {
|
||||
record.observed_time_unix_nano
|
||||
} else {
|
||||
record.time_unix_nano
|
||||
};
|
||||
let missing_trace = absent_id(&record.trace_id, 16);
|
||||
let missing_parent = absent_id(&record.span_id, 8);
|
||||
let mut hash = Sha256::new();
|
||||
hash.update(b"litellm.claude.message.v1\0");
|
||||
hash.update(&record.trace_id);
|
||||
hash.update(&record.span_id);
|
||||
if !missing_trace {
|
||||
hash.update(&record.trace_id);
|
||||
}
|
||||
if !missing_parent {
|
||||
hash.update(&record.span_id);
|
||||
}
|
||||
hash.update(text(&record.attributes, "event.name"));
|
||||
let uuid = text(&record.attributes, "message.uuid");
|
||||
if uuid.is_empty() {
|
||||
|
|
@ -50,28 +65,53 @@ fn message(record: LogRecord) -> Span {
|
|||
} else {
|
||||
hash.update(uuid);
|
||||
}
|
||||
if missing_trace {
|
||||
hash.update(b"\0unassigned\0");
|
||||
hash.update(text(&record.attributes, "session.id"));
|
||||
}
|
||||
let identity = hash.finalize();
|
||||
let trace_id = if missing_trace {
|
||||
let session = text(&record.attributes, "session.id");
|
||||
if session.is_empty() {
|
||||
identity[..16].to_vec()
|
||||
} else {
|
||||
span::session_trace_id(session)
|
||||
}
|
||||
} else {
|
||||
record.trace_id
|
||||
};
|
||||
let failed = match value(&record.attributes, "success") {
|
||||
Some(Value::BoolValue(success)) => !success,
|
||||
Some(Value::StringValue(success)) => success == "false",
|
||||
_ => false,
|
||||
};
|
||||
Span {
|
||||
trace_id: record.trace_id,
|
||||
span_id: hash.finalize()[..8].to_vec(),
|
||||
parent_span_id: record.span_id,
|
||||
name: format!("claude_code.{}", text(&record.attributes, "event.name")),
|
||||
kind: 1,
|
||||
start_time_unix_nano: timestamp,
|
||||
end_time_unix_nano: timestamp,
|
||||
status: (text(&record.attributes, "event.name") == "tool_result" && failed).then(|| {
|
||||
Status {
|
||||
code: 2,
|
||||
message: text(&record.attributes, "error").to_owned(),
|
||||
}
|
||||
}),
|
||||
attributes: record.attributes,
|
||||
..Span::default()
|
||||
}
|
||||
(
|
||||
Span {
|
||||
trace_id,
|
||||
span_id: identity[..8].to_vec(),
|
||||
parent_span_id: if missing_trace || missing_parent {
|
||||
Vec::new()
|
||||
} else {
|
||||
record.span_id
|
||||
},
|
||||
name: format!("claude_code.{}", text(&record.attributes, "event.name")),
|
||||
kind: 1,
|
||||
start_time_unix_nano: timestamp,
|
||||
end_time_unix_nano: timestamp,
|
||||
status: (text(&record.attributes, "event.name") == "tool_result" && failed).then(
|
||||
|| Status {
|
||||
code: 2,
|
||||
message: text(&record.attributes, "error").to_owned(),
|
||||
},
|
||||
),
|
||||
attributes: record.attributes,
|
||||
..Span::default()
|
||||
},
|
||||
LogContext {
|
||||
missing_trace,
|
||||
missing_parent,
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
pub(super) fn flatten(
|
||||
|
|
@ -79,6 +119,7 @@ pub(super) fn flatten(
|
|||
limits: DecodeLimits,
|
||||
) -> Result<Vec<DecodedSpan>, Error> {
|
||||
let mut count = 0usize;
|
||||
let mut contexts = Vec::new();
|
||||
let mut resources = Vec::new();
|
||||
for resource in request.resource_logs {
|
||||
let mut scopes = Vec::new();
|
||||
|
|
@ -111,7 +152,11 @@ pub(super) fn flatten(
|
|||
text(&record.attributes, "query_source"),
|
||||
)))
|
||||
})
|
||||
.map(message)
|
||||
.map(|record| {
|
||||
let (span, context) = message(record);
|
||||
contexts.push(context);
|
||||
span
|
||||
})
|
||||
.collect();
|
||||
scopes.push(ScopeSpans {
|
||||
scope: scope.scope,
|
||||
|
|
@ -125,10 +170,23 @@ pub(super) fn flatten(
|
|||
schema_url: resource.schema_url,
|
||||
});
|
||||
}
|
||||
span::flatten(
|
||||
let (mut spans, mut budget) = span::flatten_with_budget(
|
||||
ExportTraceServiceRequest {
|
||||
resource_spans: resources,
|
||||
},
|
||||
limits,
|
||||
)
|
||||
)?;
|
||||
budget.consume(contexts.len() * size_of::<LogContext>())?;
|
||||
for (span, context) in spans.iter_mut().zip(contexts) {
|
||||
if context.missing_trace {
|
||||
span.attributes.remove("lens.original_trace_id");
|
||||
}
|
||||
if context.missing_trace || context.missing_parent {
|
||||
let key = "lens.capture.warning";
|
||||
let warning = "This native log has incomplete trace context. Its content is retained, but its execution parent is unconfirmed.";
|
||||
budget.consume(key.len() + warning.len() + 96)?;
|
||||
span.attributes.insert(key.to_owned(), warning.to_owned());
|
||||
}
|
||||
}
|
||||
Ok(spans)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -20,12 +20,19 @@ pub(super) fn flatten(
|
|||
request: ExportTraceServiceRequest,
|
||||
limits: DecodeLimits,
|
||||
) -> Result<Vec<DecodedSpan>, Error> {
|
||||
flatten_with_budget(request, limits).map(|(spans, _)| spans)
|
||||
}
|
||||
|
||||
pub(super) fn flatten_with_budget(
|
||||
request: ExportTraceServiceRequest,
|
||||
limits: DecodeLimits,
|
||||
) -> Result<(Vec<DecodedSpan>, Budget), Error> {
|
||||
let mut budget = Budget::new(limits);
|
||||
let mut spans = Vec::new();
|
||||
for resource in request.resource_spans {
|
||||
append_resource(resource, &mut budget, &mut spans)?;
|
||||
}
|
||||
Ok(spans)
|
||||
Ok((spans, budget))
|
||||
}
|
||||
|
||||
fn append_resource(
|
||||
|
|
@ -124,6 +131,10 @@ fn hex_bytes(bytes: &[u8]) -> String {
|
|||
bytes.iter().map(|byte| format!("{byte:02x}")).collect()
|
||||
}
|
||||
|
||||
pub(super) fn session_trace_id(session: &str) -> Vec<u8> {
|
||||
Sha256::digest(format!("litellm.claude.session.v1\0{session}"))[..16].to_vec()
|
||||
}
|
||||
|
||||
fn decoded_span(
|
||||
span: Span,
|
||||
resource_attributes: &Shared<BTreeMap<String, String>>,
|
||||
|
|
@ -145,8 +156,7 @@ fn decoded_span(
|
|||
.get("session.id")
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
let trace_id =
|
||||
hex_bytes(&Sha256::digest(format!("litellm.claude.session.v1\0{session}"))[..16]);
|
||||
let trace_id = hex_bytes(&session_trace_id(session));
|
||||
let actor = span_attributes.get("agent_id").unwrap_or(session).clone();
|
||||
budget.consume(original_trace_id.len() + actor.len() + 256)?;
|
||||
span_attributes.insert("lens.original_trace_id".to_owned(), original_trace_id);
|
||||
|
|
|
|||
|
|
@ -1673,7 +1673,7 @@ fn session_capture_joins_native_logs_and_traces_across_turns_without_changing_sp
|
|||
|
||||
#[rstest]
|
||||
#[case::short_trace(vec![1;15], vec![2;8], 1)]
|
||||
#[case::zero_parent(vec![1;16], vec![0;8], 1)]
|
||||
#[case::short_parent(vec![1;16], vec![2;7], 1)]
|
||||
#[case::timestamp(vec![1;16], vec![2;8], i64::MAX as u64 + 1)]
|
||||
fn native_logs_reject_invalid_context(
|
||||
#[case] trace: Vec<u8>,
|
||||
|
|
@ -1692,6 +1692,106 @@ fn native_logs_reject_invalid_context(
|
|||
));
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::json_empty(true, false, true, true)]
|
||||
#[case::protobuf_empty(false, false, true, true)]
|
||||
#[case::json_zero(true, true, true, true)]
|
||||
#[case::protobuf_zero(false, true, true, true)]
|
||||
#[case::trace_only(true, false, true, false)]
|
||||
#[case::parent_only(false, false, false, true)]
|
||||
fn contextless_native_logs_preserve_the_entire_batch(
|
||||
#[case] json: bool,
|
||||
#[case] zero: bool,
|
||||
#[case] missing_trace: bool,
|
||||
#[case] missing_parent: bool,
|
||||
) {
|
||||
use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value::Value};
|
||||
use prost::Message;
|
||||
let mut request = log_request("repl_main_thread");
|
||||
let valid = request.resource_logs[0].scope_logs[0].log_records[0].clone();
|
||||
let mut uncorrelated = valid.clone();
|
||||
if missing_trace {
|
||||
uncorrelated.trace_id = if zero { vec![0; 16] } else { Vec::new() };
|
||||
}
|
||||
if missing_parent {
|
||||
uncorrelated.span_id = if zero { vec![0; 8] } else { Vec::new() };
|
||||
}
|
||||
let mut standalone = uncorrelated.clone();
|
||||
standalone
|
||||
.attributes
|
||||
.retain(|attribute| attribute.key != "session.id");
|
||||
request.resource_logs[0].scope_logs[0].log_records = vec![valid, uncorrelated, standalone];
|
||||
request.resource_logs[0].resource = Some(opentelemetry_proto::tonic::resource::v1::Resource {
|
||||
attributes: vec![KeyValue {
|
||||
key: "lens.session.capture".into(),
|
||||
value: Some(AnyValue {
|
||||
value: Some(Value::StringValue("true".into())),
|
||||
}),
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
});
|
||||
let bytes = if json {
|
||||
serde_json::to_vec(&request).unwrap()
|
||||
} else {
|
||||
request.encode_to_vec()
|
||||
};
|
||||
let media = json.then_some("application/json");
|
||||
let limits = litellm_traces::DecodeLimits {
|
||||
attributes: 6,
|
||||
..Default::default()
|
||||
};
|
||||
let spans = litellm_traces::decode_otlp_logs_with_limits(&bytes, media, limits).unwrap();
|
||||
let replayed = litellm_traces::decode_otlp_logs_with_limits(&bytes, media, limits).unwrap();
|
||||
assert_eq!(spans.len(), 3);
|
||||
assert_eq!(
|
||||
serde_json::to_value(&spans).unwrap(),
|
||||
serde_json::to_value(&replayed).unwrap()
|
||||
);
|
||||
assert_eq!(spans[0].trace_id, spans[1].trace_id);
|
||||
assert_ne!(spans[1].trace_id, spans[2].trace_id);
|
||||
assert_ne!(spans[0].span_id, spans[1].span_id);
|
||||
if missing_trace {
|
||||
assert_ne!(spans[1].span_id, spans[2].span_id);
|
||||
assert!(!spans[1].attributes.contains_key("lens.original_trace_id"));
|
||||
}
|
||||
assert_eq!(spans[0].parent_span_id, "02".repeat(8));
|
||||
assert!(spans[1].parent_span_id.is_empty());
|
||||
assert!(spans[2].parent_span_id.is_empty());
|
||||
assert!(!spans[0].attributes.contains_key("lens.capture.warning"));
|
||||
assert!(spans[1].attributes["lens.capture.warning"].contains("unconfirmed"));
|
||||
assert!(spans[2].attributes["lens.capture.warning"].contains("unconfirmed"));
|
||||
assert_eq!(spans[0].normalized.output, spans[1].normalized.output);
|
||||
assert_eq!(spans[0].normalized.output, spans[2].normalized.output);
|
||||
assert_eq!(spans[1].normalized.observation_type, ObservationType::Chain);
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::session(true)]
|
||||
#[case::standalone(false)]
|
||||
fn absent_native_log_context_has_encoding_independent_identity(#[case] session: bool) {
|
||||
use prost::Message;
|
||||
let mut request = log_request("repl_main_thread");
|
||||
let record = &mut request.resource_logs[0].scope_logs[0].log_records[0];
|
||||
record.trace_id.clear();
|
||||
record.span_id.clear();
|
||||
if !session {
|
||||
record.attributes.retain(|entry| entry.key != "session.id");
|
||||
}
|
||||
let omitted = litellm_traces::decode_otlp_logs(
|
||||
&serde_json::to_vec(&request).unwrap(),
|
||||
Some("application/json"),
|
||||
)
|
||||
.unwrap();
|
||||
let record = &mut request.resource_logs[0].scope_logs[0].log_records[0];
|
||||
record.trace_id = vec![0; 16];
|
||||
record.span_id = vec![0; 8];
|
||||
let zeroed = litellm_traces::decode_otlp_logs(&request.encode_to_vec(), None).unwrap();
|
||||
assert_eq!(omitted[0].trace_id, zeroed[0].trace_id);
|
||||
assert_eq!(omitted[0].span_id, zeroed[0].span_id);
|
||||
assert_eq!(omitted[0].normalized.output, zeroed[0].normalized.output);
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::nodes(litellm_traces::DecodeLimits { nodes: 4, ..Default::default() })]
|
||||
#[case::depth(litellm_traces::DecodeLimits { depth: 2, ..Default::default() })]
|
||||
|
|
@ -1778,6 +1878,18 @@ fn interactive_claude_exports_join_replies_with_native_child_execution_context()
|
|||
#[rstest]
|
||||
#[case::tool_result("tool_result", "", false)]
|
||||
#[case::complete_body("api_request_body", r#"{"messages":[{"role":"user","content":[{"type":"tool_result","tool_use_id":"call-1","is_error":true,"content":[{"type":"text","text":"exit 3 output"},{"type":"image","source":{"data":"PRIVATE_IMAGE"}}]}]}],"system":"PRIVATE_SYSTEM"}"#, false)]
|
||||
#[case::headless_body("api_request_body", r#"{"messages":[{"role":"user","content":[{"type":"tool_result","tool_use_id":"call-1","is_error":true,"content":"exit 3 output"}]},{"role":"system","content":"PRIVATE_SYSTEM"}]}"#, false)]
|
||||
#[case::missing_messages("api_request_body", r#"{}"#, true)]
|
||||
#[case::wrong_content(
|
||||
"api_request_body",
|
||||
r#"{"messages":[{"role":"user","content":{}}]}"#,
|
||||
true
|
||||
)]
|
||||
#[case::unexpected_last_message(
|
||||
"api_request_body",
|
||||
r#"{"messages":[{"role":"assistant","content":"unexpected"}]}"#,
|
||||
true
|
||||
)]
|
||||
#[case::truncated_body("api_request_body", "{truncated", true)]
|
||||
fn native_tool_logs_supply_arguments_and_results_without_fake_calls(
|
||||
#[case] event: &str,
|
||||
|
|
@ -1849,3 +1961,39 @@ fn interactive_claude_body_export_retains_failed_command_stdout() {
|
|||
.unwrap();
|
||||
assert_eq!(failed["content"], "Exit code 3\nRAW-EXPECTED");
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::new_prompt(serde_json::json!({"role":"user", "content":"Continue"}))]
|
||||
#[case::new_blocks(serde_json::json!({"role":"user", "content":[{"type":"text", "text":"Continue"}]}))]
|
||||
fn native_body_exports_do_not_replay_old_tool_results(#[case] final_message: serde_json::Value) {
|
||||
use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value::Value};
|
||||
let mut request = log_request("repl_main_thread");
|
||||
let body = serde_json::json!({"messages": [
|
||||
{"role": "user", "content": [{"type": "tool_result", "tool_use_id": "old-call", "content": "OLD_RESULT"}]},
|
||||
{"role": "assistant", "content": "Done"}, final_message,
|
||||
{"role": "system", "content": "PRIVATE_SYSTEM"}
|
||||
]}).to_string();
|
||||
request.resource_logs[0].scope_logs[0].log_records[0].attributes = [
|
||||
("event.name", "api_request_body"),
|
||||
("query_source", "repl_main_thread"),
|
||||
("body", body.as_str()),
|
||||
]
|
||||
.into_iter()
|
||||
.map(|(key, value)| KeyValue {
|
||||
key: key.into(),
|
||||
value: Some(AnyValue {
|
||||
value: Some(Value::StringValue(value.into())),
|
||||
}),
|
||||
..Default::default()
|
||||
})
|
||||
.collect();
|
||||
let spans = litellm_traces::decode_otlp_logs(
|
||||
&serde_json::to_vec(&request).unwrap(),
|
||||
Some("application/json"),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
serde_json::from_str::<serde_json::Value>(&spans[0].normalized.output).unwrap(),
|
||||
serde_json::json!({"tool_results":[]})
|
||||
);
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue