Skip to main content

microsandbox_agentd/
session.rs

1//! Exec session management: spawning processes with PTY or pipe I/O.
2
3use std::collections::VecDeque;
4use std::ffi::{CStr, CString};
5use std::mem::MaybeUninit;
6use std::os::fd::{AsRawFd, FromRawFd, OwnedFd, RawFd};
7use std::os::unix::process::CommandExt;
8use std::process::{Command, Stdio};
9use std::sync::Arc;
10use std::task::{Context, Poll};
11use std::{iter, mem, ptr};
12
13use nix::pty;
14use nix::sys::signal::Signal;
15use tokio::io::AsyncReadExt;
16use tokio::io::unix::AsyncFd;
17use tokio::sync::{Semaphore, mpsc, oneshot};
18
19use microsandbox_protocol::bulk::BulkRecord;
20use microsandbox_protocol::exec::{ExecFailed, ExecFailureKind, ExecRequest};
21use microsandbox_protocol::transport::ClientIncarnation;
22
23use crate::config::SecurityProfile;
24use crate::error::{AgentdError, AgentdResult};
25use crate::process::{ProcessExitWatcher, ProcessIdentity, ProcessManager};
26use crate::rlimit;
27use crate::serial::InputCharge;
28use crate::workload::WorkloadPlacement;
29
30//--------------------------------------------------------------------------------------------------
31// Constants
32//--------------------------------------------------------------------------------------------------
33
34const LINUX_CAPABILITY_VERSION_3: u32 = 0x20080522;
35const CAP_SYS_ADMIN: u32 = 21;
36const CAP_WORD_BITS: u32 = 32;
37const PR_CAPBSET_DROP: libc::c_int = 24;
38const PR_CAP_AMBIENT: libc::c_int = 47;
39const PR_CAP_AMBIENT_CLEAR_ALL: libc::c_int = 4;
40const DEFAULT_USER_SPEC: &str = "0:0";
41
42/// Aggregate guest-to-host data retained outside the serial output buffer.
43const SESSION_OUTPUT_BYTE_CAPACITY: usize = 32 * 1024 * 1024;
44
45/// Allocation granularity used by the data budget.
46const SESSION_OUTPUT_BUDGET_GRANULE: usize = 4096;
47
48/// Maximum number of data or control events waiting for the serial writer.
49const SESSION_OUTPUT_ITEM_CAPACITY: usize = 1024;
50
51/// Maximum number of records waiting for the independently scheduled bulk writer.
52const SESSION_BULK_OUTPUT_ITEM_CAPACITY: usize = 256;
53
54/// Maximum lifecycle commands waiting for the dedicated bulk scheduler.
55const SESSION_BULK_COMMAND_CAPACITY: usize = 128;
56
57//--------------------------------------------------------------------------------------------------
58// Functions: classify
59//--------------------------------------------------------------------------------------------------
60
61/// Map an `errno` integer to its standard symbolic name. Returns
62/// `None` for unrecognized values; we only enumerate the ones that
63/// can plausibly come out of fork/exec/setrlimit/setuid paths.
64fn errno_name(e: i32) -> Option<&'static str> {
65    match e {
66        libc::E2BIG => Some("E2BIG"),
67        libc::EACCES => Some("EACCES"),
68        libc::EAGAIN => Some("EAGAIN"),
69        libc::EBUSY => Some("EBUSY"),
70        libc::EFAULT => Some("EFAULT"),
71        libc::EINVAL => Some("EINVAL"),
72        libc::EIO => Some("EIO"),
73        libc::EISDIR => Some("EISDIR"),
74        libc::ELOOP => Some("ELOOP"),
75        libc::EMFILE => Some("EMFILE"),
76        libc::ENAMETOOLONG => Some("ENAMETOOLONG"),
77        libc::ENFILE => Some("ENFILE"),
78        libc::ENOENT => Some("ENOENT"),
79        libc::ENOEXEC => Some("ENOEXEC"),
80        libc::ENOMEM => Some("ENOMEM"),
81        libc::ENOSYS => Some("ENOSYS"),
82        libc::ENOTDIR => Some("ENOTDIR"),
83        libc::ENXIO => Some("ENXIO"),
84        libc::EPERM => Some("EPERM"),
85        libc::ETXTBSY => Some("ETXTBSY"),
86        _ => None,
87    }
88}
89
90/// Classify a fork/exec-time `errno` into one of the
91/// `ExecFailureKind` buckets.
92///
93/// ENOENT is ambiguous in principle (missing binary vs. missing
94/// cwd), but in practice it's overwhelmingly the binary — the cwd
95/// is set in `pre_exec` *before* execvp, and a bad cwd would more
96/// commonly produce ENOTDIR (path component isn't a directory) or
97/// EACCES (no permission to chdir). We classify ENOENT as
98/// `NotFound` and ENOTDIR as `BadCwd`. Edge cases of "bad cwd that
99/// happens to ENOENT" fall through with the message "spawn 'cmd':
100/// No such file or directory" which is still understandable.
101fn classify_spawn_errno(errno: i32) -> ExecFailureKind {
102    match errno {
103        libc::ENOENT => ExecFailureKind::NotFound,
104        libc::ENOTDIR => ExecFailureKind::BadCwd,
105        libc::EACCES | libc::EPERM => ExecFailureKind::PermissionDenied,
106        libc::ENOEXEC => ExecFailureKind::NotExecutable,
107        libc::EISDIR => ExecFailureKind::NotExecutable,
108        libc::ETXTBSY => ExecFailureKind::NotExecutable,
109        libc::E2BIG | libc::ELOOP | libc::ENAMETOOLONG | libc::EFAULT => ExecFailureKind::BadArgs,
110        libc::EMFILE | libc::ENFILE => ExecFailureKind::ResourceLimit,
111        libc::EAGAIN => ExecFailureKind::ResourceLimit,
112        libc::ENOMEM => ExecFailureKind::OutOfMemory,
113        libc::EINVAL => ExecFailureKind::Other,
114        _ => ExecFailureKind::Other,
115    }
116}
117
118/// Build a `ExecFailed` payload from a spawn-time `io::Error`.
119fn exec_failed_from_io_error(err: &std::io::Error, cmd: &str, stage: &str) -> ExecFailed {
120    let errno = err.raw_os_error();
121    let kind = errno
122        .map(classify_spawn_errno)
123        .unwrap_or(ExecFailureKind::Other);
124    let errno_name = errno.and_then(errno_name).map(str::to_string);
125    let message = format!("spawn {cmd:?}: {err}");
126    ExecFailed {
127        kind,
128        errno,
129        errno_name,
130        message,
131        stage: Some(stage.to_string()),
132    }
133}
134
135//--------------------------------------------------------------------------------------------------
136// Types
137//--------------------------------------------------------------------------------------------------
138
139/// An active exec session handle for sending input to a running process.
140///
141/// Output reading is handled by a background task that sends events
142/// via the `mpsc` channel provided at spawn time.
143#[derive(Debug)]
144pub struct ExecSession {
145    /// Stable identity for the spawned process registration.
146    process_identity: ProcessIdentity,
147
148    /// Owns process status and serializes signals with PID reuse.
149    process_manager: Arc<ProcessManager>,
150
151    /// The PTY master fd (only for PTY mode, used for writing and resize).
152    pty_master: Option<AsyncFd<OwnedFd>>,
153
154    /// The child's stdin (only for pipe mode).
155    stdin: Option<AsyncFd<OwnedFd>>,
156
157    /// Accepted input stays ordered, including EOF, while a pipe or PTY backpressures.
158    pending_stdin: VecDeque<PendingStdin>,
159}
160
161#[derive(Debug)]
162struct PendingStdin {
163    data: Vec<u8>,
164    written: usize,
165    _charge: Option<InputCharge>,
166}
167
168/// Output from a session that the agent loop should forward to the host.
169pub enum SessionOutput {
170    /// Data from stdout (or PTY master).
171    Stdout(Vec<u8>),
172
173    /// Data from stderr (pipe mode only).
174    Stderr(Vec<u8>),
175
176    /// The process has exited with the given code.
177    Exited(i32),
178
179    /// Pre-encoded frame bytes to write directly to the serial output buffer.
180    Raw(RawSessionOutput),
181
182    /// Generation-8 raw bulk record whose payload remains separately owned.
183    Bulk(BulkSessionOutput),
184}
185
186/// One queued session event and the data-budget capacity owned by its buffer.
187pub struct SessionOutputEnvelope {
188    /// Host attachment generation captured when the producer was created.
189    pub generation: u64,
190    /// Correlation ID for the session event.
191    pub id: u32,
192
193    /// Dual-port range owner captured when the session was created.
194    pub incarnation: Option<ClientIncarnation>,
195
196    /// Event consumed by the main serial loop.
197    pub output: SessionOutput,
198
199    /// Capacity follows the allocation and is released only after serial output consumes it.
200    _permit: Option<tokio::sync::OwnedSemaphorePermit>,
201}
202
203/// Lifecycle commands processed ahead of queued dedicated-lane output.
204pub enum BulkOutputCommand {
205    /// Park after the currently written complete record, retaining all queued source output.
206    Park {
207        /// Cumulative complete dedicated-lane wire bytes at the cut.
208        completion: oneshot::Sender<u64>,
209    },
210    /// Release the parked source or restored output generation after its thaw reply.
211    Resume {
212        /// Resolves once ordinary output is eligible again.
213        completion: oneshot::Sender<()>,
214    },
215    /// Discard inherited transfer output before acknowledging restore activation.
216    Restore {
217        /// New attachment generation; late output from previous generations is discarded.
218        generation: u64,
219        /// Resolves after the scheduler has crossed the cut.
220        completion: oneshot::Sender<()>,
221    },
222    /// Release queued records for one cancelled operation.
223    DropFlow {
224        /// Range owner that opened the operation.
225        incarnation: ClientIncarnation,
226        /// Correlation being cancelled.
227        id: u32,
228        /// Resolves after matching queued records and their permits are dropped.
229        completion: oneshot::Sender<()>,
230    },
231
232    /// Release every queued record owned by one disconnected SDK client.
233    DropIncarnation {
234        /// Range ownership period being removed.
235        incarnation: ClientIncarnation,
236        /// Resolves after matching queued records and their permits are dropped.
237        completion: oneshot::Sender<()>,
238    },
239}
240
241/// Capacity reserved before a producer reads or encodes a data-bearing event.
242pub struct SessionOutputPermit(tokio::sync::OwnedSemaphorePermit);
243
244/// Cloneable producer for the byte-bounded session output queue.
245#[derive(Clone)]
246pub struct SessionOutputSender {
247    generation: u64,
248    control_tx: mpsc::Sender<SessionOutputEnvelope>,
249    bulk_tx: Option<mpsc::Sender<SessionOutputEnvelope>>,
250    bulk_command_tx: Option<mpsc::Sender<BulkOutputCommand>>,
251    control_budget: Arc<Semaphore>,
252    bulk_budget: Arc<Semaphore>,
253    incarnation: Option<ClientIncarnation>,
254}
255
256/// Pre-encoded session output plus the accounting metadata known by its producer.
257pub struct RawSessionOutput {
258    /// Encoded protocol frame bytes.
259    pub frame: Vec<u8>,
260
261    /// Activity represented by the frame.
262    pub activity: RawActivity,
263
264    /// Session table entry completed by the frame, if any.
265    pub completion: Option<RawSessionCompletion>,
266}
267
268/// Raw bulk output plus activity metadata known by its producer.
269pub struct BulkSessionOutput {
270    /// Validated record whose payload is written with the fixed header via `writev`.
271    pub record: BulkRecord,
272
273    /// Activity represented by the record.
274    pub activity: RawActivity,
275}
276
277/// Activity represented by a pre-encoded session frame.
278#[derive(Debug, Clone, Copy, Default)]
279pub struct RawActivity {
280    /// Meaningful guest-to-host protocol messages represented by this update.
281    pub guest_messages: usize,
282
283    /// Filesystem bytes moved by this frame.
284    pub fs_bytes: usize,
285
286    /// TCP bytes moved by this frame.
287    pub tcp_bytes: usize,
288}
289
290/// Session table entry completed by a pre-encoded session frame.
291#[derive(Debug, Clone, Copy)]
292pub enum RawSessionCompletion {
293    /// A filesystem read stream completed.
294    FsRead,
295
296    /// A filesystem write worker completed.
297    FsWrite,
298
299    /// A TCP stream completed.
300    Tcp,
301}
302
303struct ResolvedUser {
304    uid: libc::uid_t,
305    gid: libc::gid_t,
306    initgroups_user: Option<CString>,
307    home_dir: Option<CString>,
308}
309
310struct PasswdEntry {
311    name: String,
312    uid: libc::uid_t,
313    gid: libc::gid_t,
314    home_dir: Option<String>,
315}
316
317struct GroupEntry {
318    gid: libc::gid_t,
319}
320
321struct ExecErrorPipe {
322    read_end: OwnedFd,
323    write_end: OwnedFd,
324}
325
326/// A piped process whose exit status is observed by [`ProcessManager`].
327struct PipedProcess {
328    stdin: Option<tokio::process::ChildStdin>,
329    stdout: Option<tokio::process::ChildStdout>,
330    stderr: Option<tokio::process::ChildStderr>,
331    exit_watcher: ProcessExitWatcher,
332}
333
334#[repr(C)]
335#[derive(Clone, Copy)]
336struct CapUserHeader {
337    version: u32,
338    pid: libc::c_int,
339}
340
341#[repr(C)]
342#[derive(Clone, Copy)]
343struct CapUserData {
344    effective: u32,
345    permitted: u32,
346    inheritable: u32,
347}
348
349//--------------------------------------------------------------------------------------------------
350// Methods
351//--------------------------------------------------------------------------------------------------
352
353impl RawSessionOutput {
354    /// Creates pre-encoded output with activity metadata.
355    pub fn new(
356        frame: Vec<u8>,
357        activity: RawActivity,
358        completion: Option<RawSessionCompletion>,
359    ) -> Self {
360        Self {
361            frame,
362            activity,
363            completion,
364        }
365    }
366}
367
368impl BulkSessionOutput {
369    /// Creates a raw bulk output event.
370    pub fn new(record: BulkRecord, activity: RawActivity) -> Self {
371        Self { record, activity }
372    }
373}
374
375impl SessionOutput {
376    /// Bytes retained by this event that count against bulk output capacity.
377    fn budget_bytes(&self) -> usize {
378        let allocation = match self {
379            Self::Stdout(data) | Self::Stderr(data) => data.capacity(),
380            Self::Bulk(output) => output.record.payload.len(),
381            Self::Raw(output)
382                if output.activity.fs_bytes != 0 || output.activity.tcp_bytes != 0 =>
383            {
384                output.frame.capacity()
385            }
386            Self::Exited(_) | Self::Raw(_) => 0,
387        };
388
389        allocation
390            .checked_add(SESSION_OUTPUT_BUDGET_GRANULE - 1)
391            .map(|bytes| bytes / SESSION_OUTPUT_BUDGET_GRANULE * SESSION_OUTPUT_BUDGET_GRANULE)
392            .unwrap_or(usize::MAX)
393    }
394}
395
396impl SessionOutputSender {
397    /// Create one ordered queue with a separate byte budget for data-bearing events.
398    pub fn channel() -> (Self, mpsc::Receiver<SessionOutputEnvelope>) {
399        let (tx, rx) = mpsc::channel(SESSION_OUTPUT_ITEM_CAPACITY);
400        let budget = Arc::new(Semaphore::new(SESSION_OUTPUT_BYTE_CAPACITY));
401        (
402            Self {
403                generation: 0,
404                control_tx: tx,
405                bulk_tx: None,
406                bulk_command_tx: None,
407                control_budget: Arc::clone(&budget),
408                bulk_budget: budget,
409                incarnation: None,
410            },
411            rx,
412        )
413    }
414
415    /// Create independently bounded control and raw-bulk producer queues.
416    pub fn split_channel() -> (
417        Self,
418        mpsc::Receiver<SessionOutputEnvelope>,
419        mpsc::Receiver<SessionOutputEnvelope>,
420        mpsc::Receiver<BulkOutputCommand>,
421    ) {
422        let (control_tx, control_rx) = mpsc::channel(SESSION_OUTPUT_ITEM_CAPACITY);
423        let (bulk_tx, bulk_rx) = mpsc::channel(SESSION_BULK_OUTPUT_ITEM_CAPACITY);
424        let (bulk_command_tx, bulk_command_rx) = mpsc::channel(SESSION_BULK_COMMAND_CAPACITY);
425        (
426            Self {
427                generation: 0,
428                control_tx,
429                bulk_tx: Some(bulk_tx),
430                bulk_command_tx: Some(bulk_command_tx),
431                control_budget: Arc::new(Semaphore::new(SESSION_OUTPUT_BYTE_CAPACITY)),
432                bulk_budget: Arc::new(Semaphore::new(SESSION_OUTPUT_BYTE_CAPACITY)),
433                incarnation: None,
434            },
435            control_rx,
436            bulk_rx,
437            bulk_command_rx,
438        )
439    }
440
441    /// Scope future producer events to the client incarnation that opened their session.
442    pub fn with_incarnation(&self, incarnation: Option<ClientIncarnation>) -> Self {
443        Self {
444            generation: self.generation,
445            control_tx: self.control_tx.clone(),
446            bulk_tx: self.bulk_tx.clone(),
447            bulk_command_tx: self.bulk_command_tx.clone(),
448            control_budget: Arc::clone(&self.control_budget),
449            bulk_budget: Arc::clone(&self.bulk_budget),
450            incarnation,
451        }
452    }
453
454    /// Current local attachment generation, independent of the wire protocol version.
455    pub(crate) fn generation(&self) -> u64 {
456        self.generation
457    }
458
459    pub(crate) async fn park_bulk_output(&self) -> Result<u64, &'static str> {
460        let Some(commands) = &self.bulk_command_tx else {
461            return Ok(0);
462        };
463        let (completion, completed) = oneshot::channel();
464        commands
465            .send(BulkOutputCommand::Park { completion })
466            .await
467            .map_err(|_| "bulk scheduler closed while parking")?;
468        completed.await.map_err(|_| "bulk output park failed")
469    }
470
471    pub(crate) async fn resume_bulk_output(&self) -> Result<(), &'static str> {
472        let Some(commands) = &self.bulk_command_tx else {
473            return Ok(());
474        };
475        let (completion, completed) = oneshot::channel();
476        commands
477            .send(BulkOutputCommand::Resume { completion })
478            .await
479            .map_err(|_| "bulk scheduler closed while resuming")?;
480        completed.await.map_err(|_| "bulk output resume failed")
481    }
482
483    /// Change only the root sender. Existing producers keep their old generation while they
484    /// drain inherited pipes, so their output can never complete a new client's correlation.
485    pub(crate) async fn restore_generation(&mut self) -> Result<(), &'static str> {
486        self.generation = self
487            .generation
488            .checked_add(1)
489            .ok_or("attachment generation exhausted")?;
490        if let Some(commands) = &self.bulk_command_tx {
491            let (completion, completed) = oneshot::channel();
492            commands
493                .send(BulkOutputCommand::Restore {
494                    generation: self.generation,
495                    completion,
496                })
497                .await
498                .map_err(|_| "bulk scheduler closed during restore")?;
499            completed
500                .await
501                .map_err(|_| "bulk scheduler restore barrier failed")?;
502        }
503        Ok(())
504    }
505
506    /// Restore one ordered producer queue after boot selected combined mode.
507    pub fn disable_bulk_scheduler(&mut self) {
508        // Combined mode has no cross-lane merger. Raw records, finish markers, and terminal output
509        // must therefore enter one FIFO before sharing the physical console stream.
510        self.bulk_tx = None;
511        self.bulk_command_tx = None;
512    }
513
514    /// Queue a high-priority purge for one operation without awaiting bulk-port progress.
515    pub fn drop_bulk_flow(&self, id: u32) -> Result<Option<oneshot::Receiver<()>>, &'static str> {
516        let Some(incarnation) = self.incarnation else {
517            return Ok(None);
518        };
519        let Some(commands) = self.bulk_command_tx.as_ref() else {
520            return Ok(None);
521        };
522        let (completion, completed) = oneshot::channel();
523        commands
524            .try_send(BulkOutputCommand::DropFlow {
525                incarnation,
526                id,
527                completion,
528            })
529            .map_err(|_| "dedicated bulk scheduler command queue is unavailable")?;
530        Ok(Some(completed))
531    }
532
533    /// Queue a high-priority purge for a disconnected range owner.
534    pub fn drop_bulk_incarnation(
535        &self,
536        incarnation: ClientIncarnation,
537    ) -> Result<Option<oneshot::Receiver<()>>, &'static str> {
538        let Some(commands) = self.bulk_command_tx.as_ref() else {
539            return Ok(None);
540        };
541        let (completion, completed) = oneshot::channel();
542        commands
543            .try_send(BulkOutputCommand::DropIncarnation {
544                incarnation,
545                completion,
546            })
547            .map_err(|_| "dedicated bulk scheduler command queue is unavailable")?;
548        Ok(Some(completed))
549    }
550
551    /// Queue an event after its retained allocation has acquired aggregate capacity.
552    pub async fn send(&self, id: u32, output: SessionOutput) -> bool {
553        let budget_bytes = output.budget_bytes();
554        let permit = if matches!(&output, SessionOutput::Bulk(_)) {
555            self.reserve_bulk(budget_bytes).await
556        } else {
557            self.reserve(budget_bytes).await
558        };
559        let Some(permit) = permit else {
560            eprintln!("agentd session output {id} exceeds byte budget: {budget_bytes} bytes");
561            return false;
562        };
563
564        self.send_reserved(id, output, permit).await
565    }
566
567    /// Reserve capacity before reading or encoding up to `max_bytes` of output.
568    pub async fn reserve(&self, max_bytes: usize) -> Option<SessionOutputPermit> {
569        self.reserve_from(&self.control_budget, max_bytes).await
570    }
571
572    /// Reserve capacity from the raw-bulk budget before allocating a record payload.
573    pub async fn reserve_bulk(&self, max_bytes: usize) -> Option<SessionOutputPermit> {
574        self.reserve_from(&self.bulk_budget, max_bytes).await
575    }
576
577    async fn reserve_from(
578        &self,
579        budget: &Arc<Semaphore>,
580        max_bytes: usize,
581    ) -> Option<SessionOutputPermit> {
582        let budget_bytes = max_bytes.checked_add(SESSION_OUTPUT_BUDGET_GRANULE - 1)?
583            / SESSION_OUTPUT_BUDGET_GRANULE
584            * SESSION_OUTPUT_BUDGET_GRANULE;
585        if budget_bytes > SESSION_OUTPUT_BYTE_CAPACITY {
586            return None;
587        }
588        let permit_count = u32::try_from(budget_bytes).ok()?;
589        Arc::clone(budget)
590            .acquire_many_owned(permit_count)
591            .await
592            .ok()
593            .map(SessionOutputPermit)
594    }
595
596    /// Queue output using capacity obtained before the producer created its allocation.
597    pub async fn send_reserved(
598        &self,
599        id: u32,
600        output: SessionOutput,
601        permit: SessionOutputPermit,
602    ) -> bool {
603        let charged = output.budget_bytes();
604        if charged > permit.0.num_permits() {
605            eprintln!(
606                "agentd session output {id} exceeded its reservation: {charged} > {} bytes",
607                permit.0.num_permits()
608            );
609            return false;
610        }
611
612        let tx = if matches!(&output, SessionOutput::Bulk(_)) {
613            self.bulk_tx.as_ref().unwrap_or(&self.control_tx)
614        } else {
615            &self.control_tx
616        };
617        tx.send(SessionOutputEnvelope {
618            generation: self.generation,
619            id,
620            incarnation: self.incarnation,
621            output,
622            _permit: (charged != 0).then_some(permit.0),
623        })
624        .await
625        .is_ok()
626    }
627}
628
629impl RawActivity {
630    /// A guest-to-host frame with no byte counter.
631    pub fn guest_message() -> Self {
632        Self {
633            guest_messages: 1,
634            ..Self::default()
635        }
636    }
637
638    /// A guest-to-host filesystem data frame.
639    pub fn fs_bytes(len: usize) -> Self {
640        Self {
641            guest_messages: 1,
642            fs_bytes: len,
643            tcp_bytes: 0,
644        }
645    }
646
647    /// A guest-to-host TCP data frame.
648    pub fn tcp_bytes(len: usize) -> Self {
649        Self {
650            guest_messages: 1,
651            fs_bytes: 0,
652            tcp_bytes: len,
653        }
654    }
655}
656
657impl ExecSession {
658    /// Spawns a new exec session.
659    ///
660    /// If `req.tty` is true, uses a PTY. Otherwise, uses piped stdin/stdout/stderr.
661    /// A background task is spawned to read output and send events via `tx`.
662    pub(crate) fn spawn(
663        id: u32,
664        req: &ExecRequest,
665        tx: SessionOutputSender,
666        default_user: Option<&str>,
667        security_profile: SecurityProfile,
668        workload_placement: Option<WorkloadPlacement>,
669    ) -> AgentdResult<Self> {
670        let process_manager = ProcessManager::get()?;
671        if req.tty {
672            Self::spawn_pty(
673                id,
674                req,
675                tx,
676                default_user,
677                security_profile,
678                &process_manager,
679                workload_placement,
680            )
681        } else {
682            Self::spawn_pipe(
683                id,
684                req,
685                tx,
686                default_user,
687                security_profile,
688                &process_manager,
689                workload_placement,
690            )
691        }
692    }
693
694    /// Returns the PID of the spawned process (as u32 for the protocol).
695    pub fn pid(&self) -> u32 {
696        self.process_identity.pid() as u32
697    }
698
699    /// Writes data to the process's stdin (or PTY master).
700    pub async fn write_stdin(&self, data: &[u8]) -> AgentdResult<()> {
701        let mut written = 0;
702        while written < data.len() {
703            let count =
704                std::future::poll_fn(|cx| self.poll_write_stdin(cx, &data[written..])).await?;
705            if count == 0 {
706                return Err(std::io::Error::from(std::io::ErrorKind::WriteZero).into());
707            }
708            written += count;
709        }
710        Ok(())
711    }
712
713    /// Try the common writable-stdin path without a copy, task hop, or readiness registration.
714    pub(crate) fn try_write_stdin(&self, data: &[u8]) -> std::io::Result<usize> {
715        match self.pty_master.as_ref().or(self.stdin.as_ref()) {
716            Some(input) => write_nonblocking_fd(input.as_raw_fd(), data),
717            None => Ok(data.len()),
718        }
719    }
720
721    /// Poll a previously blocked input without preventing the agent from reading lifecycle frames.
722    pub(crate) fn poll_write_stdin(
723        &self,
724        cx: &mut Context<'_>,
725        data: &[u8],
726    ) -> Poll<std::io::Result<usize>> {
727        let Some(input) = self.pty_master.as_ref().or(self.stdin.as_ref()) else {
728            return Poll::Ready(Ok(data.len()));
729        };
730        loop {
731            let mut ready = std::task::ready!(input.poll_write_ready(cx))?;
732            match ready.try_io(|inner| write_nonblocking_fd(inner.as_raw_fd(), data)) {
733                Ok(result) => return Poll::Ready(result),
734                Err(_would_block) => continue,
735            }
736        }
737    }
738
739    /// Retain only an already-admitted frame. The caller's aggregate wire credit bounds both
740    /// this allocation and zero-length EOF cardinality across every active and detached session.
741    pub(crate) fn enqueue_stdin(
742        &mut self,
743        data: Vec<u8>,
744        charge: Option<InputCharge>,
745    ) -> std::io::Result<()> {
746        let mut written = 0;
747        if self.pending_stdin.is_empty() {
748            if data.is_empty() {
749                self.close_stdin();
750                return Ok(());
751            }
752            match self.try_write_stdin(&data) {
753                Ok(count) if count == data.len() => return Ok(()),
754                Ok(0) => return Err(std::io::ErrorKind::WriteZero.into()),
755                Ok(count) => written = count,
756                Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {}
757                Err(error) => return Err(error),
758            }
759        }
760        self.pending_stdin.push_back(PendingStdin {
761            data,
762            written,
763            _charge: charge,
764        });
765        Ok(())
766    }
767
768    pub(crate) fn has_pending_stdin(&self) -> bool {
769        !self.pending_stdin.is_empty()
770    }
771
772    /// The restored process no longer has a host-side stdin owner. Drain everything accepted
773    /// before the cut, then close pipe input. A PTY has no separate write half to close safely.
774    pub(crate) fn detach_stdin(&mut self) {
775        if self.stdin.is_none() {
776            return;
777        }
778        if self.pending_stdin.is_empty() {
779            self.close_stdin();
780        } else if !self
781            .pending_stdin
782            .iter()
783            .any(|pending| pending.data.is_empty())
784        {
785            // At most one marker per already-admitted nonempty queue: detachment cannot create
786            // an unbounded stream of uncharged empty messages.
787            self.pending_stdin.push_back(PendingStdin {
788                data: Vec::new(),
789                written: 0,
790                _charge: None,
791            });
792        }
793    }
794
795    /// Consume one bounded turn without losing a partial-write cursor when a lifecycle event
796    /// cancels this poll. Restored detached sessions use the same path for accepted input only.
797    pub(crate) fn poll_pending_stdin(&mut self, cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
798        let mut progressed = false;
799        for _ in 0..16 {
800            let Some(pending) = self.pending_stdin.front() else {
801                break;
802            };
803            if pending.data.is_empty() {
804                self.close_stdin();
805                self.pending_stdin.pop_front();
806                progressed = true;
807                continue;
808            }
809            match self.poll_write_stdin(cx, &pending.data[pending.written..]) {
810                Poll::Ready(Ok(0)) => {
811                    self.pending_stdin.pop_front();
812                    return Poll::Ready(Err(std::io::ErrorKind::WriteZero.into()));
813                }
814                Poll::Ready(Ok(count)) => {
815                    let pending = self
816                        .pending_stdin
817                        .front_mut()
818                        .expect("polled pending input");
819                    pending.written += count;
820                    if pending.written == pending.data.len() {
821                        self.pending_stdin.pop_front();
822                    }
823                    progressed = true;
824                }
825                Poll::Ready(Err(error)) => {
826                    self.pending_stdin.pop_front();
827                    return Poll::Ready(Err(error));
828                }
829                Poll::Pending => break,
830            }
831        }
832        if progressed {
833            Poll::Ready(Ok(()))
834        } else {
835            Poll::Pending
836        }
837    }
838
839    /// Resizes the PTY (only applicable for TTY sessions).
840    pub fn resize(&self, rows: u16, cols: u16) -> AgentdResult<()> {
841        if let Some(ref master) = self.pty_master {
842            let ws = libc::winsize {
843                ws_row: rows,
844                ws_col: cols,
845                ws_xpixel: 0,
846                ws_ypixel: 0,
847            };
848            let ret = unsafe { libc::ioctl(master.as_raw_fd(), libc::TIOCSWINSZ, &ws) };
849            if ret < 0 {
850                return Err(std::io::Error::last_os_error().into());
851            }
852        }
853        Ok(())
854    }
855
856    /// Sends a signal to the spawned process and everything it started.
857    ///
858    /// The child is made a session leader at spawn (both pipe and PTY modes),
859    /// so signalling the negative pid reaches its whole process group. A bare
860    /// kill(pid) here leaked orphans: killing `sh -c "job &"` took out the
861    /// shell while its backgrounded children survived reparented to init,
862    /// silently accumulating load in the guest.
863    pub fn send_signal(&self, signum: i32) -> AgentdResult<()> {
864        let sig = Signal::try_from(signum)
865            .map_err(|e| AgentdError::ExecSession(format!("invalid signal {signum}: {e}")))?;
866        self.process_manager
867            .signal_process_group(self.process_identity, sig as i32)
868    }
869
870    /// Closes the process's stdin.
871    ///
872    /// For pipe mode, drops the `ChildStdin` handle which closes the fd.
873    /// For PTY mode, this is a no-op (the PTY master stays open for output).
874    pub fn close_stdin(&mut self) {
875        self.stdin.take();
876    }
877}
878
879impl ExecSession {
880    /// Spawns a process with a PTY.
881    fn spawn_pty(
882        id: u32,
883        req: &ExecRequest,
884        tx: SessionOutputSender,
885        default_user: Option<&str>,
886        security_profile: SecurityProfile,
887        process_manager: &Arc<ProcessManager>,
888        workload_placement: Option<WorkloadPlacement>,
889    ) -> AgentdResult<Self> {
890        let pty = pty::openpty(None, None)?;
891        let err_pipe = new_exec_error_pipe()?;
892
893        // Set initial window size.
894        let ws = libc::winsize {
895            ws_row: req.rows,
896            ws_col: req.cols,
897            ws_xpixel: 0,
898            ws_ypixel: 0,
899        };
900        let ret = unsafe { libc::ioctl(pty.master.as_raw_fd(), libc::TIOCSWINSZ, &ws) };
901        if ret < 0 {
902            return Err(std::io::Error::last_os_error().into());
903        }
904
905        let slave_fd = pty.slave.as_raw_fd();
906
907        // Pre-build all strings before fork to avoid allocating in the child.
908        let c_cmd = CString::new(req.cmd.as_str())
909            .map_err(|e| AgentdError::ExecSession(format!("invalid command: {e}")))?;
910        let mut c_args: Vec<CString> = vec![c_cmd.clone()];
911        for arg in &req.args {
912            c_args.push(
913                CString::new(arg.as_str())
914                    .map_err(|e| AgentdError::ExecSession(format!("invalid arg: {e}")))?,
915            );
916        }
917
918        // Build argv pointer array (null-terminated).
919        let argv_ptrs: Vec<*const libc::c_char> = c_args
920            .iter()
921            .map(|s| s.as_ptr())
922            .chain(iter::once(ptr::null()))
923            .collect();
924
925        // Pre-parse environment variables into CStrings.
926        let c_env: Vec<(CString, CString)> = req
927            .env
928            .iter()
929            .filter_map(|var| {
930                let (key, val) = var.split_once('=')?;
931                let k = CString::new(key).ok()?;
932                let v = CString::new(val).ok()?;
933                Some((k, v))
934            })
935            .collect();
936
937        // Pre-build cwd CString.
938        let c_cwd = req
939            .cwd
940            .as_ref()
941            .map(|dir| CString::new(dir.as_str()))
942            .transpose()
943            .map_err(|e| AgentdError::ExecSession(format!("invalid cwd: {e}")))?;
944
945        let resolved_user = resolve_requested_user(req, default_user)?;
946        let default_home = default_home_dir(req, resolved_user.as_ref())?;
947        let home_key = default_home
948            .as_ref()
949            .map(|_| {
950                CString::new("HOME")
951                    .map_err(|e| AgentdError::ExecSession(format!("invalid home env key: {e}")))
952            })
953            .transpose()?;
954
955        // Pre-parse rlimits before fork (no allocations in child).
956        let parsed_rlimits = rlimit::to_libc(&req.rlimits);
957
958        // Prevent the central reaper from observing this child before its PID
959        // and generation are registered.
960        let spawn_guard = process_manager.spawn_guard()?;
961
962        // Fork.
963        let pid = unsafe { libc::fork() };
964        if pid < 0 {
965            let io_err = std::io::Error::last_os_error();
966            return Err(AgentdError::ExecSpawnFailed(exec_failed_from_io_error(
967                &io_err, &req.cmd, "fork",
968            )));
969        }
970
971        #[allow(unreachable_code)]
972        if pid == 0 {
973            // Child process — only async-signal-safe operations from here.
974            drop(pty.master);
975            drop(err_pipe.read_end);
976
977            // Join the workload cgroup before any user code can execute. This
978            // closes the post-spawn PID-assignment race with checkpoint freeze.
979            if let Some(ref placement) = workload_placement
980                && placement.place_current().is_err()
981            {
982                write_exec_error_and_exit(err_pipe.write_end.as_raw_fd());
983            }
984
985            // Create new session.
986            if unsafe { libc::setsid() } < 0 {
987                unsafe { libc::_exit(1) };
988            }
989
990            // Set controlling terminal.
991            if unsafe { libc::ioctl(slave_fd, libc::TIOCSCTTY, 0) } < 0 {
992                unsafe { libc::_exit(1) };
993            }
994
995            // Dup slave to stdin/stdout/stderr.
996            unsafe {
997                if libc::dup2(slave_fd, 0) < 0 {
998                    libc::_exit(1);
999                }
1000                if libc::dup2(slave_fd, 1) < 0 {
1001                    libc::_exit(1);
1002                }
1003                if libc::dup2(slave_fd, 2) < 0 {
1004                    libc::_exit(1);
1005                }
1006                if slave_fd > 2 {
1007                    libc::close(slave_fd);
1008                }
1009            }
1010
1011            // Set environment variables using pre-built CStrings.
1012            for (key, val) in &c_env {
1013                unsafe {
1014                    libc::setenv(key.as_ptr(), val.as_ptr(), 1);
1015                }
1016            }
1017
1018            // Set working directory.
1019            if let Some(ref dir) = c_cwd {
1020                unsafe {
1021                    libc::chdir(dir.as_ptr());
1022                }
1023            }
1024
1025            if apply_exec_security_profile(security_profile).is_err() {
1026                unsafe { libc::_exit(1) };
1027            }
1028
1029            if let Some(ref user) = resolved_user
1030                && apply_resolved_user(user).is_err()
1031            {
1032                unsafe { libc::_exit(1) };
1033            }
1034
1035            if let (Some(key), Some(home)) = (&home_key, &default_home) {
1036                unsafe {
1037                    libc::setenv(key.as_ptr(), home.as_ptr(), 1);
1038                }
1039            }
1040
1041            // Apply resource limits.
1042            for (resource, limit) in &parsed_rlimits {
1043                if unsafe { libc::setrlimit(*resource as _, limit) } != 0 {
1044                    unsafe { libc::_exit(1) };
1045                }
1046            }
1047
1048            // execvp — on success this never returns.
1049            unsafe {
1050                libc::execvp(argv_ptrs[0], argv_ptrs.as_ptr());
1051            }
1052
1053            // If execvp returns, it failed.
1054            write_exec_error_and_exit(err_pipe.write_end.as_raw_fd());
1055        }
1056
1057        // Parent process.
1058        drop(pty.slave);
1059        drop(err_pipe.write_end);
1060        let exit_watcher = spawn_guard.track(pid)?;
1061        let process_identity = exit_watcher.identity();
1062
1063        match read_exec_error(err_pipe.read_end.as_raw_fd()) {
1064            Ok(Some(exec_errno)) => {
1065                drop(exit_watcher);
1066                process_manager.release(process_identity);
1067                let io_err = std::io::Error::from_raw_os_error(exec_errno);
1068                return Err(AgentdError::ExecSpawnFailed(exec_failed_from_io_error(
1069                    &io_err, &req.cmd, "execvp",
1070                )));
1071            }
1072            Ok(None) => {}
1073            Err(error) => {
1074                let _ =
1075                    process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1076                process_manager.release(process_identity);
1077                return Err(error);
1078            }
1079        }
1080
1081        // Dup the master fd for the reader task.
1082        let reader_fd = unsafe { libc::dup(pty.master.as_raw_fd()) };
1083        if reader_fd < 0 {
1084            let _ = process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1085            process_manager.release(process_identity);
1086            return Err(std::io::Error::last_os_error().into());
1087        }
1088        let reader_fd = unsafe { OwnedFd::from_raw_fd(reader_fd) };
1089        let pty_master = nonblocking_input(pty.master).inspect_err(|_| {
1090            let _ = process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1091            process_manager.release(process_identity);
1092        })?;
1093
1094        // Spawn background reader task.
1095        tokio::spawn(pty_reader_task(id, reader_fd, exit_watcher, tx));
1096
1097        Ok(Self {
1098            process_identity,
1099            process_manager: Arc::clone(process_manager),
1100            pty_master: Some(pty_master),
1101            stdin: None,
1102            pending_stdin: VecDeque::new(),
1103        })
1104    }
1105
1106    /// Spawns a process with piped stdio.
1107    fn spawn_pipe(
1108        id: u32,
1109        req: &ExecRequest,
1110        tx: SessionOutputSender,
1111        default_user: Option<&str>,
1112        security_profile: SecurityProfile,
1113        process_manager: &Arc<ProcessManager>,
1114        workload_placement: Option<WorkloadPlacement>,
1115    ) -> AgentdResult<Self> {
1116        let mut cmd = Command::new(&req.cmd);
1117        cmd.args(&req.args)
1118            .stdin(Stdio::piped())
1119            .stdout(Stdio::piped())
1120            .stderr(Stdio::piped());
1121
1122        for var in &req.env {
1123            if let Some((key, val)) = var.split_once('=') {
1124                cmd.env(key, val);
1125            }
1126        }
1127
1128        if let Some(ref dir) = req.cwd {
1129            cmd.current_dir(dir);
1130        }
1131
1132        let resolved_user = resolve_requested_user(req, default_user)?;
1133        if let Some(home) = default_home_dir(req, resolved_user.as_ref())? {
1134            cmd.env("HOME", home.to_string_lossy().into_owned());
1135        }
1136
1137        // Apply the security profile and resource limits in the child before exec.
1138        let parsed_rlimits = rlimit::to_libc(&req.rlimits);
1139        unsafe {
1140            cmd.pre_exec(move || {
1141                // This uses only write(2) in the child and therefore remains
1142                // safe in the fork-to-exec window.
1143                if let Some(ref placement) = workload_placement {
1144                    placement.place_current()?;
1145                }
1146                // Become a session (and process-group) leader so signals sent
1147                // to the group reach every descendant the command spawns, not
1148                // just the direct child. The PTY path does the same for its
1149                // controlling terminal; here it exists purely for group kills.
1150                if libc::setsid() < 0 {
1151                    return Err(std::io::Error::last_os_error());
1152                }
1153                apply_exec_security_profile(security_profile).map_err(agentd_to_io_error)?;
1154                if let Some(ref user) = resolved_user {
1155                    apply_resolved_user(user).map_err(agentd_to_io_error)?;
1156                }
1157                for (resource, limit) in &parsed_rlimits {
1158                    if libc::setrlimit(*resource as _, limit) != 0 {
1159                        return Err(std::io::Error::last_os_error());
1160                    }
1161                }
1162                Ok(())
1163            });
1164        }
1165
1166        let PipedProcess {
1167            stdin,
1168            stdout,
1169            stderr,
1170            exit_watcher,
1171        } = spawn_piped_process(cmd, process_manager)?;
1172        let process_identity = exit_watcher.identity();
1173        let stdin = stdin
1174            .map(|input| input.into_owned_fd().and_then(nonblocking_input))
1175            .transpose()
1176            .inspect_err(|_| {
1177                let _ =
1178                    process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1179                process_manager.release(process_identity);
1180            })?;
1181
1182        // Spawn background reader task.
1183        tokio::spawn(pipe_reader_task(id, stdout, stderr, exit_watcher, tx));
1184
1185        Ok(Self {
1186            process_identity,
1187            process_manager: Arc::clone(process_manager),
1188            pty_master: None,
1189            stdin,
1190            pending_stdin: VecDeque::new(),
1191        })
1192    }
1193}
1194
1195//--------------------------------------------------------------------------------------------------
1196// Trait Implementations
1197//--------------------------------------------------------------------------------------------------
1198
1199impl Drop for ExecSession {
1200    fn drop(&mut self) {
1201        // The registration deliberately outlives the direct child so signals
1202        // can still reach descendants while their output is being drained.
1203        self.process_manager.release(self.process_identity);
1204    }
1205}
1206
1207//--------------------------------------------------------------------------------------------------
1208// Functions
1209//--------------------------------------------------------------------------------------------------
1210
1211fn spawn_piped_process(
1212    mut command: Command,
1213    process_manager: &ProcessManager,
1214) -> AgentdResult<PipedProcess> {
1215    let cmd_label = command.get_program().to_string_lossy().into_owned();
1216
1217    // Prevent the central reaper from observing this child before its PID and
1218    // generation are registered.
1219    let spawn_guard = process_manager.spawn_guard()?;
1220    let mut child = command.spawn().map_err(|error| {
1221        AgentdError::ExecSpawnFailed(exec_failed_from_io_error(
1222            &error,
1223            &cmd_label,
1224            "Command::spawn",
1225        ))
1226    })?;
1227    let pid = child.id() as i32;
1228    let exit_watcher = spawn_guard.track(pid)?;
1229    let process_identity = exit_watcher.identity();
1230
1231    let stdio = (|| {
1232        let stdin = child
1233            .stdin
1234            .take()
1235            .map(tokio::process::ChildStdin::from_std)
1236            .transpose()?;
1237        let stdout = child
1238            .stdout
1239            .take()
1240            .map(tokio::process::ChildStdout::from_std)
1241            .transpose()?;
1242        let stderr = child
1243            .stderr
1244            .take()
1245            .map(tokio::process::ChildStderr::from_std)
1246            .transpose()?;
1247        Ok::<_, std::io::Error>((stdin, stdout, stderr))
1248    })();
1249    let (stdin, stdout, stderr) = stdio.map_err(|error| {
1250        // The command has already exec'd successfully. If an async stdio
1251        // adapter cannot be registered, do not leave an unreported process
1252        // group running after the host receives ExecFailed.
1253        let _ = process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1254        process_manager.release(process_identity);
1255        AgentdError::ExecSpawnFailed(exec_failed_from_io_error(
1256            &error,
1257            &cmd_label,
1258            "Command::spawn",
1259        ))
1260    })?;
1261
1262    // `std::process::Child` has no asynchronous reaper on drop. Once the PID
1263    // is tracked, the process manager owns its exit status during normal
1264    // operation; terminal teardown may reap it directly as a fallback.
1265    drop(child);
1266
1267    Ok(PipedProcess {
1268        stdin,
1269        stdout,
1270        stderr,
1271        exit_watcher,
1272    })
1273}
1274
1275fn new_exec_error_pipe() -> AgentdResult<ExecErrorPipe> {
1276    let mut fds = [0; 2];
1277    let ret = unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_CLOEXEC) };
1278    if ret != 0 {
1279        return Err(std::io::Error::last_os_error().into());
1280    }
1281
1282    Ok(ExecErrorPipe {
1283        read_end: unsafe { OwnedFd::from_raw_fd(fds[0]) },
1284        write_end: unsafe { OwnedFd::from_raw_fd(fds[1]) },
1285    })
1286}
1287
1288fn write_exec_error_and_exit(err_fd: RawFd) -> ! {
1289    let errno = unsafe { *libc::__errno_location() };
1290    let bytes = errno.to_ne_bytes();
1291    let _ = unsafe { libc::write(err_fd, bytes.as_ptr() as *const libc::c_void, bytes.len()) };
1292    unsafe { libc::_exit(127) }
1293}
1294
1295fn read_exec_error(err_fd: RawFd) -> AgentdResult<Option<i32>> {
1296    let mut buf = [0u8; mem::size_of::<i32>()];
1297    let n = unsafe { libc::read(err_fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
1298    if n < 0 {
1299        return Err(std::io::Error::last_os_error().into());
1300    }
1301    if n == 0 {
1302        return Ok(None);
1303    }
1304    if n as usize != buf.len() {
1305        return Err(AgentdError::ExecSession(format!(
1306            "short exec error report: expected {} bytes, got {n}",
1307            buf.len()
1308        )));
1309    }
1310    Ok(Some(i32::from_ne_bytes(buf)))
1311}
1312
1313fn apply_exec_security_profile(profile: SecurityProfile) -> AgentdResult<()> {
1314    match profile {
1315        SecurityProfile::Default => Ok(()),
1316        SecurityProfile::Restricted => drop_mount_admin_privileges(),
1317    }
1318}
1319
1320fn drop_mount_admin_privileges() -> AgentdResult<()> {
1321    if unsafe { libc::prctl(libc::PR_SET_NO_NEW_PRIVS, 1, 0, 0, 0) } != 0 {
1322        return Err(std::io::Error::last_os_error().into());
1323    }
1324
1325    let ret = unsafe { libc::prctl(PR_CAP_AMBIENT, PR_CAP_AMBIENT_CLEAR_ALL, 0, 0, 0) };
1326    if ret != 0 {
1327        let err = std::io::Error::last_os_error();
1328        if err.raw_os_error() != Some(libc::EINVAL) {
1329            return Err(err.into());
1330        }
1331    }
1332
1333    let mut header = CapUserHeader {
1334        version: LINUX_CAPABILITY_VERSION_3,
1335        pid: 0,
1336    };
1337    let mut data = [CapUserData {
1338        effective: 0,
1339        permitted: 0,
1340        inheritable: 0,
1341    }; 2];
1342
1343    if unsafe { libc::syscall(libc::SYS_capget, &mut header, data.as_mut_ptr()) } != 0 {
1344        return Err(std::io::Error::last_os_error().into());
1345    }
1346
1347    let index = (CAP_SYS_ADMIN / CAP_WORD_BITS) as usize;
1348    let mask = 1u32 << (CAP_SYS_ADMIN % CAP_WORD_BITS);
1349    let had_sys_admin = data[index].effective & mask != 0
1350        || data[index].permitted & mask != 0
1351        || data[index].inheritable & mask != 0;
1352
1353    if had_sys_admin {
1354        data[index].effective &= !mask;
1355        data[index].permitted &= !mask;
1356        data[index].inheritable &= !mask;
1357
1358        if unsafe { libc::syscall(libc::SYS_capset, &mut header, data.as_ptr()) } != 0 {
1359            return Err(std::io::Error::last_os_error().into());
1360        }
1361    }
1362
1363    let ret = unsafe { libc::prctl(PR_CAPBSET_DROP, CAP_SYS_ADMIN, 0, 0, 0) };
1364    if ret != 0 {
1365        let err = std::io::Error::last_os_error();
1366        let errno = err.raw_os_error();
1367        // Already-unprivileged callers may also lack CAP_SETPCAP for the bounding-set drop.
1368        let already_unprivileged = !had_sys_admin && errno == Some(libc::EPERM);
1369        if errno != Some(libc::EINVAL) && !already_unprivileged {
1370            return Err(err.into());
1371        }
1372    }
1373
1374    Ok(())
1375}
1376
1377pub(crate) fn resolve_default_user(default_user: Option<&str>) -> AgentdResult<(u32, u32)> {
1378    let Some(spec) = default_user
1379        .map(str::trim)
1380        .filter(|value| !value.is_empty())
1381    else {
1382        return Ok((0, 0));
1383    };
1384
1385    let resolved = resolve_user_spec(spec)?;
1386    Ok((resolved.uid, resolved.gid))
1387}
1388
1389fn resolve_requested_user(
1390    req: &ExecRequest,
1391    default_user: Option<&str>,
1392) -> AgentdResult<Option<ResolvedUser>> {
1393    let default_user = default_user
1394        .map(str::trim)
1395        .filter(|value| !value.is_empty());
1396    let requested = req
1397        .user
1398        .as_deref()
1399        .map(str::trim)
1400        .filter(|value| !value.is_empty())
1401        .or(default_user);
1402
1403    requested.map(resolve_user_spec).transpose()
1404}
1405
1406fn resolve_user_spec(spec: &str) -> AgentdResult<ResolvedUser> {
1407    let (user_part, group_part) = match spec.split_once(':') {
1408        Some((user, group)) => (user.trim(), Some(group.trim())),
1409        None => (spec.trim(), None),
1410    };
1411
1412    if user_part.is_empty() {
1413        return Err(AgentdError::ExecSession("user spec has empty user".into()));
1414    }
1415
1416    let passwd = if let Ok(uid) = parse_id(user_part) {
1417        lookup_passwd_by_uid(uid)?
1418    } else {
1419        lookup_passwd_by_name(user_part)?
1420            .ok_or_else(|| AgentdError::ExecSession(format!("guest user not found: {user_part}")))?
1421            .into()
1422    };
1423
1424    let (uid, passwd_entry) = match passwd {
1425        ResolvedUserLookup::Known(entry) => (entry.uid, Some(entry)),
1426        ResolvedUserLookup::Numeric(uid) => (uid, None),
1427    };
1428
1429    let gid = match group_part {
1430        Some("") => {
1431            return Err(AgentdError::ExecSession("user spec has empty group".into()));
1432        }
1433        Some(group) => resolve_group_spec(group)?,
1434        None => passwd_entry
1435            .as_ref()
1436            .map(|entry| entry.gid)
1437            .unwrap_or_else(|| unsafe { libc::getgid() }),
1438    };
1439
1440    let initgroups_user = passwd_entry
1441        .as_ref()
1442        .map(|entry| CString::new(entry.name.as_str()))
1443        .transpose()
1444        .map_err(|e| AgentdError::ExecSession(format!("invalid guest user name: {e}")))?;
1445
1446    Ok(ResolvedUser {
1447        uid,
1448        gid,
1449        initgroups_user,
1450        home_dir: passwd_entry
1451            .as_ref()
1452            .and_then(|entry| entry.home_dir.as_deref())
1453            .map(CString::new)
1454            .transpose()
1455            .map_err(|e| AgentdError::ExecSession(format!("invalid guest home directory: {e}")))?,
1456    })
1457}
1458
1459enum ResolvedUserLookup {
1460    Known(PasswdEntry),
1461    Numeric(libc::uid_t),
1462}
1463
1464impl From<PasswdEntry> for ResolvedUserLookup {
1465    fn from(value: PasswdEntry) -> Self {
1466        Self::Known(value)
1467    }
1468}
1469
1470fn resolve_group_spec(spec: &str) -> AgentdResult<libc::gid_t> {
1471    if let Ok(gid) = parse_id(spec) {
1472        return Ok(gid);
1473    }
1474
1475    lookup_group_by_name(spec)?
1476        .map(|entry| entry.gid)
1477        .ok_or_else(|| AgentdError::ExecSession(format!("guest group not found: {spec}")))
1478}
1479
1480fn parse_id(value: &str) -> Result<u32, std::num::ParseIntError> {
1481    value.parse::<u32>()
1482}
1483
1484fn lookup_passwd_by_name(name: &str) -> AgentdResult<Option<PasswdEntry>> {
1485    let name = CString::new(name)
1486        .map_err(|e| AgentdError::ExecSession(format!("invalid guest user name: {e}")))?;
1487    let mut pwd = MaybeUninit::<libc::passwd>::uninit();
1488    let mut result = ptr::null_mut();
1489    let mut buf = vec![0u8; lookup_buffer_len()];
1490    let rc = unsafe {
1491        libc::getpwnam_r(
1492            name.as_ptr(),
1493            pwd.as_mut_ptr(),
1494            buf.as_mut_ptr().cast(),
1495            buf.len(),
1496            &mut result,
1497        )
1498    };
1499    if rc != 0 {
1500        return Err(AgentdError::ExecSession(format!(
1501            "failed to resolve guest user {name:?}: {}",
1502            std::io::Error::from_raw_os_error(rc)
1503        )));
1504    }
1505    if result.is_null() {
1506        return Ok(None);
1507    }
1508
1509    let pwd = unsafe { pwd.assume_init() };
1510    let name = unsafe { CStr::from_ptr(pwd.pw_name) }
1511        .to_string_lossy()
1512        .into_owned();
1513    let home_dir = unsafe { CStr::from_ptr(pwd.pw_dir) }
1514        .to_string_lossy()
1515        .into_owned();
1516    Ok(Some(PasswdEntry {
1517        name,
1518        uid: pwd.pw_uid,
1519        gid: pwd.pw_gid,
1520        home_dir: (!home_dir.is_empty()).then_some(home_dir),
1521    }))
1522}
1523
1524fn lookup_passwd_by_uid(uid: libc::uid_t) -> AgentdResult<ResolvedUserLookup> {
1525    let mut pwd = MaybeUninit::<libc::passwd>::uninit();
1526    let mut result = ptr::null_mut();
1527    let mut buf = vec![0u8; lookup_buffer_len()];
1528    let rc = unsafe {
1529        libc::getpwuid_r(
1530            uid,
1531            pwd.as_mut_ptr(),
1532            buf.as_mut_ptr().cast(),
1533            buf.len(),
1534            &mut result,
1535        )
1536    };
1537    if rc != 0 {
1538        return Err(AgentdError::ExecSession(format!(
1539            "failed to resolve guest uid {uid}: {}",
1540            std::io::Error::from_raw_os_error(rc)
1541        )));
1542    }
1543    if result.is_null() {
1544        return Ok(ResolvedUserLookup::Numeric(uid));
1545    }
1546
1547    let pwd = unsafe { pwd.assume_init() };
1548    let name = unsafe { CStr::from_ptr(pwd.pw_name) }
1549        .to_string_lossy()
1550        .into_owned();
1551    let home_dir = unsafe { CStr::from_ptr(pwd.pw_dir) }
1552        .to_string_lossy()
1553        .into_owned();
1554    Ok(ResolvedUserLookup::Known(PasswdEntry {
1555        name,
1556        uid: pwd.pw_uid,
1557        gid: pwd.pw_gid,
1558        home_dir: (!home_dir.is_empty()).then_some(home_dir),
1559    }))
1560}
1561
1562fn lookup_group_by_name(name: &str) -> AgentdResult<Option<GroupEntry>> {
1563    let name = CString::new(name)
1564        .map_err(|e| AgentdError::ExecSession(format!("invalid guest group name: {e}")))?;
1565    let mut grp = MaybeUninit::<libc::group>::uninit();
1566    let mut result = ptr::null_mut();
1567    let mut buf = vec![0u8; lookup_buffer_len()];
1568    let rc = unsafe {
1569        libc::getgrnam_r(
1570            name.as_ptr(),
1571            grp.as_mut_ptr(),
1572            buf.as_mut_ptr().cast(),
1573            buf.len(),
1574            &mut result,
1575        )
1576    };
1577    if rc != 0 {
1578        return Err(AgentdError::ExecSession(format!(
1579            "failed to resolve guest group {name:?}: {}",
1580            std::io::Error::from_raw_os_error(rc)
1581        )));
1582    }
1583    if result.is_null() {
1584        return Ok(None);
1585    }
1586
1587    let grp = unsafe { grp.assume_init() };
1588    Ok(Some(GroupEntry { gid: grp.gr_gid }))
1589}
1590
1591fn lookup_buffer_len() -> usize {
1592    let size = unsafe { libc::sysconf(libc::_SC_GETPW_R_SIZE_MAX) };
1593    if size > 0 { size as usize } else { 16 * 1024 }
1594}
1595
1596fn apply_resolved_user(user: &ResolvedUser) -> AgentdResult<()> {
1597    if let Some(ref name) = user.initgroups_user {
1598        if unsafe { libc::initgroups(name.as_ptr(), user.gid) } != 0 {
1599            return Err(std::io::Error::last_os_error().into());
1600        }
1601    } else if unsafe { libc::setgroups(0, ptr::null()) } != 0 {
1602        return Err(std::io::Error::last_os_error().into());
1603    }
1604
1605    if unsafe { libc::setgid(user.gid) } != 0 {
1606        return Err(std::io::Error::last_os_error().into());
1607    }
1608    if unsafe { libc::setuid(user.uid) } != 0 {
1609        return Err(std::io::Error::last_os_error().into());
1610    }
1611
1612    Ok(())
1613}
1614
1615fn default_home_dir(
1616    req: &ExecRequest,
1617    user: Option<&ResolvedUser>,
1618) -> AgentdResult<Option<CString>> {
1619    if env_contains_key(&req.env, "HOME") {
1620        return Ok(None);
1621    }
1622
1623    if let Some(user) = user {
1624        return Ok(user.home_dir.clone());
1625    }
1626
1627    Ok(resolve_user_spec(DEFAULT_USER_SPEC)?.home_dir)
1628}
1629
1630fn env_contains_key(env: &[String], key: &str) -> bool {
1631    env.iter().any(|entry| {
1632        entry
1633            .split_once('=')
1634            .map(|(entry_key, _)| entry_key == key)
1635            .unwrap_or(false)
1636    })
1637}
1638
1639fn agentd_to_io_error(err: AgentdError) -> std::io::Error {
1640    std::io::Error::other(err.to_string())
1641}
1642
1643/// Keep one owned descriptor across readiness waits; cancellation cannot leave a blocking task
1644/// writing through a borrowed fd after its session has been removed or restored.
1645fn nonblocking_input(fd: OwnedFd) -> std::io::Result<AsyncFd<OwnedFd>> {
1646    let flags = unsafe { libc::fcntl(fd.as_raw_fd(), libc::F_GETFL) };
1647    if flags < 0
1648        || unsafe { libc::fcntl(fd.as_raw_fd(), libc::F_SETFL, flags | libc::O_NONBLOCK) } < 0
1649    {
1650        return Err(std::io::Error::last_os_error());
1651    }
1652    AsyncFd::new(fd)
1653}
1654
1655fn write_nonblocking_fd(fd: RawFd, data: &[u8]) -> std::io::Result<usize> {
1656    loop {
1657        let written = unsafe { libc::write(fd, data.as_ptr().cast(), data.len()) };
1658        if written >= 0 {
1659            return Ok(written as usize);
1660        }
1661        let error = std::io::Error::last_os_error();
1662        if error.kind() != std::io::ErrorKind::Interrupted {
1663            return Err(error);
1664        }
1665    }
1666}
1667
1668fn wait_fd_readable(fd: RawFd) -> AgentdResult<()> {
1669    let mut pollfd = libc::pollfd {
1670        fd,
1671        events: libc::POLLIN,
1672        revents: 0,
1673    };
1674
1675    loop {
1676        let ret = unsafe { libc::poll(&mut pollfd, 1, -1) };
1677        if ret < 0 {
1678            let err = std::io::Error::last_os_error();
1679            if err.raw_os_error() == Some(libc::EINTR) {
1680                continue;
1681            }
1682            return Err(AgentdError::Io(err));
1683        }
1684        if ret == 0 {
1685            continue;
1686        }
1687        // Always retry read on HUP/ERR too: a PTY may still contain final output before EIO.
1688        return Ok(());
1689    }
1690}
1691
1692/// Background task that reads from a PTY master fd and sends output events.
1693async fn pty_reader_task(
1694    id: u32,
1695    master_fd: OwnedFd,
1696    exit_watcher: ProcessExitWatcher,
1697    tx: SessionOutputSender,
1698) {
1699    let tx_output = tx.clone();
1700    let runtime_handle = tokio::runtime::Handle::current();
1701    let read_result = tokio::task::spawn_blocking(move || {
1702        // PTY masters are safer with a dedicated blocking read loop than with
1703        // edge-driven readiness. Fast writers followed by process exit can
1704        // strand the tail behind a missed wakeup/HUP transition.
1705        let raw = master_fd.as_raw_fd();
1706        // The duplicated master shares O_NONBLOCK with stdin. Never clear that flag here:
1707        // a blocked write would otherwise strand lifecycle handling on the agent actor.
1708
1709        loop {
1710            let mut buf = [0u8; 4096];
1711            let n = unsafe { libc::read(raw, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
1712
1713            if n > 0 {
1714                let n = n as usize;
1715                let sent = runtime_handle.block_on(async {
1716                    let Some(permit) = tx_output.reserve(n).await else {
1717                        return false;
1718                    };
1719                    tx_output
1720                        .send_reserved(id, SessionOutput::Stdout(buf[..n].to_vec()), permit)
1721                        .await
1722                });
1723                if !sent {
1724                    break;
1725                }
1726                continue;
1727            }
1728
1729            if n == 0 {
1730                break;
1731            }
1732
1733            let err = std::io::Error::last_os_error();
1734            match err.raw_os_error() {
1735                Some(libc::EINTR) => continue,
1736                Some(libc::EAGAIN) => {
1737                    if wait_fd_readable(raw).is_err() {
1738                        break;
1739                    }
1740                }
1741                Some(libc::EIO) => break,
1742                _ => break,
1743            }
1744        }
1745    })
1746    .await;
1747
1748    let _ = read_result;
1749
1750    let code = exit_watcher.await;
1751    let _ = tx.send(id, SessionOutput::Exited(code)).await;
1752}
1753
1754/// Background task that reads from piped stdout/stderr and sends output events.
1755async fn pipe_reader_task(
1756    id: u32,
1757    stdout: Option<tokio::process::ChildStdout>,
1758    stderr: Option<tokio::process::ChildStderr>,
1759    exit_watcher: ProcessExitWatcher,
1760    tx: SessionOutputSender,
1761) {
1762    let mut stdout = stdout;
1763    let mut stderr = stderr;
1764    let mut stdout_eof = stdout.is_none();
1765    let mut stderr_eof = stderr.is_none();
1766
1767    while !stdout_eof || !stderr_eof {
1768        let mut stdout_buf = [0u8; 4096];
1769        let mut stderr_buf = [0u8; 4096];
1770
1771        tokio::select! {
1772            result = async {
1773                match stdout.as_mut() {
1774                    Some(out) => out.read(&mut stdout_buf).await,
1775                    None => std::future::pending().await,
1776                }
1777            }, if !stdout_eof => {
1778                match result {
1779                    Ok(0) | Err(_) => {
1780                        stdout = None;
1781                        stdout_eof = true;
1782                    }
1783                    Ok(n) => {
1784                        let Some(permit) = tx.reserve(n).await else {
1785                            break;
1786                        };
1787                        if !tx
1788                            .send_reserved(
1789                                id,
1790                                SessionOutput::Stdout(stdout_buf[..n].to_vec()),
1791                                permit,
1792                            )
1793                            .await
1794                        {
1795                            break;
1796                        }
1797                    }
1798                }
1799            }
1800            result = async {
1801                match stderr.as_mut() {
1802                    Some(err) => err.read(&mut stderr_buf).await,
1803                    None => std::future::pending().await,
1804                }
1805            }, if !stderr_eof => {
1806                match result {
1807                    Ok(0) | Err(_) => {
1808                        stderr = None;
1809                        stderr_eof = true;
1810                    }
1811                    Ok(n) => {
1812                        let Some(permit) = tx.reserve(n).await else {
1813                            break;
1814                        };
1815                        if !tx
1816                            .send_reserved(
1817                                id,
1818                                SessionOutput::Stderr(stderr_buf[..n].to_vec()),
1819                                permit,
1820                            )
1821                            .await
1822                        {
1823                            break;
1824                        }
1825                    }
1826                }
1827            }
1828        }
1829    }
1830
1831    let code = exit_watcher.await;
1832
1833    let _ = tx.send(id, SessionOutput::Exited(code)).await;
1834}
1835
1836//--------------------------------------------------------------------------------------------------
1837// Tests
1838//--------------------------------------------------------------------------------------------------
1839
1840#[cfg(test)]
1841mod tests {
1842    use std::collections::HashMap;
1843    use std::io::Read;
1844    use std::process::{Command as StdCommand, Stdio as StdStdio};
1845    use std::time::Duration;
1846
1847    use tokio::time;
1848
1849    use microsandbox_protocol::exec::ExecRequest;
1850
1851    use super::*;
1852
1853    const REAP_HELPER_ENV: &str = "MSB_AGENTD_SESSION_REAP_HELPER";
1854    const REAP_HELPER_SENTINEL: &str = "session-reap-helper-passed";
1855    const REAP_TEST_NAME: &str = "session::tests::test_spawn_reaps_adopted_descendant";
1856    const CONCURRENT_HELPER_ENV: &str = "MSB_AGENTD_CONCURRENT_SPAWN_HELPER";
1857    const CONCURRENT_HELPER_SENTINEL: &str = "concurrent-spawn-helper-passed";
1858    const CONCURRENT_TEST_NAME: &str = "session::tests::test_concurrent_spawn_exit_codes";
1859    const RUNTIME_HELPER_ENV: &str = "MSB_AGENTD_RUNTIME_REPLACEMENT_HELPER";
1860    const RUNTIME_HELPER_SENTINEL: &str = "runtime-replacement-helper-passed";
1861    const RUNTIME_TEST_NAME: &str = "session::tests::test_spawn_survives_runtime_replacement";
1862    const PIPE_OWNER_HELPER_ENV: &str = "MSB_AGENTD_PIPE_OWNER_HELPER";
1863    const PIPE_OWNER_HELPER_SENTINEL: &str = "pipe-owner-helper-passed";
1864
1865    #[tokio::test]
1866    async fn session_output_permit_lives_until_envelope_is_consumed() {
1867        let (tx, mut rx) = SessionOutputSender::channel();
1868        assert!(tx.send(7, SessionOutput::Stdout(vec![0; 4096])).await);
1869        assert_eq!(
1870            tx.control_budget.available_permits(),
1871            SESSION_OUTPUT_BYTE_CAPACITY - 4096
1872        );
1873
1874        let envelope = rx.recv().await.unwrap();
1875        assert_eq!(
1876            tx.control_budget.available_permits(),
1877            SESSION_OUTPUT_BYTE_CAPACITY - 4096
1878        );
1879        drop(envelope);
1880        assert_eq!(
1881            tx.control_budget.available_permits(),
1882            SESSION_OUTPUT_BYTE_CAPACITY
1883        );
1884    }
1885
1886    #[tokio::test]
1887    async fn scoped_output_sender_captures_client_incarnation() {
1888        let incarnation = [0x44; 16];
1889        let (tx, mut rx) = SessionOutputSender::channel();
1890        let scoped = tx.with_incarnation(Some(incarnation));
1891
1892        assert!(scoped.send(7, SessionOutput::Exited(0)).await);
1893        let envelope = rx.recv().await.unwrap();
1894
1895        assert_eq!(envelope.id, 7);
1896        assert_eq!(envelope.incarnation, Some(incarnation));
1897    }
1898
1899    #[tokio::test]
1900    async fn split_output_queues_keep_bulk_lifecycle_commands_independent() {
1901        let incarnation = [0x55; 16];
1902        let (tx, mut control_rx, mut bulk_rx, mut command_rx) =
1903            SessionOutputSender::split_channel();
1904        let scoped = tx.with_incarnation(Some(incarnation));
1905        let record = BulkRecord {
1906            id: 9,
1907            kind: microsandbox_protocol::bulk::BulkKind::Filesystem,
1908            flow: microsandbox_protocol::bulk::BulkFlow::GuestToHost,
1909            offset: 0,
1910            payload: b"bulk".as_slice().into(),
1911        };
1912        assert!(
1913            scoped
1914                .send(
1915                    9,
1916                    SessionOutput::Bulk(BulkSessionOutput::new(record, RawActivity::fs_bytes(4),)),
1917                )
1918                .await
1919        );
1920        assert!(scoped.send(10, SessionOutput::Exited(0)).await);
1921        let mut completed = scoped.drop_bulk_flow(9).unwrap().unwrap();
1922
1923        assert!(matches!(
1924            bulk_rx.recv().await.unwrap().output,
1925            SessionOutput::Bulk(_)
1926        ));
1927        assert!(matches!(
1928            control_rx.recv().await.unwrap().output,
1929            SessionOutput::Exited(0)
1930        ));
1931        let BulkOutputCommand::DropFlow {
1932            incarnation: command_incarnation,
1933            id,
1934            completion,
1935        } = command_rx.recv().await.unwrap()
1936        else {
1937            panic!("expected flow cleanup command");
1938        };
1939        assert_eq!(command_incarnation, incarnation);
1940        assert_eq!(id, 9);
1941        completion.send(()).unwrap();
1942        assert_eq!(completed.try_recv(), Ok(()));
1943    }
1944
1945    #[tokio::test]
1946    async fn combined_mode_restores_one_ordered_output_queue() {
1947        let (mut tx, mut control_rx, mut bulk_rx, _command_rx) =
1948            SessionOutputSender::split_channel();
1949        tx.disable_bulk_scheduler();
1950        let record = BulkRecord {
1951            id: 9,
1952            kind: microsandbox_protocol::bulk::BulkKind::Filesystem,
1953            flow: microsandbox_protocol::bulk::BulkFlow::GuestToHost,
1954            offset: 0,
1955            payload: b"bulk".as_slice().into(),
1956        };
1957        assert!(
1958            tx.send(
1959                9,
1960                SessionOutput::Bulk(BulkSessionOutput::new(record, RawActivity::fs_bytes(4),)),
1961            )
1962            .await
1963        );
1964        assert!(tx.send(9, SessionOutput::Exited(0)).await);
1965
1966        assert!(matches!(
1967            control_rx.recv().await.unwrap().output,
1968            SessionOutput::Bulk(_)
1969        ));
1970        assert!(matches!(
1971            control_rx.recv().await.unwrap().output,
1972            SessionOutput::Exited(0)
1973        ));
1974        assert!(bulk_rx.recv().await.is_none());
1975    }
1976
1977    #[tokio::test]
1978    async fn control_output_remains_admissible_when_data_budget_is_exhausted() {
1979        let (tx, mut rx) = SessionOutputSender::channel();
1980        let full_budget = tx.reserve(SESSION_OUTPUT_BYTE_CAPACITY).await.unwrap();
1981
1982        assert!(tx.send(9, SessionOutput::Exited(0)).await);
1983        let envelope = rx.recv().await.unwrap();
1984        assert!(matches!(envelope.output, SessionOutput::Exited(0)));
1985        drop(full_budget);
1986    }
1987    const PIPE_OWNER_TEST_NAME: &str =
1988        "session::tests::test_piped_process_exit_outlives_spawning_runtime";
1989
1990    #[test]
1991    fn test_spawn_reaps_adopted_descendant() {
1992        if std::env::var_os(REAP_HELPER_ENV).is_some() {
1993            let runtime = tokio::runtime::Builder::new_current_thread()
1994                .enable_all()
1995                .build()
1996                .expect("session reap test runtime");
1997            runtime.block_on(run_adopted_descendant_scenario());
1998            println!("{REAP_HELPER_SENTINEL}");
1999            return;
2000        }
2001
2002        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2003            .args(["--exact", REAP_TEST_NAME, "--nocapture"])
2004            .env(REAP_HELPER_ENV, "1")
2005            .stdout(StdStdio::piped())
2006            .spawn()
2007            .expect("spawn isolated session reap test");
2008        let mut output = String::new();
2009        helper
2010            .stdout
2011            .take()
2012            .expect("helper stdout")
2013            .read_to_string(&mut output)
2014            .expect("read helper stdout");
2015
2016        match helper.wait() {
2017            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2018            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2019            Err(error) => panic!("wait for helper: {error}"),
2020        }
2021        assert!(
2022            output.contains(REAP_HELPER_SENTINEL),
2023            "helper did not complete the session reap scenario:\n{output}"
2024        );
2025    }
2026
2027    async fn run_adopted_descendant_scenario() {
2028        let ret = unsafe { libc::prctl(libc::PR_SET_CHILD_SUBREAPER, 1) };
2029        assert_eq!(
2030            ret,
2031            0,
2032            "set child subreaper: {}",
2033            std::io::Error::last_os_error()
2034        );
2035
2036        let (tx, mut rx) = SessionOutputSender::channel();
2037        let req = ExecRequest {
2038            cmd: "/bin/sh".to_string(),
2039            args: vec!["-c".to_string(), "sleep 30 & echo $!".to_string()],
2040            env: Vec::new(),
2041            cwd: None,
2042            user: None,
2043            tty: false,
2044            rows: 24,
2045            cols: 80,
2046            rlimits: Vec::new(),
2047        };
2048
2049        let session = ExecSession::spawn(17, &req, tx, None, SecurityProfile::Default, None)
2050            .expect("spawn background descendant session");
2051        let leader_pid = session.pid() as i32;
2052        let mut stdout = Vec::new();
2053        time::timeout(Duration::from_secs(10), async {
2054            while !stdout.contains(&b'\n') {
2055                let envelope = rx.recv().await.expect("session output");
2056                assert_eq!(envelope.id, 17);
2057                match envelope.output {
2058                    SessionOutput::Stdout(data) => stdout.extend_from_slice(&data),
2059                    SessionOutput::Exited(code) => panic!("session exited early with {code}"),
2060                    SessionOutput::Stderr(_) | SessionOutput::Raw(_) | SessionOutput::Bulk(_) => {}
2061                }
2062            }
2063        })
2064        .await
2065        .expect("wait for background descendant session");
2066
2067        let descendant_pid: i32 = String::from_utf8(stdout)
2068            .expect("descendant PID is UTF-8")
2069            .trim()
2070            .parse()
2071            .expect("parse descendant PID");
2072        let expected_parent = std::process::id().to_string();
2073        let status_path = format!("/proc/{descendant_pid}/status");
2074        time::timeout(Duration::from_secs(5), async {
2075            loop {
2076                if let Ok(status) = std::fs::read_to_string(&status_path)
2077                    && status
2078                        .lines()
2079                        .find_map(|line| line.strip_prefix("PPid:"))
2080                        .is_some_and(|ppid| ppid.trim() == expected_parent)
2081                {
2082                    break;
2083                }
2084                time::sleep(Duration::from_millis(10)).await;
2085            }
2086        })
2087        .await
2088        .expect("descendant should be adopted by the helper subreaper");
2089
2090        let leader_path = format!("/proc/{leader_pid}");
2091        time::timeout(Duration::from_secs(5), async {
2092            while std::path::Path::new(&leader_path).exists() {
2093                time::sleep(Duration::from_millis(10)).await;
2094            }
2095        })
2096        .await
2097        .expect("direct child should be reaped before signalling its descendants");
2098
2099        session
2100            .send_signal(libc::SIGTERM)
2101            .expect("signal descendants through completed process registration");
2102        let exit = time::timeout(Duration::from_secs(5), async {
2103            loop {
2104                let envelope = rx.recv().await.expect("session output after signal");
2105                assert_eq!(envelope.id, 17);
2106                if let SessionOutput::Exited(code) = envelope.output {
2107                    break code;
2108                }
2109            }
2110        })
2111        .await
2112        .expect("session should finish after its descendant is signalled");
2113        assert_eq!(exit, 0);
2114
2115        let proc_path = format!("/proc/{descendant_pid}");
2116        time::timeout(Duration::from_secs(5), async {
2117            while std::path::Path::new(&proc_path).exists() {
2118                time::sleep(Duration::from_millis(10)).await;
2119            }
2120        })
2121        .await
2122        .expect("descendant should be reaped");
2123
2124        let ret = unsafe { libc::waitpid(descendant_pid, ptr::null_mut(), libc::WNOHANG) };
2125        assert_eq!(ret, -1, "descendant {descendant_pid} was not reaped");
2126        assert_eq!(
2127            std::io::Error::last_os_error().raw_os_error(),
2128            Some(libc::ECHILD)
2129        );
2130    }
2131
2132    #[test]
2133    fn test_concurrent_spawn_exit_codes() {
2134        if std::env::var_os(CONCURRENT_HELPER_ENV).is_some() {
2135            let runtime = tokio::runtime::Builder::new_current_thread()
2136                .enable_all()
2137                .build()
2138                .expect("concurrent spawn test runtime");
2139            runtime.block_on(run_concurrent_spawn_scenario());
2140            println!("{CONCURRENT_HELPER_SENTINEL}");
2141            return;
2142        }
2143
2144        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2145            .args(["--exact", CONCURRENT_TEST_NAME, "--nocapture"])
2146            .env(CONCURRENT_HELPER_ENV, "1")
2147            .stdout(StdStdio::piped())
2148            .spawn()
2149            .expect("spawn isolated concurrent session test");
2150        let mut output = String::new();
2151        helper
2152            .stdout
2153            .take()
2154            .expect("helper stdout")
2155            .read_to_string(&mut output)
2156            .expect("read helper stdout");
2157
2158        match helper.wait() {
2159            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2160            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2161            Err(error) => panic!("wait for helper: {error}"),
2162        }
2163        assert!(
2164            output.contains(CONCURRENT_HELPER_SENTINEL),
2165            "helper did not complete the concurrent spawn scenario:\n{output}"
2166        );
2167    }
2168
2169    async fn run_concurrent_spawn_scenario() {
2170        const PROCESS_COUNT: u32 = 12;
2171
2172        let runtime_handle = tokio::runtime::Handle::current();
2173        let (tx, mut rx) = SessionOutputSender::channel();
2174        let mut spawn_threads = Vec::new();
2175        for offset in 0..PROCESS_COUNT {
2176            let handle = runtime_handle.clone();
2177            let tx = tx.clone();
2178            spawn_threads.push(std::thread::spawn(move || {
2179                let _runtime = handle.enter();
2180                let code = 20 + offset as i32;
2181                let req = ExecRequest {
2182                    cmd: "/bin/sh".to_string(),
2183                    args: vec!["-c".to_string(), format!("exit {code}")],
2184                    env: Vec::new(),
2185                    cwd: None,
2186                    user: None,
2187                    tty: offset % 2 == 1,
2188                    rows: 24,
2189                    cols: 80,
2190                    rlimits: Vec::new(),
2191                };
2192                ExecSession::spawn(100 + offset, &req, tx, None, SecurityProfile::Default, None)
2193            }));
2194        }
2195        drop(tx);
2196
2197        let mut sessions = Vec::new();
2198        for thread in spawn_threads {
2199            sessions.push(
2200                thread
2201                    .join()
2202                    .expect("concurrent spawn thread")
2203                    .expect("concurrent process spawn"),
2204            );
2205        }
2206
2207        let mut exits = HashMap::new();
2208        time::timeout(Duration::from_secs(15), async {
2209            while exits.len() < PROCESS_COUNT as usize {
2210                let envelope = rx.recv().await.expect("session output");
2211                if let SessionOutput::Exited(code) = envelope.output {
2212                    exits.insert(envelope.id, code);
2213                }
2214            }
2215        })
2216        .await
2217        .expect("wait for concurrent exits");
2218
2219        for offset in 0..PROCESS_COUNT {
2220            assert_eq!(exits.get(&(100 + offset)), Some(&(20 + offset as i32)));
2221        }
2222        drop(sessions);
2223    }
2224
2225    #[test]
2226    fn test_spawn_survives_runtime_replacement() {
2227        if std::env::var_os(RUNTIME_HELPER_ENV).is_some() {
2228            run_runtime_replacement_scenario();
2229            println!("{RUNTIME_HELPER_SENTINEL}");
2230            return;
2231        }
2232
2233        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2234            .args(["--exact", RUNTIME_TEST_NAME, "--nocapture"])
2235            .env(RUNTIME_HELPER_ENV, "1")
2236            .stdout(StdStdio::piped())
2237            .spawn()
2238            .expect("spawn isolated runtime replacement test");
2239        let mut output = String::new();
2240        helper
2241            .stdout
2242            .take()
2243            .expect("helper stdout")
2244            .read_to_string(&mut output)
2245            .expect("read helper stdout");
2246
2247        match helper.wait() {
2248            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2249            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2250            Err(error) => panic!("wait for helper: {error}"),
2251        }
2252        assert!(
2253            output.contains(RUNTIME_HELPER_SENTINEL),
2254            "helper did not complete the runtime replacement scenario:\n{output}"
2255        );
2256    }
2257
2258    fn run_runtime_replacement_scenario() {
2259        for (id, code) in [(201, 51), (202, 52)] {
2260            let runtime = tokio::runtime::Builder::new_current_thread()
2261                .enable_all()
2262                .build()
2263                .expect("replacement test runtime");
2264            runtime.block_on(run_single_pipe_spawn(id, code));
2265        }
2266    }
2267
2268    async fn run_single_pipe_spawn(id: u32, code: i32) {
2269        let (tx, mut rx) = SessionOutputSender::channel();
2270        let req = ExecRequest {
2271            cmd: "/bin/sh".to_string(),
2272            args: vec!["-c".to_string(), format!("exit {code}")],
2273            env: Vec::new(),
2274            cwd: None,
2275            user: None,
2276            tty: false,
2277            rows: 24,
2278            cols: 80,
2279            rlimits: Vec::new(),
2280        };
2281        let _session = ExecSession::spawn(id, &req, tx, None, SecurityProfile::Default, None)
2282            .expect("spawn session on replacement runtime");
2283
2284        let actual = time::timeout(Duration::from_secs(5), async {
2285            loop {
2286                let envelope = rx.recv().await.expect("session output");
2287                assert_eq!(envelope.id, id);
2288                if let SessionOutput::Exited(actual) = envelope.output {
2289                    break actual;
2290                }
2291            }
2292        })
2293        .await
2294        .expect("wait for exit on replacement runtime");
2295        assert_eq!(actual, code);
2296    }
2297
2298    #[test]
2299    fn test_piped_process_exit_outlives_spawning_runtime() {
2300        if std::env::var_os(PIPE_OWNER_HELPER_ENV).is_some() {
2301            run_piped_process_exit_scenario();
2302            println!("{PIPE_OWNER_HELPER_SENTINEL}");
2303            return;
2304        }
2305
2306        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2307            .args(["--exact", PIPE_OWNER_TEST_NAME, "--nocapture"])
2308            .env(PIPE_OWNER_HELPER_ENV, "1")
2309            .stdout(StdStdio::piped())
2310            .spawn()
2311            .expect("spawn isolated pipe owner test");
2312        let mut output = String::new();
2313        helper
2314            .stdout
2315            .take()
2316            .expect("helper stdout")
2317            .read_to_string(&mut output)
2318            .expect("read helper stdout");
2319
2320        match helper.wait() {
2321            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2322            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2323            Err(error) => panic!("wait for helper: {error}"),
2324        }
2325        assert!(
2326            output.contains(PIPE_OWNER_HELPER_SENTINEL),
2327            "helper did not complete the pipe owner scenario:\n{output}"
2328        );
2329    }
2330
2331    fn run_piped_process_exit_scenario() {
2332        let process_manager = ProcessManager::get().expect("get process manager");
2333        let spawning_runtime = tokio::runtime::Builder::new_current_thread()
2334            .enable_all()
2335            .build()
2336            .expect("spawning runtime");
2337        let exit_watcher = {
2338            let _runtime_guard = spawning_runtime.enter();
2339            let mut command = Command::new("/bin/sh");
2340            command
2341                .args(["-c", "exit 63"])
2342                .stdin(Stdio::piped())
2343                .stdout(Stdio::piped())
2344                .stderr(Stdio::piped());
2345            let process =
2346                spawn_piped_process(command, &process_manager).expect("spawn piped process");
2347            let PipedProcess { exit_watcher, .. } = process;
2348            exit_watcher
2349        };
2350        drop(spawning_runtime);
2351
2352        let waiting_runtime = tokio::runtime::Builder::new_current_thread()
2353            .enable_all()
2354            .build()
2355            .expect("waiting runtime");
2356        let code = waiting_runtime.block_on(async {
2357            time::timeout(Duration::from_secs(5), exit_watcher)
2358                .await
2359                .expect("wait for piped process exit")
2360        });
2361        assert_eq!(code, 63);
2362    }
2363
2364    #[tokio::test]
2365    async fn test_pty_reader_drains_ready_fd() {
2366        let (tx, mut rx) = SessionOutputSender::channel();
2367        let req = ExecRequest {
2368            cmd: "/bin/sh".to_string(),
2369            args: vec![
2370                "-c".to_string(),
2371                "i=0; while [ $i -lt 256 ]; do printf AAAA; i=$((i+1)); done; printf SECOND; sleep 0.1; printf '<END>\\n'; sleep 0.1; exit 0"
2372                    .to_string(),
2373            ],
2374            env: vec!["PATH=/usr/local/bin:/usr/bin:/bin".to_string()],
2375            cwd: None,
2376            user: None,
2377            tty: true,
2378            rows: 24,
2379            cols: 80,
2380            rlimits: Vec::new(),
2381        };
2382
2383        let session = ExecSession::spawn(7, &req, tx, None, SecurityProfile::Default, None)
2384            .expect("spawn pty session");
2385        let mut stdout = Vec::new();
2386        let mut exit = None;
2387
2388        let recv_result = time::timeout(Duration::from_secs(15), async {
2389            while let Some(envelope) = rx.recv().await {
2390                assert_eq!(envelope.id, 7);
2391                match envelope.output {
2392                    SessionOutput::Stdout(data) => stdout.extend_from_slice(&data),
2393                    SessionOutput::Exited(code) => {
2394                        exit = Some(code);
2395                        break;
2396                    }
2397                    SessionOutput::Stderr(_) | SessionOutput::Raw(_) | SessionOutput::Bulk(_) => {}
2398                }
2399            }
2400        })
2401        .await;
2402
2403        if recv_result.is_err() {
2404            let _ = session.send_signal(libc::SIGKILL);
2405            panic!("timed out waiting for PTY output");
2406        }
2407
2408        assert_eq!(exit, Some(0));
2409
2410        let second = stdout
2411            .windows(b"SECOND".len())
2412            .position(|window| window == b"SECOND");
2413        let end = stdout
2414            .windows(b"<END>".len())
2415            .position(|window| window == b"<END>");
2416
2417        assert!(
2418            matches!((second, end), (Some(second), Some(end)) if second < end),
2419            "expected immediate PTY write to arrive before later output; got {:?}",
2420            String::from_utf8_lossy(&stdout),
2421        );
2422    }
2423
2424    #[test]
2425    fn test_resolve_user_spec_for_current_uid_gid() {
2426        let uid = unsafe { libc::getuid() };
2427        let gid = unsafe { libc::getgid() };
2428        let resolved = resolve_user_spec(&format!("{uid}:{gid}")).expect("resolve numeric user");
2429        assert_eq!(resolved.uid, uid);
2430        assert_eq!(resolved.gid, gid);
2431    }
2432
2433    #[test]
2434    fn test_request_user_overrides_config_default() {
2435        let req = ExecRequest {
2436            cmd: "/bin/true".to_string(),
2437            args: Vec::new(),
2438            env: Vec::new(),
2439            cwd: None,
2440            user: Some("1:1".to_string()),
2441            tty: false,
2442            rows: 24,
2443            cols: 80,
2444            rlimits: Vec::new(),
2445        };
2446
2447        let resolved = resolve_requested_user(&req, Some("0:0")).expect("resolve requested user");
2448        assert_eq!(resolved.unwrap().uid, 1);
2449    }
2450
2451    #[test]
2452    fn test_config_default_user_used_when_request_has_none() {
2453        let req = ExecRequest {
2454            cmd: "/bin/true".to_string(),
2455            args: Vec::new(),
2456            env: Vec::new(),
2457            cwd: None,
2458            user: None,
2459            tty: false,
2460            rows: 24,
2461            cols: 80,
2462            rlimits: Vec::new(),
2463        };
2464
2465        let uid = unsafe { libc::getuid() };
2466        let gid = unsafe { libc::getgid() };
2467        let resolved = resolve_requested_user(&req, Some(&format!("{uid}:{gid}")))
2468            .expect("resolve with config default");
2469        let resolved = resolved.expect("should resolve to a user");
2470        assert_eq!(resolved.uid, uid);
2471        assert_eq!(resolved.gid, gid);
2472    }
2473
2474    #[test]
2475    fn test_request_without_user_does_not_apply_user_switch() {
2476        let req = ExecRequest {
2477            cmd: "/bin/true".to_string(),
2478            args: Vec::new(),
2479            env: Vec::new(),
2480            cwd: None,
2481            user: None,
2482            tty: false,
2483            rows: 24,
2484            cols: 80,
2485            rlimits: Vec::new(),
2486        };
2487
2488        let resolved = resolve_requested_user(&req, None).expect("resolve absent user");
2489        assert!(resolved.is_none());
2490    }
2491
2492    #[test]
2493    fn test_default_user_absent_resolves_to_root() {
2494        let resolved = resolve_default_user(None).expect("resolve absent default user");
2495        assert_eq!(resolved, (0, 0));
2496    }
2497
2498    #[test]
2499    fn test_default_home_dir_uses_resolved_user_home() {
2500        let req = ExecRequest {
2501            cmd: "/bin/true".to_string(),
2502            args: Vec::new(),
2503            env: Vec::new(),
2504            cwd: None,
2505            user: None,
2506            tty: false,
2507            rows: 24,
2508            cols: 80,
2509            rlimits: Vec::new(),
2510        };
2511        let user = ResolvedUser {
2512            uid: 1000,
2513            gid: 1000,
2514            initgroups_user: None,
2515            home_dir: Some(CString::new("/home/tester").unwrap()),
2516        };
2517
2518        assert_eq!(
2519            default_home_dir(&req, Some(&user))
2520                .expect("resolve default home")
2521                .as_deref()
2522                .map(CStr::to_string_lossy),
2523            Some("/home/tester".into()),
2524        );
2525    }
2526
2527    #[test]
2528    fn test_default_home_dir_uses_root_when_user_absent() {
2529        let req = ExecRequest {
2530            cmd: "/bin/true".to_string(),
2531            args: Vec::new(),
2532            env: Vec::new(),
2533            cwd: None,
2534            user: None,
2535            tty: false,
2536            rows: 24,
2537            cols: 80,
2538            rlimits: Vec::new(),
2539        };
2540        let root = resolve_user_spec(DEFAULT_USER_SPEC).expect("resolve implicit root");
2541
2542        assert_eq!(
2543            default_home_dir(&req, None)
2544                .expect("resolve default home")
2545                .as_deref()
2546                .map(CStr::to_string_lossy),
2547            root.home_dir.as_deref().map(CStr::to_string_lossy),
2548        );
2549    }
2550
2551    #[test]
2552    fn test_default_home_dir_respects_explicit_home_env() {
2553        let req = ExecRequest {
2554            cmd: "/bin/true".to_string(),
2555            args: Vec::new(),
2556            env: vec!["HOME=/tmp/custom".to_string()],
2557            cwd: None,
2558            user: None,
2559            tty: false,
2560            rows: 24,
2561            cols: 80,
2562            rlimits: Vec::new(),
2563        };
2564        let user = ResolvedUser {
2565            uid: 1000,
2566            gid: 1000,
2567            initgroups_user: None,
2568            home_dir: Some(CString::new("/home/tester").unwrap()),
2569        };
2570
2571        assert!(
2572            default_home_dir(&req, Some(&user))
2573                .expect("resolve default home")
2574                .is_none()
2575        );
2576    }
2577
2578    #[tokio::test]
2579    async fn test_spawn_pipe_error_does_not_include_probe_details() {
2580        let (tx, _rx) = SessionOutputSender::channel();
2581        let req = ExecRequest {
2582            cmd: "/definitely/not/a/real/binary".to_string(),
2583            args: Vec::new(),
2584            env: Vec::new(),
2585            cwd: None,
2586            user: None,
2587            tty: false,
2588            rows: 24,
2589            cols: 80,
2590            rlimits: Vec::new(),
2591        };
2592
2593        // Use the process-wide manager because other tests may have already
2594        // started its reaper thread. A private manager cannot guard this spawn
2595        // from the global `waitpid(-1, ...)` owner.
2596        let process_manager = ProcessManager::get().expect("get process manager");
2597        let err = ExecSession::spawn_pipe(
2598            9,
2599            &req,
2600            tx,
2601            None,
2602            SecurityProfile::Default,
2603            &process_manager,
2604            None,
2605        )
2606        .expect_err("spawn should fail");
2607
2608        // Spawn failures now produce the typed `ExecSpawnFailed` so
2609        // the host can render a useful message + hint. The classifier
2610        // maps ENOENT on the binary path to `NotFound`.
2611        let payload = match &err {
2612            AgentdError::ExecSpawnFailed(p) => p,
2613            other => panic!("expected ExecSpawnFailed, got: {other:?}"),
2614        };
2615        assert_eq!(payload.kind, ExecFailureKind::NotFound);
2616        assert_eq!(payload.errno, Some(libc::ENOENT));
2617        assert_eq!(payload.errno_name.as_deref(), Some("ENOENT"));
2618
2619        // The original intent of the test: probe internals leak into
2620        // the error message. The format is now
2621        // `spawn "<cmd>": <io::Error>` from
2622        // `exec_failed_from_io_error`. Verify that none of the old
2623        // probe-detail keys snuck back into the message.
2624        let message = &payload.message;
2625        assert!(message.contains("spawn"));
2626        assert!(!message.contains("symlink_metadata="));
2627        assert!(!message.contains("metadata="));
2628        assert!(!message.contains("magic="));
2629        assert!(!message.contains("path_probe="));
2630        assert!(!message.contains("cwd_probe="));
2631        assert!(!message.contains("target_probe="));
2632    }
2633}