fix(lens): reconstruct native coding agent conversations (#44711)

* fix(lens): reconstruct native coding agent conversations

* fix(lens): display timestamps used for conversation ordering

* fix(lens): preserve coding trace identity and message provenance

* fix(lens): complete capture checks after trace pagination

* fix(lens): show loaded replies during trace pagination
This commit is contained in:
moe-berri 2026-10-05 17:51:05 -07:00 • committed by GitHub
parent 6e75f28289
commit 2190003bb8
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
36 changed files with 4662 additions and 157 deletions

View file

@ -4451,6 +4451,7 @@ dependencies = [
"serde",
"serde_json",
"serde_with",
"sha2 0.10.9",
"strum",
"thiserror 2.0.19",
"time",

View file

@ -186,12 +186,14 @@ impl NativeTraceStorage {
)
}
#[pyo3(signature = (payload, content_type, tenant, logs=false))]
fn ingest<'py>(
&self,
py: Python<'py>,
payload: &[u8],
content_type: Option<String>,
#[pyo3(from_py_with = litellm_host_python::from_py_argument)] tenant: Tenant,
logs: bool,
) -> PyResult<Bound<'py, PyAny>> {
let payload = payload.to_vec();
let max_value_bytes = self.config.max_attribute_value_bytes();
@ -202,7 +204,12 @@ impl NativeTraceStorage {
py,
async move {
let rows = tokio::task::spawn_blocking(move || {
litellm_traces::decode_otlp(&payload, content_type.as_deref()).map(|spans| {
let decode = if logs {
litellm_traces::decode_otlp_logs
} else {
litellm_traces::decode_otlp
};
decode(&payload, content_type.as_deref()).map(|spans| {
litellm_traces_clickhouse::span_rows(spans, &tenant, max_value_bytes)
})
})

View file

@ -14,11 +14,12 @@ macro_rules_attribute.workspace = true
schemars = { workspace = true, optional = true }
indexmap = { version = "2", features = ["serde"] }
litellm-llms-types.workspace = true
opentelemetry-proto = { workspace = true, features = ["gen-tonic-messages", "trace", "with-serde"] }
opentelemetry-proto = { workspace = true, features = ["gen-tonic-messages", "trace", "logs", "with-serde"] }
prost.workspace = true
serde = { workspace = true, features = ["rc"] }
serde_json = { workspace = true, features = ["preserve_order"] }
serde_with.workspace = true
sha2.workspace = true
strum.workspace = true
thiserror.workspace = true
time.workspace = true

View file

@ -32,7 +32,10 @@ pub use normalize::{
AgentMetadata, AgentType, CallEvidence, CallEvidenceKind, CallKey, Integration, NormalizedSpan,
ObservationType,
};
pub use otlp::{DecodeLimits, DecodedEvent, DecodedSpan, decode_otlp, decode_otlp_with_limits};
pub use otlp::{
DecodeLimits, DecodedEvent, DecodedSpan, decode_otlp, decode_otlp_logs,
decode_otlp_logs_with_limits, decode_otlp_with_limits,
};
pub use query::ReadQuery;
pub use query_access::QueryScope;
pub use resolve::{SpendLookup, iso_time, listed_summary, resolve_trace};

View file

@ -6,8 +6,8 @@ use super::{Extraction, Format, SpanFacts};
use crate::{
Error,
normalize::{
CLAUDE_CODE_AGENT, CLAUDE_CODE_SCOPE, CallEvidence, CallKey, ObservationType, RoleEvidence,
SpanContext, attr, present, tokens,
CLAUDE_CODE_AGENT, CLAUDE_CODE_EVENTS_SCOPE, CLAUDE_CODE_SCOPE, CallEvidence, CallKey,
ObservationType, RoleEvidence, SpanContext, attr, present, tokens,
},
otlp::DecodedEvent,
};
@ -16,6 +16,10 @@ use crate::{
pub(crate) struct ClaudeCode;
enum SpanType {
AssistantResponse,
ToolResult,
ApiRequestBody,
Compaction,
Interaction,
LlmRequest,
Tool,
@ -30,6 +34,10 @@ fn span_type(name: &str, attributes: &BTreeMap<String, String>) -> SpanType {
kind
};
match kind {
"assistant_response" => SpanType::AssistantResponse,
"tool_result" => SpanType::ToolResult,
"api_request_body" => SpanType::ApiRequestBody,
"compaction" => SpanType::Compaction,
"interaction" => SpanType::Interaction,
"llm_request" => SpanType::LlmRequest,
"tool" => SpanType::Tool,
@ -147,6 +155,30 @@ fn llm_output(attributes: &BTreeMap<String, String>) -> String {
}
}
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))
.into_iter()
.flatten()
.filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_result"))
.map(|block| {
let content = match block.get("content") {
Some(Value::String(text)) => text.clone(),
Some(Value::Array(blocks)) => blocks.iter().map(|block| {
block.get("text").and_then(Value::as_str).unwrap_or("[Non-text tool output omitted by Claude export]")
}).collect::<Vec<_>>().join("\n"),
_ => String::new(),
};
json!({"id": block.get("tool_use_id"), "content": content, "is_error": block.get("is_error").and_then(Value::as_bool).unwrap_or(false)})
}).collect();
json!({"tool_results": results}).to_string()
}
fn input_tokens(attributes: &BTreeMap<String, String>) -> Result<u32, Error> {
["input_tokens", "cache_read_tokens", "cache_creation_tokens"]
.into_iter()
@ -159,7 +191,7 @@ fn input_tokens(attributes: &BTreeMap<String, String>) -> Result<u32, Error> {
impl Format for ClaudeCode {
fn matches(&self, context: &SpanContext<'_>) -> bool {
context.scope == CLAUDE_CODE_SCOPE
matches!(context.scope, CLAUDE_CODE_SCOPE | CLAUDE_CODE_EVENTS_SCOPE)
}
fn extract(&self, context: &SpanContext<'_>) -> Result<Extraction, Error> {
@ -172,6 +204,48 @@ impl Format for ClaudeCode {
..SpanFacts::default()
};
let (facts, consumed): (SpanFacts, Vec<&'static str>) = match kind {
SpanType::AssistantResponse => (
SpanFacts {
role: Some(RoleEvidence::Declared(ObservationType::Chain)),
agent_name: Some(subagent(attributes).unwrap_or(CLAUDE_CODE_AGENT).to_owned()),
model: present(attributes, &["model"]),
output: json!({"role": "assistant", "content": attr(attributes, "response")})
.to_string(),
..base
},
vec!["response"],
),
SpanType::ToolResult => (
SpanFacts {
input: tool_input(attributes),
tool_call_id: present(attributes, &["tool_use_id"]),
..base
},
if tool_arguments(attributes).is_some() {
vec!["tool_input"]
} else {
Vec::new()
},
),
SpanType::Compaction => (
SpanFacts {
role: Some(RoleEvidence::Declared(ObservationType::Chain)),
output: json!({"role": "system", "content": if attr(attributes, "success") == "true" {
"Context compacted"
} else {
"Context compaction failed"
}}).to_string(),
..base
},
Vec::new(),
),
SpanType::ApiRequestBody => (
SpanFacts {
output: exported_tool_results(attributes),
..base
},
vec!["body"],
),
SpanType::Interaction => (
SpanFacts {
role: Some(RoleEvidence::Declared(ObservationType::Agent)),
@ -268,6 +342,33 @@ mod tests {
.collect()
}
#[rstest]
fn notification_prompts_keep_user_provenance_and_compaction_is_system() {
let prompt_text =
"<task-notification><summary>Agent Reader completed</summary></task-notification>";
let notification = normalize(
"claude_code.interaction",
&attributes(&[("user_prompt", prompt_text)]),
&[],
)
.unwrap();
let prompt: Value = serde_json::from_str(&notification.input).unwrap();
assert_eq!(
prompt[0],
serde_json::json!({"role":"user","content":prompt_text})
);
let compaction = normalize(
"claude_code.compaction",
&attributes(&[("success", "true")]),
&[],
)
.unwrap();
assert_eq!(
serde_json::from_str::<Value>(&compaction.output).unwrap(),
serde_json::json!({"role":"system","content":"Context compacted"})
);
}
#[rstest]
fn tool_without_detailed_input_lists_known_arguments() {
let span = normalize(

View file

@ -1,7 +1,7 @@
use super::{
Integration, ObservationType, RoleEvidence, Rule, SpanContext, SpanFacts, attr, present,
};
use crate::normalize::{CLAUDE_CODE_AGENT, CLAUDE_CODE_SCOPE};
use crate::normalize::{CLAUDE_CODE_AGENT, CLAUDE_CODE_EVENTS_SCOPE, CLAUDE_CODE_SCOPE};
use std::collections::BTreeMap;
pub(super) const SCOPE: &str = CLAUDE_CODE_SCOPE;
@ -32,7 +32,7 @@ pub(super) struct ClaudeCode;
impl Rule for ClaudeCode {
fn matches(&self, context: &SpanContext<'_>) -> bool {
context.scope == SCOPE
matches!(context.scope, SCOPE | CLAUDE_CODE_EVENTS_SCOPE)
}
fn integration(&self, context: &SpanContext<'_>) -> Option<Integration> {
Some(framework(context.attributes))

View file

@ -18,6 +18,12 @@ mod messages;
mod metadata;
pub(crate) const CLAUDE_CODE_SCOPE: &str = "com.anthropic.claude_code.tracing";
pub(crate) const CLAUDE_CODE_EVENTS_SCOPE: &str = "com.anthropic.claude_code.events";
pub(crate) fn visible_claude_response(event: &str, source: &str) -> bool {
event == "assistant_response"
&& (matches!(source, "repl_main_thread" | "sdk" | "sdk_main_thread")
|| source.starts_with("agent:"))
}
pub(crate) const CLAUDE_CODE_AGENT: &str = "claude-code";
use instrumentation::Instrumentation;
pub(crate) use messages::{HIDDEN_BLOCK_TYPES, MessagePayload, encode};

View file

@ -160,6 +160,10 @@ impl<'de> Visitor<'de> for JsonBudget<'_> {
#[derive(Clone, Copy)]
enum MessageKind {
Export,
ExportLogs,
ResourceLogs,
ScopeLogs,
LogRecord,
ResourceSpans,
Resource,
ScopeSpans,
@ -177,6 +181,13 @@ enum MessageKind {
impl MessageKind {
fn child(self, tag: u32) -> Option<Self> {
match (self, tag) {
(Self::ExportLogs, 1) => Some(Self::ResourceLogs),
(Self::ResourceLogs, 1) => Some(Self::Resource),
(Self::ResourceLogs, 2) => Some(Self::ScopeLogs),
(Self::ScopeLogs, 1) => Some(Self::Scope),
(Self::ScopeLogs, 2) => Some(Self::LogRecord),
(Self::LogRecord, 5) => Some(Self::AnyValue),
(Self::LogRecord, 6) => Some(Self::KeyValue),
(Self::Export, 1) => Some(Self::ResourceSpans),
(Self::ResourceSpans, 1) => Some(Self::Resource),
(Self::ResourceSpans, 2) => Some(Self::ScopeSpans),
@ -203,6 +214,10 @@ pub(super) fn protobuf_preflight(payload: &[u8], limits: &DecodeLimits) -> Resul
scan_message(payload, MessageKind::Export, 0, &mut 0, limits)
}
pub(super) fn protobuf_logs_preflight(payload: &[u8], limits: &DecodeLimits) -> Result<(), Error> {
scan_message(payload, MessageKind::ExportLogs, 0, &mut 0, limits)
}
fn scan_message(
mut payload: &[u8],
kind: MessageKind,

View file

@ -0,0 +1,134 @@
use opentelemetry_proto::tonic::{
collector::{logs::v1::ExportLogsServiceRequest, trace::v1::ExportTraceServiceRequest},
common::v1::{KeyValue, any_value::Value},
logs::v1::LogRecord,
trace::v1::{ResourceSpans, ScopeSpans, Span, Status},
};
use sha2::{Digest, Sha256};
use super::{DecodeLimits, DecodedSpan, span};
use crate::{
Error,
normalize::{CLAUDE_CODE_EVENTS_SCOPE, visible_claude_response},
};
fn value<'a>(attributes: &'a [KeyValue], key: &str) -> Option<&'a Value> {
attributes
.iter()
.rev()
.find(|entry| entry.key == key)
.and_then(|entry| entry.value.as_ref())
.and_then(|value| value.value.as_ref())
}
fn text<'a>(attributes: &'a [KeyValue], key: &str) -> &'a str {
match value(attributes, key) {
Some(Value::StringValue(text)) => text,
_ => "",
}
}
fn message(record: LogRecord) -> Span {
let timestamp = if record.time_unix_nano == 0 {
record.observed_time_unix_nano
} else {
record.time_unix_nano
};
let mut hash = Sha256::new();
hash.update(b"litellm.claude.message.v1\0");
hash.update(&record.trace_id);
hash.update(&record.span_id);
hash.update(text(&record.attributes, "event.name"));
let uuid = text(&record.attributes, "message.uuid");
if uuid.is_empty() {
hash.update(timestamp.to_be_bytes());
match value(&record.attributes, "event.sequence") {
Some(Value::IntValue(sequence)) => hash.update(sequence.to_string()),
_ => hash.update(text(&record.attributes, "event.sequence")),
}
hash.update(text(&record.attributes, "response"));
} else {
hash.update(uuid);
}
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()
}
}
pub(super) fn flatten(
request: ExportLogsServiceRequest,
limits: DecodeLimits,
) -> Result<Vec<DecodedSpan>, Error> {
let mut count = 0usize;
let mut resources = Vec::new();
for resource in request.resource_logs {
let mut scopes = Vec::new();
for scope in resource.scope_logs {
count = count
.checked_add(scope.log_records.len())
.ok_or(Error::TooLarge)?;
if count > limits.spans
|| scope
.log_records
.iter()
.any(|record| record.attributes.len() > limits.attributes)
{
return Err(Error::TooLarge);
}
let supported = scope
.scope
.as_ref()
.is_some_and(|scope| scope.name == CLAUDE_CODE_EVENTS_SCOPE);
let spans = scope
.log_records
.into_iter()
.filter(|record| {
let event = text(&record.attributes, "event.name");
supported
&& (matches!(event, "tool_result" | "compaction")
|| (matches!(event, "assistant_response" | "api_request_body")
&& visible_claude_response(
"assistant_response",
text(&record.attributes, "query_source"),
)))
})
.map(message)
.collect();
scopes.push(ScopeSpans {
scope: scope.scope,
spans,
schema_url: scope.schema_url,
});
}
resources.push(ResourceSpans {
resource: resource.resource,
scope_spans: scopes,
schema_url: resource.schema_url,
});
}
span::flatten(
ExportTraceServiceRequest {
resource_spans: resources,
},
limits,
)
}

View file

@ -1,5 +1,6 @@
mod attributes;
mod limits;
mod logs;
mod span;
mod wire;
@ -49,3 +50,18 @@ pub fn decode_otlp_with_limits(
let request = wire::decode(body, content_type, &limits)?;
span::flatten(request, limits)
}
pub fn decode_otlp_logs(
body: &[u8],
content_type: Option<&str>,
) -> Result<Vec<DecodedSpan>, Error> {
decode_otlp_logs_with_limits(body, content_type, DecodeLimits::from_env()?)
}
pub fn decode_otlp_logs_with_limits(
body: &[u8],
content_type: Option<&str>,
limits: DecodeLimits,
) -> Result<Vec<DecodedSpan>, Error> {
logs::flatten(wire::decode_logs(body, content_type, &limits)?, limits)
}

View file

@ -1,3 +1,4 @@
use sha2::{Digest, Sha256};
use std::collections::BTreeMap;
use opentelemetry_proto::tonic::{
@ -12,7 +13,7 @@ use super::{
};
use crate::{
Error, Shared,
normalize::{SpanContext, normalize},
normalize::{CLAUDE_CODE_EVENTS_SCOPE, CLAUDE_CODE_SCOPE, SpanContext, normalize},
};
pub(super) fn flatten(
@ -132,7 +133,28 @@ fn decoded_span(
) -> Result<DecodedSpan, Error> {
let status = span.status.unwrap_or_default();
let parent_span_id = hex_bytes(&span.parent_span_id);
let span_attributes = attributes(span.attributes, budget)?;
let mut span_attributes = attributes(span.attributes, budget)?;
let original_trace_id = hex_bytes(&span.trace_id);
let trace_id = if matches!(
scope_name.as_str(),
CLAUDE_CODE_SCOPE | CLAUDE_CODE_EVENTS_SCOPE
) && resource_attributes
.get("lens.session.capture")
.is_some_and(|value| value == "true")
&& let Some(session) = span_attributes
.get("session.id")
.filter(|value| !value.is_empty())
{
let trace_id =
hex_bytes(&Sha256::digest(format!("litellm.claude.session.v1\0{session}"))[..16]);
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);
span_attributes.insert("gen_ai.agent.id".to_owned(), actor);
trace_id
} else {
original_trace_id
};
let events = span
.events
.into_iter()
@ -183,7 +205,7 @@ fn decoded_span(
+ normalization.display_name.as_ref().map_or(0, String::len),
)?;
Ok(DecodedSpan {
trace_id: hex_bytes(&span.trace_id),
trace_id,
span_id: hex_bytes(&span.span_id),
parent_span_id,
trace_state: span.trace_state,

View file

@ -1,7 +1,8 @@
use opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest;
use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
use prost::Message;
use super::limits::{DecodeLimits, json_preflight, protobuf_preflight};
use super::limits::{DecodeLimits, json_preflight, protobuf_logs_preflight, protobuf_preflight};
use crate::Error;
#[derive(strum::EnumString)]
@ -21,6 +22,23 @@ pub(super) fn decode(
content_type: Option<&str>,
limits: &DecodeLimits,
) -> Result<ExportTraceServiceRequest, Error> {
decode_request(body, content_type, limits, protobuf_preflight)
}
pub(super) fn decode_logs(
body: &[u8],
content_type: Option<&str>,
limits: &DecodeLimits,
) -> Result<ExportLogsServiceRequest, Error> {
decode_request(body, content_type, limits, protobuf_logs_preflight)
}
fn decode_request<T: Message + Default + serde::de::DeserializeOwned>(
body: &[u8],
content_type: Option<&str>,
limits: &DecodeLimits,
preflight: fn(&[u8], &DecodeLimits) -> Result<(), Error>,
) -> Result<T, Error> {
let media_type = content_type
.unwrap_or("application/x-protobuf")
.split(';')
@ -36,8 +54,8 @@ pub(super) fn decode(
serde_json::from_slice(body).map_err(|_| Error::InvalidPayload)?
}
OtlpMediaType::Protobuf => {
protobuf_preflight(body, limits)?;
ExportTraceServiceRequest::decode(body).map_err(|_| Error::InvalidPayload)?
preflight(body, limits)?;
T::decode(body).map_err(|_| Error::InvalidPayload)?
}
};
Ok(request)

View file

@ -25,6 +25,7 @@ pub(super) struct Resolution<'a> {
ownership: Ownership<'a>,
spend: &'a [SpendRow],
types: HashMap<&'a str, ObservationType>,
tool_failures: HashMap<&'a str, &'a TraceSpansRow>,
pub(super) model_calls: Vec<usize>,
}
@ -53,6 +54,16 @@ impl<'a> Resolution<'a> {
graph,
spend,
types,
tool_failures: rows
.iter()
.filter(|row| {
row.framework == "claude-code"
&& row.name == "claude_code.tool_result"
&& !row.tool_call_id.is_empty()
&& row.status == crate::SpanStatus::Error
})
.map(|row| (row.tool_call_id.as_str(), row))
.collect(),
model_calls,
}
}
@ -61,6 +72,28 @@ impl<'a> Resolution<'a> {
&self.graph.rows[index]
}
pub(super) fn status_source(&self, index: usize) -> &'a TraceSpansRow {
let row = self.row(index);
if row.framework != "claude-code"
|| row.kind != ObservationType::Tool
|| row.status == crate::SpanStatus::Error
{
return row;
}
self.graph
.children(index)
.into_iter()
.map(|child| self.row(child))
.find(|child| {
child.name == "claude_code.tool.execution"
&& !row.tool_call_id.is_empty()
&& child.tool_call_id == row.tool_call_id
&& child.status == crate::SpanStatus::Error
})
.or_else(|| self.tool_failures.get(row.tool_call_id.as_str()).copied())
.unwrap_or(row)
}
pub(super) fn kind(&self, index: usize) -> ObservationType {
self.types[self.graph.id(index)]
}

View file

@ -22,6 +22,7 @@ fn optional(value: &str) -> Option<String> {
fn span(resolution: &Resolution<'_>, index: usize, trace_start_ns: i64) -> Span {
let row = resolution.row(index);
let status = resolution.status_source(index);
let requests = resolution.requests(index).complete_requests();
Span {
span_id: row.span_id.clone(),
@ -33,9 +34,9 @@ fn span(resolution: &Resolution<'_>, index: usize, trace_start_ns: i64) -> Span
start_offset_ms: (i128::from(row.start_ns) - i128::from(trace_start_ns)) as f64
/ NANOS_PER_MS,
duration_ms: row.duration_ns as f64 / NANOS_PER_MS,
status: row.status,
error: optional(&row.status_message),
error_truncated: row.error_truncated,
status: status.status,
error: optional(&status.status_message),
error_truncated: status.error_truncated,
input_preview: row.input_preview.clone(),
model: optional(&row.model),
input_tokens: row.input_tokens,
@ -192,10 +193,18 @@ pub fn resolve_trace(
agent_invocations: agents.iter().map(|agent| agent.invocations).sum(),
llm_calls: calls.len() as u64,
tool_calls: resolution.unique_tools().len() as u64,
error_count: spans
error_count: rows
.iter()
.filter(|span| span.status == SpanStatus::Error)
.count() as u64,
.map(|span| {
if span.framework == "claude-code" && !span.tool_call_id.is_empty() {
("claude-tool", span.tool_call_id.as_str())
} else {
("span", span.span_id.as_str())
}
})
.collect::<BTreeSet<_>>()
.len() as u64,
input_tokens: counted.iter().map(|row| u64::from(row.input_tokens)).sum(),
output_tokens: counted.iter().map(|row| u64::from(row.output_tokens)).sum(),
models: sorted_unique(calls.iter().map(|call| rows[*call].model.as_str())),

View file

@ -0,0 +1,748 @@
{
"resourceLogs": [
{
"resource": {
"attributes": [
{
"key": "service.name",
"value": {
"stringValue": "claude-code"
}
},
{
"key": "lens.session.capture",
"value": {
"stringValue": "true"
}
}
]
},
"scopeLogs": [
{
"scope": {
"name": "com.anthropic.claude_code.events",
"version": "2.1.289"
},
"logRecords": [
{
"timeUnixNano": "1791241295512000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:01:35.512Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "19"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "73f6dfd9-e1ea-4950-8027-52e6c124e6bc"
}
},
{
"key": "response_length",
"value": {
"intValue": "18"
}
},
{
"key": "response",
"value": {
"stringValue": "MINIMAL-COMMENTARY"
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "51da354c-8278-4ead-8f90-e47fd72bf9c3"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "repl_main_thread"
}
}
],
"flags": 1,
"traceId": "ced683bb1985e31a6a67d7180245ad0a",
"spanId": "7602ccb75a2365b7",
"observedTimeUnixNano": "1791241295512000000"
},
{
"timeUnixNano": "1791241295581000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:01:35.581Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "22"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "73f6dfd9-e1ea-4950-8027-52e6c124e6bc"
}
},
{
"key": "response_length",
"value": {
"intValue": "47"
}
},
{
"key": "response",
"value": {
"stringValue": "{\"title\":\"MINIMAL README and exit-3 Bash test\"}"
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "52462a69-0006-4520-a142-e2fcd53593e5"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "generate_session_title"
}
}
],
"flags": 1,
"traceId": "ced683bb1985e31a6a67d7180245ad0a",
"spanId": "7602ccb75a2365b7",
"observedTimeUnixNano": "1791241295581000000"
}
]
}
]
},
{
"resource": {
"attributes": [
{
"key": "service.name",
"value": {
"stringValue": "claude-code"
}
},
{
"key": "lens.session.capture",
"value": {
"stringValue": "true"
}
}
]
},
"scopeLogs": [
{
"scope": {
"name": "com.anthropic.claude_code.events",
"version": "2.1.289"
},
"logRecords": [
{
"timeUnixNano": "1791241300179000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:01:40.179Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "27"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "73f6dfd9-e1ea-4950-8027-52e6c124e6bc"
}
},
{
"key": "response_length",
"value": {
"intValue": "194"
}
},
{
"key": "response",
"value": {
"stringValue": "README.txt says it's a synthetic test fixture with the marker `LENS-REPLAY-ALPHA`. The command printed MINIMAL-EXPECTED and exited with code 3, as you asked, so I didn't retry it.\n\nMINIMAL-FINAL"
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "a01fab5f-ab6a-4633-9863-d7a8ea56b53d"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "repl_main_thread"
}
}
],
"flags": 1,
"traceId": "ced683bb1985e31a6a67d7180245ad0a",
"spanId": "7602ccb75a2365b7",
"observedTimeUnixNano": "1791241300179000000"
}
]
}
]
},
{
"resource": {
"attributes": [
{
"key": "service.name",
"value": {
"stringValue": "claude-code"
}
},
{
"key": "lens.session.capture",
"value": {
"stringValue": "true"
}
}
]
},
"scopeLogs": [
{
"scope": {
"name": "com.anthropic.claude_code.events",
"version": "2.1.289"
},
"logRecords": [
{
"timeUnixNano": "1791241516987000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:05:16.987Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "39"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "a700c173-6689-4a7b-8fd7-28a4640dc577"
}
},
{
"key": "response_length",
"value": {
"intValue": "218"
}
},
{
"key": "response",
"value": {
"stringValue": "I started the Reader and Checker agents in parallel. Both ended up running in the background, not just one as you asked. I'll wait for both to finish before reporting their results and closing with NATIVE-AGENTS-FINAL."
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "646639c5-36f9-4bc7-9117-500658fa4af0"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "repl_main_thread"
}
}
],
"flags": 1,
"traceId": "dc9069174bf4ed50dbcd0088f5544c4b",
"spanId": "f6836a39a95b89d0",
"observedTimeUnixNano": "1791241516987000000"
}
]
}
]
},
{
"resource": {
"attributes": [
{
"key": "service.name",
"value": {
"stringValue": "claude-code"
}
},
{
"key": "lens.session.capture",
"value": {
"stringValue": "true"
}
}
]
},
"scopeLogs": [
{
"scope": {
"name": "com.anthropic.claude_code.events",
"version": "2.1.289"
},
"logRecords": [
{
"timeUnixNano": "1791241518190000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:05:18.190Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "44"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "a700c173-6689-4a7b-8fd7-28a4640dc577"
}
},
{
"key": "response_length",
"value": {
"intValue": "31"
}
},
{
"key": "response",
"value": {
"stringValue": "NATIVE-READER LENS-REPLAY-ALPHA"
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "0f7da27c-5d13-48f4-ab2a-c94af16cc877"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "agent:builtin:general-purpose"
}
}
],
"flags": 1,
"traceId": "dc9069174bf4ed50dbcd0088f5544c4b",
"spanId": "1016ee3f27ec48a1",
"observedTimeUnixNano": "1791241518190000000"
},
{
"timeUnixNano": "1791241518342000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:05:18.342Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "48"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "06a1d3a0-95a5-405b-997b-f36539ad420a"
}
},
{
"key": "response_length",
"value": {
"intValue": "17"
}
},
{
"key": "response",
"value": {
"stringValue": "NATIVE-CHECKER 15"
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "e2b4a8ec-8b78-4c10-a0be-650545d3639f"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "agent:builtin:general-purpose"
}
}
],
"flags": 1,
"traceId": "dc9069174bf4ed50dbcd0088f5544c4b",
"spanId": "069ab83e9b2f1d80",
"observedTimeUnixNano": "1791241518342000000"
}
]
}
]
},
{
"resource": {
"attributes": [
{
"key": "service.name",
"value": {
"stringValue": "claude-code"
}
},
{
"key": "lens.session.capture",
"value": {
"stringValue": "true"
}
}
]
},
"scopeLogs": [
{
"scope": {
"name": "com.anthropic.claude_code.events",
"version": "2.1.289"
},
"logRecords": [
{
"timeUnixNano": "1791241519849000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:05:19.849Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "52"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "06a1d3a0-95a5-405b-997b-f36539ad420a"
}
},
{
"key": "response_length",
"value": {
"intValue": "75"
}
},
{
"key": "response",
"value": {
"stringValue": "Reader finished: NATIVE-READER LENS-REPLAY-ALPHA. Checker is still running."
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "905350de-6161-4013-a7bc-3dd8f404f877"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "repl_main_thread"
}
}
],
"flags": 1,
"traceId": "b08f665a8a47e055a82cf882ae69d83b",
"spanId": "23211831eeb7496c",
"observedTimeUnixNano": "1791241519849000000"
}
]
}
]
},
{
"resource": {
"attributes": [
{
"key": "service.name",
"value": {
"stringValue": "claude-code"
}
},
{
"key": "lens.session.capture",
"value": {
"stringValue": "true"
}
}
]
},
"scopeLogs": [
{
"scope": {
"name": "com.anthropic.claude_code.events",
"version": "2.1.289"
},
"logRecords": [
{
"timeUnixNano": "1791241522476000000",
"body": {
"stringValue": "claude_code.assistant_response"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-native-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "assistant_response"
}
},
{
"key": "event.timestamp",
"value": {
"stringValue": "2026-10-05T23:05:22.476Z"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "56"
}
},
{
"key": "prompt.id",
"value": {
"stringValue": "412b0cac-7610-4d76-8591-37e6787c0f63"
}
},
{
"key": "response_length",
"value": {
"intValue": "249"
}
},
{
"key": "response",
"value": {
"stringValue": "Both subagents are done. They both ran in the background, not one in the foreground as you asked.\n\n- **Reader:** NATIVE-READER LENS-REPLAY-ALPHA (the marker from README.txt)\n- **Checker:** NATIVE-CHECKER 15 (from `python3`, 7+8)\n\nNATIVE-AGENTS-FINAL"
}
},
{
"key": "message.uuid",
"value": {
"stringValue": "cac3bf80-9439-4bda-8cb7-c88414a940a0"
}
},
{
"key": "model",
"value": {
"stringValue": "claude-opus-5-5"
}
},
{
"key": "query_source",
"value": {
"stringValue": "repl_main_thread"
}
}
],
"flags": 1,
"traceId": "e5ef60574ada78a199ad0c5cc09aac4f",
"spanId": "1b43d99139cabb38",
"observedTimeUnixNano": "1791241522476000000"
}
]
}
]
}
]
}

