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