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/// The Codex conversation a launch in `cwd` started: the newest root rollout under `sessions_root`
263/// (`YYYY/MM/DD/rollout-*.jsonl`) written since `since` whose first record names `cwd`. A launch
264/// is bound by this record of its own rather than by the process, open-file and pane chain, which
265/// missed both Codex launches of 2026-10-02 within a minute though their rollouts existed in two.
266pub fn launched_session(
267    sessions_root: &Path,
268    cwd: &Path,
269    since: std::time::SystemTime,
270) -> Option<String> {
271    let newest = |dir: &Path, n: usize| -> Vec<PathBuf> {
272        let mut found: Vec<PathBuf> = std::fs::read_dir(dir)
273            .map(|entries| {
274                entries
275                    .flatten()
276                    .map(|entry| entry.path())
277                    .filter(|path| path.is_dir())
278                    .collect()
279            })
280            .unwrap_or_default();
281        found.sort();
282        found.into_iter().rev().take(n).collect()
283    };
284    let want = std::fs::canonicalize(cwd).unwrap_or_else(|_| cwd.to_path_buf());
285    let mut best: Option<(std::time::SystemTime, String)> = None;
286    // the day the launch began and the one before it (a launch across midnight, local or UTC)
287    for year in newest(sessions_root, 2) {
288        for month in newest(&year, 2) {
289            for day in newest(&month, 2) {
290                let Ok(entries) = std::fs::read_dir(&day) else {
291                    continue;
292                };
293                for entry in entries.flatten() {
294                    let path = entry.path();
295                    if path.extension().and_then(|e| e.to_str()) != Some("jsonl") {
296                        continue;
297                    }
298                    // born after the launch: a conversation already running in this folder keeps
299                    // writing its rollout, so a modified time would bind it instead of the new one
300                    let Ok(modified) = entry.metadata().and_then(|m| m.created()) else {
301                        continue;
302                    };
303                    if modified < since {
304                        continue;
305                    }
306                    let Some((session_id, None)) = rollout_lineage(&path) else {
307                        continue;
308                    };
309                    let Some(dir) = rollout_cwd(&path) else {
310                        continue;
311                    };
312                    if std::fs::canonicalize(&dir).unwrap_or(dir) != want {
313                        continue;
314                    }
315                    if best.as_ref().is_none_or(|(at, _)| modified > *at) {
316                        best = Some((modified, session_id));
317                    }
318                }
319            }
320        }
321    }
322    best.map(|(_, session_id)| session_id)
323}
324
325/// A rollout's session id and, for a subagent thread, its parent's id.
326pub fn rollout_session(path: &Path) -> Option<(String, Option<String>)> {
327    rollout_lineage(path)
328}
329
330/// The working directory a rollout's conversation runs in, from its first
331/// `session_meta` record.
332pub fn rollout_cwd(path: &Path) -> Option<PathBuf> {
333    let mut header = String::new();
334    BufReader::new(File::open(path).ok()?.take(256 * 1024))
335        .read_line(&mut header)
336        .ok()?;
337    let value = serde_json::from_str::<Value>(&header).ok()?;
338    value
339        .pointer("/payload/cwd")
340        .and_then(Value::as_str)
341        .map(PathBuf::from)
342}
343
344/// Lightweight native identity and direct parent from Codex's first
345/// `session_meta` record. Activity aggregation uses this to treat a process-
346/// owned subagent rollout as work inside its root conversation.
347pub(crate) fn rollout_lineage(path: &Path) -> Option<(String, Option<String>)> {
348    let mut header = String::new();
349    BufReader::new(File::open(path).ok()?.take(256 * 1024))
350        .read_line(&mut header)
351        .ok()?;
352    let value = serde_json::from_str::<Value>(&header).ok()?;
353    if value.get("type").and_then(Value::as_str) != Some("session_meta") {
354        return None;
355    }
356    let payload = value.get("payload")?;
357    let session_id = payload.get("id")?.as_str()?.to_string();
358    let parent_session_id = payload
359        .pointer("/source/subagent/thread_spawn/parent_thread_id")
360        .or_else(|| payload.get("parent_thread_id"))
361        .and_then(Value::as_str)
362        .map(str::to_string);
363    Some((session_id, parent_session_id))
364}
365
366fn sample_lifecycle_status(
367    path: &Path,
368    cursor: &mut CodexLifecycleCursor,
369) -> Option<CodexPeerStatus> {
370    let mut file = File::open(path).ok()?;
371    let length = file.metadata().ok()?.len();
372    if length < cursor.offset {
373        cursor.offset = 0;
374        cursor.status = None;
375    }
376    if cursor.offset == 0 {
377        cursor.status = latest_lifecycle_status_between(&mut file, 0, length);
378    } else if length > cursor.offset {
379        if let Some(status) = latest_lifecycle_status_between(&mut file, cursor.offset, length) {
380            cursor.status = Some(status);
381        }
382    }
383    cursor.offset = length;
384    cursor.status
385}
386
387fn latest_lifecycle_status_between(
388    file: &mut File,
389    floor: u64,
390    upper: u64,
391) -> Option<CodexPeerStatus> {
392    let mut end = upper;
393    while end > floor {
394        let start = end.saturating_sub(LIFECYCLE_SCAN_BYTES).max(floor);
395        file.seek(SeekFrom::Start(start)).ok()?;
396        let mut tail = vec![0; (end - start) as usize];
397        file.read_exact(&mut tail).ok()?;
398        if let Some(status) = lifecycle_status_in_tail(&tail, start == floor) {
399            return Some(status);
400        }
401        if start == floor {
402            break;
403        }
404        // A lifecycle record is small, but an adjacent tool record can be enormous. Overlap keeps
405        // a boundary record whole without ever allocating in proportion to the rollout.
406        end = start.saturating_add(LIFECYCLE_OVERLAP_BYTES);
407    }
408    None
409}
410
411fn lifecycle_status_in_tail(
412    tail: &[u8],
413    starts_at_record_boundary: bool,
414) -> Option<CodexPeerStatus> {
415    const BOUNDARIES: [&str; 3] = [
416        "\"type\":\"task_started\"",
417        "\"type\":\"task_complete\"",
418        "\"type\":\"turn_aborted\"",
419    ];
420
421    // The read can begin in the middle of a large UTF-8 JSON string. Ignore
422    // that first fragment, then use the standard library's substring search
423    // to jump directly between lifecycle candidates instead of inspecting
424    // every byte of every tool payload with a naive sliding window.
425    let complete_start = if starts_at_record_boundary {
426        0
427    } else {
428        tail.iter()
429            .position(|byte| *byte == b'\n')
430            .map_or(tail.len(), |newline| newline + 1)
431    };
432    let text = std::str::from_utf8(&tail[complete_start..]).ok()?;
433    let mut search_end = text.len();
434    while let Some(candidate) = BOUNDARIES
435        .iter()
436        .filter_map(|boundary| text[..search_end].rfind(boundary))
437        .max()
438    {
439        let line_start = text[..candidate]
440            .rfind('\n')
441            .map_or(0, |newline| newline + 1);
442        let line_end = text[candidate..]
443            .find('\n')
444            .map_or(text.len(), |newline| candidate + newline);
445        let Ok(event) = serde_json::from_str::<Value>(&text[line_start..line_end]) else {
446            search_end = candidate;
447            continue;
448        };
449        if event.get("type").and_then(Value::as_str) != Some("event_msg") {
450            search_end = candidate;
451            continue;
452        }
453        match event
454            .get("payload")
455            .and_then(|payload| payload.get("type"))
456            .and_then(Value::as_str)
457        {
458            Some("task_started") => return Some(CodexPeerStatus::Busy),
459            Some("task_complete" | "turn_aborted") => return Some(CodexPeerStatus::Idle),
460            _ => {}
461        }
462        search_end = candidate;
463    }
464    None
465}
466
467fn normalized_path(path: &Path) -> PathBuf {
468    path.canonicalize().unwrap_or_else(|_| path.to_path_buf())
469}
470
471#[cfg(target_os = "macos")]
472fn platform_open_rollouts() -> Vec<PathBuf> {
473    macos_open_rollouts_native().unwrap_or_else(macos_open_rollouts_with_commands)
474}
475
476#[cfg(target_os = "macos")]
477fn macos_open_rollouts_with_commands() -> Vec<PathBuf> {
478    use std::process::Command;
479
480    // Darwin's pgrep omits every ancestor of the caller unless `-a` is set.
481    // Discovery commonly runs underneath the very Codex session it must
482    // report (for example inside a Supercode-powered widget), so omitting
483    // ancestors makes the current session uniquely invisible.
484    let Ok(processes) = Command::new("/usr/bin/pgrep")
485        .args(["-a", "-x", "codex"])
486        .output()
487    else {
488        return Vec::new();
489    };
490    let pids = String::from_utf8_lossy(&processes.stdout)
491        .lines()
492        .filter_map(|line| line.trim().parse::<u32>().ok())
493        .take(128)
494        .map(|pid| pid.to_string())
495        .collect::<Vec<_>>();
496    if pids.is_empty() {
497        return Vec::new();
498    }
499    let Ok(files) = Command::new("/usr/sbin/lsof")
500        .args(["-Fn", "-a", "-p", &pids.join(",")])
501        .output()
502    else {
503        return Vec::new();
504    };
505    String::from_utf8_lossy(&files.stdout)
506        .lines()
507        .filter_map(|line| line.strip_prefix('n'))
508        .filter(|path| path.ends_with(".jsonl"))
509        .map(PathBuf::from)
510        .collect()
511}
512
513#[cfg(target_os = "macos")]
514fn macos_open_rollouts_native() -> Option<Vec<PathBuf>> {
515    use std::mem::size_of;
516
517    const PROCESS_NAME_BYTES: usize = 64;
518    const INITIAL_PID_CAPACITY: usize = 2_048;
519
520    let mut pids = vec![0_i32; INITIAL_PID_CAPACITY];
521    let count = loop {
522        let count = unsafe {
523            proc_listallpids(
524                pids.as_mut_ptr().cast(),
525                i32::try_from(pids.len() * size_of::<i32>()).ok()?,
526            )
527        };
528        if count <= 0 {
529            return None;
530        }
531        if usize::try_from(count).ok()? < pids.len() {
532            break count;
533        }
534        pids.resize(pids.len() * 2, 0);
535    };
536    pids.truncate(usize::try_from(count).ok()?.min(pids.len()));
537
538    let codex_pids = pids
539        .into_iter()
540        .filter(|pid| *pid > 0)
541        .filter(|pid| {
542            let mut name = [0_u8; PROCESS_NAME_BYTES];
543            let length = unsafe {
544                proc_name(
545                    *pid,
546                    name.as_mut_ptr().cast(),
547                    u32::try_from(name.len()).expect("small process-name buffer"),
548                )
549            };
550            usize::try_from(length)
551                .ok()
552                .and_then(|length| name.get(..length))
553                == Some(b"codex".as_slice())
554        })
555        .collect::<Vec<_>>();
556    if codex_pids.is_empty() {
557        return Some(Vec::new());
558    }
559
560    macos_open_jsonl_for_pids(&codex_pids)
561}
562
563#[cfg(target_os = "macos")]
564fn macos_open_jsonl_for_pids(pids: &[i32]) -> Option<Vec<PathBuf>> {
565    use std::mem::{size_of, MaybeUninit};
566    use std::os::unix::ffi::OsStrExt;
567
568    const INITIAL_FD_CAPACITY: usize = 256;
569    const MAX_FD_CAPACITY: usize = 65_536;
570    const PROC_PIDLISTFDS: i32 = 1;
571    const PROC_PIDFDVNODEPATHINFO: i32 = 2;
572    const PROX_FDTYPE_VNODE: u32 = 1;
573
574    let mut successful_fd_reads = 0_usize;
575    let mut failed_fd_reads = 0_usize;
576    let mut paths = Vec::new();
577    for pid in pids.iter().copied() {
578        let mut capacity = INITIAL_FD_CAPACITY;
579        let descriptors = loop {
580            let mut descriptors = vec![ProcFdInfo::default(); capacity];
581            let bytes = unsafe {
582                proc_pidinfo(
583                    pid,
584                    PROC_PIDLISTFDS,
585                    0,
586                    descriptors.as_mut_ptr().cast(),
587                    i32::try_from(descriptors.len() * size_of::<ProcFdInfo>()).ok()?,
588                )
589            };
590            if bytes <= 0 {
591                failed_fd_reads += 1;
592                break None;
593            }
594            let bytes = usize::try_from(bytes).ok()?;
595            if bytes < descriptors.len() * size_of::<ProcFdInfo>() {
596                descriptors.truncate(bytes / size_of::<ProcFdInfo>());
597                break Some(descriptors);
598            }
599            if capacity >= MAX_FD_CAPACITY {
600                return None;
601            }
602            capacity = (capacity * 2).min(MAX_FD_CAPACITY);
603        };
604        let Some(descriptors) = descriptors else {
605            continue;
606        };
607        successful_fd_reads += 1;
608        for descriptor in descriptors
609            .into_iter()
610            .filter(|descriptor| descriptor.proc_fdtype == PROX_FDTYPE_VNODE)
611        {
612            let mut info = MaybeUninit::<VnodeFdInfoWithPath>::zeroed();
613            let bytes = unsafe {
614                proc_pidfdinfo(
615                    pid,
616                    descriptor.proc_fd,
617                    PROC_PIDFDVNODEPATHINFO,
618                    info.as_mut_ptr().cast(),
619                    i32::try_from(size_of::<VnodeFdInfoWithPath>()).expect("fixed native struct"),
620                )
621            };
622            if usize::try_from(bytes).ok() != Some(size_of::<VnodeFdInfoWithPath>()) {
623                continue;
624            }
625            let info = unsafe { info.assume_init() };
626            let path_length = info
627                .pvip
628                .vip_path
629                .iter()
630                .position(|byte| *byte == 0)
631                .unwrap_or(info.pvip.vip_path.len());
632            let path = PathBuf::from(std::ffi::OsStr::from_bytes(
633                &info.pvip.vip_path[..path_length],
634            ));
635            if path.extension().and_then(|extension| extension.to_str()) == Some("jsonl") {
636                paths.push(path);
637            }
638        }
639    }
640    if successful_fd_reads == 0 || failed_fd_reads != 0 {
641        return None;
642    }
643    paths.sort();
644    paths.dedup();
645    Some(paths)
646}
647
648#[cfg(target_os = "macos")]
649#[repr(C)]
650#[derive(Debug, Clone, Copy, Default)]
651struct ProcFdInfo {
652    proc_fd: i32,
653    proc_fdtype: u32,
654}
655
656#[cfg(target_os = "macos")]
657#[repr(C)]
658struct ProcFileInfo {
659    fi_openflags: u32,
660    fi_status: u32,
661    fi_offset: i64,
662    fi_type: i32,
663    fi_guardflags: u32,
664}
665
666#[cfg(target_os = "macos")]
667#[repr(C)]
668struct VinfoStat {
669    vst_dev: u32,
670    vst_mode: u16,
671    vst_nlink: u16,
672    vst_ino: u64,
673    vst_uid: u32,
674    vst_gid: u32,
675    vst_atime: i64,
676    vst_atimensec: i64,
677    vst_mtime: i64,
678    vst_mtimensec: i64,
679    vst_ctime: i64,
680    vst_ctimensec: i64,
681    vst_birthtime: i64,
682    vst_birthtimensec: i64,
683    vst_size: i64,
684    vst_blocks: i64,
685    vst_blksize: i32,
686    vst_flags: u32,
687    vst_gen: u32,
688    vst_rdev: u32,
689    vst_qspare: [i64; 2],
690}
691
692#[cfg(target_os = "macos")]
693#[repr(C)]
694struct VnodeInfo {
695    vi_stat: VinfoStat,
696    vi_type: i32,
697    vi_pad: i32,
698    vi_fsid: [i32; 2],
699}
700
701#[cfg(target_os = "macos")]
702#[repr(C)]
703struct VnodeInfoPath {
704    vip_vi: VnodeInfo,
705    vip_path: [u8; 1_024],
706}
707
708#[cfg(target_os = "macos")]
709#[repr(C)]
710struct VnodeFdInfoWithPath {
711    pfi: ProcFileInfo,
712    pvip: VnodeInfoPath,
713}
714
715#[cfg(target_os = "macos")]
716#[link(name = "proc")]
717unsafe extern "C" {
718    fn proc_listallpids(buffer: *mut std::ffi::c_void, buffersize: i32) -> i32;
719    fn proc_pidinfo(
720        pid: i32,
721        flavor: i32,
722        arg: u64,
723        buffer: *mut std::ffi::c_void,
724        buffersize: i32,
725    ) -> i32;
726    fn proc_pidfdinfo(
727        pid: i32,
728        fd: i32,
729        flavor: i32,
730        buffer: *mut std::ffi::c_void,
731        buffersize: i32,
732    ) -> i32;
733    fn proc_name(pid: i32, buffer: *mut std::ffi::c_void, buffersize: u32) -> i32;
734}
735
736#[cfg(target_os = "linux")]
737fn platform_open_rollouts() -> Vec<PathBuf> {
738    let Ok(processes) = std::fs::read_dir("/proc") else {
739        return Vec::new();
740    };
741    let mut paths = Vec::new();
742    for process in processes.flatten() {
743        let pid = process.file_name();
744        if !pid.as_encoded_bytes().iter().all(u8::is_ascii_digit) {
745            continue;
746        }
747        let process_root = process.path();
748        if std::fs::read_to_string(process_root.join("comm"))
749            .ok()
750            .is_none_or(|name| name.trim() != "codex")
751        {
752            continue;
753        }
754        let Ok(descriptors) = std::fs::read_dir(process_root.join("fd")) else {
755            continue;
756        };
757        paths.extend(
758            descriptors
759                .flatten()
760                .filter_map(|descriptor| std::fs::read_link(descriptor.path()).ok())
761                .filter(|path| {
762                    path.extension().and_then(|extension| extension.to_str()) == Some("jsonl")
763                }),
764        );
765    }
766    paths
767}
768
769#[cfg(not(any(target_os = "macos", target_os = "linux")))]
770fn platform_open_rollouts() -> Vec<PathBuf> {
771    Vec::new()
772}
773
774/// The Codex session a process runs: the root rollout (one with no parent
775/// thread) among the files `pid` holds open, as `(session id, rollout path)`.
776/// `None` when `pid` holds no Codex rollout.
777///
778/// This is how a command run from a Codex shell learns which session ran it:
779/// its ancestry reaches the `codex` process, and that process owns exactly
780/// the rollout of its conversation (subagent rollouts name their parent).
781///
782/// Codex's shared app-server daemon holds the rollouts of every thread it has
783/// loaded; a process holding more than one conversation names none of them.
784pub fn session_of_process(pid: u32) -> Option<(String, PathBuf)> {
785    let mut roots: Vec<(String, PathBuf)> = Vec::new();
786    let mut child = None;
787    for path in rollouts_open_by(pid) {
788        let Some((session_id, parent)) = rollout_lineage(&path) else {
789            continue;
790        };
791        if parent.is_none() {
792            if !roots.iter().any(|(id, _)| *id == session_id) {
793                roots.push((session_id, path));
794            }
795        } else {
796            child.get_or_insert((session_id, path));
797        }
798    }
799    match roots.len() {
800        0 => child,
801        1 => roots.pop(),
802        _ => None,
803    }
804}
805
806/// The conversation a `codex resume <id>` process continues, read from its command line: a resumed
807/// conversation that has not written since holds no rollout open, so the rollout cannot name it.
808pub fn resumed_session_of_process(pid: u32) -> Option<String> {
809    let output = std::process::Command::new("ps")
810        .args(["-o", "command=", "-p", &pid.to_string()])
811        .output()
812        .ok()?;
813    let command = String::from_utf8_lossy(&output.stdout);
814    let mut words = command.split_whitespace();
815    while let Some(word) = words.next() {
816        if word == "resume" {
817            let id = words.next()?;
818            let uuid = id.len() == 36
819                && id.chars().all(|c| c.is_ascii_hexdigit() || c == '-')
820                && id.matches('-').count() == 4;
821            return uuid.then(|| id.to_string());
822        }
823    }
824    None
825}
826
827/// The rollout file of conversation `session_id` under Codex's sessions root (`YYYY/MM/DD/rollout-…-<id>.jsonl`).
828pub fn rollout_of_session(sessions_root: &Path, session_id: &str) -> Option<PathBuf> {
829    let suffix = format!("-{session_id}.jsonl");
830    let mut stack = vec![(sessions_root.to_path_buf(), 0)];
831    while let Some((dir, depth)) = stack.pop() {
832        for entry in std::fs::read_dir(&dir).ok()?.flatten() {
833            let path = entry.path();
834            if depth < 3 && path.is_dir() {
835                stack.push((path, depth + 1));
836            } else if path
837                .file_name()
838                .and_then(|n| n.to_str())
839                .is_some_and(|n| n.ends_with(&suffix))
840            {
841                return Some(path);
842            }
843        }
844    }
845    None
846}
847
848/// The machine daemon pane each Codex conversation runs in, by session id: the pane the daemon
849/// names in the environment it gives every pane (`SUPERCODE_TEAMS_PANE`), read from the `codex`
850/// process that holds the conversation's rollout. A conversation outside any daemon pane (or held
851/// only by the shared app-server) has none.
852pub fn session_panes() -> HashMap<String, String> {
853    let mut panes = HashMap::new();
854    for pid in codex_pids() {
855        // A resumed conversation idle since its resume holds no rollout open: its own command line names it.
856        let Some(session_id) = session_of_process(pid)
857            .map(|(session_id, _)| session_id)
858            .or_else(|| resumed_session_of_process(pid))
859        else {
860            continue;
861        };
862        if let Some(pane) = pane_of_process(pid) {
863            panes.insert(session_id, pane);
864        }
865    }
866    panes
867}
868
869#[cfg(unix)]
870fn codex_pids() -> Vec<u32> {
871    // `-a`: Darwin's pgrep otherwise omits the caller's ancestors (see platform_open_rollouts).
872    let args: &[&str] = if cfg!(target_os = "macos") {
873        &["-a", "-x", "codex"]
874    } else {
875        &["-x", "codex"]
876    };
877    std::process::Command::new("pgrep")
878        .args(args)
879        .output()
880        .map(|output| {
881            String::from_utf8_lossy(&output.stdout)
882                .lines()
883                .filter_map(|line| line.trim().parse().ok())
884                .take(128)
885                .collect()
886        })
887        .unwrap_or_default()
888}
889
890#[cfg(not(unix))]
891fn codex_pids() -> Vec<u32> {
892    Vec::new()
893}
894
895/// The daemon pane a process runs in, from its own environment.
896#[cfg(target_os = "linux")]
897fn pane_of_process(pid: u32) -> Option<String> {
898    let environ = std::fs::read(format!("/proc/{pid}/environ")).ok()?;
899    environ
900        .split(|byte| *byte == 0)
901        .filter_map(|entry| std::str::from_utf8(entry).ok())
902        .find_map(|entry| entry.strip_prefix("SUPERCODE_TEAMS_PANE="))
903        .filter(|pane| pane.starts_with("p_"))
904        .map(str::to_string)
905}
906
907/// The daemon pane a process runs in, from its own environment (`ps -E` prints it after the command).
908#[cfg(target_os = "macos")]
909fn pane_of_process(pid: u32) -> Option<String> {
910    let output = std::process::Command::new("ps")
911        .args(["-E", "-o", "command=", "-p", &pid.to_string()])
912        .output()
913        .ok()?;
914    String::from_utf8_lossy(&output.stdout)
915        .split_whitespace()
916        .find_map(|word| word.strip_prefix("SUPERCODE_TEAMS_PANE="))
917        .filter(|pane| pane.starts_with("p_"))
918        .map(str::to_string)
919}
920
921#[cfg(not(any(target_os = "linux", target_os = "macos")))]
922fn pane_of_process(_pid: u32) -> Option<String> {
923    None
924}
925
926/// Whether `pid` holds the rollout of Codex session `session_id` open.
927pub fn holds_session(pid: u32, session_id: &str) -> bool {
928    rollouts_open_by(pid)
929        .iter()
930        .filter_map(|path| rollout_lineage(path))
931        .any(|(id, _)| id == session_id)
932}
933
934#[cfg(target_os = "macos")]
935fn rollouts_open_by(pid: u32) -> Vec<PathBuf> {
936    i32::try_from(pid)
937        .ok()
938        .and_then(|pid| macos_open_jsonl_for_pids(&[pid]))
939        .unwrap_or_default()
940}
941
942#[cfg(target_os = "linux")]
943fn rollouts_open_by(pid: u32) -> Vec<PathBuf> {
944    let Ok(descriptors) = std::fs::read_dir(format!("/proc/{pid}/fd")) else {
945        return Vec::new();
946    };
947    descriptors
948        .flatten()
949        .filter_map(|descriptor| std::fs::read_link(descriptor.path()).ok())
950        .filter(|path| path.extension().and_then(|extension| extension.to_str()) == Some("jsonl"))
951        .collect()
952}
953
954#[cfg(not(any(target_os = "macos", target_os = "linux")))]
955fn rollouts_open_by(_pid: u32) -> Vec<PathBuf> {
956    Vec::new()
957}