View file

@ -0,0 +1,64 @@
{
"resourceLogs": [
{
"scopeLogs": [
{
"scope": {
"name": "com.anthropic.claude_code.events",
"version": "2.1.289"
},
"logRecords": [
{
"timeUnixNano": "1791242448890000000",
"body": {
"stringValue": "claude_code.api_request_body"
},
"attributes": [
{
"key": "session.id",
"value": {
"stringValue": "session-raw-fixture"
}
},
{
"key": "event.name",
"value": {
"stringValue": "api_request_body"
}
},
{
"key": "event.sequence",
"value": {
"intValue": "30"
}
},
{
"key": "body",
"value": {
"stringValue": "{\"messages\": [{\"role\": \"user\", \"content\": [{\"tool_use_id\": \"toolu_01Jkh34bKyUY3NcG7ej8yQfw\", \"type\": \"tool_result\", \"content\": \"1\\tThis is a synthetic fixture for testing coding-session capture.\\n2\\tMarker: LENS-REPLAY-ALPHA\\n3\\tNo external services or user files should be accessed.\\n4\\t\"}, {\"type\": \"tool_result\", \"content\": \"Exit code 3\\nRAW-EXPECTED\", \"is_error\": true, \"tool_use_id\": \"toolu_01NpCr9FJfGh4PyHLkRrKhu4\", \"cache_control\": {\"type\": \"ephemeral\"}}]}]}"
}
},
{
"key": "query_source",
"value": {
"stringValue": "repl_main_thread"
}
},
{
"key": "request_body_id",
"value": {
"stringValue": "be021b83-d74f-4ab9-aab1-8ecbb564dd35"
}
}
],
"flags": 1,
"traceId": "9ac8ed7ab2370baed8b34f517e356029",
"spanId": "a2d6a9a271bd730a",
"observedTimeUnixNano": "1791242448890000000"
}
]
}
]
}
]
}

