Skip to main content

supercode_harness/
codex_peer.rs

1//! Live stock-Codex session discovery.
2//!
3//! Codex does not publish a peer registry or a supported attachment endpoint,
4//! but its process keeps every rollout it currently owns open. This module
5//! joins that process-owned file descriptor back to the persisted catalog
6//! path. The rollout's last explicit lifecycle event then distinguishes an
7//! executing turn from a merely running session. No timing or CPU heuristic
8//! is used.
9
10use std::collections::{HashMap, HashSet};
11use std::fs::File;
12use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};
13use std::path::{Path, PathBuf};
14use std::time::{Duration, Instant};
15
16use serde_json::Value;
17
18const LIFECYCLE_SCAN_BYTES: u64 = 1024 * 1024;
19const LIFECYCLE_OVERLAP_BYTES: u64 = 64 * 1024;
20const OWNERSHIP_REFRESH_INTERVAL: Duration = Duration::from_secs(1);
21const OWNERSHIP_MISS_CONFIRMATIONS: u8 = 2;
22
23/// Activity proven for a rollout owned by stock Codex.
24#[derive(Debug, Clone, Copy, PartialEq, Eq)]
25pub enum CodexPeerStatus {
26    /// Codex owns the rollout, but no active turn is proven.
27    Running,
28    /// The latest lifecycle boundary completed or aborted a turn.
29    Idle,
30    /// The latest lifecycle boundary starts a task.
31    Busy,
32}
33
34impl CodexPeerStatus {
35    /// Stable wire spelling shared by the harness protocol.
36    pub const fn as_str(self) -> &'static str {
37        match self {
38            Self::Running => "running",
39            Self::Idle => "idle",
40            Self::Busy => "busy",
41        }
42    }
43}
44
45/// Every Codex rollout currently held open by a stock `codex` process, with
46/// the narrowest activity state its own event stream proves.
47pub fn live_rollouts(sessions_root: &Path) -> HashMap<PathBuf, CodexPeerStatus> {
48    CodexPeerTracker::default().sample(sessions_root)
49}
50
51/// Cached process-ownership sampler for latency-sensitive activity streams.
52///
53/// Process/file-descriptor discovery is materially more expensive than
54/// reading a bounded lifecycle tail. Ownership is therefore refreshed once a
55/// second while known open rollouts are re-read on every activity tick.
56#[derive(Debug, Default)]
57pub(crate) struct CodexPeerTracker {
58    root: Option<PathBuf>,
59    open_rollouts: Vec<PathBuf>,
60    ownership_misses: HashMap<PathBuf, u8>,
61    lifecycle: HashMap<PathBuf, CodexLifecycleCursor>,
62    refreshed_at: Option<Instant>,
63}
64
65#[derive(Debug, Default)]
66struct CodexLifecycleCursor {
67    offset: u64,
68    status: Option<CodexPeerStatus>,
69}
70
71impl CodexPeerTracker {
72    pub(crate) fn sample(&mut self, sessions_root: &Path) -> HashMap<PathBuf, CodexPeerStatus> {
73        let root = normalized_path(sessions_root);
74        let refresh = self.root.as_ref() != Some(&root)
75            || self
76                .refreshed_at
77                .is_none_or(|at| at.elapsed() >= OWNERSHIP_REFRESH_INTERVAL);
78        if refresh {
79            let observed = platform_open_rollouts()
80                .into_iter()
81                .map(|path| normalized_path(&path))
82                .collect();
83            self.open_rollouts =
84                reconcile_open_rollouts(&self.open_rollouts, observed, &mut self.ownership_misses);
85            self.root = Some(root.clone());
86            self.refreshed_at = Some(Instant::now());
87            self.lifecycle
88                .retain(|path, _| self.open_rollouts.contains(path));
89        }
90        let mut statuses = HashMap::new();
91        for path in &self.open_rollouts {
92            if !(path.starts_with(&root)
93                && path.extension().and_then(|extension| extension.to_str()) == Some("jsonl"))
94            {
95                continue;
96            }
97            let cursor = self.lifecycle.entry(path.clone()).or_default();
98            let status = sample_lifecycle_status(path, cursor).unwrap_or(CodexPeerStatus::Running);
99            statuses.insert(path.clone(), status);
100        }
101        statuses
102    }
103}
104
105/// One process-table sample is negative evidence, not a lifecycle boundary.
106/// Keep a previously open rollout through one miss; only consecutive misses
107/// retire it. This absorbs transient `pgrep`/`lsof` snapshots without putting
108/// a wall-clock guess into session state, while a genuinely exited process is
109/// removed on the next independent ownership observation.
110fn reconcile_open_rollouts(
111    previous: &[PathBuf],
112    observed: Vec<PathBuf>,
113    misses: &mut HashMap<PathBuf, u8>,
114) -> Vec<PathBuf> {
115    let observed = observed.into_iter().collect::<HashSet<_>>();
116    let previous = previous.iter().cloned().collect::<HashSet<_>>();
117    let mut reconciled = observed.clone();
118
119    for path in &observed {
120        misses.remove(path);
121    }
122    for path in previous.difference(&observed) {
123        let count = misses.entry(path.clone()).or_insert(0);
124        *count = count.saturating_add(1);
125        if *count < OWNERSHIP_MISS_CONFIRMATIONS {
126            reconciled.insert(path.clone());
127        } else {
128            misses.remove(path);
129        }
130    }
131    misses.retain(|path, _| previous.contains(path) && !observed.contains(path));
132
133    let mut reconciled = reconciled.into_iter().collect::<Vec<_>>();
134    reconciled.sort();
135    reconciled
136}
137
138/// Activity for a discovered catalog path owned by a currently running Codex.
139pub fn rollout_status(
140    live: &HashMap<PathBuf, CodexPeerStatus>,
141    path: &Path,
142) -> Option<CodexPeerStatus> {
143    live.get(&normalized_path(path)).copied()
144}
145
146/// A rollout's session id and, for a subagent thread, its parent's id.
147pub fn rollout_session(path: &Path) -> Option<(String, Option<String>)> {
148    rollout_lineage(path)
149}
150
151/// The working directory a rollout's conversation runs in, from its first
152/// `session_meta` record.
153pub fn rollout_cwd(path: &Path) -> Option<PathBuf> {
154    let mut header = String::new();
155    BufReader::new(File::open(path).ok()?.take(256 * 1024))
156        .read_line(&mut header)
157        .ok()?;
158    let value = serde_json::from_str::<Value>(&header).ok()?;
159    value
160        .pointer("/payload/cwd")
161        .and_then(Value::as_str)
162        .map(PathBuf::from)
163}
164
165/// Lightweight native identity and direct parent from Codex's first
166/// `session_meta` record. Activity aggregation uses this to treat a process-
167/// owned subagent rollout as work inside its root conversation.
168pub(crate) fn rollout_lineage(path: &Path) -> Option<(String, Option<String>)> {
169    let mut header = String::new();
170    BufReader::new(File::open(path).ok()?.take(256 * 1024))
171        .read_line(&mut header)
172        .ok()?;
173    let value = serde_json::from_str::<Value>(&header).ok()?;
174    if value.get("type").and_then(Value::as_str) != Some("session_meta") {
175        return None;
176    }
177    let payload = value.get("payload")?;
178    let session_id = payload.get("id")?.as_str()?.to_string();
179    let parent_session_id = payload
180        .pointer("/source/subagent/thread_spawn/parent_thread_id")
181        .or_else(|| payload.get("parent_thread_id"))
182        .and_then(Value::as_str)
183        .map(str::to_string);
184    Some((session_id, parent_session_id))
185}
186
187fn sample_lifecycle_status(
188    path: &Path,
189    cursor: &mut CodexLifecycleCursor,
190) -> Option<CodexPeerStatus> {
191    let mut file = File::open(path).ok()?;
192    let length = file.metadata().ok()?.len();
193    if length < cursor.offset {
194        cursor.offset = 0;
195        cursor.status = None;
196    }
197    if cursor.offset == 0 {
198        cursor.status = latest_lifecycle_status_between(&mut file, 0, length);
199    } else if length > cursor.offset {
200        if let Some(status) = latest_lifecycle_status_between(&mut file, cursor.offset, length) {
201            cursor.status = Some(status);
202        }
203    }
204    cursor.offset = length;
205    cursor.status
206}
207
208fn latest_lifecycle_status_between(
209    file: &mut File,
210    floor: u64,
211    upper: u64,
212) -> Option<CodexPeerStatus> {
213    let mut end = upper;
214    while end > floor {
215        let start = end.saturating_sub(LIFECYCLE_SCAN_BYTES).max(floor);
216        file.seek(SeekFrom::Start(start)).ok()?;
217        let mut tail = vec![0; (end - start) as usize];
218        file.read_exact(&mut tail).ok()?;
219        if let Some(status) = lifecycle_status_in_tail(&tail, start == floor) {
220            return Some(status);
221        }
222        if start == floor {
223            break;
224        }
225        // A lifecycle record is small, but an adjacent tool record can be enormous. Overlap keeps
226        // a boundary record whole without ever allocating in proportion to the rollout.
227        end = start.saturating_add(LIFECYCLE_OVERLAP_BYTES);
228    }
229    None
230}
231
232fn lifecycle_status_in_tail(
233    tail: &[u8],
234    starts_at_record_boundary: bool,
235) -> Option<CodexPeerStatus> {
236    const BOUNDARIES: [&str; 3] = [
237        "\"type\":\"task_started\"",
238        "\"type\":\"task_complete\"",
239        "\"type\":\"turn_aborted\"",
240    ];
241
242    // The read can begin in the middle of a large UTF-8 JSON string. Ignore
243    // that first fragment, then use the standard library's substring search
244    // to jump directly between lifecycle candidates instead of inspecting
245    // every byte of every tool payload with a naive sliding window.
246    let complete_start = if starts_at_record_boundary {
247        0
248    } else {
249        tail.iter()
250            .position(|byte| *byte == b'\n')
251            .map_or(tail.len(), |newline| newline + 1)
252    };
253    let text = std::str::from_utf8(&tail[complete_start..]).ok()?;
254    let mut search_end = text.len();
255    while let Some(candidate) = BOUNDARIES
256        .iter()
257        .filter_map(|boundary| text[..search_end].rfind(boundary))
258        .max()
259    {
260        let line_start = text[..candidate]
261            .rfind('\n')
262            .map_or(0, |newline| newline + 1);
263        let line_end = text[candidate..]
264            .find('\n')
265            .map_or(text.len(), |newline| candidate + newline);
266        let Ok(event) = serde_json::from_str::<Value>(&text[line_start..line_end]) else {
267            search_end = candidate;
268            continue;
269        };
270        if event.get("type").and_then(Value::as_str) != Some("event_msg") {
271            search_end = candidate;
272            continue;
273        }
274        match event
275            .get("payload")
276            .and_then(|payload| payload.get("type"))
277            .and_then(Value::as_str)
278        {
279            Some("task_started") => return Some(CodexPeerStatus::Busy),
280            Some("task_complete" | "turn_aborted") => return Some(CodexPeerStatus::Idle),
281            _ => {}
282        }
283        search_end = candidate;
284    }
285    None
286}
287
288fn normalized_path(path: &Path) -> PathBuf {
289    path.canonicalize().unwrap_or_else(|_| path.to_path_buf())
290}
291
292#[cfg(target_os = "macos")]
293fn platform_open_rollouts() -> Vec<PathBuf> {
294    macos_open_rollouts_native().unwrap_or_else(macos_open_rollouts_with_commands)
295}
296
297#[cfg(target_os = "macos")]
298fn macos_open_rollouts_with_commands() -> Vec<PathBuf> {
299    use std::process::Command;
300
301    // Darwin's pgrep omits every ancestor of the caller unless `-a` is set.
302    // Discovery commonly runs underneath the very Codex session it must
303    // report (for example inside a Supercode-powered widget), so omitting
304    // ancestors makes the current session uniquely invisible.
305    let Ok(processes) = Command::new("/usr/bin/pgrep")
306        .args(["-a", "-x", "codex"])
307        .output()
308    else {
309        return Vec::new();
310    };
311    let pids = String::from_utf8_lossy(&processes.stdout)
312        .lines()
313        .filter_map(|line| line.trim().parse::<u32>().ok())
314        .take(128)
315        .map(|pid| pid.to_string())
316        .collect::<Vec<_>>();
317    if pids.is_empty() {
318        return Vec::new();
319    }
320    let Ok(files) = Command::new("/usr/sbin/lsof")
321        .args(["-Fn", "-a", "-p", &pids.join(",")])
322        .output()
323    else {
324        return Vec::new();
325    };
326    String::from_utf8_lossy(&files.stdout)
327        .lines()
328        .filter_map(|line| line.strip_prefix('n'))
329        .filter(|path| path.ends_with(".jsonl"))
330        .map(PathBuf::from)
331        .collect()
332}
333
334#[cfg(target_os = "macos")]
335fn macos_open_rollouts_native() -> Option<Vec<PathBuf>> {
336    use std::mem::size_of;
337
338    const PROCESS_NAME_BYTES: usize = 64;
339    const INITIAL_PID_CAPACITY: usize = 2_048;
340
341    let mut pids = vec![0_i32; INITIAL_PID_CAPACITY];
342    let count = loop {
343        let count = unsafe {
344            proc_listallpids(
345                pids.as_mut_ptr().cast(),
346                i32::try_from(pids.len() * size_of::<i32>()).ok()?,
347            )
348        };
349        if count <= 0 {
350            return None;
351        }
352        if usize::try_from(count).ok()? < pids.len() {
353            break count;
354        }
355        pids.resize(pids.len() * 2, 0);
356    };
357    pids.truncate(usize::try_from(count).ok()?.min(pids.len()));
358
359    let codex_pids = pids
360        .into_iter()
361        .filter(|pid| *pid > 0)
362        .filter(|pid| {
363            let mut name = [0_u8; PROCESS_NAME_BYTES];
364            let length = unsafe {
365                proc_name(
366                    *pid,
367                    name.as_mut_ptr().cast(),
368                    u32::try_from(name.len()).expect("small process-name buffer"),
369                )
370            };
371            usize::try_from(length)
372                .ok()
373                .and_then(|length| name.get(..length))
374                == Some(b"codex".as_slice())
375        })
376        .collect::<Vec<_>>();
377    if codex_pids.is_empty() {
378        return Some(Vec::new());
379    }
380
381    macos_open_jsonl_for_pids(&codex_pids)
382}
383
384#[cfg(target_os = "macos")]
385fn macos_open_jsonl_for_pids(pids: &[i32]) -> Option<Vec<PathBuf>> {
386    use std::mem::{size_of, MaybeUninit};
387    use std::os::unix::ffi::OsStrExt;
388
389    const INITIAL_FD_CAPACITY: usize = 256;
390    const MAX_FD_CAPACITY: usize = 65_536;
391    const PROC_PIDLISTFDS: i32 = 1;
392    const PROC_PIDFDVNODEPATHINFO: i32 = 2;
393    const PROX_FDTYPE_VNODE: u32 = 1;
394
395    let mut successful_fd_reads = 0_usize;
396    let mut failed_fd_reads = 0_usize;
397    let mut paths = Vec::new();
398    for pid in pids.iter().copied() {
399        let mut capacity = INITIAL_FD_CAPACITY;
400        let descriptors = loop {
401            let mut descriptors = vec![ProcFdInfo::default(); capacity];
402            let bytes = unsafe {
403                proc_pidinfo(
404                    pid,
405                    PROC_PIDLISTFDS,
406                    0,
407                    descriptors.as_mut_ptr().cast(),
408                    i32::try_from(descriptors.len() * size_of::<ProcFdInfo>()).ok()?,
409                )
410            };
411            if bytes <= 0 {
412                failed_fd_reads += 1;
413                break None;
414            }
415            let bytes = usize::try_from(bytes).ok()?;
416            if bytes < descriptors.len() * size_of::<ProcFdInfo>() {
417                descriptors.truncate(bytes / size_of::<ProcFdInfo>());
418                break Some(descriptors);
419            }
420            if capacity >= MAX_FD_CAPACITY {
421                return None;
422            }
423            capacity = (capacity * 2).min(MAX_FD_CAPACITY);
424        };
425        let Some(descriptors) = descriptors else {
426            continue;
427        };
428        successful_fd_reads += 1;
429        for descriptor in descriptors
430            .into_iter()
431            .filter(|descriptor| descriptor.proc_fdtype == PROX_FDTYPE_VNODE)
432        {
433            let mut info = MaybeUninit::<VnodeFdInfoWithPath>::zeroed();
434            let bytes = unsafe {
435                proc_pidfdinfo(
436                    pid,
437                    descriptor.proc_fd,
438                    PROC_PIDFDVNODEPATHINFO,
439                    info.as_mut_ptr().cast(),
440                    i32::try_from(size_of::<VnodeFdInfoWithPath>()).expect("fixed native struct"),
441                )
442            };
443            if usize::try_from(bytes).ok() != Some(size_of::<VnodeFdInfoWithPath>()) {
444                continue;
445            }
446            let info = unsafe { info.assume_init() };
447            let path_length = info
448                .pvip
449                .vip_path
450                .iter()
451                .position(|byte| *byte == 0)
452                .unwrap_or(info.pvip.vip_path.len());
453            let path = PathBuf::from(std::ffi::OsStr::from_bytes(
454                &info.pvip.vip_path[..path_length],
455            ));
456            if path.extension().and_then(|extension| extension.to_str()) == Some("jsonl") {
457                paths.push(path);
458            }
459        }
460    }
461    if successful_fd_reads == 0 || failed_fd_reads != 0 {
462        return None;
463    }
464    paths.sort();
465    paths.dedup();
466    Some(paths)
467}
468
469#[cfg(target_os = "macos")]
470#[repr(C)]
471#[derive(Debug, Clone, Copy, Default)]
472struct ProcFdInfo {
473    proc_fd: i32,
474    proc_fdtype: u32,
475}
476
477#[cfg(target_os = "macos")]
478#[repr(C)]
479struct ProcFileInfo {
480    fi_openflags: u32,
481    fi_status: u32,
482    fi_offset: i64,
483    fi_type: i32,
484    fi_guardflags: u32,
485}
486
487#[cfg(target_os = "macos")]
488#[repr(C)]
489struct VinfoStat {
490    vst_dev: u32,
491    vst_mode: u16,
492    vst_nlink: u16,
493    vst_ino: u64,
494    vst_uid: u32,
495    vst_gid: u32,
496    vst_atime: i64,
497    vst_atimensec: i64,
498    vst_mtime: i64,
499    vst_mtimensec: i64,
500    vst_ctime: i64,
501    vst_ctimensec: i64,
502    vst_birthtime: i64,
503    vst_birthtimensec: i64,
504    vst_size: i64,
505    vst_blocks: i64,
506    vst_blksize: i32,
507    vst_flags: u32,
508    vst_gen: u32,
509    vst_rdev: u32,
510    vst_qspare: [i64; 2],
511}
512
513#[cfg(target_os = "macos")]
514#[repr(C)]
515struct VnodeInfo {
516    vi_stat: VinfoStat,
517    vi_type: i32,
518    vi_pad: i32,
519    vi_fsid: [i32; 2],
520}
521
522#[cfg(target_os = "macos")]
523#[repr(C)]
524struct VnodeInfoPath {
525    vip_vi: VnodeInfo,
526    vip_path: [u8; 1_024],
527}
528
529#[cfg(target_os = "macos")]
530#[repr(C)]
531struct VnodeFdInfoWithPath {
532    pfi: ProcFileInfo,
533    pvip: VnodeInfoPath,
534}
535
536#[cfg(target_os = "macos")]
537#[link(name = "proc")]
538unsafe extern "C" {
539    fn proc_listallpids(buffer: *mut std::ffi::c_void, buffersize: i32) -> i32;
540    fn proc_pidinfo(
541        pid: i32,
542        flavor: i32,
543        arg: u64,
544        buffer: *mut std::ffi::c_void,
545        buffersize: i32,
546    ) -> i32;
547    fn proc_pidfdinfo(
548        pid: i32,
549        fd: i32,
550        flavor: i32,
551        buffer: *mut std::ffi::c_void,
552        buffersize: i32,
553    ) -> i32;
554    fn proc_name(pid: i32, buffer: *mut std::ffi::c_void, buffersize: u32) -> i32;
555}
556
557#[cfg(target_os = "linux")]
558fn platform_open_rollouts() -> Vec<PathBuf> {
559    let Ok(processes) = std::fs::read_dir("/proc") else {
560        return Vec::new();
561    };
562    let mut paths = Vec::new();
563    for process in processes.flatten() {
564        let pid = process.file_name();
565        if !pid.as_encoded_bytes().iter().all(u8::is_ascii_digit) {
566            continue;
567        }
568        let process_root = process.path();
569        if std::fs::read_to_string(process_root.join("comm"))
570            .ok()
571            .is_none_or(|name| name.trim() != "codex")
572        {
573            continue;
574        }
575        let Ok(descriptors) = std::fs::read_dir(process_root.join("fd")) else {
576            continue;
577        };
578        paths.extend(
579            descriptors
580                .flatten()
581                .filter_map(|descriptor| std::fs::read_link(descriptor.path()).ok())
582                .filter(|path| {
583                    path.extension().and_then(|extension| extension.to_str()) == Some("jsonl")
584                }),
585        );
586    }
587    paths
588}
589
590#[cfg(not(any(target_os = "macos", target_os = "linux")))]
591fn platform_open_rollouts() -> Vec<PathBuf> {
592    Vec::new()
593}
594
595/// The Codex session a process runs: the root rollout (one with no parent
596/// thread) among the files `pid` holds open, as `(session id, rollout path)`.
597/// `None` when `pid` holds no Codex rollout.
598///
599/// This is how a command run from a Codex shell learns which session ran it:
600/// its ancestry reaches the `codex` process, and that process owns exactly
601/// the rollout of its conversation (subagent rollouts name their parent).
602pub fn session_of_process(pid: u32) -> Option<(String, PathBuf)> {
603    let mut root = None;
604    for path in rollouts_open_by(pid) {
605        let Some((session_id, parent)) = rollout_lineage(&path) else {
606            continue;
607        };
608        if parent.is_none() {
609            return Some((session_id, path));
610        }
611        root.get_or_insert((session_id, path));
612    }
613    root
614}
615
616#[cfg(target_os = "macos")]
617fn rollouts_open_by(pid: u32) -> Vec<PathBuf> {
618    i32::try_from(pid)
619        .ok()
620        .and_then(|pid| macos_open_jsonl_for_pids(&[pid]))
621        .unwrap_or_default()
622}
623
624#[cfg(target_os = "linux")]
625fn rollouts_open_by(pid: u32) -> Vec<PathBuf> {
626    let Ok(descriptors) = std::fs::read_dir(format!("/proc/{pid}/fd")) else {
627        return Vec::new();
628    };
629    descriptors
630        .flatten()
631        .filter_map(|descriptor| std::fs::read_link(descriptor.path()).ok())
632        .filter(|path| path.extension().and_then(|extension| extension.to_str()) == Some("jsonl"))
633        .collect()
634}
635
636#[cfg(not(any(target_os = "macos", target_os = "linux")))]
637fn rollouts_open_by(_pid: u32) -> Vec<PathBuf> {
638    Vec::new()
639}
640
641#[cfg(test)]
642mod tests {
643    use std::fs::{remove_file, OpenOptions};
644    use std::io::Write;
645
646    use super::*;
647
648    #[cfg(target_os = "macos")]
649    #[test]
650    fn native_macos_fd_layout_and_open_jsonl_discovery_match_the_sdk() {
651        assert_eq!(std::mem::size_of::<ProcFdInfo>(), 8);
652        assert_eq!(std::mem::size_of::<ProcFileInfo>(), 24);
653        assert_eq!(std::mem::size_of::<VnodeInfo>(), 152);
654        assert_eq!(std::mem::size_of::<VnodeFdInfoWithPath>(), 1_200);
655
656        let path = std::env::temp_dir().join(format!(
657            "supercode-native-open-rollout-{}.jsonl",
658            std::process::id()
659        ));
660        let file = File::create(&path).unwrap();
661        let open = macos_open_jsonl_for_pids(&[std::process::id() as i32])
662            .expect("the current process's descriptor table should be readable");
663        assert!(open.contains(&normalized_path(&path)));
664        drop(file);
665        remove_file(path).unwrap();
666    }
667
668    #[test]
669    fn long_tool_heavy_turn_is_found_once_then_followed_incrementally() {
670        let path = std::env::temp_dir().join(format!(
671            "supercode-codex-long-turn-{}-{}.jsonl",
672            std::process::id(),
673            std::thread::current().name().unwrap_or("test")
674        ));
675        let mut file = File::create(&path).unwrap();
676        writeln!(
677            file,
678            r#"{{"type":"event_msg","payload":{{"type":"task_started"}}}}"#
679        )
680        .unwrap();
681        write!(
682            file,
683            r#"{{"type":"response_item","payload":"{}"}}"#,
684            "x".repeat(6 * 1024 * 1024)
685        )
686        .unwrap();
687        writeln!(file).unwrap();
688        file.flush().unwrap();
689
690        let mut cursor = CodexLifecycleCursor::default();
691        assert_eq!(
692            sample_lifecycle_status(&path, &mut cursor),
693            Some(CodexPeerStatus::Busy)
694        );
695        let first_offset = cursor.offset;
696
697        let mut file = OpenOptions::new().append(true).open(&path).unwrap();
698        writeln!(
699            file,
700            r#"{{"type":"event_msg","payload":{{"type":"item_completed"}}}}"#
701        )
702        .unwrap();
703        file.flush().unwrap();
704        assert_eq!(
705            sample_lifecycle_status(&path, &mut cursor),
706            Some(CodexPeerStatus::Busy)
707        );
708        assert!(cursor.offset > first_offset);
709
710        writeln!(
711            file,
712            r#"{{"type":"event_msg","payload":{{"type":"task_complete"}}}}"#
713        )
714        .unwrap();
715        file.flush().unwrap();
716        assert_eq!(
717            sample_lifecycle_status(&path, &mut cursor),
718            Some(CodexPeerStatus::Idle)
719        );
720        remove_file(path).unwrap();
721    }
722
723    #[test]
724    fn open_rollout_requires_consecutive_misses_before_retirement() {
725        let rollout = PathBuf::from("/tmp/session.jsonl");
726        let mut misses = HashMap::new();
727
728        let observed = reconcile_open_rollouts(&[], vec![rollout.clone()], &mut misses);
729        assert_eq!(observed, vec![rollout.clone()]);
730
731        let retained = reconcile_open_rollouts(&observed, vec![], &mut misses);
732        assert_eq!(retained, vec![rollout.clone()]);
733        assert_eq!(misses.get(&rollout), Some(&1));
734
735        let recovered = reconcile_open_rollouts(&retained, vec![rollout.clone()], &mut misses);
736        assert_eq!(recovered, vec![rollout.clone()]);
737        assert!(misses.is_empty());
738
739        let retained = reconcile_open_rollouts(&recovered, vec![], &mut misses);
740        assert_eq!(retained, vec![rollout.clone()]);
741        let retired = reconcile_open_rollouts(&retained, vec![], &mut misses);
742        assert!(retired.is_empty());
743        assert!(misses.is_empty());
744    }
745}