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