use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use serde::Deserialize;
use serde_json::value::RawValue;
use crate::attribution::CodexSessionFile;
use super::files::{self, FileFingerprint};
use super::payload::{IngestEnvelope, KIND_INTERACTED, TranscriptPayload};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AnchorKind {
Started,
Interacted,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubAgentAnchor {
pub kind: AnchorKind,
pub thread_id: String,
pub call_id: String,
pub agent_path: Option<String>,
pub line: String,
}
impl SubAgentAnchor {
#[must_use]
pub fn agent_type(&self) -> Option<&str> {
self.agent_path
.as_deref()
.and_then(|path| path.rsplit('/').next())
.filter(|segment| !segment.is_empty())
}
#[must_use]
pub fn payload_kind(&self) -> Option<&'static str> {
match self.kind {
AnchorKind::Started => None,
AnchorKind::Interacted => Some(KIND_INTERACTED),
}
}
#[must_use]
pub fn dedup_key(&self) -> String {
match self.kind {
AnchorKind::Started => format!("started:{}", self.thread_id),
AnchorKind::Interacted => format!("interacted:{}:{}", self.call_id, self.thread_id),
}
}
}
#[derive(Deserialize)]
struct RolloutRow {
#[serde(rename = "type")]
row_type: String,
payload: Option<serde_json::Value>,
}
#[derive(Deserialize)]
struct SubAgentActivity {
#[serde(rename = "type")]
activity_type: String,
event_id: Option<String>,
agent_thread_id: Option<String>,
agent_path: Option<String>,
kind: Option<String>,
}
const ROLLOUT_KIND_STARTED: &str = "started";
#[must_use]
pub fn parse_subagent_anchors(raw: &[u8]) -> Vec<SubAgentAnchor> {
let mut out: Vec<SubAgentAnchor> = Vec::new();
let mut seen: HashSet<String> = HashSet::new();
for line in raw.split(|&byte| byte == b'\n') {
if !line
.windows(b"sub_agent_activity".len())
.any(|window| window == b"sub_agent_activity")
{
continue;
}
let Ok(text) = std::str::from_utf8(line) else {
continue;
};
let text = text.trim();
let Ok(row) = serde_json::from_str::<RolloutRow>(text) else {
continue;
};
if row.row_type != "event_msg" {
continue;
}
let Some(activity) = row
.payload
.and_then(|payload| serde_json::from_value::<SubAgentActivity>(payload).ok())
else {
continue;
};
if activity.activity_type != "sub_agent_activity" {
continue;
}
let kind = match activity.kind.as_deref() {
Some(ROLLOUT_KIND_STARTED) => AnchorKind::Started,
Some(KIND_INTERACTED) => AnchorKind::Interacted,
_ => continue,
};
let (Some(call_id), Some(thread)) = (activity.event_id, activity.agent_thread_id) else {
continue;
};
if call_id.is_empty() || thread.is_empty() {
continue;
}
let anchor = SubAgentAnchor {
kind,
thread_id: thread,
call_id,
agent_path: activity.agent_path,
line: text.to_owned(),
};
if !seen.insert(anchor.dedup_key()) {
if anchor.kind == AnchorKind::Started {
tracing::warn!(
child_thread_id = %anchor.thread_id,
call_id = %anchor.call_id,
"codex-anchors: duplicate started record for one thread; keeping the first",
);
}
continue;
}
out.push(anchor);
}
out
}
#[must_use]
pub fn build_anchor_payload<'a>(
rollout: &'a CodexSessionFile,
anchor: &'a SubAgentAnchor,
harness_id: &'a str,
records: &'a RawValue,
) -> TranscriptPayload<'a> {
TranscriptPayload {
session: IngestEnvelope {
org_id: "",
auth_subject: "",
harness_id,
harness_session_id: rollout
.root_session_id
.as_deref()
.unwrap_or(&rollout.session_id),
harness_version: rollout.cli_version.as_deref(),
cwd: rollout.cwd.as_deref(),
},
agent_id: Some(&anchor.thread_id),
agent_type: anchor.agent_type(),
description: anchor.agent_path.as_deref(),
tool_use_id: Some(&anchor.call_id),
kind: anchor.payload_kind(),
records,
}
}
#[must_use]
pub fn anchor_records(anchor: &SubAgentAnchor) -> Option<Box<RawValue>> {
RawValue::from_string(files::jsonl_to_records(anchor.line.as_bytes())).ok()
}
#[derive(Debug, Default)]
struct RolloutScan {
scanned: Option<FileFingerprint>,
delivered: HashSet<String>,
}
#[derive(Debug, Default)]
pub struct CodexAnchorScanner {
states: HashMap<PathBuf, RolloutScan>,
}
impl CodexAnchorScanner {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn retain_live<'a, I>(&mut self, live: I)
where
I: IntoIterator<Item = &'a Path>,
{
let live: HashSet<&Path> = live.into_iter().collect();
self.states.retain(|path, _| live.contains(path.as_path()));
}
#[must_use]
pub fn needs_read(&self, rollout: &Path, fingerprint: Option<FileFingerprint>) -> bool {
fingerprint.is_some() && self.states.get(rollout).and_then(|s| s.scanned) != fingerprint
}
#[must_use]
pub fn undelivered(&self, rollout: &Path, raw: &[u8]) -> Vec<SubAgentAnchor> {
let state = self.states.get(rollout);
parse_subagent_anchors(raw)
.into_iter()
.filter(|anchor| {
state.is_none_or(|state| !state.delivered.contains(&anchor.dedup_key()))
})
.collect()
}
pub fn record_delivered(&mut self, rollout: &Path, anchor: &SubAgentAnchor) {
self.states
.entry(rollout.to_path_buf())
.or_default()
.delivered
.insert(anchor.dedup_key());
}
pub fn record_clean_scan(&mut self, rollout: &Path, fingerprint: Option<FileFingerprint>) {
self.states
.entry(rollout.to_path_buf())
.or_default()
.scanned = fingerprint;
}
#[must_use]
pub fn delivered_count(&self, rollout: &Path) -> usize {
self.states
.get(rollout)
.map_or(0, |state| state.delivered.len())
}
}
pub mod fixtures {
pub const STARTED_LINE: &str = r#"{"timestamp":"2026-07-23T04:41:01.858Z","type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_J7B6r7ZdtqkECtSJV8YDQaL7","occurred_at_ms":1784781661858,"agent_thread_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_path":"/root/depth2_cli_child","kind":"started"}}"#;
pub const INTERACTED_LINE: &str = r#"{"timestamp":"2026-07-23T04:41:18.008Z","type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_cqusEjhomv5zKjZ7vodiY7Og","occurred_at_ms":1784781678008,"agent_thread_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_path":"/root/depth2_cli_child","kind":"interacted"}}"#;
pub const ROOT_SESSION_ID: &str = "019f8d46-beb1-7c40-9a1f-2e8b1c0d5a33";
pub const SESSION_ID: &str = "019f8d46-c0de-7000-8000-000000000001";
pub const CWD: &str = "/w/repo";
pub const CLI_VERSION: &str = "0.145.0";
pub const ROLLOUT: &str = concat!(
r#"{"timestamp":"2026-07-23T04:41:00.000Z","type":"session_meta","payload":{"id":"019f8d46-c0de-7000-8000-000000000001"}}"#,
"\n",
r#"{"timestamp":"2026-07-23T04:41:01.858Z","type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_J7B6r7ZdtqkECtSJV8YDQaL7","occurred_at_ms":1784781661858,"agent_thread_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_path":"/root/depth2_cli_child","kind":"started"}}"#,
"\n",
r#"{"timestamp":"2026-07-23T04:41:10.000Z","type":"event_msg","payload":{"type":"agent_message","message":"spawned via sub_agent_activity"}}"#,
"\n",
r#"{"timestamp":"2026-07-23T04:41:18.008Z","type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_cqusEjhomv5zKjZ7vodiY7Og","occurred_at_ms":1784781678008,"agent_thread_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_path":"/root/depth2_cli_child","kind":"interacted"}}"#,
"\n",
r#"{"timestamp":"2026-07-23T04:41:20.000Z","type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_J7B6r7ZdtqkECtSJV8YDQaL7","agent_thread_id":"019f8d46-e663-74e1-940c-f82e34c07618","kind":"finished"}}"#,
"\n",
);
pub const STARTED_BODY: &str = concat!(
r#"{"session":{"org_id":"","auth_subject":"","harness_id":"codex","#,
r#""harness_session_id":"019f8d46-beb1-7c40-9a1f-2e8b1c0d5a33","harness_version":"0.145.0","cwd":"/w/repo"},"#,
r#""agent_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_type":"depth2_cli_child","#,
r#""description":"/root/depth2_cli_child","tool_use_id":"call_J7B6r7ZdtqkECtSJV8YDQaL7","#,
r#""records":[{"timestamp":"2026-07-23T04:41:01.858Z","type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_J7B6r7ZdtqkECtSJV8YDQaL7","occurred_at_ms":1784781661858,"agent_thread_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_path":"/root/depth2_cli_child","kind":"started"}}]}"#,
);
pub const INTERACTED_BODY: &str = concat!(
r#"{"session":{"org_id":"","auth_subject":"","harness_id":"codex","#,
r#""harness_session_id":"019f8d46-beb1-7c40-9a1f-2e8b1c0d5a33","harness_version":"0.145.0","cwd":"/w/repo"},"#,
r#""agent_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_type":"depth2_cli_child","#,
r#""description":"/root/depth2_cli_child","tool_use_id":"call_cqusEjhomv5zKjZ7vodiY7Og","kind":"interacted","#,
r#""records":[{"timestamp":"2026-07-23T04:41:18.008Z","type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_cqusEjhomv5zKjZ7vodiY7Og","occurred_at_ms":1784781678008,"agent_thread_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_path":"/root/depth2_cli_child","kind":"interacted"}}]}"#,
);
pub const BODIES: [&str; 2] = [STARTED_BODY, INTERACTED_BODY];
#[must_use]
pub fn session_file(path: std::path::PathBuf) -> crate::attribution::CodexSessionFile {
crate::attribution::CodexSessionFile {
session_id: SESSION_ID.to_owned(),
root_session_id: Some(ROOT_SESSION_ID.to_owned()),
parent_thread_id: Some(ROOT_SESSION_ID.to_owned()),
subagent_kind: None,
timestamp: time::OffsetDateTime::UNIX_EPOCH,
modified_at: Some(time::OffsetDateTime::UNIX_EPOCH),
cwd: Some(CWD.to_owned()),
originator: Some("codex_exec".to_owned()),
cli_version: Some(CLI_VERSION.to_owned()),
source: Some("exec".to_owned()),
thread_source: Some("subagent".to_owned()),
model_provider: None,
path,
}
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::fixtures::{INTERACTED_LINE, STARTED_LINE};
use super::*;
use tapes_capture::envelope::HARNESS_ID_CODEX;
fn rollout_at(path: &Path) -> CodexSessionFile {
fixtures::session_file(path.to_path_buf())
}
fn body(harness_id: &str, rollout: &CodexSessionFile, anchor: &SubAgentAnchor) -> String {
let records = anchor_records(anchor).unwrap();
serde_json::to_string(&build_anchor_payload(rollout, anchor, harness_id, &records)).unwrap()
}
#[test]
fn parse_extracts_started_and_interacted_records_verbatim() {
let anchors = parse_subagent_anchors(fixtures::ROLLOUT.as_bytes());
assert_eq!(anchors.len(), 2, "the ignorable rows must stay ignored");
let started = &anchors[0];
assert_eq!(started.kind, AnchorKind::Started);
assert_eq!(started.thread_id, "019f8d46-e663-74e1-940c-f82e34c07618");
assert_eq!(started.call_id, "call_J7B6r7ZdtqkECtSJV8YDQaL7");
assert_eq!(
started.agent_path.as_deref(),
Some("/root/depth2_cli_child")
);
assert_eq!(started.agent_type(), Some("depth2_cli_child"));
assert_eq!(
started.line, STARTED_LINE,
"the rollout line must survive verbatim — the server dedups on its hash",
);
let interacted = &anchors[1];
assert_eq!(interacted.kind, AnchorKind::Interacted);
assert_eq!(
interacted.thread_id, "019f8d46-e663-74e1-940c-f82e34c07618",
"interacted rows carry the TARGET thread",
);
assert_eq!(interacted.call_id, "call_cqusEjhomv5zKjZ7vodiY7Og");
assert_eq!(interacted.line, INTERACTED_LINE);
}
#[test]
fn parse_tolerates_a_truncated_final_line() {
let raw = format!("{STARTED_LINE}\n{{\"type\":\"event_msg\",\"payl");
assert_eq!(parse_subagent_anchors(raw.as_bytes()).len(), 1);
}
#[test]
fn parse_keeps_first_started_record_per_child() {
let dup = STARTED_LINE.replace("call_J7B6r7ZdtqkECtSJV8YDQaL7", "call_second");
let raw = format!("{STARTED_LINE}\n{dup}\n");
let anchors = parse_subagent_anchors(raw.as_bytes());
assert_eq!(anchors.len(), 1);
assert_eq!(anchors[0].call_id, "call_J7B6r7ZdtqkECtSJV8YDQaL7");
}
#[test]
fn parse_keeps_one_interacted_anchor_per_triggering_call() {
let second_send = INTERACTED_LINE.replace("call_cqusEjhomv5zKjZ7vodiY7Og", "call_2nd");
let raw = format!("{INTERACTED_LINE}\n{second_send}\n{INTERACTED_LINE}\n");
let anchors = parse_subagent_anchors(raw.as_bytes());
assert_eq!(anchors.len(), 2);
assert_eq!(anchors[0].call_id, "call_cqusEjhomv5zKjZ7vodiY7Og");
assert_eq!(anchors[1].call_id, "call_2nd");
let raw = format!("{STARTED_LINE}\n{INTERACTED_LINE}\n");
assert_eq!(parse_subagent_anchors(raw.as_bytes()).len(), 2);
}
#[test]
fn parse_requires_call_and_thread_ids() {
let missing_thread = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_x","kind":"started"}}"#;
let missing_call = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","agent_thread_id":"child","kind":"started"}}"#;
let missing_thread_interacted = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_x","kind":"interacted"}}"#;
let missing_call_interacted = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","agent_thread_id":"child","kind":"interacted"}}"#;
let raw = format!(
"{missing_thread}\n{missing_call}\n{missing_thread_interacted}\n{missing_call_interacted}\n"
);
assert!(parse_subagent_anchors(raw.as_bytes()).is_empty());
}
#[test]
fn the_fixture_bodies_are_what_the_derivation_produces() {
let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
let anchors = parse_subagent_anchors(fixtures::ROLLOUT.as_bytes());
let bodies: Vec<String> = anchors
.iter()
.map(|anchor| body(HARNESS_ID_CODEX, &rollout, anchor))
.collect();
assert_eq!(bodies, fixtures::BODIES.to_vec());
}
#[test]
fn anchor_rows_key_to_the_root_session_not_the_spawning_thread() {
let rollout = rollout_at(Path::new("/tmp/launcher.jsonl"));
assert_ne!(rollout.session_id, fixtures::ROOT_SESSION_ID);
let anchor = &parse_subagent_anchors(fixtures::ROLLOUT.as_bytes())[0];
let got: serde_json::Value =
serde_json::from_str(&body(HARNESS_ID_CODEX, &rollout, anchor)).unwrap();
assert_eq!(
got["session"]["harness_session_id"],
fixtures::ROOT_SESSION_ID
);
let mut root = rollout;
root.root_session_id = None;
let got: serde_json::Value =
serde_json::from_str(&body(HARNESS_ID_CODEX, &root, anchor)).unwrap();
assert_eq!(got["session"]["harness_session_id"], fixtures::SESSION_ID);
}
#[test]
fn a_started_row_carries_no_kind_field_at_all() {
let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
let anchors = parse_subagent_anchors(fixtures::ROLLOUT.as_bytes());
let started: serde_json::Value =
serde_json::from_str(&body(HARNESS_ID_CODEX, &rollout, &anchors[0])).unwrap();
assert!(started.get("kind").is_none());
let interacted: serde_json::Value =
serde_json::from_str(&body(HARNESS_ID_CODEX, &rollout, &anchors[1])).unwrap();
assert_eq!(interacted["kind"], KIND_INTERACTED);
}
#[test]
fn the_row_names_the_harness_the_caller_declares() {
let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
let anchor = &parse_subagent_anchors(fixtures::ROLLOUT.as_bytes())[0];
let got = body(
tapes_capture::envelope::HARNESS_ID_CODEX_APP,
&rollout,
anchor,
);
assert!(got.contains(r#""harness_id":"codex-app""#), "got: {got}");
}
#[test]
fn an_anchor_with_no_agent_path_omits_the_optional_slots() {
let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
let anchor = SubAgentAnchor {
kind: AnchorKind::Started,
thread_id: "child".to_owned(),
call_id: "call_x".to_owned(),
agent_path: None,
line: "{}".to_owned(),
};
let got = body(HARNESS_ID_CODEX, &rollout, &anchor);
assert!(!got.contains("agent_type"), "got: {got}");
assert!(!got.contains("description"), "got: {got}");
assert!(got.contains(r#""agent_id":"child""#), "got: {got}");
}
fn write_rollout(dir: &Path, body: &str) -> PathBuf {
let path = dir.join("rollout.jsonl");
std::fs::write(&path, body).unwrap();
path
}
#[test]
fn a_fresh_rollout_offers_every_anchor_once() {
let dir = tempfile::tempdir().unwrap();
let path = write_rollout(dir.path(), fixtures::ROLLOUT);
let mut scanner = CodexAnchorScanner::new();
let fingerprint = files::fingerprint(&path);
assert!(scanner.needs_read(&path, fingerprint));
let raw = std::fs::read(&path).unwrap();
let anchors = scanner.undelivered(&path, &raw);
assert_eq!(anchors.len(), 2);
for anchor in &anchors {
scanner.record_delivered(&path, anchor);
}
scanner.record_clean_scan(&path, fingerprint);
assert_eq!(scanner.delivered_count(&path), 2);
assert!(
!scanner.needs_read(&path, files::fingerprint(&path)),
"an unchanged append-only file cannot hide a new anchor",
);
assert!(scanner.undelivered(&path, &raw).is_empty());
}
#[test]
fn a_grown_rollout_offers_only_what_is_new() {
let dir = tempfile::tempdir().unwrap();
let path = write_rollout(dir.path(), fixtures::ROLLOUT);
let mut scanner = CodexAnchorScanner::new();
let raw = std::fs::read(&path).unwrap();
for anchor in scanner.undelivered(&path, &raw) {
scanner.record_delivered(&path, &anchor);
}
scanner.record_clean_scan(&path, files::fingerprint(&path));
let second = STARTED_LINE
.replace("call_J7B6r7ZdtqkECtSJV8YDQaL7", "call_second")
.replace(
"019f8d46-e663-74e1-940c-f82e34c07618",
"019f8d47-0473-7743-a1ed-9e4c0ae92ad8",
);
std::fs::write(&path, format!("{}{second}\n", fixtures::ROLLOUT)).unwrap();
assert!(scanner.needs_read(&path, files::fingerprint(&path)));
let raw = std::fs::read(&path).unwrap();
let pending = scanner.undelivered(&path, &raw);
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].call_id, "call_second");
}
#[test]
fn an_undelivered_anchor_is_offered_again_next_read() {
let dir = tempfile::tempdir().unwrap();
let path = write_rollout(dir.path(), fixtures::ROLLOUT);
let scanner = CodexAnchorScanner::new();
let raw = std::fs::read(&path).unwrap();
assert_eq!(scanner.undelivered(&path, &raw).len(), 2);
assert_eq!(scanner.undelivered(&path, &raw).len(), 2);
assert!(
scanner.needs_read(&path, files::fingerprint(&path)),
"a scan that was never marked clean must re-run",
);
}
#[test]
fn a_vanished_rollout_is_skipped_rather_than_read() {
let scanner = CodexAnchorScanner::new();
assert!(!scanner.needs_read(Path::new("/nonexistent/rollout.jsonl"), None));
}
#[test]
fn retain_live_drops_rollouts_that_aged_out_of_the_snapshot() {
let mut scanner = CodexAnchorScanner::new();
let kept = PathBuf::from("/tmp/kept.jsonl");
let gone = PathBuf::from("/tmp/gone.jsonl");
let anchor = &parse_subagent_anchors(fixtures::ROLLOUT.as_bytes())[0];
scanner.record_delivered(&kept, anchor);
scanner.record_delivered(&gone, anchor);
scanner.retain_live([kept.as_path()]);
assert_eq!(scanner.delivered_count(&kept), 1);
assert_eq!(scanner.delivered_count(&gone), 0);
}
}