1use std::{
15 collections::{HashMap, HashSet},
16 ffi::OsString,
17 io::Write as _,
18 path::{Path, PathBuf},
19 sync::{
20 Arc, Mutex,
21 atomic::{AtomicBool, Ordering},
22 },
23 time::{Duration, SystemTime, UNIX_EPOCH},
24};
25
26use serde::{Deserialize, Serialize};
27use sha2::{Digest, Sha256};
28
29pub const PARENT_VARIABLE: &str = "SCV_PARENT";
31pub const DEPTH_VARIABLE: &str = "SCV_DELEGATION_DEPTH";
33const STOP_GRACE: Duration = Duration::from_secs(2);
35const MAX_RECORD_BYTES: u64 = 64 * 1024;
37const ZOMBIE_MIN_AGE: Duration = Duration::from_secs(10);
39
40pub fn current_depth() -> u32 {
42 std::env::var(DEPTH_VARIABLE)
43 .ok()
44 .and_then(|value| value.trim().parse().ok())
45 .unwrap_or(0)
46}
47
48#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
50pub struct ProcessIdentity {
51 pub pid: u32,
52 pub start_time: u64,
53}
54
55impl ProcessIdentity {
56 pub fn current() -> Option<Self> {
57 Self::of(std::process::id())
58 }
59
60 pub fn of(pid: u32) -> Option<Self> {
61 process_start_time(pid).map(|start_time| Self { pid, start_time })
62 }
63
64 pub fn is_alive(&self) -> bool {
67 #[cfg(target_os = "linux")]
68 {
69 linux::stat(self.pid)
70 .is_some_and(|info| info.start_time == self.start_time && info.state != 'Z')
71 }
72 #[cfg(not(target_os = "linux"))]
73 {
74 Self::of(self.pid) == Some(*self)
75 }
76 }
77}
78
79#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
81pub struct DelegationRecord {
82 pub handle: String,
83 pub agent: String,
84 pub instance: String,
85 pub session: String,
86 pub owner: ProcessIdentity,
87 pub process: ProcessIdentity,
89 pub pgid: u32,
90 pub cwd: PathBuf,
91 pub started_unix: u64,
92 pub depth: u32,
94 #[serde(default, skip_serializing_if = "Option::is_none")]
96 pub conversation: Option<String>,
97 #[serde(default, skip_serializing_if = "Option::is_none")]
98 pub turn: Option<u32>,
99}
100
101#[derive(Debug, Clone, PartialEq, Eq)]
103pub struct DelegationEntry {
104 pub record: DelegationRecord,
105 pub orphaned: bool,
107 pub processes: usize,
109}
110
111#[derive(Debug, Clone, Default, PartialEq, Eq)]
113pub struct ReconcileReport {
114 pub reaped: Vec<String>,
116 pub removed: usize,
118 pub stale_markers: usize,
120}
121
122#[derive(Debug, Default)]
123struct Inner {
124 active: HashMap<String, Arc<AtomicBool>>,
125 reaped: u64,
126}
127
128#[derive(Debug)]
130pub struct DelegationRegistry {
131 record_dir: PathBuf,
132 instance: String,
133 owner: Option<ProcessIdentity>,
134 depth: u32,
135 chain: Option<String>,
136 inner: Mutex<Inner>,
137}
138
139pub(crate) struct PendingDelegation {
141 pub handle: String,
142 pub environment: Vec<(OsString, OsString)>,
143 agent: String,
144 session: String,
145 cwd: PathBuf,
146 conversation: Option<(String, u32)>,
147}
148
149impl DelegationRegistry {
150 pub fn new(instance_home: &Path) -> Self {
152 let digest = Sha256::digest(instance_home.as_os_str().as_encoded_bytes());
153 let instance = digest[..4]
154 .iter()
155 .map(|byte| format!("{byte:02x}"))
156 .collect();
157 Self {
158 record_dir: instance_home.join("run").join("delegations"),
159 instance,
160 owner: ProcessIdentity::current(),
161 depth: current_depth(),
162 chain: std::env::var(PARENT_VARIABLE)
163 .ok()
164 .filter(|value| !value.trim().is_empty()),
165 inner: Mutex::new(Inner::default()),
166 }
167 }
168
169 pub fn depth(&self) -> u32 {
171 self.depth
172 }
173
174 pub fn record_dir(&self) -> &Path {
175 &self.record_dir
176 }
177
178 pub fn instance(&self) -> &str {
180 &self.instance
181 }
182
183 pub fn reaped_total(&self) -> u64 {
185 self.inner.lock().expect("registry lock").reaped
186 }
187
188 pub fn conversation_dir(&self) -> PathBuf {
190 self.record_dir.parent().map_or_else(
191 || self.record_dir.join("conversations"),
192 |run| run.join("conversations"),
193 )
194 }
195
196 pub(crate) fn begin(
197 &self,
198 agent: &str,
199 session: &str,
200 cwd: &Path,
201 conversation: Option<(&str, u32)>,
202 ) -> PendingDelegation {
203 let suffix = uuid::Uuid::new_v4().simple().to_string();
204 let handle = format!("{agent}-{}", &suffix[..6]);
205 let entry = format!("{}/{session}/{handle}", self.instance);
206 let chain = match &self.chain {
207 Some(chain) => format!("{chain};{entry}"),
208 None => entry,
209 };
210 PendingDelegation {
211 environment: vec![
212 (PARENT_VARIABLE.into(), chain.into()),
213 (DEPTH_VARIABLE.into(), (self.depth + 1).to_string().into()),
214 ],
215 handle,
216 agent: agent.to_owned(),
217 session: session.to_owned(),
218 cwd: cwd.to_owned(),
219 conversation: conversation.map(|(handle, turn)| (handle.to_owned(), turn)),
220 }
221 }
222
223 pub(crate) fn register(
226 self: &Arc<Self>,
227 pending: PendingDelegation,
228 pid: u32,
229 ) -> std::io::Result<DelegationGuard> {
230 let killed = Arc::new(AtomicBool::new(false));
231 let record = DelegationRecord {
232 handle: pending.handle.clone(),
233 agent: pending.agent,
234 instance: self.instance.clone(),
235 session: pending.session,
236 owner: self.owner.unwrap_or(ProcessIdentity {
237 pid: std::process::id(),
238 start_time: 0,
239 }),
240 process: ProcessIdentity::of(pid).unwrap_or(ProcessIdentity { pid, start_time: 0 }),
241 pgid: pid,
242 cwd: pending.cwd,
243 started_unix: SystemTime::now()
244 .duration_since(UNIX_EPOCH)
245 .map_or(0, |elapsed| elapsed.as_secs()),
246 depth: self.depth + 1,
247 conversation: pending
248 .conversation
249 .as_ref()
250 .map(|(handle, _)| handle.clone()),
251 turn: pending.conversation.as_ref().map(|(_, turn)| *turn),
252 };
253 self.inner
254 .lock()
255 .expect("registry lock")
256 .active
257 .insert(record.handle.clone(), Arc::clone(&killed));
258 if let Err(error) = write_record(&self.record_dir, &record) {
259 self.inner
260 .lock()
261 .expect("registry lock")
262 .active
263 .remove(&record.handle);
264 return Err(error);
265 }
266 Ok(DelegationGuard {
267 registry: Arc::clone(self),
268 handle: record.handle,
269 pgid: pid,
270 killed,
271 finished: false,
272 })
273 }
274
275 pub fn list(&self, include_orphans: bool) -> Vec<DelegationEntry> {
278 let table = ProcessTable::snapshot();
279 let mut entries: Vec<_> = self
280 .records()
281 .into_iter()
282 .filter_map(|record| {
283 let orphaned = !self.owner_alive(&record);
284 if orphaned && !include_orphans {
285 return None;
286 }
287 let processes = table.members(&record).len();
288 Some(DelegationEntry {
289 record,
290 orphaned,
291 processes,
292 })
293 })
294 .collect();
295 entries.sort_by(|a, b| {
296 a.record
297 .started_unix
298 .cmp(&b.record.started_unix)
299 .then_with(|| a.record.handle.cmp(&b.record.handle))
300 });
301 entries
302 }
303
304 pub async fn kill(&self, handle: &str) -> Result<(), String> {
306 let record = self
307 .records()
308 .into_iter()
309 .find(|record| record.handle == handle)
310 .ok_or_else(|| format!("no running delegation {handle:?}"))?;
311 let local = self
312 .inner
313 .lock()
314 .expect("registry lock")
315 .active
316 .get(handle)
317 .cloned();
318 if let Some(killed) = &local {
319 killed.store(true, Ordering::Release);
320 }
321 stop_delegation(&record).await;
322 if local.is_none() && !self.owner_alive(&record) {
323 remove_record(&self.record_dir, handle);
324 self.inner.lock().expect("registry lock").reaped += 1;
325 }
326 Ok(())
327 }
328
329 pub async fn reconcile(&self) -> ReconcileReport {
331 let mut report = ReconcileReport::default();
332 for record in self.records() {
333 if self.owner_alive(&record) {
334 continue;
335 }
336 if stop_delegation(&record).await {
337 report.reaped.push(record.handle.clone());
338 } else {
339 report.removed += 1;
340 }
341 remove_record(&self.record_dir, &record.handle);
342 }
343 self.inner.lock().expect("registry lock").reaped += report.reaped.len() as u64;
344 report.stale_markers = crate::conversation::remove_stale_markers(&self.conversation_dir());
345 report
346 }
347
348 fn owner_alive(&self, record: &DelegationRecord) -> bool {
351 if Some(record.owner) == self.owner {
352 return self
353 .inner
354 .lock()
355 .expect("registry lock")
356 .active
357 .contains_key(&record.handle);
358 }
359 record.owner.is_alive()
360 }
361
362 fn records(&self) -> Vec<DelegationRecord> {
363 let Ok(entries) = std::fs::read_dir(&self.record_dir) else {
364 return Vec::new();
365 };
366 entries
367 .filter_map(Result::ok)
368 .filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
369 .filter_map(|entry| read_record(&entry.path()))
370 .filter(|record| record.instance == self.instance)
371 .collect()
372 }
373
374 fn finish_local(&self, handle: &str) {
375 self.inner
376 .lock()
377 .expect("registry lock")
378 .active
379 .remove(handle);
380 remove_record(&self.record_dir, handle);
381 }
382}
383
384pub(crate) struct DelegationGuard {
386 registry: Arc<DelegationRegistry>,
387 handle: String,
388 pgid: u32,
389 killed: Arc<AtomicBool>,
390 finished: bool,
391}
392
393impl DelegationGuard {
394 #[cfg(test)]
395 pub(crate) fn handle(&self) -> &str {
396 &self.handle
397 }
398
399 pub(crate) fn was_killed(&self) -> bool {
401 self.killed.load(Ordering::Acquire)
402 }
403
404 pub(crate) async fn finish(mut self) {
406 self.finished = true;
407 stop_tagged(&self.handle).await;
408 self.registry.finish_local(&self.handle);
409 }
410}
411
412impl Drop for DelegationGuard {
413 fn drop(&mut self) {
414 if self.finished {
415 return;
416 }
417 signal_group(self.pgid, libc::SIGKILL);
420 self.registry.finish_local(&self.handle);
421 let handle = self.handle.clone();
422 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
423 runtime.spawn(async move { stop_tagged(&handle).await });
424 } else {
425 for identity in tagged_processes(&handle) {
426 signal(identity.pid, libc::SIGKILL);
427 }
428 }
429 }
430}
431
432async fn stop_delegation(record: &DelegationRecord) -> bool {
435 let mut stopped = false;
436 let leader = ProcessIdentity::of(record.process.pid);
440 let group_is_ours = record.pgid == record.process.pid
441 && match leader {
442 Some(leader) => leader == record.process,
443 None => group_exists(record.pgid),
444 };
445 if group_is_ours && group_exists(record.pgid) {
446 stopped = true;
447 signal_group(record.pgid, libc::SIGTERM);
448 let deadline = tokio::time::Instant::now() + STOP_GRACE;
449 while group_exists(record.pgid) && tokio::time::Instant::now() < deadline {
450 tokio::time::sleep(Duration::from_millis(50)).await;
451 }
452 signal_group(record.pgid, libc::SIGKILL);
453 }
454 stopped | stop_tagged(&record.handle).await
455}
456
457async fn stop_tagged(handle: &str) -> bool {
459 let tagged = tagged_processes(handle);
460 if tagged.is_empty() {
461 return false;
462 }
463 for identity in &tagged {
464 signal(identity.pid, libc::SIGTERM);
465 }
466 let deadline = tokio::time::Instant::now() + STOP_GRACE;
467 while tagged.iter().any(ProcessIdentity::is_alive) && tokio::time::Instant::now() < deadline {
468 tokio::time::sleep(Duration::from_millis(50)).await;
469 }
470 for identity in tagged.iter().filter(|identity| identity.is_alive()) {
471 signal(identity.pid, libc::SIGKILL);
472 }
473 true
474}
475
476fn tagged_processes(handle: &str) -> Vec<ProcessIdentity> {
478 let own = std::process::id();
479 ProcessTable::snapshot()
480 .tagged
481 .into_iter()
482 .filter(|(identity, chain)| identity.pid != own && chain_names(chain, handle))
483 .map(|(identity, _)| identity)
484 .collect()
485}
486
487fn chain_names(chain: &str, handle: &str) -> bool {
488 chain
489 .split(';')
490 .any(|entry| entry.rsplit('/').next() == Some(handle))
491}
492
493fn signal(pid: u32, signal: i32) {
494 if let Ok(pid) = i32::try_from(pid)
495 && pid > 0
496 {
497 unsafe {
498 libc::kill(pid, signal);
499 }
500 }
501}
502
503fn signal_group(pgid: u32, signal: i32) {
504 if let Ok(pgid) = i32::try_from(pgid)
506 && pgid > 1
507 {
508 unsafe {
509 libc::kill(-pgid, signal);
510 }
511 }
512}
513
514fn group_exists(pgid: u32) -> bool {
516 let Ok(group) = i32::try_from(pgid) else {
517 return false;
518 };
519 if group <= 1 {
520 return false;
521 }
522 let result = unsafe { libc::kill(-group, 0) };
523 let signalable =
524 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
525 #[cfg(target_os = "linux")]
526 {
527 signalable
528 && linux::all_stats()
529 .iter()
530 .any(|info| info.pgid == pgid && info.state != 'Z')
531 }
532 #[cfg(not(target_os = "linux"))]
533 {
534 signalable
535 }
536}
537
538fn write_record(dir: &Path, record: &DelegationRecord) -> std::io::Result<()> {
539 write_private_json(dir, &format!("{}.json", record.handle), record)
540}
541
542pub(crate) fn write_private_json(
545 dir: &Path,
546 name: &str,
547 value: &impl Serialize,
548) -> std::io::Result<()> {
549 use std::os::unix::fs::{OpenOptionsExt as _, PermissionsExt as _};
550 std::fs::create_dir_all(dir)?;
551 if let Some(run) = dir.parent() {
552 std::fs::set_permissions(run, std::fs::Permissions::from_mode(0o700))?;
553 }
554 std::fs::set_permissions(dir, std::fs::Permissions::from_mode(0o700))?;
555 let bytes = serde_json::to_vec_pretty(value).map_err(std::io::Error::other)?;
556 let temporary = dir.join(format!(".{name}.tmp"));
557 let mut file = std::fs::OpenOptions::new()
558 .write(true)
559 .create(true)
560 .truncate(true)
561 .mode(0o600)
562 .open(&temporary)?;
563 file.write_all(&bytes)?;
564 file.sync_all()?;
565 drop(file);
566 std::fs::rename(&temporary, dir.join(name))
567}
568
569fn read_record(path: &Path) -> Option<DelegationRecord> {
570 let file = std::fs::File::open(path).ok()?;
571 let mut bytes = Vec::new();
572 std::io::Read::read_to_end(&mut std::io::Read::take(file, MAX_RECORD_BYTES), &mut bytes)
573 .ok()?;
574 let record: DelegationRecord = serde_json::from_slice(&bytes).ok()?;
575 (path.file_stem().and_then(|stem| stem.to_str()) == Some(record.handle.as_str()))
577 .then_some(record)
578}
579
580fn remove_record(dir: &Path, handle: &str) {
581 let _ = std::fs::remove_file(dir.join(format!("{handle}.json")));
582}
583
584pub fn become_child_subreaper() -> bool {
587 #[cfg(target_os = "linux")]
588 {
589 unsafe { libc::prctl(libc::PR_SET_CHILD_SUBREAPER, 1, 0, 0, 0) == 0 }
590 }
591 #[cfg(not(target_os = "linux"))]
592 {
593 false
594 }
595}
596
597static SPAWNED: Mutex<Option<HashSet<u32>>> = Mutex::new(None);
598
599pub(crate) fn track_spawned(pid: u32) {
601 SPAWNED
602 .lock()
603 .expect("spawned lock")
604 .get_or_insert_with(HashSet::new)
605 .insert(pid);
606}
607
608pub(crate) fn untrack_spawned(pid: u32) {
609 if let Some(spawned) = SPAWNED.lock().expect("spawned lock").as_mut() {
610 spawned.remove(&pid);
611 }
612}
613
614pub fn reap_orphaned_zombies() -> usize {
617 #[cfg(target_os = "linux")]
618 {
619 let own = std::process::id();
620 let spawned = SPAWNED
621 .lock()
622 .expect("spawned lock")
623 .clone()
624 .unwrap_or_default();
625 let uptime = linux::uptime_ticks();
626 let mut reaped = 0;
627 for info in linux::all_stats() {
628 if info.ppid != own || info.state != 'Z' || spawned.contains(&info.pid) {
629 continue;
630 }
631 let old_enough = uptime.is_some_and(|now| {
632 now.saturating_sub(info.start_time)
633 >= ZOMBIE_MIN_AGE.as_secs() * linux::clock_ticks()
634 });
635 if !old_enough {
636 continue;
637 }
638 let mut status = 0;
639 if unsafe { libc::waitpid(info.pid as i32, &mut status, libc::WNOHANG) }
640 == info.pid as i32
641 {
642 reaped += 1;
643 }
644 }
645 reaped
646 }
647 #[cfg(not(target_os = "linux"))]
648 {
649 0
650 }
651}
652
653struct ProcessTable {
655 groups: Vec<(ProcessIdentity, u32)>,
656 tagged: Vec<(ProcessIdentity, String)>,
657}
658
659impl ProcessTable {
660 fn members(&self, record: &DelegationRecord) -> HashSet<u32> {
661 let mut members: HashSet<u32> = self
662 .groups
663 .iter()
664 .filter(|(_, pgid)| *pgid == record.pgid)
665 .map(|(identity, _)| identity.pid)
666 .collect();
667 members.extend(
668 self.tagged
669 .iter()
670 .filter(|(_, chain)| chain_names(chain, &record.handle))
671 .map(|(identity, _)| identity.pid),
672 );
673 members
674 }
675
676 #[cfg(target_os = "linux")]
677 fn snapshot() -> Self {
678 let mut groups = Vec::new();
679 let mut tagged = Vec::new();
680 for info in linux::all_stats() {
681 if info.state == 'Z' {
682 continue;
683 }
684 let identity = ProcessIdentity {
685 pid: info.pid,
686 start_time: info.start_time,
687 };
688 groups.push((identity, info.pgid));
689 if let Some(chain) = linux::parent_chain(info.pid) {
690 tagged.push((identity, chain));
691 }
692 }
693 Self { groups, tagged }
694 }
695
696 #[cfg(not(target_os = "linux"))]
697 fn snapshot() -> Self {
698 let mut groups = Vec::new();
699 let mut tagged = Vec::new();
700 let Ok(output) = std::process::Command::new("ps")
702 .args(["-E", "-ww", "-axo", "pid=,pgid=,command="])
703 .output()
704 else {
705 return Self { groups, tagged };
706 };
707 for line in String::from_utf8_lossy(&output.stdout).lines() {
708 let mut fields = line.split_whitespace();
709 let (Some(pid), Some(pgid)) = (
710 fields.next().and_then(|value| value.parse::<u32>().ok()),
711 fields.next().and_then(|value| value.parse::<u32>().ok()),
712 ) else {
713 continue;
714 };
715 let Some(identity) = ProcessIdentity::of(pid) else {
716 continue;
717 };
718 groups.push((identity, pgid));
719 if let Some(chain) = fields.find_map(|field| {
720 field
721 .strip_prefix(PARENT_VARIABLE)
722 .and_then(|rest| rest.strip_prefix('='))
723 }) {
724 tagged.push((identity, chain.to_owned()));
725 }
726 }
727 Self { groups, tagged }
728 }
729}
730
731#[cfg(target_os = "linux")]
732fn process_start_time(pid: u32) -> Option<u64> {
733 linux::stat(pid).map(|info| info.start_time)
734}
735
736#[cfg(target_os = "macos")]
737fn process_start_time(pid: u32) -> Option<u64> {
738 let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
739 let size = std::mem::size_of::<libc::proc_bsdinfo>() as i32;
740 let written = unsafe {
741 libc::proc_pidinfo(
742 pid as i32,
743 libc::PROC_PIDTBSDINFO,
744 0,
745 (&mut info as *mut libc::proc_bsdinfo).cast(),
746 size,
747 )
748 };
749 (written == size).then(|| info.pbi_start_tvsec * 1_000_000 + info.pbi_start_tvusec)
750}
751
752#[cfg(not(any(target_os = "linux", target_os = "macos")))]
753fn process_start_time(pid: u32) -> Option<u64> {
754 let alive = unsafe { libc::kill(pid as i32, 0) } == 0;
755 alive.then_some(0)
756}
757
758#[cfg(target_os = "linux")]
759mod linux {
760 pub(super) struct Stat {
761 pub pid: u32,
762 pub ppid: u32,
763 pub pgid: u32,
764 pub state: char,
765 pub start_time: u64,
766 }
767
768 pub(super) fn stat(pid: u32) -> Option<Stat> {
769 let text = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
770 let rest = &text[text.rfind(')')? + 2..];
772 let fields: Vec<&str> = rest.split_whitespace().collect();
773 Some(Stat {
775 pid,
776 state: fields.first()?.chars().next()?,
777 ppid: fields.get(1)?.parse().ok()?,
778 pgid: fields.get(2)?.parse().ok()?,
779 start_time: fields.get(19)?.parse().ok()?,
780 })
781 }
782
783 pub(super) fn all_stats() -> Vec<Stat> {
784 let Ok(entries) = std::fs::read_dir("/proc") else {
785 return Vec::new();
786 };
787 entries
788 .filter_map(Result::ok)
789 .filter_map(|entry| entry.file_name().to_str()?.parse::<u32>().ok())
790 .filter_map(stat)
791 .collect()
792 }
793
794 pub(super) fn parent_chain(pid: u32) -> Option<String> {
796 let environ = std::fs::read(format!("/proc/{pid}/environ")).ok()?;
797 let prefix = format!("{}=", super::PARENT_VARIABLE);
798 environ.split(|byte| *byte == 0).find_map(|entry| {
799 entry
800 .strip_prefix(prefix.as_bytes())
801 .map(|value| String::from_utf8_lossy(value).into_owned())
802 })
803 }
804
805 pub(super) fn clock_ticks() -> u64 {
806 let ticks = unsafe { libc::sysconf(libc::_SC_CLK_TCK) };
807 u64::try_from(ticks)
808 .ok()
809 .filter(|ticks| *ticks > 0)
810 .unwrap_or(100)
811 }
812
813 pub(super) fn uptime_ticks() -> Option<u64> {
814 let text = std::fs::read_to_string("/proc/uptime").ok()?;
815 let seconds: f64 = text.split_whitespace().next()?.parse().ok()?;
816 Some((seconds * clock_ticks() as f64) as u64)
817 }
818}
819
820#[cfg(test)]
821mod tests {
822 use super::*;
823 use std::os::unix::{fs::PermissionsExt as _, process::CommandExt as _};
824
825 fn registry(home: &Path) -> Arc<DelegationRegistry> {
826 Arc::new(DelegationRegistry::new(home))
827 }
828
829 fn spawn_tagged(script: &str, environment: &[(OsString, OsString)]) -> std::process::Child {
831 std::process::Command::new("sh")
832 .args(["-c", script])
833 .envs(environment.iter().map(|(key, value)| (key, value)))
834 .stdin(std::process::Stdio::null())
835 .stdout(std::process::Stdio::null())
836 .stderr(std::process::Stdio::null())
837 .process_group(0)
838 .spawn()
839 .unwrap()
840 }
841
842 async fn wait_for(mut condition: impl FnMut() -> bool) -> bool {
843 for _ in 0..200 {
844 if condition() {
845 return true;
846 }
847 tokio::time::sleep(Duration::from_millis(25)).await;
848 }
849 false
850 }
851
852 #[test]
853 fn chains_match_only_their_own_handle() {
854 assert!(chain_names("abcd/s1/codex-1a2b3c", "codex-1a2b3c"));
855 assert!(chain_names(
856 "x/s/claude-000000;abcd/s1/codex-1a2b3c",
857 "codex-1a2b3c"
858 ));
859 assert!(!chain_names("abcd/s1/codex-1a2b3c", "codex-1a2b3"));
860 assert!(!chain_names("abcd/s1/codex-1a2b3c", "1a2b3c"));
861 }
862
863 #[test]
864 fn nested_tags_extend_the_chain_and_depth() {
865 let home = tempfile::tempdir().unwrap();
866 let mut registry = DelegationRegistry::new(home.path());
867 registry.chain = Some("aaaa/s0/codex-111111".into());
868 registry.depth = 1;
869 let pending = registry.begin("claude", "s1", home.path(), Some(("claude-1", 2)));
870 let value = |name: &str| {
871 pending
872 .environment
873 .iter()
874 .find(|(key, _)| key == name)
875 .map(|(_, value)| value.to_str().unwrap().to_owned())
876 .unwrap()
877 };
878 assert_eq!(
879 value(PARENT_VARIABLE),
880 format!(
881 "aaaa/s0/codex-111111;{}/s1/{}",
882 registry.instance, pending.handle
883 )
884 );
885 assert_eq!(value(DEPTH_VARIABLE), "2");
886 assert!(pending.handle.starts_with("claude-"));
887 }
888
889 #[tokio::test]
890 async fn records_are_private_and_removed_when_the_run_finishes() {
891 let home = tempfile::tempdir().unwrap();
892 let registry = registry(home.path());
893 let pending = registry.begin("codex", "session", home.path(), None);
894 let environment = pending.environment.clone();
895 let mut child = spawn_tagged("sleep 30", &environment);
896 let guard = registry.register(pending, child.id()).unwrap();
897 let path = registry
898 .record_dir()
899 .join(format!("{}.json", guard.handle()));
900 let mode = |path: &Path| std::fs::metadata(path).unwrap().permissions().mode() & 0o777;
901 assert_eq!(mode(&path), 0o600);
902 assert_eq!(mode(registry.record_dir()), 0o700);
903 assert_eq!(mode(registry.record_dir().parent().unwrap()), 0o700);
904
905 let listed = registry.list(false);
906 assert_eq!(listed.len(), 1);
907 assert!(!listed[0].orphaned);
908 assert_eq!(listed[0].record.process.pid, child.id());
909 assert!(listed[0].processes >= 1);
910
911 signal_group(child.id(), libc::SIGKILL);
912 child.wait().unwrap();
913 guard.finish().await;
914 assert!(!path.exists());
915 assert!(registry.list(true).is_empty());
916 }
917
918 #[tokio::test]
919 async fn kill_stops_a_local_run_and_marks_it_killed() {
920 let home = tempfile::tempdir().unwrap();
921 let registry = registry(home.path());
922 let pending = registry.begin("claude", "session", home.path(), None);
923 let environment = pending.environment.clone();
924 let mut child = spawn_tagged("trap '' TERM; sleep 30", &environment);
925 let guard = registry.register(pending, child.id()).unwrap();
926 registry.kill(guard.handle()).await.unwrap();
927 assert!(guard.was_killed());
928 assert!(child.wait().unwrap().code().is_none(), "killed by a signal");
929 assert!(registry.kill("claude-nosuch").await.is_err());
930 guard.finish().await;
931 }
932
933 #[cfg(target_os = "linux")]
934 #[tokio::test]
935 async fn reconcile_removes_conversation_markers_of_exited_processes() {
936 let home = tempfile::tempdir().unwrap();
937 let daemon = registry(home.path());
938 let markers = daemon.conversation_dir();
939 let mut gone = std::process::Command::new("true").spawn().unwrap();
940 let gone_pid = gone.id();
941 gone.wait().unwrap();
942 let dead = ProcessIdentity {
943 pid: gone_pid,
944 start_time: 1,
945 };
946 let live = ProcessIdentity::current().unwrap();
947 for (id, owner) in [("dead-id", dead), ("live-id", live)] {
948 let marker = serde_json::json!({"owner": owner, "agent": "codex", "handle": "codex-1"});
949 write_private_json(&markers, &format!("{id}.json"), &marker).unwrap();
950 }
951 let report = daemon.reconcile().await;
952 assert_eq!(report.stale_markers, 1);
953 assert!(!markers.join("dead-id.json").exists());
954 assert!(markers.join("live-id.json").is_file());
955 assert_eq!(daemon.reconcile().await, ReconcileReport::default());
956 }
957
958 #[cfg(target_os = "linux")]
959 #[tokio::test]
960 async fn reconcile_reaps_an_orphan_and_its_detached_descendants() {
961 let home = tempfile::tempdir().unwrap();
962 let owner = registry(home.path());
963 let pending = owner.begin("codex", "session", home.path(), None);
964 let environment = pending.environment.clone();
965 let mut child = spawn_tagged("setsid sleep 60 & exec sleep 60", &environment);
967 let guard = owner.register(pending, child.id()).unwrap();
968 let handle = guard.handle().to_owned();
969 assert!(wait_for(|| tagged_processes(&handle).len() >= 2).await);
970 let path = owner.record_dir().join(format!("{handle}.json"));
971 let mut record = read_record(&path).unwrap();
973 let mut gone = std::process::Command::new("true").spawn().unwrap();
974 let gone_pid = gone.id();
975 gone.wait().unwrap();
976 record.owner = ProcessIdentity {
977 pid: gone_pid,
978 start_time: 1,
979 };
980 write_record(owner.record_dir(), &record).unwrap();
981 std::mem::forget(guard);
982
983 let daemon = registry(home.path());
985 assert!(daemon.list(false).is_empty());
986 let orphans = daemon.list(true);
987 assert_eq!(orphans.len(), 1);
988 assert!(orphans[0].orphaned);
989 assert!(orphans[0].processes >= 2);
990 let report = daemon.reconcile().await;
991 assert_eq!(report.reaped, vec![handle.clone()]);
992 assert_eq!(daemon.reaped_total(), 1);
993 assert!(child.wait().unwrap().code().is_none());
994 assert!(wait_for(|| tagged_processes(&handle).is_empty()).await);
995 assert!(!path.exists());
996 assert_eq!(daemon.reconcile().await, ReconcileReport::default());
997 }
998
999 #[tokio::test]
1000 async fn an_abandoned_run_is_cleaned_up_when_its_guard_drops() {
1001 let home = tempfile::tempdir().unwrap();
1002 let registry = registry(home.path());
1003 let pending = registry.begin("pi", "session", home.path(), None);
1004 let environment = pending.environment.clone();
1005 let mut child = spawn_tagged("sleep 30", &environment);
1006 let guard = registry.register(pending, child.id()).unwrap();
1007 let path = registry
1008 .record_dir()
1009 .join(format!("{}.json", guard.handle()));
1010 drop(guard);
1011 assert!(child.wait().unwrap().code().is_none());
1012 assert!(!path.exists());
1013 }
1014
1015 #[test]
1016 fn records_for_another_instance_or_under_the_wrong_name_are_ignored() {
1017 let home = tempfile::tempdir().unwrap();
1018 let registry = DelegationRegistry::new(home.path());
1019 let dir = registry.record_dir().to_owned();
1020 let record = DelegationRecord {
1021 handle: "codex-abcdef".into(),
1022 agent: "codex".into(),
1023 instance: "other".into(),
1024 session: "s".into(),
1025 owner: ProcessIdentity {
1026 pid: 1,
1027 start_time: 1,
1028 },
1029 process: ProcessIdentity {
1030 pid: 1,
1031 start_time: 1,
1032 },
1033 pgid: 1,
1034 cwd: "/".into(),
1035 started_unix: 0,
1036 depth: 1,
1037 conversation: None,
1038 turn: None,
1039 };
1040 write_record(&dir, &record).unwrap();
1041 std::fs::copy(
1042 dir.join("codex-abcdef.json"),
1043 dir.join("codex-renamed.json"),
1044 )
1045 .unwrap();
1046 assert!(registry.list(true).is_empty());
1047 }
1048}