1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum CodexPeerStatus {
33 Running,
35 Idle,
37 Busy,
39 Waiting,
41}
42
43impl CodexPeerStatus {
44 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
55pub fn live_rollouts(sessions_root: &Path) -> HashMap<PathBuf, CodexPeerStatus> {
58 CodexPeerTracker::default().sample(sessions_root)
59}
60
61#[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 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#[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 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
221fn 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
254pub fn rollout_status(
256 live: &HashMap<PathBuf, CodexPeerStatus>,
257 path: &Path,
258) -> Option<CodexPeerStatus> {
259 live.get(&normalized_path(path)).copied()
260}
261
262pub 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 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 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
325pub fn rollout_session(path: &Path) -> Option<(String, Option<String>)> {
327 rollout_lineage(path)
328}
329
330pub 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
344pub(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 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 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 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
774pub 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
806pub 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
827pub 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
848pub fn session_panes() -> HashMap<String, String> {
853 let mut panes = HashMap::new();
854 for pid in codex_pids() {
855 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 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#[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#[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
926pub 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}