Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 35 additions & 12 deletions crates/tracedecay-capture/src/claude/canonical.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,21 +42,37 @@ pub struct ClaudeSpawnParent<'a> {
pub fn normalize(
native: &Value,
session_id: &str,
transcript_path: Option<&str>,
stable_record_id: ObservationId,
range: ObservationSourceRangeV1,
) -> Result<CanonicalObservationEnvelopeV1, ObservationRecordParseErrorV1> {
normalize_spawned(native, session_id, None, stable_record_id, range)
normalize_spawned(
native,
session_id,
None,
transcript_path,
stable_record_id,
range,
)
}

/// [`normalize`] for a record of a subagent transcript spawned by `parent`.
pub fn normalize_spawned(
native: &Value,
session_id: &str,
parent: Option<ClaudeSpawnParent<'_>>,
transcript_path: Option<&str>,
stable_record_id: ObservationId,
range: ObservationSourceRangeV1,
) -> Result<CanonicalObservationEnvelopeV1, ObservationRecordParseErrorV1> {
normalize_record(native, session_id, parent, stable_record_id, range)
normalize_record(
native,
session_id,
parent,
transcript_path,
stable_record_id,
range,
)
}

/// One source-record canonicalization, not a per-block walk.
Expand All @@ -65,6 +81,7 @@ fn normalize_record(
native: &Value,
session_id: &str,
parent: Option<ClaudeSpawnParent<'_>>,
transcript_path: Option<&str>,
stable_record_id: ObservationId,
range: ObservationSourceRangeV1,
) -> Result<CanonicalObservationEnvelopeV1, ObservationRecordParseErrorV1> {
Expand All @@ -80,7 +97,7 @@ fn normalize_record(
.and_then(parse_rfc3339_timestamp);
let mut facts = Vec::new();

append_session_location_fact(&mut facts, native);
append_session_location_fact(&mut facts, native, transcript_path);
if matches!(record_kind, "user" | "assistant") {
let message = native.get("message").unwrap_or(native);
// Message facts carry only provider-authored visible text. Thinking,
Expand Down Expand Up @@ -261,19 +278,23 @@ fn tool_result_only_message_content(message: &Value) -> Option<Value> {
(!parts.is_empty()).then(|| Value::String(parts.join("\n\n")))
}

fn append_session_location_fact(facts: &mut Vec<CanonicalObservationFactV1>, native: &Value) {
let Some(cwd) = native
fn append_session_location_fact(
facts: &mut Vec<CanonicalObservationFactV1>,
native: &Value,
transcript_path: Option<&str>,
) {
let cwd = native
.get("cwd")
.and_then(Value::as_str)
.filter(|cwd| !cwd.is_empty())
.map(str::to_owned)
else {
.map(str::to_owned);
if cwd.is_none() && transcript_path.is_none() {
return;
};
}
facts.push(CanonicalObservationFactV1::Session {
project_path: Some(cwd.clone()),
location_path: Some(cwd),
transcript_path: None,
project_path: cwd.clone(),
location_path: cwd,
transcript_path: transcript_path.map(str::to_owned),
title: None,
started_at: None,
ended_at: None,
Expand Down Expand Up @@ -721,7 +742,7 @@ mod provider_usage_tests {
ObservationSourceRangeV1::new(20, 30).unwrap(),
),
] {
let envelope = normalize(native, "session.fixture", stable_id, range).unwrap();
let envelope = normalize(native, "session.fixture", None, stable_id, range).unwrap();
assert_eq!(
envelope.relations().message_id().map(ObservationId::as_str),
Some("message.shared")
Expand All @@ -747,6 +768,7 @@ mod provider_usage_tests {
let envelope = normalize(
&native,
"session.fixture",
None,
ObservationId::new("message.fixture").unwrap(),
ObservationSourceRangeV1::new(10, 20).unwrap(),
)
Expand Down Expand Up @@ -781,6 +803,7 @@ mod provider_usage_tests {
let envelope = normalize(
&native,
"session.fixture",
None,
ObservationId::new("message.user-fixture").unwrap(),
ObservationSourceRangeV1::new(20, 30).unwrap(),
)
Expand Down
84 changes: 58 additions & 26 deletions crates/tracedecay-capture/src/codex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -125,27 +125,31 @@ pub fn normalize_codex_observation(
)
}

/// The rollout context in effect for one record: where the session projects,
/// which file it was read from, and the model its turn runs under.
#[derive(Clone, Copy)]
pub struct CodexObservationLocation<'a> {
pub struct CodexObservationContext<'a> {
pub project_path: Option<&'a Path>,
pub location_path: Option<&'a Path>,
pub transcript_path: Option<&'a Path>,
pub model: Option<&'a str>,
}

pub fn normalize_codex_observation_with_location(
pub fn normalize_codex_observation_with_context(
native: &Value,
session_id: &str,
native_thread_id: Option<&str>,
stable_record_id: ObservationId,
range: tracedecay_domain::ObservationSourceRangeV1,
location: CodexObservationLocation<'_>,
context: CodexObservationContext<'_>,
) -> Result<CanonicalObservationEnvelopeV1, ObservationRecordParseErrorV1> {
normalize_codex_observation_inner(
native,
session_id,
native_thread_id,
stable_record_id,
range,
Some(location),
Some(context),
)
}

Expand All @@ -158,15 +162,15 @@ fn normalize_codex_observation_inner(
native_thread_id: Option<&str>,
stable_record_id: ObservationId,
range: tracedecay_domain::ObservationSourceRangeV1,
location: Option<CodexObservationLocation<'_>>,
context: Option<CodexObservationContext<'_>>,
) -> Result<CanonicalObservationEnvelopeV1, ObservationRecordParseErrorV1> {
normalize_codex_record(
native,
session_id,
native_thread_id,
stable_record_id,
range,
location,
context,
)
}

Expand All @@ -178,7 +182,7 @@ fn normalize_codex_record(
native_thread_id: Option<&str>,
stable_record_id: ObservationId,
range: tracedecay_domain::ObservationSourceRangeV1,
location: Option<CodexObservationLocation<'_>>,
context: Option<CodexObservationContext<'_>>,
) -> Result<CanonicalObservationEnvelopeV1, ObservationRecordParseErrorV1> {
let native_kind = native
.get("type")
Expand Down Expand Up @@ -218,15 +222,17 @@ fn normalize_codex_record(
relations = append_codex_session_meta_agent_relations(relations, payload, native_thread_id);
}
let mut facts = Vec::new();
if let Some(location) = location {
if let Some(context) = context {
facts.push(CanonicalObservationFactV1::Session {
project_path: location
project_path: context
.project_path
.map(|path| path.to_string_lossy().into_owned()),
location_path: location
location_path: context
.location_path
.map(|path| path.to_string_lossy().into_owned()),
transcript_path: None,
transcript_path: context
.transcript_path
.map(|path| path.to_string_lossy().into_owned()),
title: None,
started_at: None,
ended_at: None,
Expand Down Expand Up @@ -256,9 +262,19 @@ fn normalize_codex_record(
native_kind: "turn_context".to_string(),
state: CanonicalUnknownStateV1::Unsupported,
}),
"event_msg" => append_codex_event_facts(payload, timestamp, &mut facts),
"event_msg" => append_codex_event_facts(
payload,
timestamp,
context.and_then(|context| context.model),
&mut facts,
),
"response_item" => {
append_codex_response_item_facts(payload, timestamp, &mut facts);
append_codex_response_item_facts(
payload,
timestamp,
context.and_then(|context| context.model),
&mut facts,
);
}
"compacted" => {
facts.push(CanonicalObservationFactV1::Compaction {
Expand Down Expand Up @@ -359,6 +375,7 @@ fn append_codex_session_meta_agent_relations(
fn append_codex_event_facts(
payload: &Value,
timestamp: Option<i64>,
context_model: Option<&str>,
facts: &mut Vec<CanonicalObservationFactV1>,
) {
match payload.get("type").and_then(Value::as_str) {
Expand All @@ -375,7 +392,8 @@ fn append_codex_event_facts(
model: payload
.get("model")
.and_then(Value::as_str)
.map(str::to_string),
.map(str::to_string)
.or_else(|| context_model.map(str::to_string)),
timestamp,
});
}
Expand Down Expand Up @@ -419,7 +437,8 @@ fn append_codex_event_facts(
.get("model")
.or_else(|| payload.get("model"))
.and_then(Value::as_str)
.map(str::to_string),
.map(str::to_string)
.or_else(|| context_model.map(str::to_string)),
timestamp,
});
}
Expand Down Expand Up @@ -809,6 +828,7 @@ fn append_codex_update_plan_lifecycle_fact(
fn append_codex_response_item_facts(
payload: &Value,
timestamp: Option<i64>,
context_model: Option<&str>,
facts: &mut Vec<CanonicalObservationFactV1>,
) {
let Some(item_kind) = payload.get("type").and_then(Value::as_str) else {
Expand Down Expand Up @@ -842,7 +862,8 @@ fn append_codex_response_item_facts(
model: payload
.get("model")
.and_then(Value::as_str)
.map(str::to_string),
.map(str::to_string)
.or_else(|| context_model.map(str::to_string)),
timestamp,
});
}
Expand Down Expand Up @@ -1028,7 +1049,7 @@ mod provider_usage_tests {
ProviderUsageModelV1, ProviderUsageScopeV1,
};

use super::{CodexObservationLocation, normalize_codex_observation_with_location};
use super::{CodexObservationContext, normalize_codex_observation_with_context};

#[test]
fn canonical_payload_uses_only_bound_source_location() {
Expand All @@ -1045,21 +1066,26 @@ mod provider_usage_tests {
}
}
});
let canonical = normalize_codex_observation_with_location(
let canonical = normalize_codex_observation_with_context(
&native,
"session.fixture",
Some("thread.fixture"),
ObservationId::new("record.fixture").unwrap(),
ObservationSourceRangeV1::new(10, 20).unwrap(),
CodexObservationLocation {
CodexObservationContext {
project_path: Some(Path::new("/redacted/project")),
location_path: Some(Path::new("/redacted/project")),
transcript_path: Some(Path::new("/redacted/project/rollout.jsonl")),
model: Some("gpt-5.5"),
},
)
.unwrap();

let encoded = serde_json::to_value(canonical).unwrap();
assert!(encoded["facts"][0].get("transcript_path").is_none());
assert_eq!(
encoded["facts"][0]["transcript_path"].as_str(),
Some("/redacted/project/rollout.jsonl")
);
assert!(encoded["relations"].get("turn_id").is_none());
assert_eq!(encoded["facts"][1]["model"]["state"], "unknown");
}
Expand Down Expand Up @@ -1089,15 +1115,17 @@ mod provider_usage_tests {
}
}
});
let envelope = normalize_codex_observation_with_location(
let envelope = normalize_codex_observation_with_context(
&native,
"session.fixture",
Some("thread.fixture"),
ObservationId::new("record.fixture").unwrap(),
ObservationSourceRangeV1::new(10, 20).unwrap(),
CodexObservationLocation {
CodexObservationContext {
project_path: None,
location_path: None,
transcript_path: None,
model: None,
},
)
.unwrap();
Expand Down Expand Up @@ -1151,15 +1179,17 @@ mod provider_usage_tests {
}
}
});
let envelope = normalize_codex_observation_with_location(
let envelope = normalize_codex_observation_with_context(
&native,
"session.fixture",
Some("thread.fixture"),
ObservationId::new("record.fixture").unwrap(),
ObservationSourceRangeV1::new(10, 20).unwrap(),
CodexObservationLocation {
CodexObservationContext {
project_path: None,
location_path: None,
transcript_path: None,
model: None,
},
)
.unwrap();
Expand Down Expand Up @@ -1211,15 +1241,17 @@ mod provider_usage_tests {
}
}
});
let envelope = normalize_codex_observation_with_location(
let envelope = normalize_codex_observation_with_context(
&native,
"session.fixture",
Some("thread-1"),
ObservationId::new("record.fixture").unwrap(),
ObservationSourceRangeV1::new(10, 20).unwrap(),
CodexObservationLocation {
CodexObservationContext {
project_path: None,
location_path: None,
transcript_path: None,
model: Some("gpt-ctx"),
},
)
.unwrap();
Expand Down
Loading
Loading