Skip to main content

scv_tools/delegate/
records.rs

1//! What SCV started, so it can list, stop, and clean up delegated agents.
2//!
3//! Every delegated process is tagged through its environment
4//! (`SCV_PARENT=<instance>/<session>/<handle>`, chained through nested SCVs,
5//! and `SCV_DELEGATION_DEPTH`) and recorded in
6//! `$SCV_HOME/state/delegations/<handle>.json` while it runs. A record whose
7//! owning SCV process died is an orphan: the daemon's reconciliation kills its
8//! process group and anything still carrying its tag, then removes it.
9//!
10//! This is cooperative bookkeeping. Delegated agents run as the user, so one
11//! that deliberately clears its environment or leaves its process group can
12//! escape it; the tags and records exist to clean up accidental leaks.
13
14use std::{
15    collections::{HashMap, HashSet},
16    ffi::OsString,
17    path::{Path, PathBuf},
18    sync::{
19        Arc, Mutex,
20        atomic::{AtomicBool, Ordering},
21    },
22    time::{Duration, SystemTime, UNIX_EPOCH},
23};
24
25use scv_client::Layout;
26use serde::{Deserialize, Serialize};
27use sha2::{Digest, Sha256};
28
29use crate::{process::ProcessGroup, sync::lock};
30
31/// Environment variable carrying the delegation chain.
32pub const PARENT_VARIABLE: &str = "SCV_PARENT";
33/// Environment variable carrying how deeply this process is delegated.
34pub use scv_client::DELEGATION_DEPTH_VARIABLE as DEPTH_VARIABLE;
35/// Grace between TERM and KILL when stopping delegated processes.
36pub(crate) const STOP_GRACE: Duration = Duration::from_secs(2);
37/// Largest record file read.
38const MAX_RECORD_BYTES: u64 = 64 * 1024;
39/// A zombie child younger than this may still be awaited by its spawner.
40const ZOMBIE_MIN_AGE: Duration = Duration::from_secs(10);
41
42/// Delegation depth of the current process: 0 unless an SCV started it.
43pub fn current_depth() -> u32 {
44    #[cfg(test)]
45    {
46        0
47    }
48    #[cfg(not(test))]
49    {
50        scv_client::inherited_delegation_depth().unwrap_or(0)
51    }
52}
53
54/// A process, identified by PID plus start time so a reused PID never matches.
55#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
56pub struct ProcessIdentity {
57    pub pid: u32,
58    pub(crate) start_time: u64,
59}
60
61impl ProcessIdentity {
62    pub fn current() -> Option<Self> {
63        Self::of(std::process::id())
64    }
65
66    pub fn of(pid: u32) -> Option<Self> {
67        process_start_time(pid).map(|start_time| Self { pid, start_time })
68    }
69
70    /// Whether this exact process still runs. An exited process that its
71    /// parent has not yet collected (a zombie) does not count.
72    pub fn is_alive(&self) -> bool {
73        #[cfg(target_os = "linux")]
74        {
75            linux::stat(self.pid)
76                .is_some_and(|info| info.start_time == self.start_time && info.state != 'Z')
77        }
78        #[cfg(not(target_os = "linux"))]
79        {
80            Self::of(self.pid) == Some(*self) && !process_is_zombie(self.pid)
81        }
82    }
83}
84
85/// One delegated run, as recorded on disk.
86#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
87pub struct DelegationRecord {
88    pub handle: String,
89    pub agent: String,
90    pub instance: String,
91    pub session: String,
92    pub owner: ProcessIdentity,
93    /// The agent process, which leads its own process group.
94    pub process: ProcessIdentity,
95    pub pgid: u32,
96    pub cwd: PathBuf,
97    pub started_unix: u64,
98    /// Depth of the delegated process (the owner's depth plus one).
99    pub depth: u32,
100    /// The conversation this run is a turn of, and which turn.
101    #[serde(default, skip_serializing_if = "Option::is_none")]
102    pub conversation: Option<String>,
103    #[serde(default, skip_serializing_if = "Option::is_none")]
104    pub turn: Option<u32>,
105    /// A live child that keeps its process between turns (a nested SCV or an
106    /// ACP agent) and has no turn running: when its last turn ended. Absent
107    /// while a turn runs, and always for a per-turn run.
108    #[serde(default, skip_serializing_if = "Option::is_none")]
109    pub idle_since_unix: Option<u64>,
110    /// A nested SCV's own background jobs that still run or whose results
111    /// its model has not yet taken in: finished and not yet reported, or in
112    /// the turn it started to report them, until that turn ends. Absent when
113    /// there are none, for every other agent, and from older releases.
114    #[serde(default, skip_serializing_if = "Option::is_none")]
115    pub background_jobs: Option<u32>,
116}
117
118/// A record plus what SCV currently observes about it.
119#[derive(Debug, Clone, PartialEq, Eq)]
120pub struct DelegationEntry {
121    pub record: DelegationRecord,
122    /// The owning SCV process is gone; reconciliation will clean it up.
123    pub orphaned: bool,
124    /// Live processes in its group plus tagged processes outside it.
125    pub processes: usize,
126}
127
128impl DelegationEntry {
129    /// Whether the run is still at work: it has live processes and is not a
130    /// live child waiting between turns of its conversation with no
131    /// background jobs of its own.
132    pub fn working(&self) -> bool {
133        self.processes > 0
134            && (self.record.idle_since_unix.is_none()
135                || self.record.background_jobs.is_some_and(|jobs| jobs > 0))
136    }
137}
138
139/// A delegation named in an `SCV_PARENT` chain, and the session that
140/// started it.
141#[derive(Debug, Clone, PartialEq, Eq)]
142pub struct ChainRun {
143    /// The delegation's handle, such as `codex-3f9a2c`.
144    pub handle: String,
145    /// The SCV session that started it.
146    pub session: String,
147}
148
149/// What a reconciliation pass did.
150#[derive(Debug, Clone, Default, PartialEq, Eq)]
151pub struct ReconcileReport {
152    /// Orphaned delegations whose processes were stopped.
153    pub reaped: Vec<String>,
154    /// Orphaned records whose processes had already exited.
155    pub removed: usize,
156    /// Conversation markers left by SCV processes that no longer run.
157    pub stale_markers: usize,
158}
159
160#[derive(Debug, Default)]
161struct Inner {
162    active: HashMap<String, Arc<AtomicBool>>,
163    reaped: u64,
164}
165
166/// The delegations one SCV process started, backed by the instance's records.
167#[derive(Debug)]
168pub struct DelegationRegistry {
169    record_dir: PathBuf,
170    conversation_dir: PathBuf,
171    instance: String,
172    owner: Option<ProcessIdentity>,
173    depth: u32,
174    chain: Option<String>,
175    inner: Mutex<Inner>,
176}
177
178/// A delegation about to start: its handle and the environment tagging it.
179pub(crate) struct PendingDelegation {
180    pub(crate) handle: String,
181    pub(crate) environment: Vec<(OsString, OsString)>,
182    agent: String,
183    session: String,
184    cwd: PathBuf,
185    conversation: Option<(String, u32)>,
186    /// Depth of the delegated process.
187    depth: u32,
188}
189
190impl DelegationRegistry {
191    /// The registry for the SCV instance at `layout`: records in
192    /// [`Layout::delegations`], conversation markers in
193    /// [`Layout::conversations`], and an instance ID hashed from its home.
194    pub fn new(layout: &Layout) -> Self {
195        let digest = Sha256::digest(layout.home().as_os_str().as_encoded_bytes());
196        let instance = digest[..4]
197            .iter()
198            .map(|byte| format!("{byte:02x}"))
199            .collect();
200        Self {
201            record_dir: layout.delegations(),
202            conversation_dir: layout.conversations(),
203            instance,
204            owner: ProcessIdentity::current(),
205            depth: current_depth(),
206            chain: {
207                #[cfg(test)]
208                {
209                    None
210                }
211                #[cfg(not(test))]
212                {
213                    std::env::var(PARENT_VARIABLE)
214                        .ok()
215                        .filter(|value| !value.trim().is_empty())
216                }
217            },
218            inner: Mutex::new(Inner::default()),
219        }
220    }
221
222    /// A unit-test fixture whose caller is independent of the test runner.
223    /// Production construction always reads the inherited depth and chain.
224    #[cfg(test)]
225    pub(crate) fn for_test(layout: &Layout) -> Self {
226        Self {
227            depth: 0,
228            chain: None,
229            ..Self::new(layout)
230        }
231    }
232
233    /// This process's own delegation depth.
234    pub(crate) fn depth(&self) -> u32 {
235        self.depth
236    }
237
238    pub fn record_dir(&self) -> &Path {
239        &self.record_dir
240    }
241
242    /// Short identifier of the SCV instance, shared by all its processes.
243    pub fn instance(&self) -> &str {
244        &self.instance
245    }
246
247    /// Delegations this process stopped as orphans since it started.
248    pub fn reaped_total(&self) -> u64 {
249        lock(&self.inner).reaped
250    }
251
252    /// Where live conversations leave markers for `scv agents gc`.
253    pub fn conversation_dir(&self) -> &Path {
254        &self.conversation_dir
255    }
256
257    #[cfg(test)]
258    pub(crate) fn begin(
259        &self,
260        agent: &str,
261        session: &str,
262        cwd: &Path,
263        conversation: Option<(&str, u32)>,
264    ) -> PendingDelegation {
265        self.begin_at(self.depth, agent, session, cwd, conversation)
266    }
267
268    /// Start recording a delegation whose owner is at `owner_depth`: the
269    /// process's own depth, or more when its client is itself delegated.
270    pub(crate) fn begin_at(
271        &self,
272        owner_depth: u32,
273        agent: &str,
274        session: &str,
275        cwd: &Path,
276        conversation: Option<(&str, u32)>,
277    ) -> PendingDelegation {
278        let suffix = uuid::Uuid::new_v4().simple().to_string();
279        let handle = format!("{agent}-{}", &suffix[..6]);
280        let entry = format!("{}/{session}/{handle}", self.instance);
281        let chain = match &self.chain {
282            Some(chain) => format!("{chain};{entry}"),
283            None => entry,
284        };
285        PendingDelegation {
286            environment: vec![
287                (PARENT_VARIABLE.into(), chain.into()),
288                (
289                    DEPTH_VARIABLE.into(),
290                    owner_depth.saturating_add(1).to_string().into(),
291                ),
292            ],
293            handle,
294            agent: agent.to_owned(),
295            session: session.to_owned(),
296            cwd: cwd.to_owned(),
297            conversation: conversation.map(|(handle, turn)| (handle.to_owned(), turn)),
298            depth: owner_depth.saturating_add(1),
299        }
300    }
301
302    /// Record a spawned delegation. The returned guard removes the record and
303    /// stops leftovers when the run ends, even if the run is abandoned.
304    pub(crate) fn register(
305        self: &Arc<Self>,
306        pending: PendingDelegation,
307        pid: u32,
308    ) -> std::io::Result<DelegationGuard> {
309        let killed = Arc::new(AtomicBool::new(false));
310        let record = DelegationRecord {
311            handle: pending.handle.clone(),
312            agent: pending.agent,
313            instance: self.instance.clone(),
314            session: pending.session,
315            owner: self.owner.unwrap_or(ProcessIdentity {
316                pid: std::process::id(),
317                start_time: 0,
318            }),
319            process: ProcessIdentity::of(pid).unwrap_or(ProcessIdentity { pid, start_time: 0 }),
320            pgid: pid,
321            cwd: pending.cwd,
322            started_unix: unix_now(),
323            depth: pending.depth,
324            conversation: pending
325                .conversation
326                .as_ref()
327                .map(|(handle, _)| handle.clone()),
328            turn: pending.conversation.as_ref().map(|(_, turn)| *turn),
329            idle_since_unix: None,
330            background_jobs: None,
331        };
332        lock(&self.inner)
333            .active
334            .insert(record.handle.clone(), Arc::clone(&killed));
335        if let Err(error) = write_record(&self.record_dir, &record) {
336            lock(&self.inner).active.remove(&record.handle);
337            return Err(error);
338        }
339        Ok(DelegationGuard {
340            registry: Arc::clone(self),
341            handle: record.handle,
342            pgid: pid,
343            killed,
344            finished: false,
345        })
346    }
347
348    /// Delegations of this instance that are still running. With
349    /// `include_orphans`, also records whose owner died and await cleanup.
350    pub fn list(&self, include_orphans: bool) -> Vec<DelegationEntry> {
351        let table = ProcessTable::snapshot();
352        let mut entries: Vec<_> = self
353            .records()
354            .into_iter()
355            .filter_map(|record| {
356                let orphaned = !self.owner_alive(&record);
357                if orphaned && !include_orphans {
358                    return None;
359                }
360                let processes = table.members(&record).len();
361                Some(DelegationEntry {
362                    record,
363                    orphaned,
364                    processes,
365                })
366            })
367            .collect();
368        entries.sort_by(|a, b| {
369            a.record
370                .started_unix
371                .cmp(&b.record.started_unix)
372                .then_with(|| a.record.handle.cmp(&b.record.handle))
373        });
374        entries
375    }
376
377    /// The running delegation this process started that an `SCV_PARENT`
378    /// `chain` names: the one a caller runs inside, whatever nested SCVs lie
379    /// between. Entries of other instances, runs other processes own, and
380    /// malformed entries are skipped.
381    pub fn own_run(&self, chain: &str) -> Option<ChainRun> {
382        let own = std::process::id();
383        let entries = self.list(true);
384        chain.split(';').find_map(|entry| {
385            let mut parts = entry.splitn(3, '/');
386            let (instance, session, handle) = (parts.next()?, parts.next()?, parts.next()?);
387            if instance != self.instance {
388                return None;
389            }
390            entries
391                .iter()
392                .find(|running| running.record.handle == handle && running.record.owner.pid == own)
393                .map(|_| ChainRun {
394                    handle: handle.to_owned(),
395                    session: session.to_owned(),
396                })
397        })
398    }
399
400    /// Stop one delegation of this instance, whichever process owns it.
401    pub async fn kill(&self, handle: &str) -> Result<(), String> {
402        let record = self
403            .records()
404            .into_iter()
405            .find(|record| record.handle == handle)
406            .ok_or_else(|| format!("no running delegation {handle:?}"))?;
407        let local = lock(&self.inner).active.get(handle).cloned();
408        if let Some(killed) = &local {
409            killed.store(true, Ordering::Release);
410        }
411        stop_delegation(&record).await;
412        if local.is_none() && !self.owner_alive(&record) {
413            remove_record(&self.record_dir, handle);
414            lock(&self.inner).reaped += 1;
415        }
416        Ok(())
417    }
418
419    /// Stop and remove every orphaned delegation of this instance.
420    pub async fn reconcile(&self) -> ReconcileReport {
421        let mut report = ReconcileReport::default();
422        for record in self.records() {
423            if self.owner_alive(&record) {
424                continue;
425            }
426            if stop_delegation(&record).await {
427                report.reaped.push(record.handle.clone());
428            } else {
429                report.removed += 1;
430            }
431            remove_record(&self.record_dir, &record.handle);
432        }
433        lock(&self.inner).reaped += report.reaped.len() as u64;
434        report.stale_markers =
435            crate::delegate::conversation::remove_stale_markers(self.conversation_dir());
436        report
437    }
438
439    /// Whether the process that owns `record` still runs it. A record this
440    /// process owns counts only while its run is active here.
441    fn owner_alive(&self, record: &DelegationRecord) -> bool {
442        if Some(record.owner) == self.owner {
443            return lock(&self.inner).active.contains_key(&record.handle);
444        }
445        record.owner.is_alive()
446    }
447
448    fn records(&self) -> Vec<DelegationRecord> {
449        let Ok(entries) = std::fs::read_dir(&self.record_dir) else {
450            return Vec::new();
451        };
452        entries
453            .filter_map(Result::ok)
454            .filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
455            .filter_map(|entry| read_record(&entry.path()))
456            .filter(|record| record.instance == self.instance)
457            .collect()
458    }
459
460    fn finish_local(&self, handle: &str) {
461        lock(&self.inner).active.remove(handle);
462        remove_record(&self.record_dir, handle);
463    }
464}
465
466/// Keeps a delegation recorded while it runs.
467pub(crate) struct DelegationGuard {
468    registry: Arc<DelegationRegistry>,
469    handle: String,
470    pgid: u32,
471    killed: Arc<AtomicBool>,
472    finished: bool,
473}
474
475impl DelegationGuard {
476    #[cfg(test)]
477    pub(crate) fn handle(&self) -> &str {
478        &self.handle
479    }
480
481    /// Record that the run moved on to `turn` of its conversation and is at
482    /// work, for a live child that serves every turn.
483    pub(crate) fn set_turn(&self, turn: u32) {
484        self.update(|record| {
485            record.turn = Some(turn);
486            record.idle_since_unix = None;
487        });
488    }
489
490    /// Record that a live child ended its turn and waits for the next one.
491    pub(crate) fn set_idle(&self) {
492        let now = unix_now();
493        self.update(|record| record.idle_since_unix = Some(now));
494    }
495
496    /// Record how many of a live child's own background jobs still run or
497    /// wait to be reported to it.
498    pub(crate) fn set_background_jobs(&self, jobs: usize) {
499        let jobs = (jobs > 0).then(|| u32::try_from(jobs).unwrap_or(u32::MAX));
500        self.update(|record| record.background_jobs = jobs);
501    }
502
503    /// Rewrite the record. Bookkeeping only: a failure never fails the turn,
504    /// and a record already removed stays removed.
505    fn update(&self, change: impl FnOnce(&mut DelegationRecord)) {
506        let dir = &self.registry.record_dir;
507        if let Some(mut record) = read_record(&dir.join(format!("{}.json", self.handle))) {
508            change(&mut record);
509            if let Err(error) = write_record(dir, &record) {
510                tracing::debug!(handle = %record.handle, %error, "could not update a delegation record");
511            }
512        }
513    }
514
515    /// Whether `scv agents kill` stopped this run.
516    pub(crate) fn was_killed(&self) -> bool {
517        self.killed.load(Ordering::Acquire)
518    }
519
520    /// The run ended: stop anything still tagged with it, then forget it.
521    pub(crate) async fn finish(mut self) {
522        self.finished = true;
523        stop_tagged(&self.handle).await;
524        self.registry.finish_local(&self.handle);
525    }
526}
527
528impl Drop for DelegationGuard {
529    fn drop(&mut self) {
530        if self.finished {
531            return;
532        }
533        // The run was abandoned mid-flight: kill its group now and sweep
534        // tagged leftovers in the background.
535        if let Some(group) = ProcessGroup::new(self.pgid) {
536            group.signal(libc::SIGKILL);
537        }
538        self.registry.finish_local(&self.handle);
539        let handle = self.handle.clone();
540        if let Ok(runtime) = tokio::runtime::Handle::try_current() {
541            runtime.spawn(async move { stop_tagged(&handle).await });
542        } else {
543            for identity in tagged_processes(&handle) {
544                signal(identity.pid, libc::SIGKILL);
545            }
546        }
547    }
548}
549
550/// Stop a delegation's process group and tagged processes: TERM, then KILL
551/// after a short grace. Returns whether anything was still running.
552async fn stop_delegation(record: &DelegationRecord) -> bool {
553    let mut stopped = false;
554    // The group ID is the leader's PID, which the kernel does not reuse while
555    // the group has members. A live leader with a different start time means
556    // the PID was reused, so the group is not ours.
557    let leader = ProcessIdentity::of(record.process.pid);
558    let group_is_ours = record.pgid == record.process.pid
559        && match leader {
560            Some(leader) => leader == record.process,
561            None => group_exists(record.pgid),
562        };
563    if group_is_ours && group_exists(record.pgid) {
564        stopped = true;
565        let group = ProcessGroup::new(record.pgid);
566        if let Some(group) = group {
567            group.signal(libc::SIGTERM);
568        }
569        let deadline = tokio::time::Instant::now() + STOP_GRACE;
570        while group_exists(record.pgid) && tokio::time::Instant::now() < deadline {
571            tokio::time::sleep(Duration::from_millis(50)).await;
572        }
573        if let Some(group) = group {
574            group.signal(libc::SIGKILL);
575        }
576    }
577    stopped | stop_tagged(&record.handle).await
578}
579
580/// TERM, then KILL, every process tagged with `handle`. Returns whether any was found.
581async fn stop_tagged(handle: &str) -> bool {
582    let tagged = tagged_processes(handle);
583    if tagged.is_empty() {
584        return false;
585    }
586    for identity in &tagged {
587        signal(identity.pid, libc::SIGTERM);
588    }
589    let deadline = tokio::time::Instant::now() + STOP_GRACE;
590    while tagged.iter().any(ProcessIdentity::is_alive) && tokio::time::Instant::now() < deadline {
591        tokio::time::sleep(Duration::from_millis(50)).await;
592    }
593    for identity in tagged.iter().filter(|identity| identity.is_alive()) {
594        signal(identity.pid, libc::SIGKILL);
595    }
596    true
597}
598
599/// Processes whose `SCV_PARENT` chain names `handle`.
600fn tagged_processes(handle: &str) -> Vec<ProcessIdentity> {
601    let own = std::process::id();
602    ProcessTable::snapshot()
603        .tagged
604        .into_iter()
605        .filter(|(identity, chain)| identity.pid != own && chain_names(chain, handle))
606        .map(|(identity, _)| identity)
607        .collect()
608}
609
610fn chain_names(chain: &str, handle: &str) -> bool {
611    chain
612        .split(';')
613        .any(|entry| entry.rsplit('/').next() == Some(handle))
614}
615
616fn signal(pid: u32, signal: i32) {
617    if let Ok(pid) = i32::try_from(pid)
618        && pid > 0
619    {
620        // SAFETY: kill(2) takes plain integers and touches no memory of
621        // ours; a positive PID addresses exactly one process.
622        unsafe {
623            libc::kill(pid, signal);
624        }
625    }
626}
627
628/// Whether the process group still has a running member (zombies excluded on Linux).
629pub(crate) fn group_exists(pgid: u32) -> bool {
630    let Some(group) = ProcessGroup::new(pgid) else {
631        return false;
632    };
633    let signalable = group.is_signalable();
634    #[cfg(target_os = "linux")]
635    {
636        signalable
637            && linux::all_stats()
638                .iter()
639                .any(|info| info.pgid == pgid && info.state != 'Z')
640    }
641    #[cfg(not(target_os = "linux"))]
642    {
643        signalable
644    }
645}
646
647fn unix_now() -> u64 {
648    SystemTime::now()
649        .duration_since(UNIX_EPOCH)
650        .map_or(0, |elapsed| elapsed.as_secs())
651}
652
653fn write_record(dir: &Path, record: &DelegationRecord) -> std::io::Result<()> {
654    write_private_json(dir, &format!("{}.json", record.handle), record)
655}
656
657/// Atomically write `value` as `dir/name` with mode 0600, creating `dir` and
658/// keeping it and its parent (`run/`) private.
659pub(crate) fn write_private_json(
660    dir: &Path,
661    name: &str,
662    value: &impl Serialize,
663) -> std::io::Result<()> {
664    use std::os::unix::fs::PermissionsExt as _;
665    std::fs::create_dir_all(dir)?;
666    if let Some(run) = dir.parent() {
667        std::fs::set_permissions(run, std::fs::Permissions::from_mode(0o700))?;
668    }
669    std::fs::set_permissions(dir, std::fs::Permissions::from_mode(0o700))?;
670    let bytes = serde_json::to_vec_pretty(value).map_err(std::io::Error::other)?;
671    scv_client::fs::replace_private(&dir.join(name), &bytes)
672}
673
674fn read_record(path: &Path) -> Option<DelegationRecord> {
675    let bytes = match std::fs::File::open(path).and_then(|file| {
676        let mut bytes = Vec::new();
677        std::io::Read::read_to_end(&mut std::io::Read::take(file, MAX_RECORD_BYTES), &mut bytes)
678            .map(|_| bytes)
679    }) {
680        Ok(bytes) => bytes,
681        Err(error) => {
682            // A record removed between listing and reading is normal.
683            if error.kind() != std::io::ErrorKind::NotFound {
684                tracing::debug!(path = %path.display(), %error, "unreadable delegation record");
685            }
686            return None;
687        }
688    };
689    let record: DelegationRecord = match serde_json::from_slice(&bytes) {
690        Ok(record) => record,
691        Err(error) => {
692            tracing::debug!(path = %path.display(), %error, "malformed delegation record");
693            return None;
694        }
695    };
696    // Only a record named after its own handle is trusted.
697    if path.file_stem().and_then(|stem| stem.to_str()) == Some(record.handle.as_str()) {
698        Some(record)
699    } else {
700        tracing::debug!(path = %path.display(), "delegation record named for another handle");
701        None
702    }
703}
704
705fn remove_record(dir: &Path, handle: &str) {
706    let path = dir.join(format!("{handle}.json"));
707    if let Err(error) = std::fs::remove_file(&path)
708        && error.kind() != std::io::ErrorKind::NotFound
709    {
710        tracing::debug!(path = %path.display(), %error, "could not remove a delegation record");
711    }
712}
713
714/// Make this process the reaper of orphaned descendants (Linux), so processes
715/// a delegated agent leaves behind stay in SCV's process tree.
716pub fn become_child_subreaper() -> bool {
717    #[cfg(target_os = "linux")]
718    {
719        // SAFETY: PR_SET_CHILD_SUBREAPER takes integer arguments only and
720        // changes only this process's own reaping attribute.
721        unsafe { libc::prctl(libc::PR_SET_CHILD_SUBREAPER, 1, 0, 0, 0) == 0 }
722    }
723    #[cfg(not(target_os = "linux"))]
724    {
725        false
726    }
727}
728
729static SPAWNED: Mutex<Option<HashSet<u32>>> = Mutex::new(None);
730
731/// Note a child this process spawned and will wait for itself.
732pub(crate) fn track_spawned(pid: u32) {
733    lock(&SPAWNED).get_or_insert_with(HashSet::new).insert(pid);
734}
735
736pub(crate) fn untrack_spawned(pid: u32) {
737    if let Some(spawned) = lock(&SPAWNED).as_mut() {
738        spawned.remove(&pid);
739    }
740}
741
742/// Collect exited orphans reparented to this subreaper. Children SCV spawned
743/// itself are left to their own waiters.
744pub fn reap_orphaned_zombies() -> usize {
745    #[cfg(target_os = "linux")]
746    {
747        let own = std::process::id();
748        let spawned = lock(&SPAWNED).clone().unwrap_or_default();
749        let uptime = linux::uptime_ticks();
750        let mut reaped = 0;
751        for info in linux::all_stats() {
752            if info.ppid != own || info.state != 'Z' || spawned.contains(&info.pid) {
753                continue;
754            }
755            let old_enough = uptime.is_some_and(|now| {
756                now.saturating_sub(info.start_time)
757                    >= ZOMBIE_MIN_AGE.as_secs() * linux::clock_ticks()
758            });
759            if !old_enough {
760                continue;
761            }
762            let mut status = 0;
763            // SAFETY: `status` is a live local that waitpid(2) writes one
764            // int into; WNOHANG keeps the call from blocking.
765            if unsafe { libc::waitpid(info.pid as i32, &raw mut status, libc::WNOHANG) }
766                == info.pid as i32
767            {
768                reaped += 1;
769            }
770        }
771        reaped
772    }
773    #[cfg(not(target_os = "linux"))]
774    {
775        0
776    }
777}
778
779/// Processes of interest at one moment: group membership and tags.
780struct ProcessTable {
781    groups: Vec<(ProcessIdentity, u32)>,
782    tagged: Vec<(ProcessIdentity, String)>,
783}
784
785impl ProcessTable {
786    fn members(&self, record: &DelegationRecord) -> HashSet<u32> {
787        let mut members: HashSet<u32> = self
788            .groups
789            .iter()
790            .filter(|(_, pgid)| *pgid == record.pgid)
791            .map(|(identity, _)| identity.pid)
792            .collect();
793        members.extend(
794            self.tagged
795                .iter()
796                .filter(|(_, chain)| chain_names(chain, &record.handle))
797                .map(|(identity, _)| identity.pid),
798        );
799        members
800    }
801
802    #[cfg(target_os = "linux")]
803    fn snapshot() -> Self {
804        let mut groups = Vec::new();
805        let mut tagged = Vec::new();
806        for info in linux::all_stats() {
807            if info.state == 'Z' {
808                continue;
809            }
810            let identity = ProcessIdentity {
811                pid: info.pid,
812                start_time: info.start_time,
813            };
814            groups.push((identity, info.pgid));
815            if let Some(chain) = linux::parent_chain(info.pid) {
816                tagged.push((identity, chain));
817            }
818        }
819        Self { groups, tagged }
820    }
821
822    #[cfg(not(target_os = "linux"))]
823    fn snapshot() -> Self {
824        let mut groups = Vec::new();
825        let mut tagged = Vec::new();
826        // `ps -E` appends each process's environment to its command line.
827        let Ok(output) = std::process::Command::new("ps")
828            .args(["-E", "-ww", "-axo", "pid=,pgid=,command="])
829            .output()
830        else {
831            return Self { groups, tagged };
832        };
833        for line in String::from_utf8_lossy(&output.stdout).lines() {
834            let mut fields = line.split_whitespace();
835            let (Some(pid), Some(pgid)) = (
836                fields.next().and_then(|value| value.parse::<u32>().ok()),
837                fields.next().and_then(|value| value.parse::<u32>().ok()),
838            ) else {
839                continue;
840            };
841            let Some(identity) = ProcessIdentity::of(pid) else {
842                continue;
843            };
844            groups.push((identity, pgid));
845            if let Some(chain) = fields.find_map(|field| {
846                field
847                    .strip_prefix(PARENT_VARIABLE)
848                    .and_then(|rest| rest.strip_prefix('='))
849            }) {
850                tagged.push((identity, chain.to_owned()));
851            }
852        }
853        Self { groups, tagged }
854    }
855}
856
857#[cfg(target_os = "linux")]
858fn process_start_time(pid: u32) -> Option<u64> {
859    linux::stat(pid).map(|info| info.start_time)
860}
861
862#[cfg(target_os = "macos")]
863fn process_start_time(pid: u32) -> Option<u64> {
864    // SAFETY: proc_bsdinfo is a plain C struct of integers and byte arrays,
865    // for which all-zero bytes are a valid value.
866    let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
867    let size = std::mem::size_of::<libc::proc_bsdinfo>() as i32;
868    // SAFETY: the buffer is `info` itself and `size` is its exact size, so
869    // proc_pidinfo(2) writes at most that many bytes into it.
870    let written = unsafe {
871        libc::proc_pidinfo(
872            pid as i32,
873            libc::PROC_PIDTBSDINFO,
874            0,
875            (&mut info as *mut libc::proc_bsdinfo).cast(),
876            size,
877        )
878    };
879    (written == size).then(|| info.pbi_start_tvsec * 1_000_000 + info.pbi_start_tvusec)
880}
881
882#[cfg(not(target_os = "linux"))]
883fn process_is_zombie(pid: u32) -> bool {
884    std::process::Command::new("ps")
885        .args(["-o", "stat=", "-p", &pid.to_string()])
886        .output()
887        .ok()
888        .filter(|output| output.status.success())
889        .is_some_and(|output| {
890            String::from_utf8_lossy(&output.stdout)
891                .trim_start()
892                .starts_with('Z')
893        })
894}
895
896#[cfg(not(any(target_os = "linux", target_os = "macos")))]
897fn process_start_time(pid: u32) -> Option<u64> {
898    // SAFETY: signal 0 only checks that the process exists; kill(2) takes
899    // plain integers and touches no memory of ours.
900    let alive = unsafe { libc::kill(pid as i32, 0) } == 0;
901    alive.then_some(0)
902}
903
904#[cfg(target_os = "linux")]
905mod linux {
906    pub(super) struct Stat {
907        pub(crate) pid: u32,
908        pub(crate) ppid: u32,
909        pub(crate) pgid: u32,
910        pub(crate) state: char,
911        pub(crate) start_time: u64,
912    }
913
914    pub(super) fn stat(pid: u32) -> Option<Stat> {
915        let text = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
916        // The command name is parenthesized and may contain spaces or ')'.
917        let rest = &text[text.rfind(')')? + 2..];
918        let fields: Vec<&str> = rest.split_whitespace().collect();
919        // After the name: state(3) ppid(4) pgrp(5) ... starttime(22).
920        Some(Stat {
921            pid,
922            state: fields.first()?.chars().next()?,
923            ppid: fields.get(1)?.parse().ok()?,
924            pgid: fields.get(2)?.parse().ok()?,
925            start_time: fields.get(19)?.parse().ok()?,
926        })
927    }
928
929    pub(super) fn all_stats() -> Vec<Stat> {
930        let Ok(entries) = std::fs::read_dir("/proc") else {
931            return Vec::new();
932        };
933        entries
934            .filter_map(Result::ok)
935            .filter_map(|entry| entry.file_name().to_str()?.parse::<u32>().ok())
936            .filter_map(stat)
937            .collect()
938    }
939
940    /// `SCV_PARENT` from a process's environment, when readable.
941    pub(super) fn parent_chain(pid: u32) -> Option<String> {
942        let environ = std::fs::read(format!("/proc/{pid}/environ")).ok()?;
943        let prefix = format!("{}=", super::PARENT_VARIABLE);
944        environ.split(|byte| *byte == 0).find_map(|entry| {
945            entry
946                .strip_prefix(prefix.as_bytes())
947                .map(|value| String::from_utf8_lossy(value).into_owned())
948        })
949    }
950
951    pub(super) fn clock_ticks() -> u64 {
952        // SAFETY: sysconf(3) only reads a system constant.
953        let ticks = unsafe { libc::sysconf(libc::_SC_CLK_TCK) };
954        u64::try_from(ticks)
955            .ok()
956            .filter(|ticks| *ticks > 0)
957            .unwrap_or(100)
958    }
959
960    pub(super) fn uptime_ticks() -> Option<u64> {
961        let text = std::fs::read_to_string("/proc/uptime").ok()?;
962        let seconds: f64 = text.split_whitespace().next()?.parse().ok()?;
963        Some((seconds * clock_ticks() as f64) as u64)
964    }
965}
966
967#[cfg(test)]
968mod tests;