1use std::collections::{HashMap, HashSet};
104use std::path::{Path, PathBuf};
105
106use serde::Deserialize;
107use serde_json::value::RawValue;
108
109use crate::attribution::CodexSessionFile;
110
111use super::files::{self, FileFingerprint};
112use super::payload::{IngestEnvelope, KIND_INTERACTED, TranscriptPayload};
113
114#[derive(Debug, Clone, Copy, PartialEq, Eq)]
116pub enum AnchorKind {
117 Started,
121 Interacted,
126}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
134pub struct SubAgentAnchor {
135 pub kind: AnchorKind,
137 pub thread_id: String,
140 pub call_id: String,
143 pub agent_path: Option<String>,
146 pub line: String,
148}
149
150impl SubAgentAnchor {
151 #[must_use]
154 pub fn agent_type(&self) -> Option<&str> {
155 self.agent_path
156 .as_deref()
157 .and_then(|path| path.rsplit('/').next())
158 .filter(|segment| !segment.is_empty())
159 }
160
161 #[must_use]
168 pub fn payload_kind(&self) -> Option<&'static str> {
169 match self.kind {
170 AnchorKind::Started => None,
171 AnchorKind::Interacted => Some(KIND_INTERACTED),
172 }
173 }
174
175 #[must_use]
182 pub fn dedup_key(&self) -> String {
183 match self.kind {
184 AnchorKind::Started => format!("started:{}", self.thread_id),
185 AnchorKind::Interacted => format!("interacted:{}:{}", self.call_id, self.thread_id),
186 }
187 }
188}
189
190#[derive(Deserialize)]
191struct RolloutRow {
192 #[serde(rename = "type")]
193 row_type: String,
194 payload: Option<serde_json::Value>,
195}
196
197#[derive(Deserialize)]
198struct SubAgentActivity {
199 #[serde(rename = "type")]
200 activity_type: String,
201 event_id: Option<String>,
202 agent_thread_id: Option<String>,
203 agent_path: Option<String>,
204 kind: Option<String>,
205}
206
207const ROLLOUT_KIND_STARTED: &str = "started";
209
210#[must_use]
224pub fn parse_subagent_anchors(raw: &[u8]) -> Vec<SubAgentAnchor> {
225 let mut out: Vec<SubAgentAnchor> = Vec::new();
226 let mut seen: HashSet<String> = HashSet::new();
227 for line in raw.split(|&byte| byte == b'\n') {
228 if !line
230 .windows(b"sub_agent_activity".len())
231 .any(|window| window == b"sub_agent_activity")
232 {
233 continue;
234 }
235 let Ok(text) = std::str::from_utf8(line) else {
236 continue;
237 };
238 let text = text.trim();
239 let Ok(row) = serde_json::from_str::<RolloutRow>(text) else {
240 continue;
241 };
242 if row.row_type != "event_msg" {
243 continue;
244 }
245 let Some(activity) = row
246 .payload
247 .and_then(|payload| serde_json::from_value::<SubAgentActivity>(payload).ok())
248 else {
249 continue;
250 };
251 if activity.activity_type != "sub_agent_activity" {
252 continue;
253 }
254 let kind = match activity.kind.as_deref() {
255 Some(ROLLOUT_KIND_STARTED) => AnchorKind::Started,
256 Some(KIND_INTERACTED) => AnchorKind::Interacted,
257 _ => continue,
260 };
261 let (Some(call_id), Some(thread)) = (activity.event_id, activity.agent_thread_id) else {
262 continue;
263 };
264 if call_id.is_empty() || thread.is_empty() {
265 continue;
266 }
267 let anchor = SubAgentAnchor {
268 kind,
269 thread_id: thread,
270 call_id,
271 agent_path: activity.agent_path,
272 line: text.to_owned(),
273 };
274 if !seen.insert(anchor.dedup_key()) {
275 if anchor.kind == AnchorKind::Started {
276 tracing::warn!(
277 child_thread_id = %anchor.thread_id,
278 call_id = %anchor.call_id,
279 "codex-anchors: duplicate started record for one thread; keeping the first",
280 );
281 }
282 continue;
283 }
284 out.push(anchor);
285 }
286 out
287}
288
289#[must_use]
308pub fn build_anchor_payload<'a>(
309 rollout: &'a CodexSessionFile,
310 anchor: &'a SubAgentAnchor,
311 harness_id: &'a str,
312 records: &'a RawValue,
313) -> TranscriptPayload<'a> {
314 TranscriptPayload {
315 session: IngestEnvelope {
316 org_id: "",
317 auth_subject: "",
318 harness_id,
319 harness_session_id: rollout
323 .root_session_id
324 .as_deref()
325 .unwrap_or(&rollout.session_id),
326 harness_version: rollout.cli_version.as_deref(),
327 cwd: rollout.cwd.as_deref(),
328 },
329 agent_id: Some(&anchor.thread_id),
330 agent_type: anchor.agent_type(),
331 description: anchor.agent_path.as_deref(),
332 tool_use_id: Some(&anchor.call_id),
333 kind: anchor.payload_kind(),
334 records,
335 }
336}
337
338#[must_use]
345pub fn anchor_records(anchor: &SubAgentAnchor) -> Option<Box<RawValue>> {
346 RawValue::from_string(files::jsonl_to_records(anchor.line.as_bytes())).ok()
347}
348
349#[derive(Debug, Default)]
351struct RolloutScan {
352 scanned: Option<FileFingerprint>,
355 delivered: HashSet<String>,
357}
358
359#[derive(Debug, Default)]
375pub struct CodexAnchorScanner {
376 states: HashMap<PathBuf, RolloutScan>,
377}
378
379impl CodexAnchorScanner {
380 #[must_use]
382 pub fn new() -> Self {
383 Self::default()
384 }
385
386 pub fn retain_live<'a, I>(&mut self, live: I)
392 where
393 I: IntoIterator<Item = &'a Path>,
394 {
395 let live: HashSet<&Path> = live.into_iter().collect();
396 self.states.retain(|path, _| live.contains(path.as_path()));
397 }
398
399 #[must_use]
406 pub fn needs_read(&self, rollout: &Path, fingerprint: Option<FileFingerprint>) -> bool {
407 fingerprint.is_some() && self.states.get(rollout).and_then(|s| s.scanned) != fingerprint
408 }
409
410 #[must_use]
413 pub fn undelivered(&self, rollout: &Path, raw: &[u8]) -> Vec<SubAgentAnchor> {
414 let state = self.states.get(rollout);
415 parse_subagent_anchors(raw)
416 .into_iter()
417 .filter(|anchor| {
418 state.is_none_or(|state| !state.delivered.contains(&anchor.dedup_key()))
419 })
420 .collect()
421 }
422
423 pub fn record_delivered(&mut self, rollout: &Path, anchor: &SubAgentAnchor) {
426 self.states
427 .entry(rollout.to_path_buf())
428 .or_default()
429 .delivered
430 .insert(anchor.dedup_key());
431 }
432
433 pub fn record_clean_scan(&mut self, rollout: &Path, fingerprint: Option<FileFingerprint>) {
439 self.states
440 .entry(rollout.to_path_buf())
441 .or_default()
442 .scanned = fingerprint;
443 }
444
445 #[must_use]
447 pub fn delivered_count(&self, rollout: &Path) -> usize {
448 self.states
449 .get(rollout)
450 .map_or(0, |state| state.delivered.len())
451 }
452}
453
454pub mod fixtures {
462 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"}}"#;
466
467 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"}}"#;
472
473 pub const ROOT_SESSION_ID: &str = "019f8d46-beb1-7c40-9a1f-2e8b1c0d5a33";
475
476 pub const SESSION_ID: &str = "019f8d46-c0de-7000-8000-000000000001";
480
481 pub const CWD: &str = "/w/repo";
483
484 pub const CLI_VERSION: &str = "0.145.0";
486
487 pub const ROLLOUT: &str = concat!(
491 r#"{"timestamp":"2026-07-23T04:41:00.000Z","type":"session_meta","payload":{"id":"019f8d46-c0de-7000-8000-000000000001"}}"#,
492 "\n",
493 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"}}"#,
494 "\n",
495 r#"{"timestamp":"2026-07-23T04:41:10.000Z","type":"event_msg","payload":{"type":"agent_message","message":"spawned via sub_agent_activity"}}"#,
496 "\n",
497 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"}}"#,
498 "\n",
499 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"}}"#,
500 "\n",
501 );
502
503 pub const STARTED_BODY: &str = concat!(
513 r#"{"session":{"org_id":"","auth_subject":"","harness_id":"codex","#,
514 r#""harness_session_id":"019f8d46-beb1-7c40-9a1f-2e8b1c0d5a33","harness_version":"0.145.0","cwd":"/w/repo"},"#,
515 r#""agent_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_type":"depth2_cli_child","#,
516 r#""description":"/root/depth2_cli_child","tool_use_id":"call_J7B6r7ZdtqkECtSJV8YDQaL7","#,
517 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"}}]}"#,
518 );
519
520 pub const INTERACTED_BODY: &str = concat!(
524 r#"{"session":{"org_id":"","auth_subject":"","harness_id":"codex","#,
525 r#""harness_session_id":"019f8d46-beb1-7c40-9a1f-2e8b1c0d5a33","harness_version":"0.145.0","cwd":"/w/repo"},"#,
526 r#""agent_id":"019f8d46-e663-74e1-940c-f82e34c07618","agent_type":"depth2_cli_child","#,
527 r#""description":"/root/depth2_cli_child","tool_use_id":"call_cqusEjhomv5zKjZ7vodiY7Og","kind":"interacted","#,
528 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"}}]}"#,
529 );
530
531 pub const BODIES: [&str; 2] = [STARTED_BODY, INTERACTED_BODY];
534
535 #[must_use]
541 pub fn session_file(path: std::path::PathBuf) -> crate::attribution::CodexSessionFile {
542 crate::attribution::CodexSessionFile {
543 session_id: SESSION_ID.to_owned(),
544 root_session_id: Some(ROOT_SESSION_ID.to_owned()),
545 parent_thread_id: Some(ROOT_SESSION_ID.to_owned()),
546 subagent_kind: None,
547 timestamp: time::OffsetDateTime::UNIX_EPOCH,
548 modified_at: Some(time::OffsetDateTime::UNIX_EPOCH),
549 cwd: Some(CWD.to_owned()),
550 originator: Some("codex_exec".to_owned()),
551 cli_version: Some(CLI_VERSION.to_owned()),
552 source: Some("exec".to_owned()),
553 thread_source: Some("subagent".to_owned()),
554 model_provider: None,
555 path,
556 }
557 }
558}
559
560#[cfg(test)]
561#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
562mod tests {
563 use super::fixtures::{INTERACTED_LINE, STARTED_LINE};
564 use super::*;
565 use tapes_capture::envelope::HARNESS_ID_CODEX;
566
567 fn rollout_at(path: &Path) -> CodexSessionFile {
568 fixtures::session_file(path.to_path_buf())
569 }
570
571 fn body(harness_id: &str, rollout: &CodexSessionFile, anchor: &SubAgentAnchor) -> String {
572 let records = anchor_records(anchor).unwrap();
573 serde_json::to_string(&build_anchor_payload(rollout, anchor, harness_id, &records)).unwrap()
574 }
575
576 #[test]
579 fn parse_extracts_started_and_interacted_records_verbatim() {
580 let anchors = parse_subagent_anchors(fixtures::ROLLOUT.as_bytes());
581 assert_eq!(anchors.len(), 2, "the ignorable rows must stay ignored");
582
583 let started = &anchors[0];
584 assert_eq!(started.kind, AnchorKind::Started);
585 assert_eq!(started.thread_id, "019f8d46-e663-74e1-940c-f82e34c07618");
586 assert_eq!(started.call_id, "call_J7B6r7ZdtqkECtSJV8YDQaL7");
587 assert_eq!(
588 started.agent_path.as_deref(),
589 Some("/root/depth2_cli_child")
590 );
591 assert_eq!(started.agent_type(), Some("depth2_cli_child"));
592 assert_eq!(
593 started.line, STARTED_LINE,
594 "the rollout line must survive verbatim — the server dedups on its hash",
595 );
596
597 let interacted = &anchors[1];
598 assert_eq!(interacted.kind, AnchorKind::Interacted);
599 assert_eq!(
600 interacted.thread_id, "019f8d46-e663-74e1-940c-f82e34c07618",
601 "interacted rows carry the TARGET thread",
602 );
603 assert_eq!(interacted.call_id, "call_cqusEjhomv5zKjZ7vodiY7Og");
604 assert_eq!(interacted.line, INTERACTED_LINE);
605 }
606
607 #[test]
608 fn parse_tolerates_a_truncated_final_line() {
609 let raw = format!("{STARTED_LINE}\n{{\"type\":\"event_msg\",\"payl");
611 assert_eq!(parse_subagent_anchors(raw.as_bytes()).len(), 1);
612 }
613
614 #[test]
615 fn parse_keeps_first_started_record_per_child() {
616 let dup = STARTED_LINE.replace("call_J7B6r7ZdtqkECtSJV8YDQaL7", "call_second");
617 let raw = format!("{STARTED_LINE}\n{dup}\n");
618 let anchors = parse_subagent_anchors(raw.as_bytes());
619 assert_eq!(anchors.len(), 1);
620 assert_eq!(anchors[0].call_id, "call_J7B6r7ZdtqkECtSJV8YDQaL7");
621 }
622
623 #[test]
624 fn parse_keeps_one_interacted_anchor_per_triggering_call() {
625 let second_send = INTERACTED_LINE.replace("call_cqusEjhomv5zKjZ7vodiY7Og", "call_2nd");
628 let raw = format!("{INTERACTED_LINE}\n{second_send}\n{INTERACTED_LINE}\n");
629 let anchors = parse_subagent_anchors(raw.as_bytes());
630 assert_eq!(anchors.len(), 2);
631 assert_eq!(anchors[0].call_id, "call_cqusEjhomv5zKjZ7vodiY7Og");
632 assert_eq!(anchors[1].call_id, "call_2nd");
633
634 let raw = format!("{STARTED_LINE}\n{INTERACTED_LINE}\n");
637 assert_eq!(parse_subagent_anchors(raw.as_bytes()).len(), 2);
638 }
639
640 #[test]
641 fn parse_requires_call_and_thread_ids() {
642 let missing_thread = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_x","kind":"started"}}"#;
643 let missing_call = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","agent_thread_id":"child","kind":"started"}}"#;
644 let missing_thread_interacted = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","event_id":"call_x","kind":"interacted"}}"#;
645 let missing_call_interacted = r#"{"type":"event_msg","payload":{"type":"sub_agent_activity","agent_thread_id":"child","kind":"interacted"}}"#;
646 let raw = format!(
647 "{missing_thread}\n{missing_call}\n{missing_thread_interacted}\n{missing_call_interacted}\n"
648 );
649 assert!(parse_subagent_anchors(raw.as_bytes()).is_empty());
650 }
651
652 #[test]
657 fn the_fixture_bodies_are_what_the_derivation_produces() {
658 let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
659 let anchors = parse_subagent_anchors(fixtures::ROLLOUT.as_bytes());
660 let bodies: Vec<String> = anchors
661 .iter()
662 .map(|anchor| body(HARNESS_ID_CODEX, &rollout, anchor))
663 .collect();
664 assert_eq!(bodies, fixtures::BODIES.to_vec());
665 }
666
667 #[test]
668 fn anchor_rows_key_to_the_root_session_not_the_spawning_thread() {
669 let rollout = rollout_at(Path::new("/tmp/launcher.jsonl"));
673 assert_ne!(rollout.session_id, fixtures::ROOT_SESSION_ID);
674 let anchor = &parse_subagent_anchors(fixtures::ROLLOUT.as_bytes())[0];
675 let got: serde_json::Value =
676 serde_json::from_str(&body(HARNESS_ID_CODEX, &rollout, anchor)).unwrap();
677 assert_eq!(
678 got["session"]["harness_session_id"],
679 fixtures::ROOT_SESSION_ID
680 );
681
682 let mut root = rollout;
684 root.root_session_id = None;
685 let got: serde_json::Value =
686 serde_json::from_str(&body(HARNESS_ID_CODEX, &root, anchor)).unwrap();
687 assert_eq!(got["session"]["harness_session_id"], fixtures::SESSION_ID);
688 }
689
690 #[test]
691 fn a_started_row_carries_no_kind_field_at_all() {
692 let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
696 let anchors = parse_subagent_anchors(fixtures::ROLLOUT.as_bytes());
697 let started: serde_json::Value =
700 serde_json::from_str(&body(HARNESS_ID_CODEX, &rollout, &anchors[0])).unwrap();
701 assert!(started.get("kind").is_none());
702 let interacted: serde_json::Value =
703 serde_json::from_str(&body(HARNESS_ID_CODEX, &rollout, &anchors[1])).unwrap();
704 assert_eq!(interacted["kind"], KIND_INTERACTED);
705 }
706
707 #[test]
708 fn the_row_names_the_harness_the_caller_declares() {
709 let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
713 let anchor = &parse_subagent_anchors(fixtures::ROLLOUT.as_bytes())[0];
714 let got = body(
715 tapes_capture::envelope::HARNESS_ID_CODEX_APP,
716 &rollout,
717 anchor,
718 );
719 assert!(got.contains(r#""harness_id":"codex-app""#), "got: {got}");
720 }
721
722 #[test]
723 fn an_anchor_with_no_agent_path_omits_the_optional_slots() {
724 let rollout = rollout_at(Path::new("/tmp/rollout.jsonl"));
725 let anchor = SubAgentAnchor {
726 kind: AnchorKind::Started,
727 thread_id: "child".to_owned(),
728 call_id: "call_x".to_owned(),
729 agent_path: None,
730 line: "{}".to_owned(),
731 };
732 let got = body(HARNESS_ID_CODEX, &rollout, &anchor);
733 assert!(!got.contains("agent_type"), "got: {got}");
734 assert!(!got.contains("description"), "got: {got}");
735 assert!(got.contains(r#""agent_id":"child""#), "got: {got}");
736 }
737
738 fn write_rollout(dir: &Path, body: &str) -> PathBuf {
741 let path = dir.join("rollout.jsonl");
742 std::fs::write(&path, body).unwrap();
743 path
744 }
745
746 #[test]
747 fn a_fresh_rollout_offers_every_anchor_once() {
748 let dir = tempfile::tempdir().unwrap();
749 let path = write_rollout(dir.path(), fixtures::ROLLOUT);
750 let mut scanner = CodexAnchorScanner::new();
751
752 let fingerprint = files::fingerprint(&path);
753 assert!(scanner.needs_read(&path, fingerprint));
754 let raw = std::fs::read(&path).unwrap();
755 let anchors = scanner.undelivered(&path, &raw);
756 assert_eq!(anchors.len(), 2);
757
758 for anchor in &anchors {
759 scanner.record_delivered(&path, anchor);
760 }
761 scanner.record_clean_scan(&path, fingerprint);
762
763 assert_eq!(scanner.delivered_count(&path), 2);
764 assert!(
765 !scanner.needs_read(&path, files::fingerprint(&path)),
766 "an unchanged append-only file cannot hide a new anchor",
767 );
768 assert!(scanner.undelivered(&path, &raw).is_empty());
769 }
770
771 #[test]
772 fn a_grown_rollout_offers_only_what_is_new() {
773 let dir = tempfile::tempdir().unwrap();
774 let path = write_rollout(dir.path(), fixtures::ROLLOUT);
775 let mut scanner = CodexAnchorScanner::new();
776
777 let raw = std::fs::read(&path).unwrap();
778 for anchor in scanner.undelivered(&path, &raw) {
779 scanner.record_delivered(&path, &anchor);
780 }
781 scanner.record_clean_scan(&path, files::fingerprint(&path));
782
783 let second = STARTED_LINE
785 .replace("call_J7B6r7ZdtqkECtSJV8YDQaL7", "call_second")
786 .replace(
787 "019f8d46-e663-74e1-940c-f82e34c07618",
788 "019f8d47-0473-7743-a1ed-9e4c0ae92ad8",
789 );
790 std::fs::write(&path, format!("{}{second}\n", fixtures::ROLLOUT)).unwrap();
791
792 assert!(scanner.needs_read(&path, files::fingerprint(&path)));
793 let raw = std::fs::read(&path).unwrap();
794 let pending = scanner.undelivered(&path, &raw);
795 assert_eq!(pending.len(), 1);
796 assert_eq!(pending[0].call_id, "call_second");
797 }
798
799 #[test]
800 fn an_undelivered_anchor_is_offered_again_next_read() {
801 let dir = tempfile::tempdir().unwrap();
804 let path = write_rollout(dir.path(), fixtures::ROLLOUT);
805 let scanner = CodexAnchorScanner::new();
806 let raw = std::fs::read(&path).unwrap();
807 assert_eq!(scanner.undelivered(&path, &raw).len(), 2);
808 assert_eq!(scanner.undelivered(&path, &raw).len(), 2);
809 assert!(
810 scanner.needs_read(&path, files::fingerprint(&path)),
811 "a scan that was never marked clean must re-run",
812 );
813 }
814
815 #[test]
816 fn a_vanished_rollout_is_skipped_rather_than_read() {
817 let scanner = CodexAnchorScanner::new();
818 assert!(!scanner.needs_read(Path::new("/nonexistent/rollout.jsonl"), None));
819 }
820
821 #[test]
822 fn retain_live_drops_rollouts_that_aged_out_of_the_snapshot() {
823 let mut scanner = CodexAnchorScanner::new();
824 let kept = PathBuf::from("/tmp/kept.jsonl");
825 let gone = PathBuf::from("/tmp/gone.jsonl");
826 let anchor = &parse_subagent_anchors(fixtures::ROLLOUT.as_bytes())[0];
827 scanner.record_delivered(&kept, anchor);
828 scanner.record_delivered(&gone, anchor);
829
830 scanner.retain_live([kept.as_path()]);
831
832 assert_eq!(scanner.delivered_count(&kept), 1);
833 assert_eq!(scanner.delivered_count(&gone), 0);
834 }
835}