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 rollout_session(path: &Path) -> Option<(String, Option<String>)> {
264 rollout_lineage(path)
265}
266
267pub 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
281pub(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 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 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 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
711pub 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
743pub 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 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#[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#[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
817pub 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}