1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
25pub enum CodexPeerStatus {
26 Running,
28 Idle,
30 Busy,
32}
33
34impl CodexPeerStatus {
35 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
45pub fn live_rollouts(sessions_root: &Path) -> HashMap<PathBuf, CodexPeerStatus> {
48 CodexPeerTracker::default().sample(sessions_root)
49}
50
51#[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
105fn 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
138pub fn rollout_status(
140 live: &HashMap<PathBuf, CodexPeerStatus>,
141 path: &Path,
142) -> Option<CodexPeerStatus> {
143 live.get(&normalized_path(path)).copied()
144}
145
146pub fn rollout_session(path: &Path) -> Option<(String, Option<String>)> {
148 rollout_lineage(path)
149}
150
151pub 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
165pub(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 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 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 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
595pub fn session_of_process(pid: u32) -> Option<(String, PathBuf)> {
606 let mut roots: Vec<(String, PathBuf)> = Vec::new();
607 let mut child = None;
608 for path in rollouts_open_by(pid) {
609 let Some((session_id, parent)) = rollout_lineage(&path) else {
610 continue;
611 };
612 if parent.is_none() {
613 if !roots.iter().any(|(id, _)| *id == session_id) {
614 roots.push((session_id, path));
615 }
616 } else {
617 child.get_or_insert((session_id, path));
618 }
619 }
620 match roots.len() {
621 0 => child,
622 1 => roots.pop(),
623 _ => None,
624 }
625}
626
627pub fn holds_session(pid: u32, session_id: &str) -> bool {
629 rollouts_open_by(pid)
630 .iter()
631 .filter_map(|path| rollout_lineage(path))
632 .any(|(id, _)| id == session_id)
633}
634
635#[cfg(target_os = "macos")]
636fn rollouts_open_by(pid: u32) -> Vec<PathBuf> {
637 i32::try_from(pid)
638 .ok()
639 .and_then(|pid| macos_open_jsonl_for_pids(&[pid]))
640 .unwrap_or_default()
641}
642
643#[cfg(target_os = "linux")]
644fn rollouts_open_by(pid: u32) -> Vec<PathBuf> {
645 let Ok(descriptors) = std::fs::read_dir(format!("/proc/{pid}/fd")) else {
646 return Vec::new();
647 };
648 descriptors
649 .flatten()
650 .filter_map(|descriptor| std::fs::read_link(descriptor.path()).ok())
651 .filter(|path| path.extension().and_then(|extension| extension.to_str()) == Some("jsonl"))
652 .collect()
653}
654
655#[cfg(not(any(target_os = "macos", target_os = "linux")))]
656fn rollouts_open_by(_pid: u32) -> Vec<PathBuf> {
657 Vec::new()
658}