diff --git a/litellm-rust/crates/traces/src/normalize/format/claude_code.rs b/litellm-rust/crates/traces/src/normalize/format/claude_code.rs index d24b9e6d8d2..37fb2d01330 100644 --- a/litellm-rust/crates/traces/src/normalize/format/claude_code.rs +++ b/litellm-rust/crates/traces/src/normalize/format/claude_code.rs @@ -159,10 +159,23 @@ fn exported_tool_results(attributes: &BTreeMap) -> String { let Ok(body) = serde_json::from_str::(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 = 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 = content.as_array() .into_iter() .flatten() .filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_result")) diff --git a/litellm-rust/crates/traces/src/otlp/logs.rs b/litellm-rust/crates/traces/src/otlp/logs.rs index 061ec3f6662..089c5edb798 100644 --- a/litellm-rust/crates/traces/src/otlp/logs.rs +++ b/litellm-rust/crates/traces/src/otlp/logs.rs @@ -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, 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::())?; + 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) } diff --git a/litellm-rust/crates/traces/src/otlp/span.rs b/litellm-rust/crates/traces/src/otlp/span.rs index 58ccb9ea3b4..7e01191cbdc 100644 --- a/litellm-rust/crates/traces/src/otlp/span.rs +++ b/litellm-rust/crates/traces/src/otlp/span.rs @@ -20,12 +20,19 @@ pub(super) fn flatten( request: ExportTraceServiceRequest, limits: DecodeLimits, ) -> Result, Error> { + flatten_with_budget(request, limits).map(|(spans, _)| spans) +} + +pub(super) fn flatten_with_budget( + request: ExportTraceServiceRequest, + limits: DecodeLimits, +) -> Result<(Vec, 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 { + Sha256::digest(format!("litellm.claude.session.v1\0{session}"))[..16].to_vec() +} + fn decoded_span( span: Span, resource_attributes: &Shared>, @@ -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); diff --git a/litellm-rust/crates/traces/tests/otlp.rs b/litellm-rust/crates/traces/tests/otlp.rs index d4fd3a248a9..68a25fb028b 100644 --- a/litellm-rust/crates/traces/tests/otlp.rs +++ b/litellm-rust/crates/traces/tests/otlp.rs @@ -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, @@ -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::(&spans[0].normalized.output).unwrap(), + serde_json::json!({"tool_results":[]}) + ); +}