File diff suppressed because it is too large Load diff

View file

@ -72,6 +72,28 @@ fn decode(
.unwrap())
}
#[rstest]
#[case::interaction("interaction", "user_prompt", "")]
#[case::model_context("llm_request", "new_context", "[USER]\n")]
fn native_claude_prompts_preserve_notification_text_and_user_role(
span: Span,
#[case] kind: &str,
#[case] key: &str,
#[case] prefix: &str,
) {
let prompt = "<task-notification><summary>Quoted summary</summary><result>Keep this result</result></task-notification>\nExplain this example";
let payload = format!("{prefix}{prompt}");
let decoded = decode(
span,
"com.anthropic.claude_code.tracing",
&[("span.type", kind), (key, &payload)],
vec![],
)
.unwrap();
let messages: Value = serde_json::from_str(&decoded.normalized.input).unwrap();
assert_eq!(messages, json!([{"role": "user", "content": prompt}]));
}
#[rstest]
#[case::agent("agent", ObservationType::Agent)]
#[case::workflow("workflow", ObservationType::Chain)]

View file

@ -1428,3 +1428,424 @@ fn environment_decode_limits_child() {
)),
}
}
fn log_request(
source: &str,
) -> opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest {
use opentelemetry_proto::tonic::{
collector::logs::v1::ExportLogsServiceRequest,
common::v1::{AnyValue, InstrumentationScope, KeyValue, any_value::Value},
logs::v1::{LogRecord, ResourceLogs, ScopeLogs},
};
let attributes = [
("event.name", "assistant_response"),
("query_source", source),
("response", "Visible reply"),
("message.uuid", "message-one"),
("model", "test-model"),
("session.id", "session-one"),
]
.into_iter()
.map(|(key, value)| KeyValue {
key: key.to_owned(),
value: Some(AnyValue {
value: Some(Value::StringValue(value.to_owned())),
}),
..Default::default()
})
.collect();
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
scope_logs: vec![ScopeLogs {
scope: Some(InstrumentationScope {
name: "com.anthropic.claude_code.events".to_owned(),
..Default::default()
}),
log_records: vec![LogRecord {
trace_id: vec![1; 16],
span_id: vec![2; 8],
time_unix_nano: 100,
attributes,
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
#[rstest]
#[case::main("repl_main_thread", 1)]
#[case::subagent("agent:builtin:general-purpose", 1)]
#[case::title("generate_session_title", 0)]
#[case::suggestion("prompt_suggestion", 0)]
fn native_assistant_logs_preserve_visible_messages_without_counting_model_calls(
#[case] source: &str,
#[case] count: usize,
) {
use prost::Message;
let request = log_request(source);
let json = litellm_traces::decode_otlp_logs(
&serde_json::to_vec(&request).unwrap(),
Some("application/json"),
)
.unwrap();
let binary = litellm_traces::decode_otlp_logs(&request.encode_to_vec(), None).unwrap();
assert_eq!(
serde_json::to_value(&json).unwrap(),
serde_json::to_value(&binary).unwrap()
);
assert_eq!(json.len(), count);
if let Some(span) = json.first() {
assert_eq!(span.trace_id, "01".repeat(16));
assert_eq!(span.parent_span_id, "02".repeat(8));
assert_ne!(span.span_id, span.parent_span_id);
assert_eq!(span.normalized.observation_type, ObservationType::Chain);
assert_eq!(span.normalized.framework, Some(Integration::ClaudeCode));
assert_eq!(span.normalized.model.as_deref(), Some("test-model"));
assert_eq!(span.normalized.output_tokens, 0);
assert_eq!(span.normalized.input_tokens, 0);
assert_eq!(
serde_json::from_str::<serde_json::Value>(&span.normalized.output).unwrap()["content"],
"Visible reply"
);
}
}
#[rstest]
#[case::json(true)]
#[case::protobuf(false)]
fn simultaneous_native_tool_logs_keep_distinct_sequence_ids(#[case] json: bool) {
use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value::Value};
use prost::Message;
let mut request = log_request("repl_main_thread");
let template = request.resource_logs[0].scope_logs[0].log_records[0].clone();
request.resource_logs[0].scope_logs[0].log_records = [1, 2]
.into_iter()
.map(|sequence| {
let mut record = template.clone();
record.attributes = [
("event.name", Value::StringValue("tool_result".into())),
("event.sequence", Value::IntValue(sequence)),
(
"tool_use_id",
Value::StringValue(format!("call-{sequence}")),
),
]
.into_iter()
.map(|(key, value)| KeyValue {
key: key.into(),
value: Some(AnyValue { value: Some(value) }),
..Default::default()
})
.collect();
record
})
.collect();
let bytes = if json {
serde_json::to_vec(&request).unwrap()
} else {
request.encode_to_vec()
};
let content_type = json.then_some("application/json");
let spans = litellm_traces::decode_otlp_logs(&bytes, content_type).unwrap();
assert_eq!(spans.len(), 2);
assert_ne!(spans[0].span_id, spans[1].span_id);
let replayed = litellm_traces::decode_otlp_logs(&bytes, content_type).unwrap();
assert_eq!(spans[0].span_id, replayed[0].span_id);
assert_eq!(spans[1].span_id, replayed[1].span_id);
}
#[rstest]
#[case::boolean_failure(false, true)]
#[case::boolean_success(true, true)]
#[case::string_failure(false, false)]
#[case::string_success(true, false)]
fn native_tool_log_status_accepts_boolean_and_string_values(
#[case] success: bool,
#[case] typed: bool,
) {
use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value::Value};
use prost::Message;
let mut request = log_request("repl_main_thread");
request.resource_logs[0].scope_logs[0].log_records[0].attributes = [
("event.name", Value::StringValue("tool_result".into())),
("error", Value::StringValue("Command failed".into())),
(
"success",
if typed {
Value::BoolValue(success)
} else {
Value::StringValue(success.to_string())
},
),
]
.into_iter()
.map(|(key, value)| KeyValue {
key: key.into(),
value: Some(AnyValue { value: Some(value) }),
..Default::default()
})
.collect();
let binary = litellm_traces::decode_otlp_logs(&request.encode_to_vec(), None).unwrap();
let json = litellm_traces::decode_otlp_logs(
&serde_json::to_vec(&request).unwrap(),
Some("application/json"),
)
.unwrap();
assert_eq!(binary[0].status_code == "STATUS_CODE_ERROR", !success);
assert_eq!(json[0].status_code, binary[0].status_code);
if !success {
assert_eq!(binary[0].status_message, "Command failed");
}
}
#[rstest]
fn session_capture_joins_native_logs_and_traces_across_turns_without_changing_span_parents() {
use opentelemetry_proto::tonic::{
common::v1::{AnyValue, KeyValue, any_value::Value},
resource::v1::Resource,
};
use prost::Message;
let mut logs = log_request("repl_main_thread");
let resource = Resource {
attributes: [
("lens.session.capture", "true"),
("gen_ai.agent.name", "custom-claude"),
]
.into_iter()
.map(|(key, value)| KeyValue {
key: key.to_owned(),
value: Some(AnyValue {
value: Some(Value::StringValue(value.to_owned())),
}),
..Default::default()
})
.collect(),
..Default::default()
};
logs.resource_logs[0].resource = Some(resource.clone());
let mut request = request_with(Span {
trace_id: vec![3; 16],
span_id: vec![4; 8],
name: "claude_code.interaction".to_owned(),
attributes: logs.resource_logs[0].scope_logs[0].log_records[0]
.attributes
.iter()
.filter(|attr| attr.key == "session.id")
.cloned()
.collect(),
start_time_unix_nano: 100,
end_time_unix_nano: 200,
..Default::default()
});
request.resource_spans[0].resource = Some(resource);
request.resource_spans[0].scope_spans[0].scope = Some(
opentelemetry_proto::tonic::common::v1::InstrumentationScope {
name: "com.anthropic.claude_code.tracing".to_owned(),
..Default::default()
},
);
let first = litellm_traces::decode_otlp_logs(&logs.encode_to_vec(), None).unwrap();
let second = decode_otlp(&request.encode_to_vec(), None).unwrap();
assert_eq!(first[0].trace_id, second[0].trace_id);
assert_eq!(
first[0].attributes["lens.original_trace_id"],
"01".repeat(16)
);
assert_eq!(
second[0].attributes["lens.original_trace_id"],
"03".repeat(16)
);
assert_eq!(first[0].parent_span_id, "02".repeat(8));
assert_eq!(second[0].attributes["gen_ai.agent.id"], "session-one");
assert_eq!(
first[0].normalized.agent_name.as_deref(),
Some("custom-claude")
);
request.resource_spans[0].resource = None;
assert_eq!(
decode_otlp(&request.encode_to_vec(), None).unwrap()[0].trace_id,
"03".repeat(16)
);
}
#[rstest]
#[case::short_trace(vec![1;15], vec![2;8], 1)]
#[case::zero_parent(vec![1;16], vec![0;8], 1)]
#[case::timestamp(vec![1;16], vec![2;8], i64::MAX as u64 + 1)]
fn native_logs_reject_invalid_context(
#[case] trace: Vec<u8>,
#[case] parent: Vec<u8>,
#[case] time: u64,
) {
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 = trace;
record.span_id = parent;
record.time_unix_nano = time;
assert!(matches!(
litellm_traces::decode_otlp_logs(&request.encode_to_vec(), None),
Err(litellm_traces::Error::InvalidPayload)
));
}
#[rstest]
#[case::nodes(litellm_traces::DecodeLimits { nodes: 4, ..Default::default() })]
#[case::depth(litellm_traces::DecodeLimits { depth: 2, ..Default::default() })]
#[case::bytes(litellm_traces::DecodeLimits { decoded_span_bytes: 20, ..Default::default() })]
#[case::attributes(litellm_traces::DecodeLimits { attributes: 2, ..Default::default() })]
fn native_logs_enforce_budgets_for_both_encodings(#[case] limits: litellm_traces::DecodeLimits) {
use prost::Message;
let request = log_request("repl_main_thread");
assert!(matches!(
litellm_traces::decode_otlp_logs_with_limits(&request.encode_to_vec(), None, limits),
Err(litellm_traces::Error::TooLarge)
));
assert!(matches!(
litellm_traces::decode_otlp_logs_with_limits(
&serde_json::to_vec(&request).unwrap(),
Some("application/json"),
limits
),
Err(litellm_traces::Error::TooLarge)
));
}
#[rstest]
fn interactive_claude_exports_join_replies_with_native_child_execution_context() {
let traces = decode_otlp(
include_bytes!("fixtures/claude_code_native_traces.json"),
Some("application/json"),
)
.unwrap();
let logs = litellm_traces::decode_otlp_logs(
include_bytes!("fixtures/claude_code_native_logs.json"),
Some("application/json"),
)
.unwrap();
assert!(
logs.iter()
.any(|span| span.normalized.output.contains("MINIMAL-COMMENTARY"))
);
assert!(
logs.iter()
.any(|span| span.normalized.output.contains("MINIMAL-FINAL"))
);
assert!(
logs.iter()
.any(|span| span.normalized.output.contains("NATIVE-AGENTS-FINAL"))
);
assert!(logs.iter().all(|span| span.trace_id == traces[0].trace_id));
assert!(logs.iter().all(|span| {
traces
.iter()
.any(|parent| parent.span_id == span.parent_span_id)
}));
let child = logs
.iter()
.find(|span| {
span.normalized.output.contains("NATIVE-READER")
&& span
.attributes
.get("query_source")
.is_some_and(|source| source.starts_with("agent:"))
})
.unwrap();
let execution = traces
.iter()
.find(|span| span.span_id == child.parent_span_id)
.unwrap();
assert_eq!(execution.name, "claude_code.tool.execution");
assert!(
traces
.iter()
.any(|span| span.span_id == execution.parent_span_id && span.name == "Agent")
);
assert!(
logs.iter().all(|span| span
.attributes
.get("query_source")
.is_none_or(|source| !matches!(
source.as_str(),
"prompt_suggestion" | "generate_session_title"
)))
);
}
#[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::truncated_body("api_request_body", "{truncated", true)]
fn native_tool_logs_supply_arguments_and_results_without_fake_calls(
#[case] event: &str,
#[case] body: &str,
#[case] warning: bool,
) {
use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value::Value};
use prost::Message;
let mut request = log_request("repl_main_thread");
request.resource_logs[0].scope_logs[0].log_records[0].attributes = [
("event.name", event),
("query_source", "repl_main_thread"),
("body", body),
("tool_use_id", "call-1"),
("success", "false"),
("error", "exit 3"),
(
"tool_input",
r#"{"command":"exit 3","description":"Expected failure"}"#,
),
]
.into_iter()
.map(|(key, text)| KeyValue {
key: key.into(),
value: Some(AnyValue {
value: Some(Value::StringValue(text.into())),
}),
..Default::default()
})
.collect();
let spans = litellm_traces::decode_otlp_logs(&request.encode_to_vec(), None).unwrap();
let span = &spans[0];
assert_eq!(span.normalized.observation_type, ObservationType::Framework);
assert_eq!(span.normalized.input_tokens, 0);
if event == "tool_result" {
assert_eq!(span.normalized.tool_call_id.as_deref(), Some("call-1"));
assert!(span.normalized.input.contains("Expected failure"));
assert_eq!(span.status_code, "STATUS_CODE_ERROR");
} else {
let output: serde_json::Value = serde_json::from_str(&span.normalized.output).unwrap();
assert_eq!(output.get("warning").is_some(), warning);
assert!(span.consumed_attributes.contains(&"body"));
if !warning {
assert_eq!(output["tool_results"][0]["id"], "call-1");
assert!(
output["tool_results"][0]["content"]
.as_str()
.unwrap()
.contains("exit 3 output")
);
assert!(!span.normalized.output.contains("PRIVATE"));
}
}
}
#[rstest]
fn interactive_claude_body_export_retains_failed_command_stdout() {
let spans = litellm_traces::decode_otlp_logs(
include_bytes!("fixtures/claude_code_native_tool_result.json"),
Some("application/json"),
)
.unwrap();
let output: serde_json::Value = serde_json::from_str(&spans[0].normalized.output).unwrap();
let failed = output["tool_results"]
.as_array()
.unwrap()
.iter()
.find(|result| result["is_error"] == true)
.unwrap();
assert_eq!(failed["content"], "Exit code 3\nRAW-EXPECTED");
}

View file

@ -1327,3 +1327,65 @@ fn gateway_lookup_respects_legacy_fallback_and_ownership(
expected
);
}
#[rstest]
#[case::matching("call-one", "claude_code.tool.execution", SpanStatus::Error)]
#[case::other_tool("other-call", "claude_code.tool.execution", SpanStatus::Ok)]
#[case::child_agent("call-one", "child agent", SpanStatus::Ok)]
fn native_tool_status_uses_only_its_own_execution_error(
#[case] call: &str,
#[case] name: &str,
#[case] expected: SpanStatus,
) {
let tool = TraceSpansRow {
framework: "claude-code".to_owned(),
tool_call_id: "call-one".to_owned(),
..row("tool", "", "Bash", "tool", "claude-code")
};
let execution = TraceSpansRow {
status: SpanStatus::Error,
status_message: "exit 3".to_owned(),
tool_call_id: call.to_owned(),
..row("execution", "tool", name, "framework", "claude-code")
};
let trace = resolve_trace("trace", "", &[tool, execution], &[]).unwrap();
assert_eq!(trace.spans[0].status, expected);
assert_eq!(
trace.spans[0].error.as_deref(),
if expected == SpanStatus::Error {
Some("exit 3")
} else {
None
}
);
}
#[rstest]
#[case::matching("call-one", SpanStatus::Error)]
#[case::other_tool("other-call", SpanStatus::Ok)]
fn native_tool_failure_log_matches_by_call_id_without_double_counting(
#[case] call: &str,
#[case] expected: SpanStatus,
) {
let tool = TraceSpansRow {
framework: "claude-code".into(),
tool_call_id: "call-one".into(),
..row("tool", "root", "Bash", "tool", "claude-code")
};
let log = TraceSpansRow {
framework: "claude-code".into(),
status: SpanStatus::Error,
status_message: "Permission denied".into(),
tool_call_id: call.into(),
..row(
"log",
"root",
"claude_code.tool_result",
"framework",
"claude-code",
)
};
let trace = resolve_trace("trace", "", &[tool, log], &[]).unwrap();
assert_eq!(trace.spans[0].status, expected);
assert_eq!(trace.summary.error_count, 1);
}

View file

@ -545,6 +545,7 @@ class LiteLLMRoutes(enum.Enum):
"/lens/workers/register",
"/lens/workers/{worker_id}",
"/v1/traces",
"/v1/logs",
"/v1/traces/query",
"/v1/traces/query/help",
"/v1/traces/{trace_id}",

View file

@ -183,7 +183,7 @@ def _mark_body_received(byte_count: int | None) -> None:
def is_otlp_trace_request(request: Request) -> bool:
return request.method == "POST" and get_route_path(request.scope) == "/v1/traces"
return request.method == "POST" and get_route_path(request.scope) in {"/v1/traces", "/v1/logs"}
async def _read_request_body(request: Request | None) -> dict:

View file

@ -122,6 +122,7 @@ def _otlp_error(content_type: str | None, status_code: int, message: str, retry:
)
@router.post("/v1/logs", include_in_schema=False)
@router.post("/v1/traces", include_in_schema=False)
async def ingest_otlp_traces(
request: Request,
@ -135,6 +136,7 @@ async def ingest_otlp_traces(
content_type=content_type,
content_encoding=request.headers.get("content-encoding"),
tenant=tenant,
logs=request.url.path.endswith("/v1/logs"),
)
except TracingPayloadTooLargeError as e:
return _otlp_error(content_type, 413, str(e))

View file

@ -41,7 +41,7 @@ class NativeTraceStorage:
def __new__(cls, config: NativeTraceConfig) -> NativeTraceStorage: ...
def ensure_schema(self) -> Future[None]: ...
def insert_rows(self, table: str, rows: Sequence[Mapping[str, object]]) -> Future[None]: ...
def ingest(self, payload: bytes, content_type: str | None, tenant: Mapping[str, str]) -> Future[int]: ...
def ingest(self, payload: bytes, content_type: str | None, tenant: Mapping[str, str], logs: bool = False) -> Future[int]: ...
def list_traces(
self, scope: TraceScope, start_ms: int, end_ms: int, cursor: str | None, limit: int
) -> Future[JsonValue]: ...

View file

@ -62,7 +62,9 @@ class NativeStore(Protocol):
def insert_rows(self, table: str, rows: Sequence[Mapping[str, object]]) -> Awaitable[None]: ...
def ingest(self, payload: bytes, content_type: str | None, tenant: Mapping[str, str]) -> Awaitable[int]: ...
def ingest(
self, payload: bytes, content_type: str | None, tenant: Mapping[str, str], logs: bool = False
) -> Awaitable[int]: ...
def list_traces(
self, scope: TraceScope, start_ms: int, end_ms: int, cursor: str | None, limit: int
@ -178,8 +180,8 @@ class ClickHouseStorage:
async def insert_rows(self, table: str, rows: Sequence[Mapping[str, object]]) -> None:
await self._native.insert_rows(table, rows)
async def ingest(self, payload: bytes, content_type: str | None, tenant: Tenant) -> int:
return await self._native.ingest(payload, content_type, asdict(tenant))
async def ingest(self, payload: bytes, content_type: str | None, tenant: Tenant, logs: bool = False) -> int:
return await self._native.ingest(payload, content_type, asdict(tenant), logs)
async def list_traces(
self,

View file

@ -61,10 +61,11 @@ class TraceReceiver:
content_type: str | None,
content_encoding: str | None,
tenant: Tenant,
logs: bool = False,
) -> int:
if not self._ingest_slots.acquire(blocking=False):
raise TracingOverloadedError("OTLP ingestion is at capacity")
task: Final = asyncio.create_task(self._ingest(body, content_type, content_encoding, tenant))
task: Final = asyncio.create_task(self._ingest(body, content_type, content_encoding, tenant, logs))
task.add_done_callback(self._release_ingest)
return await asyncio.shield(task)
@ -79,6 +80,7 @@ class TraceReceiver:
content_type: str | None,
content_encoding: str | None,
tenant: Tenant,
logs: bool = False,
) -> int:
try:
received: Final = (
@ -90,7 +92,7 @@ class TraceReceiver:
raise TracingOverloadedError("OTLP body upload timed out") from error
payload: Final = await asyncio.to_thread(self._decompressor, received, content_encoding)
try:
return await self.storage.ingest(payload, content_type, tenant)
return await self.storage.ingest(payload, content_type, tenant, logs)
except OverflowError as error:
raise TracingPayloadTooLargeError(str(error)) from error
except ValueError as error:

View file

@ -28,11 +28,14 @@ def _fake_storage() -> MagicMock:
@pytest.mark.asyncio
async def test_ingest_decompresses_and_passes_the_authenticated_tenant() -> None:
@pytest.mark.parametrize("logs", (False, True))
async def test_ingest_decompresses_and_passes_the_authenticated_tenant(logs: bool) -> None:
storage: Final = _fake_storage()
count: Final = await TraceReceiver(storage).ingest(gzip.compress(b"export"), "application/json", "gzip", TENANT)
count: Final = await TraceReceiver(storage).ingest(
gzip.compress(b"export"), "application/json", "gzip", TENANT, logs=logs
)
assert count == 6
storage.ingest.assert_awaited_once_with(b"export", "application/json", TENANT)
storage.ingest.assert_awaited_once_with(b"export", "application/json", TENANT, logs)
@pytest.mark.asyncio
@ -85,7 +88,7 @@ async def test_cancelled_request_keeps_its_worker_slot_until_decompression_finis
storage: Final = _fake_storage()
async def store(payload: bytes, content_type: str | None, tenant: Tenant) -> int:
async def store(payload: bytes, content_type: str | None, tenant: Tenant, logs: bool) -> int:
stored.set()
return 0

View file

@ -12,7 +12,6 @@ from starlette.datastructures import FormData
from starlette.requests import Request
import litellm
import litellm.proxy.common_utils.http_parsing_utils as http_parsing_utils
from litellm.proxy._types import ProxyException
@ -478,9 +477,7 @@ async def test_circular_reference_handling():
# Second parse using the same request - will use the modified cached value
result2 = await _read_request_body(mock_request)
assert (
"proxy_server_request" not in result2
) # This will pass, showing the cache pollution
assert "proxy_server_request" not in result2 # This will pass, showing the cache pollution
@pytest.mark.asyncio
@ -591,9 +588,7 @@ async def test_surrogate_repair_skipped_above_size_limit(monkeypatch):
import litellm.proxy.common_utils.http_parsing_utils as http_parsing_utils
# Cap the repair at ~100 bytes so the test stays fast and independent of the default.
monkeypatch.setattr(
http_parsing_utils, "MAX_REQUEST_BODY_SIZE_TO_REPAIR_MB", 100 / (1024 * 1024)
)
monkeypatch.setattr(http_parsing_utils, "MAX_REQUEST_BODY_SIZE_TO_REPAIR_MB", 100 / (1024 * 1024))
small_body = b'{"model":"gpt-4o","x":NaN}'
assert len(small_body) <= 100
@ -601,9 +596,7 @@ async def test_surrogate_repair_skipped_above_size_limit(monkeypatch):
assert repaired["model"] == "gpt-4o"
padding = "a" * 200
large_body = (
b'{"model":"gpt-4o","pad":"' + padding.encode() + b'","x":NaN}'
)
large_body = b'{"model":"gpt-4o","pad":"' + padding.encode() + b'","x":NaN}'
assert len(large_body) > 100
with pytest.raises(ProxyException) as exc_info:
await _read_request_body(_make_json_request(large_body))
@ -612,9 +605,7 @@ async def test_surrogate_repair_skipped_above_size_limit(monkeypatch):
# Disabling the cap (0) restores repair for the same large body, proving the cap
# — not the malformed content — is what short-circuits the repair.
monkeypatch.setattr(
http_parsing_utils, "MAX_REQUEST_BODY_SIZE_TO_REPAIR_MB", 0
)
monkeypatch.setattr(http_parsing_utils, "MAX_REQUEST_BODY_SIZE_TO_REPAIR_MB", 0)
repaired_large = await _read_request_body(_make_json_request(large_body))
assert repaired_large["model"] == "gpt-4o"
@ -643,7 +634,7 @@ async def test_lone_surrogate_escape_is_rejected_with_400(content: bytes):
paired = body.replace(content, b"say ok \\ud83d\\ude00")
parsed = await _read_request_body(_make_json_request(paired))
assert parsed["messages"][0]["content"] == "say ok \U0001F600"
assert parsed["messages"][0]["content"] == "say ok \U0001f600"
@pytest.mark.asyncio
@ -865,9 +856,7 @@ def test_populate_request_with_path_params_does_not_overwrite_existing_values():
# Verify existing values were NOT overwritten
assert result["model"] == "gpt-4" # Should keep original, not "gpt-3.5-turbo"
assert (
result["organization_id"] == "org-existing"
) # Should keep original, not "org-query-param"
assert result["organization_id"] == "org-existing" # Should keep original, not "org-query-param"
# Verify other data is preserved
assert result["messages"] == [{"role": "user", "content": "Hello"}]
@ -1035,9 +1024,7 @@ class TestGetTagsFromRequestBodyStringCoerce:
)
# Must not raise; must yield no metadata tags but keep root tags
tags = get_tags_from_request_body(
{"metadata": "not-json", "tags": ["root-only"]}
)
tags = get_tags_from_request_body({"metadata": "not-json", "tags": ["root-only"]})
assert tags == ["root-only"]
def test_dict_metadata_still_works(self):
@ -1094,9 +1081,7 @@ class TestReadRequestBodyNonCanonicalContentType:
"multiform/anything",
],
)
async def test_json_body_with_formlike_content_type_parses_as_json(
self, content_type
):
async def test_json_body_with_formlike_content_type_parses_as_json(self, content_type):
payload = {"user_config": {"model_list": []}, "model": "x"}
mock_request = MagicMock()
@ -1329,6 +1314,8 @@ def test_shared_inference_model_selection_preserves_handler_precedence(
"method,path,skip_parse",
[
("POST", "/v1/traces", True),
("POST", "/v1/logs", True),
("GET", "/v1/logs", False),
("GET", "/v1/traces", False),
("POST", "/v1/messages", False),
("POST", "/v1/traces/other", False),
@ -1340,7 +1327,10 @@ async def test_only_trace_ingest_skips_json_body(method: str, path: str, skip_pa
receive: Final = AsyncMock(return_value={"type": "http.request", "body": body, "more_body": False})
request: Final = Request(
{
"type": "http", "method": method, "path": root_path + path, "root_path": root_path,
"type": "http",
"method": method,
"path": root_path + path,
"root_path": root_path,
"headers": [(b"content-type", b"application/json")],
},
receive,
@ -1356,9 +1346,14 @@ async def test_only_trace_ingest_skips_json_body(method: str, path: str, skip_pa
@pytest.mark.asyncio
@pytest.mark.parametrize("content_type, encoding", [
("application/json", ""), ("application/x-protobuf", ""), ("application/json", "gzip"),
])
@pytest.mark.parametrize(
"content_type, encoding",
[
("application/json", ""),
("application/x-protobuf", ""),
("application/json", "gzip"),
],
)
async def test_otlp_auth_does_not_consume_chunked_bodies_before_the_receiver_limit(content_type, encoding):
from litellm.constants import OTLP_MAX_BODY_BYTES
from litellm.tracing import Tenant, TraceReceiver, TracingPayloadTooLargeError
@ -1371,9 +1366,18 @@ async def test_otlp_auth_does_not_consume_chunked_bodies_before_the_receiver_lim
assert len(received) <= 2, "receiver must reject without consuming subsequent chunks"
return {"type": "http.request", "body": chunk, "more_body": True}
request = Request({"type": "http", "method": "POST", "path": "/v1/traces", "headers": [
(b"content-type", content_type.encode()), (b"content-encoding", encoding.encode()),
]}, receive)
request = Request(
{
"type": "http",
"method": "POST",
"path": "/v1/traces",
"headers": [
(b"content-type", content_type.encode()),
(b"content-encoding", encoding.encode()),
],
},
receive,
)
assert await _read_request_body(request) == {}
assert received == []
storage = MagicMock()
@ -1393,9 +1397,7 @@ async def test_auth_body_read_and_trace_handler_leave_stream_for_receiver_limit(
from litellm.tracing import TraceReceiver
chunk: Final = b"x" * (OTLP_MAX_BODY_BYTES // 2 + 1)
receive: Final = AsyncMock(
side_effect=[{"type": "http.request", "body": chunk, "more_body": True}] * 2
)
receive: Final = AsyncMock(side_effect=[{"type": "http.request", "body": chunk, "more_body": True}] * 2)
request: Final = Request(
{"type": "http", "method": "POST", "path": "/v1/traces", "headers": [(b"content-type", b"application/json")]},
receive,

View file

@ -43,8 +43,7 @@ QUERY_HELP: Final[Mapping[str, object]] = {
"dialect": "test SQL",
"access": "authenticated scope",
"response": (
'JSON object {"data": [rows]}; each row maps selected columns to values; '
"64-bit integers may be strings"
'JSON object {"data": [rows]}; each row maps selected columns to values; 64-bit integers may be strings'
),
"tables": [{"name": "otel_traces", "columns": [{"name": "value", "type": "String", "comment": "label"}]}],
"normalized_fields": [],
@ -236,9 +235,10 @@ def test_501_when_tracing_not_enabled(
assert client.get("/v1/traces").status_code == 501
def test_post_protobuf_returns_empty_protobuf(client, receiver):
@pytest.mark.parametrize("endpoint", ("/v1/traces", "/v1/logs"))
def test_post_protobuf_returns_empty_protobuf(client, receiver, endpoint):
response = client.post(
"/v1/traces",
endpoint,
content=b"\x0a\x00",
headers={"content-type": "application/x-protobuf", "content-encoding": "gzip"},
)
@ -247,27 +247,31 @@ def test_post_protobuf_returns_empty_protobuf(client, receiver):
assert response.headers["content-type"] == "application/x-protobuf"
kwargs = receiver.ingest.call_args.kwargs
assert kwargs["body"] is not None
assert kwargs["logs"] is (endpoint == "/v1/logs")
assert kwargs["content_type"] == "application/x-protobuf"
assert kwargs["content_encoding"] == "gzip"
assert kwargs["tenant"].team_id == "team-research"
def test_post_json_returns_empty_json(client, receiver):
response = client.post("/v1/traces", content=b"{}", headers={"content-type": "application/json"})
@pytest.mark.parametrize("endpoint", ("/v1/traces", "/v1/logs"))
def test_post_json_returns_empty_json(client, receiver, endpoint):
response = client.post(endpoint, content=b"{}", headers={"content-type": "application/json"})
assert response.status_code == 200
assert response.json() == {}
def test_post_clickhouse_failure_is_503_with_retry_after(client, receiver):
@pytest.mark.parametrize("endpoint", ("/v1/traces", "/v1/logs"))
def test_post_clickhouse_failure_is_503_with_retry_after(client, receiver, endpoint):
receiver.ingest.side_effect = RuntimeError("ClickHouse unavailable")
response = client.post("/v1/traces", content=b"", headers={"content-type": "application/x-protobuf"})
response = client.post(endpoint, content=b"", headers={"content-type": "application/x-protobuf"})
assert response.status_code == 503
assert response.headers["retry-after"] == str(tracing_endpoints.OTLP_RETRY_AFTER_SECONDS)
def test_post_too_large_is_413(client, receiver):
@pytest.mark.parametrize("endpoint", ("/v1/traces", "/v1/logs"))
def test_post_too_large_is_413(client, receiver, endpoint):
receiver.ingest.side_effect = TracingPayloadTooLargeError("OTLP body exceeds 10 bytes")
response = client.post("/v1/traces", content=b"x" * 20)
response = client.post(endpoint, content=b"x" * 20)
assert response.status_code == 413
from google.rpc.status_pb2 import Status
@ -379,9 +383,7 @@ def test_trace_detail_reports_page_size_validation(
receiver.get_trace.assert_not_awaited()
def test_trace_read_routes_accept_and_forward_512_character_cursors(
client: TestClient, receiver: MagicMock
) -> None:
def test_trace_read_routes_accept_and_forward_512_character_cursors(client: TestClient, receiver: MagicMock) -> None:
cursor: Final = "x" * 512
receiver.get_trace.return_value = TRACE_RESPONSE
receiver.get_span_error = AsyncMock(return_value=SPAN_ERROR_RESPONSE)
@ -414,9 +416,7 @@ def test_trace_read_routes_accept_and_forward_512_character_cursors(
"path",
("/v1/traces", "/v1/traces/t1", "/v1/traces/t1/spans/s1/error"),
)
def test_trace_read_routes_reject_513_character_cursors(
client: TestClient, receiver: MagicMock, path: str
) -> None:
def test_trace_read_routes_reject_513_character_cursors(client: TestClient, receiver: MagicMock, path: str) -> None:
response: Final = client.get(path, params={"cursor": "x" * 513})
_assert_validation_error(response, "string_too_long", ("query", "cursor"))
@ -433,14 +433,10 @@ def test_trace_read_routes_ignore_unknown_query_parameters(client: TestClient, r
detail_response: Final = client.get("/v1/traces/t1", params=detail_params)
detail_unknown_response: Final = client.get("/v1/traces/t1", params={**detail_params, "foo": "bar"})
span_response: Final = client.get("/v1/traces/t1/spans/s1", params={"trace_ref": "run-one"})
span_unknown_response: Final = client.get(
"/v1/traces/t1/spans/s1", params={"trace_ref": "run-one", "foo": "bar"}
)
span_unknown_response: Final = client.get("/v1/traces/t1/spans/s1", params={"trace_ref": "run-one", "foo": "bar"})
error_params: Final = {"trace_ref": "run-one", "cursor": "error-cursor"}
error_response: Final = client.get("/v1/traces/t1/spans/s1/error", params=error_params)
error_unknown_response: Final = client.get(
"/v1/traces/t1/spans/s1/error", params={**error_params, "foo": "bar"}
)
error_unknown_response: Final = client.get("/v1/traces/t1/spans/s1/error", params={**error_params, "foo": "bar"})
assert list_response.status_code == 200, list_response.text
assert list_unknown_response.status_code == 200, list_unknown_response.text
@ -574,11 +570,12 @@ def test_key_without_user_cannot_read_traces(client: TestClient, auth: UserAPIKe
storage.query_help.assert_not_called()
def test_view_only_admin_cannot_ingest_traces(client, receiver):
@pytest.mark.parametrize("endpoint", ("/v1/traces", "/v1/logs"))
def test_view_only_admin_cannot_ingest_traces(client, receiver, endpoint):
client.app.dependency_overrides[user_api_key_auth] = lambda: UserAPIKeyAuth(
token="admin-key", user_role=LitellmUserRoles.PROXY_ADMIN_VIEW_ONLY
)
response = client.post("/v1/traces", content=b"{}")
response = client.post(endpoint, content=b"{}")
assert response.status_code == 403
receiver.ingest.assert_not_called()
@ -634,6 +631,7 @@ def test_injected_receiver_ingests_with_the_authenticated_tenant(client: TestCli
org_id=TEAM_KEY.org_id or "",
user_id=TEAM_KEY.user_id or "",
),
False,
)
@ -818,9 +816,9 @@ def test_sql_query_returns_empty_data(client: TestClient, receiver: MagicMock) -
def test_sql_query_openapi_declares_a_closed_response_object(client: TestClient) -> None:
openapi: Final = client.app.openapi()
response: Final = openapi["paths"]["/v1/traces/query"]["post"]["responses"]["200"]["content"][
"application/json"
]["schema"]
response: Final = openapi["paths"]["/v1/traces/query"]["post"]["responses"]["200"]["content"]["application/json"][
"schema"
]
component_name: Final = response["$ref"].rsplit("/", 1)[-1]
component: Final = openapi["components"]["schemas"][component_name]
@ -844,9 +842,7 @@ def test_sql_query_rejects_invalid_request_bodies(
) -> None:
client.app.dependency_overrides[tracing_endpoints.provide_trace_query_secret] = lambda: "test-secret"
receiver.storage.query_sql = AsyncMock()
response: Final = client.post(
"/v1/traces/query", content=body, headers={"content-type": "application/json"}
)
response: Final = client.post("/v1/traces/query", content=body, headers={"content-type": "application/json"})
_assert_validation_error(response, error_type, location)
receiver.storage.query_sql.assert_not_awaited()

View file

@ -138,6 +138,63 @@ describe("TraceConversation", () => {
expect(screen.queryByRole("button", { name: /Load next/ })).not.toBeInTheDocument();
});
it("shows a loaded agent reply while later trace pages remain available", async () => {
const user = userEvent.setup();
vi.mocked(agentTraceCall).mockResolvedValue({ ...trace, spans: [root], next_cursor: "next-page" });
renderWithProviders(
<RoutedRunView traceId={trace.summary.trace_id} accessToken="test" onBack={vi.fn()} embedded />,
);
await user.click(await screen.findByRole("tab", { name: "Conversation" }));
expect(await screen.findByText("The release is ready")).toBeVisible();
expect(screen.getByRole("button", { name: "Load more steps" })).toBeVisible();
expect(screen.queryByText("End of conversation")).not.toBeInTheDocument();
});
it.each([false, true])(
"checks for missing replies after the final trace page and details load (reply: %s)",
async (hasReply) => {
const user = userEvent.setup();
const laterDetail = Promise.withResolvers<SpanDetail>();
const summary = { ...trace.summary, span_count: 2 };
const first: Trace = { ...trace, summary, spans: [root], next_cursor: "last-page" };
const last: Trace = {
...trace,
summary,
spans: [{ ...tool, span_id: "later", name: "later response", type: "llm" }],
next_cursor: null,
};
vi.mocked(agentTraceCall).mockImplementation(async (_token, _trace, _ref, cursor) => (cursor ? last : first));
vi.mocked(agentTraceSpanCall).mockImplementation(async (_token, _trace, id) =>
id === "root" ? { ...rootDetail, output: "", attributes: { "span.type": "llm_request" } } : laterDetail.promise,
);
renderWithProviders(
<RoutedRunView traceId={trace.summary.trace_id} accessToken="test" onBack={vi.fn()} embedded />,
);
await user.click(await screen.findByRole("tab", { name: "Conversation" }));
expect(await screen.findByText("Read the release notes")).toBeVisible();
expect(screen.queryByText(/no recorded assistant replies/)).not.toBeInTheDocument();
expect(screen.queryByText("End of conversation")).not.toBeInTheDocument();
await user.click(screen.getByRole("button", { name: "Load more steps" }));
expect(await screen.findByText("Loading conversation…")).toBeVisible();
expect(screen.queryByText(/no recorded assistant replies/)).not.toBeInTheDocument();
expect(screen.queryByText("End of conversation")).not.toBeInTheDocument();
const resolved: SpanDetail = {
span_id: "later",
input: "",
output: hasReply ? rootDetail.output : "",
attributes: hasReply ? { "event.name": "assistant_response" } : {},
};
await act(async () => laterDetail.resolve(resolved));
expect(await screen.findByText("End of conversation")).toBeVisible();
if (hasReply) {
expect(screen.getByText("The release is ready")).toBeVisible();
expect(screen.queryByText(/no recorded assistant replies/)).not.toBeInTheDocument();
} else {
expect(screen.getByText(/no recorded assistant replies/)).toBeVisible();
}
},
);
it("shows distinct agent invocation labels together with each step's time", async () => {
const first = { ...root, name: "reviewer", start_offset_ms: 1000 };
const second = { ...first, span_id: "second", start_offset_ms: 2000 };
@ -249,6 +306,7 @@ describe("TraceConversation", () => {
it.each(["Partial investigation", ""])(
"shows a failed child agent's error once after its work, with output %j",
async (output) => {
const user = userEvent.setup();
const agent = {
...root,
span_id: "child",
@ -269,6 +327,7 @@ describe("TraceConversation", () => {
renderWithProviders(<TraceConversation trace={traced} accessToken="test" onOpenStep={vi.fn()} />);
expect(await screen.findByText("End of conversation")).toBeVisible();
await user.click(screen.getByText("Subagent: Investigate release", { exact: true }));
expect(screen.getAllByText("Investigation timed out")).toHaveLength(1);
const entries = screen.getAllByRole("region", { name: "Conversation step Investigate release" });
expect(within(entries[0]).getByText("Investigate failed checks")).toBeVisible();

View file

@ -9,7 +9,15 @@ import { ToolArguments, ToolOutput } from "../content/ToolContent";
import { toolSummary } from "../content/payload";
import { useTracesApi } from "../../api";
import { Button } from "@/components/ui/button";
import { buildConversation, conversationSteps, CONVERSATION_PAGE_SIZE, type ConversationItem } from "./conversation";
import {
buildConversation,
conversationSteps,
conversationWarnings,
groupConversation,
CONVERSATION_PAGE_SIZE,
type ConversationItem,
type ConversationGroup,
} from "./conversation";
import { ErrorBlock } from "../content/SpanError";
import { Markdown } from "../content/Markdown";
import { ToolCallBlock, ToolResultCard } from "../content/Messages";
@ -46,7 +54,9 @@ export function TraceConversation({
queries.slice(0, loadedCount).map((query, index) => [visible[index].span_id, query.data!] as const),
);
const complete = loadedCount === steps.length;
const traceComplete = complete && !trace.next_cursor;
const items = buildConversation(trace.spans, details, complete);
const warnings = conversationWarnings(details, traceComplete);
const multipleAgents = new Set(items.map((item) => item.agentId).filter(Boolean)).size > 1;
const inlineErrorIds = new Set(items.filter((item) => item.showError).map((item) => item.span.span_id));
const rootErrors = trace.spans.filter((span) => {
@ -59,38 +69,16 @@ export function TraceConversation({
{rootErrors.map((span) => (
<ErrorBlock key={span.span_id} span={span} />
))}
{items.map((item) => (
<section key={item.id} className="min-w-0 space-y-2" aria-label={`Conversation step ${item.span.name}`}>
{(multipleAgents || item.toolResult === undefined) && (
<div className="flex min-w-0 items-center justify-between gap-2 text-xs text-muted-foreground">
<div className="flex min-w-0 items-center gap-2">
{multipleAgents && (
<span className="truncate" title={item.agentName}>
{item.agentName}
</span>
)}
<span className="shrink-0 tabular-nums">{fmtMs(item.span.start_offset_ms)}</span>
</div>
{item.toolResult === undefined && (
<Button
variant="ghost"
size="xs"
aria-label={`Inspect step ${item.span.name}`}
onClick={() => onOpenStep(item.span.span_id)}
className="shrink-0 text-muted-foreground"
>
Inspect step
</Button>
)}
</div>
)}
{item.showError && <ErrorBlock span={item.span} />}
{item.messages.map((message, index) => (
<ConversationMessage key={index} message={message} />
))}
{item.toolResult !== undefined && <ConversationTool item={item} onOpenStep={onOpenStep} />}
</section>
{warnings.map((warning) => (
<p key={warning} role="status" className="rounded-md border p-3 text-sm text-muted-foreground">
{warning}
</p>
))}
<ConversationGroups
groups={groupConversation(items, trace.spans)}
multipleAgents={multipleAgents}
onOpenStep={onOpenStep}
/>
{queries.map(
(query, index) =>
query.isError && (
@ -107,11 +95,11 @@ export function TraceConversation({
Loading conversation…
</p>
)}
{complete && items.length === 0 && (
{traceComplete && items.length === 0 && (
<p className="text-sm text-muted-foreground">No conversation content recorded.</p>
)}
<div className="flex items-center justify-between gap-3 border-t pt-4 text-xs text-muted-foreground">
<span>{complete ? "End of conversation" : `${loadedCount} of ${steps.length} steps loaded`}</span>
<span>{traceComplete ? "End of conversation" : `${loadedCount} of ${steps.length} steps loaded`}</span>
{visible.length < steps.length && (
<Button
variant="outline"
@ -128,6 +116,86 @@ export function TraceConversation({
);
}
function ConversationGroups({
groups,
multipleAgents,
onOpenStep,
}: {
groups: ConversationGroup[];
multipleAgents: boolean;
onOpenStep: (id: string) => void;
}) {
return (
<>
{groups.map((group) =>
group.kind === "item" ? (
<ConversationStep
key={group.item.id}
item={group.item}
multipleAgents={multipleAgents}
onOpenStep={onOpenStep}
/>
) : (
<details key={group.id} className="min-w-0 rounded-md border p-3">
<summary className="cursor-pointer text-sm font-medium">Subagent: {group.name}</summary>
<div className="mt-4 min-w-0 space-y-5 border-l pl-3">
<ConversationGroups groups={group.children} multipleAgents={multipleAgents} onOpenStep={onOpenStep} />
</div>
</details>
),
)}
</>
);
}
function ConversationStep({
item,
multipleAgents,
onOpenStep,
}: {
item: ConversationItem;
multipleAgents: boolean;
onOpenStep: (id: string) => void;
}) {
return (
<section className="min-w-0 space-y-2" aria-label={`Conversation step ${item.span.name}`}>
{(multipleAgents || item.toolResult === undefined) && (
<div className="flex min-w-0 items-center justify-between gap-2 text-xs text-muted-foreground">
<div className="flex min-w-0 items-center gap-2">
{multipleAgents && (
<span className="truncate" title={item.agentName}>
{item.agentName}
</span>
)}
{item.model && (
<span className="truncate" title={item.model}>
{item.model}
</span>
)}
<span className="shrink-0 tabular-nums">{fmtMs(item.time ?? item.span.start_offset_ms)}</span>
</div>
{item.toolResult === undefined && (
<Button
variant="ghost"
size="xs"
aria-label={`Inspect step ${item.span.name}`}
onClick={() => onOpenStep(item.span.span_id)}
className="shrink-0 text-muted-foreground"
>
Inspect step
</Button>
)}
</div>
)}
{item.showError && <ErrorBlock span={item.span} />}
{item.messages.map((message, index) => (
<ConversationMessage key={index} message={message} />
))}
{item.toolResult !== undefined && <ConversationTool item={item} onOpenStep={onOpenStep} />}
</section>
);
}
function ConversationTool({ item, onOpenStep }: { item: ConversationItem; onOpenStep: (id: string) => void }) {
const [open, setOpen] = useState(item.span.status === "error");
const failed = item.span.status === "error";
@ -183,7 +251,11 @@ function ConversationTool({ item, onOpenStep }: { item: ConversationItem; onOpen
iconOnly
/>
</div>
<ToolOutput result={item.toolResult ?? ""} failed={failed} />
{item.toolResult ? (
<ToolOutput result={item.toolResult} failed={failed} />
) : (
<p className="text-sm text-muted-foreground">No result content recorded.</p>
)}
</div>
)}
</div>
@ -191,10 +263,12 @@ function ConversationTool({ item, onOpenStep }: { item: ConversationItem; onOpen
}
function ConversationMessage({ message }: { message: TraceMessage }) {
if (message.role === "system" && (message.content.startsWith("Agent ") || message.content === "Context compacted"))
return <p className="text-xs text-muted-foreground">{message.content}</p>;
if (message.role === "system")
return (
<details className="text-sm text-muted-foreground">
<summary className="cursor-pointer">System instructions</summary>
<summary className="cursor-pointer">Session context</summary>
<div className="pt-3">
<Markdown text={message.content} />
</div>

View file

@ -1,5 +1,11 @@
import { describe, expect, it } from "vitest";
import { buildConversation, conversationSteps, newConversationMessages } from "./conversation";
import {
buildConversation,
conversationSteps,
newConversationMessages,
groupConversation,
conversationWarnings,
} from "./conversation";
import type { Span, SpanDetail, TraceMessage } from "../../types";
import research from "../../__fixtures__/research_trace.json";
@ -339,3 +345,273 @@ describe("recorded tool summaries", () => {
).toHaveLength(1);
});
});
describe("coding sessions", () => {
it("preserves identical user messages and assistant replies in transcript events", () => {
const spans = [
root,
...[1, 2, 3, 4].map((n) => ({
...root,
span_id: `m${n}`,
parent_span_id: root.span_id,
type: "chain" as const,
start_offset_ms: n,
})),
];
const details = new Map([
[root.span_id, { ...detail(root.span_id, [user], []), attributes: { "lens.capture.messages_separate": "true" } }],
...spans.slice(1).map(
(span, i) =>
[
span.span_id,
{
...detail(span.span_id, i % 2 === 0 ? [user] : [], i % 2 === 1 ? [answer] : []),
attributes: { "lens.capture.source": "session_transcript" },
},
] as const,
),
]);
expect(buildConversation(spans, details, true).flatMap((item) => item.messages)).toEqual([
user,
answer,
user,
answer,
]);
});
it("uses one actor across resumed turns and keeps child branches together", () => {
const resumed = { ...root, span_id: "resumed", start_offset_ms: 10 };
const child = { ...root, span_id: "child", parent_span_id: root.span_id, start_offset_ms: 2, name: "reader" };
const nested = { ...child, span_id: "nested", parent_span_id: "child", start_offset_ms: 3, name: "checker" };
const details = new Map(
[root, resumed, child, nested].map((span) => [
span.span_id,
{
...detail(span.span_id, span.name, "Done"),
attributes: { "gen_ai.agent.id": span === root || span === resumed ? "main-session" : span.span_id },
},
]),
);
const items = buildConversation([root, resumed, child, nested], details, true);
expect(new Set(items.filter((item) => !item.parentBranchId).map((item) => item.agentId))).toEqual(
new Set(["main-session"]),
);
const groups = groupConversation(items, [root, resumed, child, nested]);
const branch = groups.find((group) => group.kind === "branch");
expect(branch).toMatchObject({ kind: "branch", id: "child", name: "reader" });
if (branch?.kind === "branch")
expect(branch.children.some((group) => group.kind === "branch" && group.id === "nested")).toBe(true);
});
it("shows child conversations when their parent has no recorded messages", () => {
const child = { ...root, span_id: "child", parent_span_id: root.span_id, start_offset_ms: 1, name: "reader" };
const details = new Map([
[root.span_id, detail(root.span_id, [], [])],
[child.span_id, detail(child.span_id, "Read the file", "File contents")],
]);
const groups = groupConversation(buildConversation([root, child], details, true), [root, child]);
expect(groups).toHaveLength(1);
expect(groups[0]).toMatchObject({ kind: "branch", id: "child", name: "reader" });
if (groups[0].kind === "branch") {
expect(groups[0].children.flatMap((group) => (group.kind === "item" ? group.item.messages : []))).toEqual([
{ role: "user", content: "Read the file" },
{ role: "assistant", content: "File contents" },
]);
}
});
it("retains silent intermediate agents in nested branches", () => {
const child = { ...root, span_id: "child", parent_span_id: root.span_id, start_offset_ms: 1, name: "reader" };
const nested = { ...child, span_id: "nested", parent_span_id: child.span_id, start_offset_ms: 2, name: "checker" };
const spans = [root, child, nested];
const details = new Map([
[root.span_id, detail(root.span_id, [user], [])],
[child.span_id, detail(child.span_id, [], [])],
[nested.span_id, detail(nested.span_id, "Check the result", "Checked")],
]);
const expectedBranch = {
kind: "branch",
id: child.span_id,
name: "reader",
children: [expect.objectContaining({ kind: "branch", id: nested.span_id, name: "checker" })],
};
expect(groupConversation(buildConversation(spans, details, true), spans)).toEqual([
expect.objectContaining({ kind: "item" }),
expect.objectContaining(expectedBranch),
]);
});
it("folds native Claude Agent descendants and omits auxiliary suggestions", () => {
const agent = {
...root,
span_id: "agent-tool",
parent_span_id: root.span_id,
type: "tool" as const,
name: "Agent",
framework: "claude-code",
start_offset_ms: 1,
};
const execution = {
...agent,
span_id: "execution",
parent_span_id: agent.span_id,
type: "framework" as const,
name: "claude_code.tool.execution",
};
const response = {
...root,
span_id: "response",
parent_span_id: execution.span_id,
type: "chain" as const,
start_offset_ms: 3,
};
const suggestion = { ...response, span_id: "suggestion", parent_span_id: root.span_id, type: "llm" as const };
const details = new Map([
[root.span_id, { ...detail(root.span_id, [user], []), attributes: { "gen_ai.agent.id": "session" } }],
[
agent.span_id,
{
...detail(agent.span_id, { prompt: "Read", description: "Reader" }, "Started"),
attributes: { "gen_ai.agent.id": "session" },
},
],
[
response.span_id,
{ ...detail(response.span_id, [], "Child reply"), attributes: { "event.name": "assistant_response" } },
],
[
suggestion.span_id,
{
...detail(suggestion.span_id, [], "HIDDEN SUGGESTION"),
attributes: { query_source_safe: "prompt_suggestion" },
},
],
]);
const items = buildConversation([root, agent, execution, response, suggestion], details, true);
expect(items.flatMap((item) => item.messages).map((message) => message.content)).not.toContain("HIDDEN SUGGESTION");
expect(groupConversation(items, [root, agent, execution, response, suggestion])).toContainEqual(
expect.objectContaining({ kind: "branch", id: "agent-tool", name: "Reader" }),
);
expect(items.find((item) => item.span.span_id === response.span_id)?.agentId).toBe("agent-tool");
});
it("fills native tool arguments and missing results from logs without duplicating calls or replacing recorded output", () => {
const tool = {
...root,
span_id: "tool",
parent_span_id: root.span_id,
type: "tool" as const,
name: "Bash",
start_offset_ms: 1,
};
const log = {
...tool,
span_id: "log",
type: "framework" as const,
framework: "claude-code",
name: "claude_code.tool_result",
start_offset_ms: 2,
};
const body = { ...log, span_id: "body", name: "claude_code.api_request_body", start_offset_ms: 3 };
const toolDetail = {
...detail(tool.span_id, { command: "exit 3" }, ""),
output: "",
attributes: { "gen_ai.tool.call.id": "call-1" },
};
const details = new Map([
[root.span_id, detail(root.span_id, [user], [])],
[tool.span_id, toolDetail],
[
log.span_id,
{
...detail(log.span_id, { command: "exit 3", description: "Expected failure" }, ""),
attributes: { tool_use_id: "call-1" },
},
],
[
body.span_id,
detail(body.span_id, "", {
tool_results: [
{ id: "call-1", content: "Expected stdout" },
{ id: "unrelated", content: "Other stdout" },
],
}),
],
]);
const spans = [root, tool, log, body];
const items = buildConversation(spans, details, true);
expect(items.filter((item) => item.toolCall)).toHaveLength(1);
expect(items.find((item) => item.toolCall)).toMatchObject({
toolCall: { args: { command: "exit 3", description: "Expected failure" } },
toolResult: "Expected stdout",
});
details.set(tool.span_id, { ...toolDetail, output: "Recorded output" });
expect(buildConversation(spans, details, true).find((item) => item.toolCall)?.toolResult).toBe("Recorded output");
expect(
conversationWarnings(
new Map([
[
body.span_id,
{
...detail(body.span_id, "", { warning: "Export truncated" }),
attributes: { "event.name": "api_request_body" },
},
],
]),
true,
),
).toEqual(["Export truncated"]);
});
it.each([
{ recorded: "repl_main_thread", source: "repl_main_thread" },
{ recorded: "agent", source: "agent:builtin:general-purpose" },
{ recorded: "agent.builtin:general-purpose", source: "agent:builtin:general-purpose" },
{ recorded: "agent.custom:reader", source: "agent:custom:reader" },
])("positions $source commentary before tools with native source $recorded", ({ recorded, source }) => {
const llm = {
...root,
span_id: "llm",
parent_span_id: root.span_id,
type: "llm" as const,
start_offset_ms: 1,
duration_ms: 10,
model: "model-a",
};
const tool = { ...llm, span_id: "tool", name: "Read", type: "tool" as const, start_offset_ms: 8, duration_ms: 1 };
const response = { ...llm, span_id: "reply", type: "chain" as const, start_offset_ms: 11, duration_ms: 0 };
const details = new Map([
[root.span_id, detail(root.span_id, [], [])],
[
llm.span_id,
{
...detail(llm.span_id, [], []),
attributes: { query_source_safe: recorded, first_content_ms: "2" },
},
],
[tool.span_id, detail(tool.span_id, {}, "File contents")],
[
response.span_id,
{
...detail(response.span_id, [], "Checking"),
attributes: { "event.name": "assistant_response", query_source: source },
},
],
]);
expect(buildConversation([root, llm, tool, response], details, true)[0].time).toBe(3);
expect(buildConversation([root, llm, tool, response], details, true).map((item) => item.span.span_id)).toEqual([
"reply",
"tool",
]);
});
it("warns about old incomplete native captures and clears the warning when reply logs arrive", () => {
const details = new Map([
["llm", { ...detail("llm", [], []), output: "", attributes: { "span.type": "llm_request" } }],
]);
expect(conversationWarnings(details, false)).toEqual([]);
expect(conversationWarnings(details, true)).toHaveLength(1);
details.set("reply", { ...detail("reply", [], "Hello"), attributes: { "event.name": "assistant_response" } });
expect(conversationWarnings(details, true)).toEqual([]);
});
});

View file

@ -1,6 +1,8 @@
import { groupBy, orderBy } from "es-toolkit";
import type { Span, SpanDetail, TraceMessage, TraceToolCall, UIContent } from "../../types";
import { isFrameworkSpan, parseAssistantSummary, parseMessages, prettyPayload } from "../../utils";
import { toTraceMessage, toolInput } from "../content/payload";
import { claudeCaptureWarnings, claudeToolDetails, isClaudeSupplement } from "./nativeClaude";
export const CONVERSATION_PAGE_SIZE = 20;
@ -9,7 +11,8 @@ export function conversationSteps(spans: readonly Span[]): Span[] {
return spans
.filter((span) => {
const isEvent = ["agent", "llm", "tool"].includes(span.type) || !parents.has(span.span_id);
return !isFrameworkSpan(span) && (span.parent_span_id === null || isEvent);
const visible = !isFrameworkSpan(span) && (span.parent_span_id === null || isEvent);
return isClaudeSupplement(span) || visible;
})
.sort((a, b) => a.start_offset_ms - b.start_offset_ms);
}
@ -31,6 +34,12 @@ function messages(value: string, content: UIContent | undefined, role: string):
return text ? [{ role, content: text }] : [];
}
function inputMessages(detail: SpanDetail): TraceMessage[] {
return detail.attributes["lens.capture.messages_separate"] === "true"
? []
: messages(detail.input, detail.input_ui, "user");
}
function stableValue(value: unknown): unknown {
if (Array.isArray(value)) return value.map(stableValue);
if (value !== null && typeof value === "object")
@ -72,6 +81,10 @@ export interface ConversationItem {
toolResult?: string;
agentId?: string;
agentName?: string;
branchId?: string;
parentBranchId?: string;
model?: string;
time?: number;
showError?: boolean;
}
@ -110,6 +123,26 @@ interface ConversationEvent {
output: boolean;
}
function messageTime(span: Span, steps: readonly Span[], details: ReadonlyMap<string, SpanDetail>): number {
const attributes = details.get(span.span_id)?.attributes;
if (attributes?.["event.name"] !== "assistant_response") return span.start_offset_ms;
const source = attributes["query_source"];
const matches = steps.filter((candidate) => {
const recorded = details.get(candidate.span_id)?.attributes["query_source_safe"];
const unnamedAgent = recorded === "agent" && source?.startsWith("agent:");
const sameSource = recorded === source || recorded === source?.replace(/^agent:/, "agent.") || unnamedAgent;
const sameRequest = candidate.parent_span_id === span.parent_span_id && candidate.model === span.model;
const matchingCall = candidate.type === "llm" && sameRequest && sameSource;
return matchingCall && Math.abs(candidate.start_offset_ms + candidate.duration_ms - span.start_offset_ms) < 2;
});
if (matches.length !== 1) return span.start_offset_ms;
const request = matches[0];
const firstContent = Number(details.get(request.span_id)?.attributes["first_content_ms"]);
return Number.isFinite(firstContent) && firstContent >= 0 && firstContent <= request.duration_ms
? request.start_offset_ms + firstContent
: span.start_offset_ms;
}
function conversationEvents(
steps: readonly Span[],
byId: ReadonlyMap<string, Span>,
@ -134,7 +167,7 @@ function conversationEvents(
}
return loaded
.flatMap((span) => {
const start = { span, time: span.start_offset_ms, output: false };
const start = { span, time: messageTime(span, loaded, details), output: false };
if (span.type === "tool" || (span.type !== "agent" && span.parent_span_id !== null)) return [start];
const end = ends.get(span.span_id)!;
return (complete && missingIndex < 0) || end < boundary ? [start, { span, time: end, output: true }] : [start];
@ -171,12 +204,45 @@ function withoutForwardedAnswers(
);
}
function agentLabels(agents: readonly Span[]): ReadonlyMap<string, string> {
const named = agents.map((agent) => ({ agent, name: agent.name || agent.agent || "Agent" }));
function isNativeAgent(span: Span): boolean {
return span.framework === "claude-code" && span.type === "tool" && ["Agent", "Task"].includes(span.name);
}
function conversationBranch(span: Span, byId: ReadonlyMap<string, Span>): string {
if (span.type === "agent" || span.parent_span_id === null) return span.span_id;
let parent = span.parent_span_id ? byId.get(span.parent_span_id) : undefined;
const visited = new Set<string>();
while (parent && !visited.has(parent.span_id)) {
visited.add(parent.span_id);
if (parent.type === "agent" || isNativeAgent(parent)) return parent.span_id;
parent = parent.parent_span_id ? byId.get(parent.parent_span_id) : undefined;
}
return span.parent_span_id ?? span.span_id;
}
function agentIdentity(span: Span, details: ReadonlyMap<string, SpanDetail>): string {
return isNativeAgent(span) ? span.span_id : details.get(span.span_id)?.attributes["gen_ai.agent.id"] || span.span_id;
}
function agentLabels(agents: readonly Span[], details: ReadonlyMap<string, SpanDetail>): ReadonlyMap<string, string> {
const identity = (agent: Span): string => agentIdentity(agent, details);
const unique = agents.filter(
(agent, index) => agents.findIndex((other) => identity(other) === identity(agent)) === index,
);
const named = unique.map((agent) => {
const detail = details.get(agent.span_id);
const args = detail && isNativeAgent(agent) ? toolInput(detail.input, detail.input_ui) : undefined;
const description = args && typeof args === "object" && "description" in args ? args.description : undefined;
const name =
agent.framework === "claude-code" && !isNativeAgent(agent)
? agent.agent || agent.name || "Agent"
: agent.name || agent.agent || "Agent";
return { agent, name: typeof description === "string" && description ? description : name };
});
return new Map(
Object.values(groupBy(named, ({ name }) => JSON.stringify(name))).flatMap((group) =>
orderBy(group, [({ agent }) => agent.start_offset_ms, ({ agent }) => agent.span_id], ["asc", "asc"]).map(
({ agent, name }, index) => [agent.span_id, group.length > 1 ? `${name} (${index + 1})` : name] as const,
({ agent, name }, index) => [identity(agent), group.length > 1 ? `${name} (${index + 1})` : name] as const,
),
),
);
@ -184,29 +250,22 @@ function agentLabels(agents: readonly Span[]): ReadonlyMap<string, string> {
export function buildConversation(
spans: readonly Span[],
details: ReadonlyMap<string, SpanDetail>,
recordedDetails: ReadonlyMap<string, SpanDetail>,
complete: boolean,
): ConversationItem[] {
const details = claudeToolDetails(spans, recordedDetails);
const byId = new Map(spans.map((span) => [span.span_id, span]));
const histories = new Map<string, TraceMessage[]>();
const completedOutputs = new Map<string, TraceMessage[]>();
const pendingCalls = new Map<string, TraceToolCall[]>();
const items: ConversationItem[] = [];
const events = conversationEvents(conversationSteps(spans), byId, details, complete);
const branch = (span: Span): string => {
if (span.type === "agent" || span.parent_span_id === null) return span.span_id;
let parent = span.parent_span_id ? byId.get(span.parent_span_id) : undefined;
const visited = new Set<string>();
while (parent && !visited.has(parent.span_id)) {
visited.add(parent.span_id);
if (parent.type === "agent") return parent.span_id;
parent = parent.parent_span_id ? byId.get(parent.parent_span_id) : undefined;
}
return span.parent_span_id ?? span.span_id;
};
const branch = (span: Span): string => conversationBranch(span, byId);
for (const event of events) {
const { span } = event;
if (isClaudeSupplement(span)) continue;
const detail = details.get(span.span_id)!;
if (["generate_session_title", "prompt_suggestion"].includes(detail.attributes["query_source_safe"])) continue;
const key = branch(span);
const history = histories.get(key) ?? [];
if (event.output) {
@ -217,7 +276,13 @@ export function buildConversation(
completedOutputs,
byId,
);
const item = { id: `${span.span_id}-output`, span, messages: fresh, showError: span.status === "error" };
const item = {
id: `${span.span_id}-output`,
span,
time: event.time,
messages: fresh,
showError: span.status === "error",
};
if (fresh.length || item.showError) items.push(item);
completedOutputs.set(span.span_id, output);
histories.set(key, [...history, ...fresh]);
@ -229,12 +294,16 @@ export function buildConversation(
histories.set(key, [...history, { role: "tool", name: span.name, content: item.toolResult ?? "" }]);
continue;
}
const input = messages(detail.input, detail.input_ui, "user");
const input = inputMessages(detail);
const output = messages(detail.output, detail.output_ui, "assistant");
const fresh = newConversationMessages(history, input);
const fresh =
detail.attributes["lens.capture.source"] === "session_transcript"
? input
: newConversationMessages(history, input);
if (span.type === "agent" || span.parent_span_id === null) {
histories.set(key, input);
if (fresh.length) items.push({ id: span.span_id, span, messages: fresh });
const item = { id: span.span_id, span, time: event.time, messages: fresh };
if (fresh.length) items.push(item);
continue;
}
const combined = [...fresh, ...output];
@ -247,6 +316,7 @@ export function buildConversation(
const item = {
id: span.span_id,
span,
time: event.time,
showError: span.status === "error",
messages: combined.map((message) => ({
...message,
@ -259,14 +329,93 @@ export function buildConversation(
const span = byId.get(id);
return span ? [span] : [];
});
const labels = agentLabels(agents);
const labels = agentLabels(agents, details);
const actor = (id: string): string => {
const span = byId.get(id);
return span ? agentIdentity(span, details) : id;
};
return items
.map((item) => ({
...item,
agentId: branch(item.span),
agentName: labels.get(branch(item.span)) || item.span.agent,
agentId: actor(branch(item.span)),
agentName: labels.get(actor(branch(item.span))) || item.span.agent,
branchId: branch(item.span),
parentBranchId: (() => {
const parentId = byId.get(branch(item.span))?.parent_span_id;
const parent = parentId ? byId.get(parentId) : undefined;
return parent ? branch(parent) : undefined;
})(),
model: item.span.model || details.get(item.span.span_id)?.attributes["gen_ai.request.model"] || undefined,
messages: item.messages.filter((message) => Boolean(message.content) || Boolean(message.tool_calls?.length)),
}))
.filter((item) => item.messages.length || item.toolResult !== undefined || item.showError);
}
import { groupBy, orderBy } from "es-toolkit";
export type ConversationGroup =
| { kind: "item"; item: ConversationItem }
| { kind: "branch"; id: string; name: string; children: ConversationGroup[] };
export function groupConversation(items: readonly ConversationItem[], spans: readonly Span[]): ConversationGroup[] {
const byId = new Map(spans.map((span) => [span.span_id, span]));
const parentById = new Map(
conversationSteps(spans).flatMap((span) => {
const branch = byId.get(conversationBranch(span, byId));
if (!branch) return [];
const parent = branch.parent_span_id ? byId.get(branch.parent_span_id) : undefined;
return [[branch.span_id, parent ? conversationBranch(parent, byId) : undefined] as const];
}),
);
const roots = new Set(
[...parentById].flatMap(([id, parent]) => {
if (!parent) return [id];
return parentById.has(parent) ? [] : [parent];
}),
);
const directBranch = (item: ConversationItem, parent?: string): string | undefined => {
let id = item.branchId;
const visited = new Set<string>();
while (id && !visited.has(id)) {
visited.add(id);
const ancestor = parentById.get(id);
if (parent ? ancestor === parent : ancestor !== undefined && roots.has(ancestor)) return id;
id = ancestor;
}
return undefined;
};
const build = (parent?: string, ancestors = new Set<string>()): ConversationGroup[] => {
const seen = new Set<string>();
return items.flatMap((item): ConversationGroup[] => {
if (parent ? item.branchId === parent : !item.parentBranchId) return [{ kind: "item", item }];
const id = directBranch(item, parent);
if (!id || seen.has(id) || ancestors.has(id)) return [];
seen.add(id);
const first = items.find((candidate) => candidate.branchId === id);
const name = first?.agentName || byId.get(id)?.name || "Subagent";
return [{ kind: "branch", id, name, children: build(id, new Set([...ancestors, id])) }];
});
};
return build();
}
export function conversationWarnings(details: ReadonlyMap<string, SpanDetail>, complete: boolean): string[] {
const warnings = [
...claudeCaptureWarnings(details),
...[...details.values()].flatMap((detail) =>
detail.attributes["lens.capture.warning"] ? [detail.attributes["lens.capture.warning"]] : [],
),
];
if (
complete &&
![...details.values()].some((detail) => detail.attributes["event.name"] === "assistant_response") &&
[...details.values()].some(
(detail) =>
detail.attributes["span.type"] === "llm_request" &&
!detail.output &&
!messages(detail.output, detail.output_ui, "assistant").length,
)
) {
warnings.push(
"This Claude Code trace has no recorded assistant replies. Enable assistant response logs for future sessions.",
);
}
return [...new Set(warnings)];
}

View file

@ -0,0 +1,68 @@
import type { Span, SpanDetail } from "../../types";
export function isClaudeSupplement(span: Span): boolean {
return (
span.framework === "claude-code" && ["claude_code.tool_result", "claude_code.api_request_body"].includes(span.name)
);
}
function object(value: unknown): Record<string, unknown> | undefined {
return value !== null && typeof value === "object" && !Array.isArray(value)
? (value as Record<string, unknown>)
: undefined;
}
function payload(detail: SpanDetail): Record<string, unknown> | undefined {
try {
return object(JSON.parse(detail.output));
} catch {
return undefined;
}
}
export function claudeToolDetails(
spans: readonly Span[],
details: ReadonlyMap<string, SpanDetail>,
): ReadonlyMap<string, SpanDetail> {
const inputs = new Map<string, SpanDetail>();
const outputs = new Map<string, string>();
for (const span of spans.filter(isClaudeSupplement)) {
const detail = details.get(span.span_id);
if (!detail) continue;
const id = detail.attributes["tool_use_id"];
if (span.name === "claude_code.tool_result" && id && detail.input) inputs.set(id, detail);
const results = payload(detail)?.tool_results;
if (!Array.isArray(results)) continue;
for (const value of results) {
const result = object(value);
if (typeof result?.id === "string" && typeof result.content === "string") {
outputs.set(result.id, result.content);
}
}
}
return new Map(
[...details].map(([id, detail]) => {
const callId = detail.attributes["gen_ai.tool.call.id"] || detail.attributes["tool_use_id"];
const input = inputs.get(callId);
const output = outputs.get(callId);
return [
id,
{
...detail,
...(input ? { input: input.input, input_ui: input.input_ui } : {}),
...(!detail.output && output !== undefined
? { output, output_ui: { kind: "text" as const, text: output } }
: {}),
},
];
}),
);
}
export function claudeCaptureWarnings(details: ReadonlyMap<string, SpanDetail>): string[] {
return [...details.values()].flatMap((detail) => {
if (detail.attributes["event.name"] !== "api_request_body") return [];
const warning = payload(detail)?.warning;
return typeof warning === "string" ? [warning] : [];
});
}

View file

@ -122,8 +122,12 @@ function LoadedRun({
const refreshTrace = () => queryClient.resetQueries({ queryKey, exact: true });
const failure = traceQuery.error ? classifyTraceReadFailure(traceQuery.error) : null;
const trace = useMemo(() => {
const [first, ...rest] = traceQuery.data.pages;
return { ...first, spans: [first, ...rest].flatMap((page) => page.spans) };
const pages = traceQuery.data.pages;
return {
...pages[0],
spans: pages.flatMap((page) => page.spans),
next_cursor: pages[pages.length - 1].next_cursor,
};
}, [traceQuery.data]);
const seekingSpan = !switching && selectedSpanMissing(trace, selection.spanId);
const { hasNextPage, isFetching, isError, fetchNextPage } = traceQuery;

View file

@ -20385,6 +20385,23 @@ export interface paths {
patch?: never;
trace?: never;
};
"/v1/logs": {
parameters: {
query?: never;
header?: never;
path?: never;
cookie?: never;
};
get?: never;
put?: never;
/** Ingest Otlp Traces */
post: operations["ingest_otlp_traces_v1_logs_post"];
delete?: never;
options?: never;
head?: never;
patch?: never;
trace?: never;
};
"/v1/mcp/access_groups": {
parameters: {
query?: never;
@ -77422,6 +77439,26 @@ export interface operations {
};
};
};
ingest_otlp_traces_v1_logs_post: {
parameters: {
query?: never;
header?: never;
path?: never;
cookie?: never;
};
requestBody?: never;
responses: {
/** @description Successful Response */
200: {
headers: {
[name: string]: unknown;
};
content: {
"application/json": unknown;
};
};
};
};
get_mcp_access_groups_v1_mcp_access_groups_get: {
parameters: {
query?: never;