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            // Group lookup can open /etc/group; finish it before restrictive limits.
1030            if let Some(ref user) = resolved_user
1031                && apply_resolved_groups(user).is_err()
1032            {
1033                unsafe { libc::_exit(1) };
1034            }
1035
1036            // Apply limits while privileged so non-root commands can raise hard limits.
1037            for (resource, limit) in &parsed_rlimits {
1038                if unsafe { libc::setrlimit(*resource as _, limit) } != 0 {
1039                    unsafe { libc::_exit(1) };
1040                }
1041            }
1042
1043            if let Some(ref user) = resolved_user
1044                && apply_resolved_user(user).is_err()
1045            {
1046                unsafe { libc::_exit(1) };
1047            }
1048
1049            if let (Some(key), Some(home)) = (&home_key, &default_home) {
1050                unsafe {
1051                    libc::setenv(key.as_ptr(), home.as_ptr(), 1);
1052                }
1053            }
1054
1055            // execvp — on success this never returns.
1056            unsafe {
1057                libc::execvp(argv_ptrs[0], argv_ptrs.as_ptr());
1058            }
1059
1060            // If execvp returns, it failed.
1061            write_exec_error_and_exit(err_pipe.write_end.as_raw_fd());
1062        }
1063
1064        // Parent process.
1065        drop(pty.slave);
1066        drop(err_pipe.write_end);
1067        let exit_watcher = spawn_guard.track(pid)?;
1068        let process_identity = exit_watcher.identity();
1069
1070        match read_exec_error(err_pipe.read_end.as_raw_fd()) {
1071            Ok(Some(exec_errno)) => {
1072                drop(exit_watcher);
1073                process_manager.release(process_identity);
1074                let io_err = std::io::Error::from_raw_os_error(exec_errno);
1075                return Err(AgentdError::ExecSpawnFailed(exec_failed_from_io_error(
1076                    &io_err, &req.cmd, "execvp",
1077                )));
1078            }
1079            Ok(None) => {}
1080            Err(error) => {
1081                let _ =
1082                    process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1083                process_manager.release(process_identity);
1084                return Err(error);
1085            }
1086        }
1087
1088        // Dup the master fd for the reader task.
1089        let reader_fd = unsafe { libc::dup(pty.master.as_raw_fd()) };
1090        if reader_fd < 0 {
1091            let _ = process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1092            process_manager.release(process_identity);
1093            return Err(std::io::Error::last_os_error().into());
1094        }
1095        let reader_fd = unsafe { OwnedFd::from_raw_fd(reader_fd) };
1096        let pty_master = nonblocking_input(pty.master).inspect_err(|_| {
1097            let _ = process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1098            process_manager.release(process_identity);
1099        })?;
1100
1101        // Spawn background reader task.
1102        tokio::spawn(pty_reader_task(id, reader_fd, exit_watcher, tx));
1103
1104        Ok(Self {
1105            process_identity,
1106            process_manager: Arc::clone(process_manager),
1107            pty_master: Some(pty_master),
1108            stdin: None,
1109            pending_stdin: VecDeque::new(),
1110        })
1111    }
1112
1113    /// Spawns a process with piped stdio.
1114    fn spawn_pipe(
1115        id: u32,
1116        req: &ExecRequest,
1117        tx: SessionOutputSender,
1118        default_user: Option<&str>,
1119        security_profile: SecurityProfile,
1120        process_manager: &Arc<ProcessManager>,
1121        workload_placement: Option<WorkloadPlacement>,
1122    ) -> AgentdResult<Self> {
1123        let mut cmd = Command::new(&req.cmd);
1124        cmd.args(&req.args)
1125            .stdin(Stdio::piped())
1126            .stdout(Stdio::piped())
1127            .stderr(Stdio::piped());
1128
1129        for var in &req.env {
1130            if let Some((key, val)) = var.split_once('=') {
1131                cmd.env(key, val);
1132            }
1133        }
1134
1135        if let Some(ref dir) = req.cwd {
1136            cmd.current_dir(dir);
1137        }
1138
1139        let resolved_user = resolve_requested_user(req, default_user)?;
1140        if let Some(home) = default_home_dir(req, resolved_user.as_ref())? {
1141            cmd.env("HOME", home.to_string_lossy().into_owned());
1142        }
1143
1144        // Apply the security profile and resource limits in the child before exec.
1145        let parsed_rlimits = rlimit::to_libc(&req.rlimits);
1146        unsafe {
1147            cmd.pre_exec(move || {
1148                // This uses only write(2) in the child and therefore remains
1149                // safe in the fork-to-exec window.
1150                if let Some(ref placement) = workload_placement {
1151                    placement.place_current()?;
1152                }
1153                // Become a session (and process-group) leader so signals sent
1154                // to the group reach every descendant the command spawns, not
1155                // just the direct child. The PTY path does the same for its
1156                // controlling terminal; here it exists purely for group kills.
1157                if libc::setsid() < 0 {
1158                    return Err(std::io::Error::last_os_error());
1159                }
1160                apply_exec_security_profile(security_profile).map_err(agentd_to_io_error)?;
1161
1162                // Group lookup must finish before limits can prevent opening /etc/group.
1163                if let Some(ref user) = resolved_user {
1164                    apply_resolved_groups(user).map_err(agentd_to_io_error)?;
1165                }
1166
1167                // Apply limits before dropping the privilege needed to raise hard limits.
1168                for (resource, limit) in &parsed_rlimits {
1169                    if libc::setrlimit(*resource as _, limit) != 0 {
1170                        return Err(std::io::Error::last_os_error());
1171                    }
1172                }
1173
1174                if let Some(ref user) = resolved_user {
1175                    apply_resolved_user(user).map_err(agentd_to_io_error)?;
1176                }
1177                Ok(())
1178            });
1179        }
1180
1181        let PipedProcess {
1182            stdin,
1183            stdout,
1184            stderr,
1185            exit_watcher,
1186        } = spawn_piped_process(cmd, process_manager)?;
1187        let process_identity = exit_watcher.identity();
1188        let stdin = stdin
1189            .map(|input| input.into_owned_fd().and_then(nonblocking_input))
1190            .transpose()
1191            .inspect_err(|_| {
1192                let _ =
1193                    process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1194                process_manager.release(process_identity);
1195            })?;
1196
1197        // Spawn background reader task.
1198        tokio::spawn(pipe_reader_task(id, stdout, stderr, exit_watcher, tx));
1199
1200        Ok(Self {
1201            process_identity,
1202            process_manager: Arc::clone(process_manager),
1203            pty_master: None,
1204            stdin,
1205            pending_stdin: VecDeque::new(),
1206        })
1207    }
1208}
1209
1210//--------------------------------------------------------------------------------------------------
1211// Trait Implementations
1212//--------------------------------------------------------------------------------------------------
1213
1214impl Drop for ExecSession {
1215    fn drop(&mut self) {
1216        // The registration deliberately outlives the direct child so signals
1217        // can still reach descendants while their output is being drained.
1218        self.process_manager.release(self.process_identity);
1219    }
1220}
1221
1222//--------------------------------------------------------------------------------------------------
1223// Functions
1224//--------------------------------------------------------------------------------------------------
1225
1226fn spawn_piped_process(
1227    mut command: Command,
1228    process_manager: &ProcessManager,
1229) -> AgentdResult<PipedProcess> {
1230    let cmd_label = command.get_program().to_string_lossy().into_owned();
1231
1232    // Prevent the central reaper from observing this child before its PID and
1233    // generation are registered.
1234    let spawn_guard = process_manager.spawn_guard()?;
1235    let mut child = command.spawn().map_err(|error| {
1236        AgentdError::ExecSpawnFailed(exec_failed_from_io_error(
1237            &error,
1238            &cmd_label,
1239            "Command::spawn",
1240        ))
1241    })?;
1242    let pid = child.id() as i32;
1243    let exit_watcher = spawn_guard.track(pid)?;
1244    let process_identity = exit_watcher.identity();
1245
1246    let stdio = (|| {
1247        let stdin = child
1248            .stdin
1249            .take()
1250            .map(tokio::process::ChildStdin::from_std)
1251            .transpose()?;
1252        let stdout = child
1253            .stdout
1254            .take()
1255            .map(tokio::process::ChildStdout::from_std)
1256            .transpose()?;
1257        let stderr = child
1258            .stderr
1259            .take()
1260            .map(tokio::process::ChildStderr::from_std)
1261            .transpose()?;
1262        Ok::<_, std::io::Error>((stdin, stdout, stderr))
1263    })();
1264    let (stdin, stdout, stderr) = stdio.map_err(|error| {
1265        // The command has already exec'd successfully. If an async stdio
1266        // adapter cannot be registered, do not leave an unreported process
1267        // group running after the host receives ExecFailed.
1268        let _ = process_manager.signal_process_group(process_identity, Signal::SIGKILL as i32);
1269        process_manager.release(process_identity);
1270        AgentdError::ExecSpawnFailed(exec_failed_from_io_error(
1271            &error,
1272            &cmd_label,
1273            "Command::spawn",
1274        ))
1275    })?;
1276
1277    // `std::process::Child` has no asynchronous reaper on drop. Once the PID
1278    // is tracked, the process manager owns its exit status during normal
1279    // operation; terminal teardown may reap it directly as a fallback.
1280    drop(child);
1281
1282    Ok(PipedProcess {
1283        stdin,
1284        stdout,
1285        stderr,
1286        exit_watcher,
1287    })
1288}
1289
1290fn new_exec_error_pipe() -> AgentdResult<ExecErrorPipe> {
1291    let mut fds = [0; 2];
1292    let ret = unsafe { libc::pipe2(fds.as_mut_ptr(), libc::O_CLOEXEC) };
1293    if ret != 0 {
1294        return Err(std::io::Error::last_os_error().into());
1295    }
1296
1297    Ok(ExecErrorPipe {
1298        read_end: unsafe { OwnedFd::from_raw_fd(fds[0]) },
1299        write_end: unsafe { OwnedFd::from_raw_fd(fds[1]) },
1300    })
1301}
1302
1303fn write_exec_error_and_exit(err_fd: RawFd) -> ! {
1304    let errno = unsafe { *libc::__errno_location() };
1305    let bytes = errno.to_ne_bytes();
1306    let _ = unsafe { libc::write(err_fd, bytes.as_ptr() as *const libc::c_void, bytes.len()) };
1307    unsafe { libc::_exit(127) }
1308}
1309
1310fn read_exec_error(err_fd: RawFd) -> AgentdResult<Option<i32>> {
1311    let mut buf = [0u8; mem::size_of::<i32>()];
1312    let n = unsafe { libc::read(err_fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
1313    if n < 0 {
1314        return Err(std::io::Error::last_os_error().into());
1315    }
1316    if n == 0 {
1317        return Ok(None);
1318    }
1319    if n as usize != buf.len() {
1320        return Err(AgentdError::ExecSession(format!(
1321            "short exec error report: expected {} bytes, got {n}",
1322            buf.len()
1323        )));
1324    }
1325    Ok(Some(i32::from_ne_bytes(buf)))
1326}
1327
1328fn apply_exec_security_profile(profile: SecurityProfile) -> AgentdResult<()> {
1329    match profile {
1330        SecurityProfile::Default => Ok(()),
1331        SecurityProfile::Restricted => drop_mount_admin_privileges(),
1332    }
1333}
1334
1335fn drop_mount_admin_privileges() -> AgentdResult<()> {
1336    if unsafe { libc::prctl(libc::PR_SET_NO_NEW_PRIVS, 1, 0, 0, 0) } != 0 {
1337        return Err(std::io::Error::last_os_error().into());
1338    }
1339
1340    let ret = unsafe { libc::prctl(PR_CAP_AMBIENT, PR_CAP_AMBIENT_CLEAR_ALL, 0, 0, 0) };
1341    if ret != 0 {
1342        let err = std::io::Error::last_os_error();
1343        if err.raw_os_error() != Some(libc::EINVAL) {
1344            return Err(err.into());
1345        }
1346    }
1347
1348    let mut header = CapUserHeader {
1349        version: LINUX_CAPABILITY_VERSION_3,
1350        pid: 0,
1351    };
1352    let mut data = [CapUserData {
1353        effective: 0,
1354        permitted: 0,
1355        inheritable: 0,
1356    }; 2];
1357
1358    if unsafe { libc::syscall(libc::SYS_capget, &mut header, data.as_mut_ptr()) } != 0 {
1359        return Err(std::io::Error::last_os_error().into());
1360    }
1361
1362    let index = (CAP_SYS_ADMIN / CAP_WORD_BITS) as usize;
1363    let mask = 1u32 << (CAP_SYS_ADMIN % CAP_WORD_BITS);
1364    let had_sys_admin = data[index].effective & mask != 0
1365        || data[index].permitted & mask != 0
1366        || data[index].inheritable & mask != 0;
1367
1368    if had_sys_admin {
1369        data[index].effective &= !mask;
1370        data[index].permitted &= !mask;
1371        data[index].inheritable &= !mask;
1372
1373        if unsafe { libc::syscall(libc::SYS_capset, &mut header, data.as_ptr()) } != 0 {
1374            return Err(std::io::Error::last_os_error().into());
1375        }
1376    }
1377
1378    let ret = unsafe { libc::prctl(PR_CAPBSET_DROP, CAP_SYS_ADMIN, 0, 0, 0) };
1379    if ret != 0 {
1380        let err = std::io::Error::last_os_error();
1381        let errno = err.raw_os_error();
1382        // Already-unprivileged callers may also lack CAP_SETPCAP for the bounding-set drop.
1383        let already_unprivileged = !had_sys_admin && errno == Some(libc::EPERM);
1384        if errno != Some(libc::EINVAL) && !already_unprivileged {
1385            return Err(err.into());
1386        }
1387    }
1388
1389    Ok(())
1390}
1391
1392pub(crate) fn resolve_default_user(default_user: Option<&str>) -> AgentdResult<(u32, u32)> {
1393    let Some(spec) = default_user
1394        .map(str::trim)
1395        .filter(|value| !value.is_empty())
1396    else {
1397        return Ok((0, 0));
1398    };
1399
1400    let resolved = resolve_user_spec(spec)?;
1401    Ok((resolved.uid, resolved.gid))
1402}
1403
1404/// Like [`resolve_default_user`], plus the supplementary groups exec would apply to that user.
1405pub(crate) fn resolve_user_groups(user: &str) -> AgentdResult<(u32, u32, Vec<libc::gid_t>)> {
1406    let resolved = resolve_user_spec(user)?;
1407    let mut groups = Vec::new();
1408    if let Some(ref name) = resolved.initgroups_user {
1409        let mut count: libc::c_int = 32;
1410        loop {
1411            groups.resize(count as usize, 0);
1412            if unsafe {
1413                libc::getgrouplist(name.as_ptr(), resolved.gid, groups.as_mut_ptr(), &mut count)
1414            } >= 0
1415            {
1416                groups.truncate(count as usize);
1417                break;
1418            }
1419            count = count.max(groups.len() as libc::c_int + 1);
1420        }
1421    }
1422    Ok((resolved.uid, resolved.gid, groups))
1423}
1424
1425fn resolve_requested_user(
1426    req: &ExecRequest,
1427    default_user: Option<&str>,
1428) -> AgentdResult<Option<ResolvedUser>> {
1429    let default_user = default_user
1430        .map(str::trim)
1431        .filter(|value| !value.is_empty());
1432    let requested = req
1433        .user
1434        .as_deref()
1435        .map(str::trim)
1436        .filter(|value| !value.is_empty())
1437        .or(default_user);
1438
1439    requested.map(resolve_user_spec).transpose()
1440}
1441
1442fn resolve_user_spec(spec: &str) -> AgentdResult<ResolvedUser> {
1443    let (user_part, group_part) = match spec.split_once(':') {
1444        Some((user, group)) => (user.trim(), Some(group.trim())),
1445        None => (spec.trim(), None),
1446    };
1447
1448    if user_part.is_empty() {
1449        return Err(AgentdError::ExecSession("user spec has empty user".into()));
1450    }
1451
1452    let passwd = if let Ok(uid) = parse_id(user_part) {
1453        lookup_passwd_by_uid(uid)?
1454    } else {
1455        lookup_passwd_by_name(user_part)?
1456            .ok_or_else(|| AgentdError::UserNotFound(user_part.to_owned()))?
1457            .into()
1458    };
1459
1460    let (uid, passwd_entry) = match passwd {
1461        ResolvedUserLookup::Known(entry) => (entry.uid, Some(entry)),
1462        ResolvedUserLookup::Numeric(uid) => (uid, None),
1463    };
1464
1465    let gid = match group_part {
1466        Some("") => {
1467            return Err(AgentdError::ExecSession("user spec has empty group".into()));
1468        }
1469        Some(group) => resolve_group_spec(group)?,
1470        None => passwd_entry
1471            .as_ref()
1472            .map(|entry| entry.gid)
1473            .unwrap_or_else(|| unsafe { libc::getgid() }),
1474    };
1475
1476    let initgroups_user = passwd_entry
1477        .as_ref()
1478        .map(|entry| CString::new(entry.name.as_str()))
1479        .transpose()
1480        .map_err(|e| AgentdError::ExecSession(format!("invalid guest user name: {e}")))?;
1481
1482    Ok(ResolvedUser {
1483        uid,
1484        gid,
1485        initgroups_user,
1486        home_dir: passwd_entry
1487            .as_ref()
1488            .and_then(|entry| entry.home_dir.as_deref())
1489            .map(CString::new)
1490            .transpose()
1491            .map_err(|e| AgentdError::ExecSession(format!("invalid guest home directory: {e}")))?,
1492    })
1493}
1494
1495enum ResolvedUserLookup {
1496    Known(PasswdEntry),
1497    Numeric(libc::uid_t),
1498}
1499
1500impl From<PasswdEntry> for ResolvedUserLookup {
1501    fn from(value: PasswdEntry) -> Self {
1502        Self::Known(value)
1503    }
1504}
1505
1506fn resolve_group_spec(spec: &str) -> AgentdResult<libc::gid_t> {
1507    if let Ok(gid) = parse_id(spec) {
1508        return Ok(gid);
1509    }
1510
1511    lookup_group_by_name(spec)?
1512        .map(|entry| entry.gid)
1513        .ok_or_else(|| AgentdError::GroupNotFound(spec.to_owned()))
1514}
1515
1516fn parse_id(value: &str) -> Result<u32, std::num::ParseIntError> {
1517    value.parse::<u32>()
1518}
1519
1520fn lookup_passwd_by_name(name: &str) -> AgentdResult<Option<PasswdEntry>> {
1521    let name = CString::new(name)
1522        .map_err(|e| AgentdError::ExecSession(format!("invalid guest user name: {e}")))?;
1523    let mut pwd = MaybeUninit::<libc::passwd>::uninit();
1524    let mut result = ptr::null_mut();
1525    let mut buf = vec![0u8; lookup_buffer_len()];
1526    let rc = unsafe {
1527        libc::getpwnam_r(
1528            name.as_ptr(),
1529            pwd.as_mut_ptr(),
1530            buf.as_mut_ptr().cast(),
1531            buf.len(),
1532            &mut result,
1533        )
1534    };
1535    if rc != 0 {
1536        return Err(AgentdError::ExecSession(format!(
1537            "failed to resolve guest user {name:?}: {}",
1538            std::io::Error::from_raw_os_error(rc)
1539        )));
1540    }
1541    if result.is_null() {
1542        return Ok(None);
1543    }
1544
1545    let pwd = unsafe { pwd.assume_init() };
1546    let name = unsafe { CStr::from_ptr(pwd.pw_name) }
1547        .to_string_lossy()
1548        .into_owned();
1549    let home_dir = unsafe { CStr::from_ptr(pwd.pw_dir) }
1550        .to_string_lossy()
1551        .into_owned();
1552    Ok(Some(PasswdEntry {
1553        name,
1554        uid: pwd.pw_uid,
1555        gid: pwd.pw_gid,
1556        home_dir: (!home_dir.is_empty()).then_some(home_dir),
1557    }))
1558}
1559
1560fn lookup_passwd_by_uid(uid: libc::uid_t) -> AgentdResult<ResolvedUserLookup> {
1561    let mut pwd = MaybeUninit::<libc::passwd>::uninit();
1562    let mut result = ptr::null_mut();
1563    let mut buf = vec![0u8; lookup_buffer_len()];
1564    let rc = unsafe {
1565        libc::getpwuid_r(
1566            uid,
1567            pwd.as_mut_ptr(),
1568            buf.as_mut_ptr().cast(),
1569            buf.len(),
1570            &mut result,
1571        )
1572    };
1573    if rc != 0 {
1574        return Err(AgentdError::ExecSession(format!(
1575            "failed to resolve guest uid {uid}: {}",
1576            std::io::Error::from_raw_os_error(rc)
1577        )));
1578    }
1579    if result.is_null() {
1580        return Ok(ResolvedUserLookup::Numeric(uid));
1581    }
1582
1583    let pwd = unsafe { pwd.assume_init() };
1584    let name = unsafe { CStr::from_ptr(pwd.pw_name) }
1585        .to_string_lossy()
1586        .into_owned();
1587    let home_dir = unsafe { CStr::from_ptr(pwd.pw_dir) }
1588        .to_string_lossy()
1589        .into_owned();
1590    Ok(ResolvedUserLookup::Known(PasswdEntry {
1591        name,
1592        uid: pwd.pw_uid,
1593        gid: pwd.pw_gid,
1594        home_dir: (!home_dir.is_empty()).then_some(home_dir),
1595    }))
1596}
1597
1598fn lookup_group_by_name(name: &str) -> AgentdResult<Option<GroupEntry>> {
1599    let name = CString::new(name)
1600        .map_err(|e| AgentdError::ExecSession(format!("invalid guest group name: {e}")))?;
1601    let mut grp = MaybeUninit::<libc::group>::uninit();
1602    let mut result = ptr::null_mut();
1603    let mut buf = vec![0u8; lookup_buffer_len()];
1604    let rc = unsafe {
1605        libc::getgrnam_r(
1606            name.as_ptr(),
1607            grp.as_mut_ptr(),
1608            buf.as_mut_ptr().cast(),
1609            buf.len(),
1610            &mut result,
1611        )
1612    };
1613    if rc != 0 {
1614        return Err(AgentdError::ExecSession(format!(
1615            "failed to resolve guest group {name:?}: {}",
1616            std::io::Error::from_raw_os_error(rc)
1617        )));
1618    }
1619    if result.is_null() {
1620        return Ok(None);
1621    }
1622
1623    let grp = unsafe { grp.assume_init() };
1624    Ok(Some(GroupEntry { gid: grp.gr_gid }))
1625}
1626
1627fn lookup_buffer_len() -> usize {
1628    let size = unsafe { libc::sysconf(libc::_SC_GETPW_R_SIZE_MAX) };
1629    if size > 0 { size as usize } else { 16 * 1024 }
1630}
1631
1632fn apply_resolved_groups(user: &ResolvedUser) -> AgentdResult<()> {
1633    if let Some(ref name) = user.initgroups_user {
1634        if unsafe { libc::initgroups(name.as_ptr(), user.gid) } != 0 {
1635            return Err(std::io::Error::last_os_error().into());
1636        }
1637    } else if unsafe { libc::setgroups(0, ptr::null()) } != 0 {
1638        return Err(std::io::Error::last_os_error().into());
1639    }
1640
1641    Ok(())
1642}
1643
1644fn apply_resolved_user(user: &ResolvedUser) -> AgentdResult<()> {
1645    if unsafe { libc::setgid(user.gid) } != 0 {
1646        return Err(std::io::Error::last_os_error().into());
1647    }
1648    if unsafe { libc::setuid(user.uid) } != 0 {
1649        return Err(std::io::Error::last_os_error().into());
1650    }
1651
1652    Ok(())
1653}
1654
1655fn default_home_dir(
1656    req: &ExecRequest,
1657    user: Option<&ResolvedUser>,
1658) -> AgentdResult<Option<CString>> {
1659    if env_contains_key(&req.env, "HOME") {
1660        return Ok(None);
1661    }
1662
1663    if let Some(user) = user {
1664        return Ok(user.home_dir.clone());
1665    }
1666
1667    Ok(resolve_user_spec(DEFAULT_USER_SPEC)?.home_dir)
1668}
1669
1670fn env_contains_key(env: &[String], key: &str) -> bool {
1671    env.iter().any(|entry| {
1672        entry
1673            .split_once('=')
1674            .map(|(entry_key, _)| entry_key == key)
1675            .unwrap_or(false)
1676    })
1677}
1678
1679fn agentd_to_io_error(err: AgentdError) -> std::io::Error {
1680    std::io::Error::other(err.to_string())
1681}
1682
1683/// Keep one owned descriptor across readiness waits; cancellation cannot leave a blocking task
1684/// writing through a borrowed fd after its session has been removed or restored.
1685fn nonblocking_input(fd: OwnedFd) -> std::io::Result<AsyncFd<OwnedFd>> {
1686    let flags = unsafe { libc::fcntl(fd.as_raw_fd(), libc::F_GETFL) };
1687    if flags < 0
1688        || unsafe { libc::fcntl(fd.as_raw_fd(), libc::F_SETFL, flags | libc::O_NONBLOCK) } < 0
1689    {
1690        return Err(std::io::Error::last_os_error());
1691    }
1692    AsyncFd::new(fd)
1693}
1694
1695fn write_nonblocking_fd(fd: RawFd, data: &[u8]) -> std::io::Result<usize> {
1696    loop {
1697        let written = unsafe { libc::write(fd, data.as_ptr().cast(), data.len()) };
1698        if written >= 0 {
1699            return Ok(written as usize);
1700        }
1701        let error = std::io::Error::last_os_error();
1702        if error.kind() != std::io::ErrorKind::Interrupted {
1703            return Err(error);
1704        }
1705    }
1706}
1707
1708fn wait_fd_readable(fd: RawFd) -> AgentdResult<()> {
1709    let mut pollfd = libc::pollfd {
1710        fd,
1711        events: libc::POLLIN,
1712        revents: 0,
1713    };
1714
1715    loop {
1716        let ret = unsafe { libc::poll(&mut pollfd, 1, -1) };
1717        if ret < 0 {
1718            let err = std::io::Error::last_os_error();
1719            if err.raw_os_error() == Some(libc::EINTR) {
1720                continue;
1721            }
1722            return Err(AgentdError::Io(err));
1723        }
1724        if ret == 0 {
1725            continue;
1726        }
1727        // Always retry read on HUP/ERR too: a PTY may still contain final output before EIO.
1728        return Ok(());
1729    }
1730}
1731
1732/// Background task that reads from a PTY master fd and sends output events.
1733async fn pty_reader_task(
1734    id: u32,
1735    master_fd: OwnedFd,
1736    exit_watcher: ProcessExitWatcher,
1737    tx: SessionOutputSender,
1738) {
1739    let tx_output = tx.clone();
1740    let runtime_handle = tokio::runtime::Handle::current();
1741    let read_result = tokio::task::spawn_blocking(move || {
1742        // PTY masters are safer with a dedicated blocking read loop than with
1743        // edge-driven readiness. Fast writers followed by process exit can
1744        // strand the tail behind a missed wakeup/HUP transition.
1745        let raw = master_fd.as_raw_fd();
1746        // The duplicated master shares O_NONBLOCK with stdin. Never clear that flag here:
1747        // a blocked write would otherwise strand lifecycle handling on the agent actor.
1748
1749        loop {
1750            let mut buf = [0u8; 4096];
1751            let n = unsafe { libc::read(raw, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
1752
1753            if n > 0 {
1754                let n = n as usize;
1755                let sent = runtime_handle.block_on(async {
1756                    let Some(permit) = tx_output.reserve(n).await else {
1757                        return false;
1758                    };
1759                    tx_output
1760                        .send_reserved(id, SessionOutput::Stdout(buf[..n].to_vec()), permit)
1761                        .await
1762                });
1763                if !sent {
1764                    break;
1765                }
1766                continue;
1767            }
1768
1769            if n == 0 {
1770                break;
1771            }
1772
1773            let err = std::io::Error::last_os_error();
1774            match err.raw_os_error() {
1775                Some(libc::EINTR) => continue,
1776                Some(libc::EAGAIN) => {
1777                    if wait_fd_readable(raw).is_err() {
1778                        break;
1779                    }
1780                }
1781                Some(libc::EIO) => break,
1782                _ => break,
1783            }
1784        }
1785    })
1786    .await;
1787
1788    let _ = read_result;
1789
1790    let code = exit_watcher.await;
1791    let _ = tx.send(id, SessionOutput::Exited(code)).await;
1792}
1793
1794/// Background task that reads from piped stdout/stderr and sends output events.
1795async fn pipe_reader_task(
1796    id: u32,
1797    stdout: Option<tokio::process::ChildStdout>,
1798    stderr: Option<tokio::process::ChildStderr>,
1799    exit_watcher: ProcessExitWatcher,
1800    tx: SessionOutputSender,
1801) {
1802    let mut stdout = stdout;
1803    let mut stderr = stderr;
1804    let mut stdout_eof = stdout.is_none();
1805    let mut stderr_eof = stderr.is_none();
1806
1807    while !stdout_eof || !stderr_eof {
1808        let mut stdout_buf = [0u8; 4096];
1809        let mut stderr_buf = [0u8; 4096];
1810
1811        tokio::select! {
1812            result = async {
1813                match stdout.as_mut() {
1814                    Some(out) => out.read(&mut stdout_buf).await,
1815                    None => std::future::pending().await,
1816                }
1817            }, if !stdout_eof => {
1818                match result {
1819                    Ok(0) | Err(_) => {
1820                        stdout = None;
1821                        stdout_eof = true;
1822                    }
1823                    Ok(n) => {
1824                        let Some(permit) = tx.reserve(n).await else {
1825                            break;
1826                        };
1827                        if !tx
1828                            .send_reserved(
1829                                id,
1830                                SessionOutput::Stdout(stdout_buf[..n].to_vec()),
1831                                permit,
1832                            )
1833                            .await
1834                        {
1835                            break;
1836                        }
1837                    }
1838                }
1839            }
1840            result = async {
1841                match stderr.as_mut() {
1842                    Some(err) => err.read(&mut stderr_buf).await,
1843                    None => std::future::pending().await,
1844                }
1845            }, if !stderr_eof => {
1846                match result {
1847                    Ok(0) | Err(_) => {
1848                        stderr = None;
1849                        stderr_eof = true;
1850                    }
1851                    Ok(n) => {
1852                        let Some(permit) = tx.reserve(n).await else {
1853                            break;
1854                        };
1855                        if !tx
1856                            .send_reserved(
1857                                id,
1858                                SessionOutput::Stderr(stderr_buf[..n].to_vec()),
1859                                permit,
1860                            )
1861                            .await
1862                        {
1863                            break;
1864                        }
1865                    }
1866                }
1867            }
1868        }
1869    }
1870
1871    let code = exit_watcher.await;
1872
1873    let _ = tx.send(id, SessionOutput::Exited(code)).await;
1874}
1875
1876//--------------------------------------------------------------------------------------------------
1877// Tests
1878//--------------------------------------------------------------------------------------------------
1879
1880#[cfg(test)]
1881mod tests {
1882    use std::collections::HashMap;
1883    use std::io::Read;
1884    use std::process::{Command as StdCommand, Stdio as StdStdio};
1885    use std::time::Duration;
1886
1887    use tokio::time;
1888
1889    use microsandbox_protocol::exec::ExecRequest;
1890
1891    use super::*;
1892
1893    const REAP_HELPER_ENV: &str = "MSB_AGENTD_SESSION_REAP_HELPER";
1894    const REAP_HELPER_SENTINEL: &str = "session-reap-helper-passed";
1895    const REAP_TEST_NAME: &str = "session::tests::test_spawn_reaps_adopted_descendant";
1896    const CONCURRENT_HELPER_ENV: &str = "MSB_AGENTD_CONCURRENT_SPAWN_HELPER";
1897    const CONCURRENT_HELPER_SENTINEL: &str = "concurrent-spawn-helper-passed";
1898    const CONCURRENT_TEST_NAME: &str = "session::tests::test_concurrent_spawn_exit_codes";
1899    const RUNTIME_HELPER_ENV: &str = "MSB_AGENTD_RUNTIME_REPLACEMENT_HELPER";
1900    const RUNTIME_HELPER_SENTINEL: &str = "runtime-replacement-helper-passed";
1901    const RUNTIME_TEST_NAME: &str = "session::tests::test_spawn_survives_runtime_replacement";
1902    const PIPE_OWNER_HELPER_ENV: &str = "MSB_AGENTD_PIPE_OWNER_HELPER";
1903    const PIPE_OWNER_HELPER_SENTINEL: &str = "pipe-owner-helper-passed";
1904
1905    #[tokio::test]
1906    async fn session_output_permit_lives_until_envelope_is_consumed() {
1907        let (tx, mut rx) = SessionOutputSender::channel();
1908        assert!(tx.send(7, SessionOutput::Stdout(vec![0; 4096])).await);
1909        assert_eq!(
1910            tx.control_budget.available_permits(),
1911            SESSION_OUTPUT_BYTE_CAPACITY - 4096
1912        );
1913
1914        let envelope = rx.recv().await.unwrap();
1915        assert_eq!(
1916            tx.control_budget.available_permits(),
1917            SESSION_OUTPUT_BYTE_CAPACITY - 4096
1918        );
1919        drop(envelope);
1920        assert_eq!(
1921            tx.control_budget.available_permits(),
1922            SESSION_OUTPUT_BYTE_CAPACITY
1923        );
1924    }
1925
1926    #[tokio::test]
1927    async fn scoped_output_sender_captures_client_incarnation() {
1928        let incarnation = [0x44; 16];
1929        let (tx, mut rx) = SessionOutputSender::channel();
1930        let scoped = tx.with_incarnation(Some(incarnation));
1931
1932        assert!(scoped.send(7, SessionOutput::Exited(0)).await);
1933        let envelope = rx.recv().await.unwrap();
1934
1935        assert_eq!(envelope.id, 7);
1936        assert_eq!(envelope.incarnation, Some(incarnation));
1937    }
1938
1939    #[tokio::test]
1940    async fn split_output_queues_keep_bulk_lifecycle_commands_independent() {
1941        let incarnation = [0x55; 16];
1942        let (tx, mut control_rx, mut bulk_rx, mut command_rx) =
1943            SessionOutputSender::split_channel();
1944        let scoped = tx.with_incarnation(Some(incarnation));
1945        let record = BulkRecord {
1946            id: 9,
1947            kind: microsandbox_protocol::bulk::BulkKind::Filesystem,
1948            flow: microsandbox_protocol::bulk::BulkFlow::GuestToHost,
1949            offset: 0,
1950            payload: b"bulk".as_slice().into(),
1951        };
1952        assert!(
1953            scoped
1954                .send(
1955                    9,
1956                    SessionOutput::Bulk(BulkSessionOutput::new(record, RawActivity::fs_bytes(4),)),
1957                )
1958                .await
1959        );
1960        assert!(scoped.send(10, SessionOutput::Exited(0)).await);
1961        let mut completed = scoped.drop_bulk_flow(9).unwrap().unwrap();
1962
1963        assert!(matches!(
1964            bulk_rx.recv().await.unwrap().output,
1965            SessionOutput::Bulk(_)
1966        ));
1967        assert!(matches!(
1968            control_rx.recv().await.unwrap().output,
1969            SessionOutput::Exited(0)
1970        ));
1971        let BulkOutputCommand::DropFlow {
1972            incarnation: command_incarnation,
1973            id,
1974            completion,
1975        } = command_rx.recv().await.unwrap()
1976        else {
1977            panic!("expected flow cleanup command");
1978        };
1979        assert_eq!(command_incarnation, incarnation);
1980        assert_eq!(id, 9);
1981        completion.send(()).unwrap();
1982        assert_eq!(completed.try_recv(), Ok(()));
1983    }
1984
1985    #[tokio::test]
1986    async fn combined_mode_restores_one_ordered_output_queue() {
1987        let (mut tx, mut control_rx, mut bulk_rx, _command_rx) =
1988            SessionOutputSender::split_channel();
1989        tx.disable_bulk_scheduler();
1990        let record = BulkRecord {
1991            id: 9,
1992            kind: microsandbox_protocol::bulk::BulkKind::Filesystem,
1993            flow: microsandbox_protocol::bulk::BulkFlow::GuestToHost,
1994            offset: 0,
1995            payload: b"bulk".as_slice().into(),
1996        };
1997        assert!(
1998            tx.send(
1999                9,
2000                SessionOutput::Bulk(BulkSessionOutput::new(record, RawActivity::fs_bytes(4),)),
2001            )
2002            .await
2003        );
2004        assert!(tx.send(9, SessionOutput::Exited(0)).await);
2005
2006        assert!(matches!(
2007            control_rx.recv().await.unwrap().output,
2008            SessionOutput::Bulk(_)
2009        ));
2010        assert!(matches!(
2011            control_rx.recv().await.unwrap().output,
2012            SessionOutput::Exited(0)
2013        ));
2014        assert!(bulk_rx.recv().await.is_none());
2015    }
2016
2017    #[tokio::test]
2018    async fn control_output_remains_admissible_when_data_budget_is_exhausted() {
2019        let (tx, mut rx) = SessionOutputSender::channel();
2020        let full_budget = tx.reserve(SESSION_OUTPUT_BYTE_CAPACITY).await.unwrap();
2021
2022        assert!(tx.send(9, SessionOutput::Exited(0)).await);
2023        let envelope = rx.recv().await.unwrap();
2024        assert!(matches!(envelope.output, SessionOutput::Exited(0)));
2025        drop(full_budget);
2026    }
2027    const PIPE_OWNER_TEST_NAME: &str =
2028        "session::tests::test_piped_process_exit_outlives_spawning_runtime";
2029
2030    #[test]
2031    fn test_spawn_reaps_adopted_descendant() {
2032        if std::env::var_os(REAP_HELPER_ENV).is_some() {
2033            let runtime = tokio::runtime::Builder::new_current_thread()
2034                .enable_all()
2035                .build()
2036                .expect("session reap test runtime");
2037            runtime.block_on(run_adopted_descendant_scenario());
2038            println!("{REAP_HELPER_SENTINEL}");
2039            return;
2040        }
2041
2042        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2043            .args(["--exact", REAP_TEST_NAME, "--nocapture"])
2044            .env(REAP_HELPER_ENV, "1")
2045            .stdout(StdStdio::piped())
2046            .spawn()
2047            .expect("spawn isolated session reap test");
2048        let mut output = String::new();
2049        helper
2050            .stdout
2051            .take()
2052            .expect("helper stdout")
2053            .read_to_string(&mut output)
2054            .expect("read helper stdout");
2055
2056        match helper.wait() {
2057            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2058            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2059            Err(error) => panic!("wait for helper: {error}"),
2060        }
2061        assert!(
2062            output.contains(REAP_HELPER_SENTINEL),
2063            "helper did not complete the session reap scenario:\n{output}"
2064        );
2065    }
2066
2067    async fn run_adopted_descendant_scenario() {
2068        let ret = unsafe { libc::prctl(libc::PR_SET_CHILD_SUBREAPER, 1) };
2069        assert_eq!(
2070            ret,
2071            0,
2072            "set child subreaper: {}",
2073            std::io::Error::last_os_error()
2074        );
2075
2076        let (tx, mut rx) = SessionOutputSender::channel();
2077        let req = ExecRequest {
2078            cmd: "/bin/sh".to_string(),
2079            args: vec!["-c".to_string(), "sleep 30 & echo $!".to_string()],
2080            env: Vec::new(),
2081            cwd: None,
2082            user: None,
2083            tty: false,
2084            rows: 24,
2085            cols: 80,
2086            rlimits: Vec::new(),
2087        };
2088
2089        let session = ExecSession::spawn(17, &req, tx, None, SecurityProfile::Default, None)
2090            .expect("spawn background descendant session");
2091        let leader_pid = session.pid() as i32;
2092        let mut stdout = Vec::new();
2093        time::timeout(Duration::from_secs(10), async {
2094            while !stdout.contains(&b'\n') {
2095                let envelope = rx.recv().await.expect("session output");
2096                assert_eq!(envelope.id, 17);
2097                match envelope.output {
2098                    SessionOutput::Stdout(data) => stdout.extend_from_slice(&data),
2099                    SessionOutput::Exited(code) => panic!("session exited early with {code}"),
2100                    SessionOutput::Stderr(_) | SessionOutput::Raw(_) | SessionOutput::Bulk(_) => {}
2101                }
2102            }
2103        })
2104        .await
2105        .expect("wait for background descendant session");
2106
2107        let descendant_pid: i32 = String::from_utf8(stdout)
2108            .expect("descendant PID is UTF-8")
2109            .trim()
2110            .parse()
2111            .expect("parse descendant PID");
2112        let expected_parent = std::process::id().to_string();
2113        let status_path = format!("/proc/{descendant_pid}/status");
2114        time::timeout(Duration::from_secs(5), async {
2115            loop {
2116                if let Ok(status) = std::fs::read_to_string(&status_path)
2117                    && status
2118                        .lines()
2119                        .find_map(|line| line.strip_prefix("PPid:"))
2120                        .is_some_and(|ppid| ppid.trim() == expected_parent)
2121                {
2122                    break;
2123                }
2124                time::sleep(Duration::from_millis(10)).await;
2125            }
2126        })
2127        .await
2128        .expect("descendant should be adopted by the helper subreaper");
2129
2130        let leader_path = format!("/proc/{leader_pid}");
2131        time::timeout(Duration::from_secs(5), async {
2132            while std::path::Path::new(&leader_path).exists() {
2133                time::sleep(Duration::from_millis(10)).await;
2134            }
2135        })
2136        .await
2137        .expect("direct child should be reaped before signalling its descendants");
2138
2139        session
2140            .send_signal(libc::SIGTERM)
2141            .expect("signal descendants through completed process registration");
2142        let exit = time::timeout(Duration::from_secs(5), async {
2143            loop {
2144                let envelope = rx.recv().await.expect("session output after signal");
2145                assert_eq!(envelope.id, 17);
2146                if let SessionOutput::Exited(code) = envelope.output {
2147                    break code;
2148                }
2149            }
2150        })
2151        .await
2152        .expect("session should finish after its descendant is signalled");
2153        assert_eq!(exit, 0);
2154
2155        let proc_path = format!("/proc/{descendant_pid}");
2156        time::timeout(Duration::from_secs(5), async {
2157            while std::path::Path::new(&proc_path).exists() {
2158                time::sleep(Duration::from_millis(10)).await;
2159            }
2160        })
2161        .await
2162        .expect("descendant should be reaped");
2163
2164        let ret = unsafe { libc::waitpid(descendant_pid, ptr::null_mut(), libc::WNOHANG) };
2165        assert_eq!(ret, -1, "descendant {descendant_pid} was not reaped");
2166        assert_eq!(
2167            std::io::Error::last_os_error().raw_os_error(),
2168            Some(libc::ECHILD)
2169        );
2170    }
2171
2172    #[test]
2173    fn test_concurrent_spawn_exit_codes() {
2174        if std::env::var_os(CONCURRENT_HELPER_ENV).is_some() {
2175            let runtime = tokio::runtime::Builder::new_current_thread()
2176                .enable_all()
2177                .build()
2178                .expect("concurrent spawn test runtime");
2179            runtime.block_on(run_concurrent_spawn_scenario());
2180            println!("{CONCURRENT_HELPER_SENTINEL}");
2181            return;
2182        }
2183
2184        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2185            .args(["--exact", CONCURRENT_TEST_NAME, "--nocapture"])
2186            .env(CONCURRENT_HELPER_ENV, "1")
2187            .stdout(StdStdio::piped())
2188            .spawn()
2189            .expect("spawn isolated concurrent session test");
2190        let mut output = String::new();
2191        helper
2192            .stdout
2193            .take()
2194            .expect("helper stdout")
2195            .read_to_string(&mut output)
2196            .expect("read helper stdout");
2197
2198        match helper.wait() {
2199            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2200            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2201            Err(error) => panic!("wait for helper: {error}"),
2202        }
2203        assert!(
2204            output.contains(CONCURRENT_HELPER_SENTINEL),
2205            "helper did not complete the concurrent spawn scenario:\n{output}"
2206        );
2207    }
2208
2209    async fn run_concurrent_spawn_scenario() {
2210        const PROCESS_COUNT: u32 = 12;
2211
2212        let runtime_handle = tokio::runtime::Handle::current();
2213        let (tx, mut rx) = SessionOutputSender::channel();
2214        let mut spawn_threads = Vec::new();
2215        for offset in 0..PROCESS_COUNT {
2216            let handle = runtime_handle.clone();
2217            let tx = tx.clone();
2218            spawn_threads.push(std::thread::spawn(move || {
2219                let _runtime = handle.enter();
2220                let code = 20 + offset as i32;
2221                let req = ExecRequest {
2222                    cmd: "/bin/sh".to_string(),
2223                    args: vec!["-c".to_string(), format!("exit {code}")],
2224                    env: Vec::new(),
2225                    cwd: None,
2226                    user: None,
2227                    tty: offset % 2 == 1,
2228                    rows: 24,
2229                    cols: 80,
2230                    rlimits: Vec::new(),
2231                };
2232                ExecSession::spawn(100 + offset, &req, tx, None, SecurityProfile::Default, None)
2233            }));
2234        }
2235        drop(tx);
2236
2237        let mut sessions = Vec::new();
2238        for thread in spawn_threads {
2239            sessions.push(
2240                thread
2241                    .join()
2242                    .expect("concurrent spawn thread")
2243                    .expect("concurrent process spawn"),
2244            );
2245        }
2246
2247        let mut exits = HashMap::new();
2248        time::timeout(Duration::from_secs(15), async {
2249            while exits.len() < PROCESS_COUNT as usize {
2250                let envelope = rx.recv().await.expect("session output");
2251                if let SessionOutput::Exited(code) = envelope.output {
2252                    exits.insert(envelope.id, code);
2253                }
2254            }
2255        })
2256        .await
2257        .expect("wait for concurrent exits");
2258
2259        for offset in 0..PROCESS_COUNT {
2260            assert_eq!(exits.get(&(100 + offset)), Some(&(20 + offset as i32)));
2261        }
2262        drop(sessions);
2263    }
2264
2265    #[test]
2266    fn test_spawn_survives_runtime_replacement() {
2267        if std::env::var_os(RUNTIME_HELPER_ENV).is_some() {
2268            run_runtime_replacement_scenario();
2269            println!("{RUNTIME_HELPER_SENTINEL}");
2270            return;
2271        }
2272
2273        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2274            .args(["--exact", RUNTIME_TEST_NAME, "--nocapture"])
2275            .env(RUNTIME_HELPER_ENV, "1")
2276            .stdout(StdStdio::piped())
2277            .spawn()
2278            .expect("spawn isolated runtime replacement test");
2279        let mut output = String::new();
2280        helper
2281            .stdout
2282            .take()
2283            .expect("helper stdout")
2284            .read_to_string(&mut output)
2285            .expect("read helper stdout");
2286
2287        match helper.wait() {
2288            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2289            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2290            Err(error) => panic!("wait for helper: {error}"),
2291        }
2292        assert!(
2293            output.contains(RUNTIME_HELPER_SENTINEL),
2294            "helper did not complete the runtime replacement scenario:\n{output}"
2295        );
2296    }
2297
2298    fn run_runtime_replacement_scenario() {
2299        for (id, code) in [(201, 51), (202, 52)] {
2300            let runtime = tokio::runtime::Builder::new_current_thread()
2301                .enable_all()
2302                .build()
2303                .expect("replacement test runtime");
2304            runtime.block_on(run_single_pipe_spawn(id, code));
2305        }
2306    }
2307
2308    async fn run_single_pipe_spawn(id: u32, code: i32) {
2309        let (tx, mut rx) = SessionOutputSender::channel();
2310        let req = ExecRequest {
2311            cmd: "/bin/sh".to_string(),
2312            args: vec!["-c".to_string(), format!("exit {code}")],
2313            env: Vec::new(),
2314            cwd: None,
2315            user: None,
2316            tty: false,
2317            rows: 24,
2318            cols: 80,
2319            rlimits: Vec::new(),
2320        };
2321        let _session = ExecSession::spawn(id, &req, tx, None, SecurityProfile::Default, None)
2322            .expect("spawn session on replacement runtime");
2323
2324        let actual = time::timeout(Duration::from_secs(5), async {
2325            loop {
2326                let envelope = rx.recv().await.expect("session output");
2327                assert_eq!(envelope.id, id);
2328                if let SessionOutput::Exited(actual) = envelope.output {
2329                    break actual;
2330                }
2331            }
2332        })
2333        .await
2334        .expect("wait for exit on replacement runtime");
2335        assert_eq!(actual, code);
2336    }
2337
2338    #[test]
2339    fn test_piped_process_exit_outlives_spawning_runtime() {
2340        if std::env::var_os(PIPE_OWNER_HELPER_ENV).is_some() {
2341            run_piped_process_exit_scenario();
2342            println!("{PIPE_OWNER_HELPER_SENTINEL}");
2343            return;
2344        }
2345
2346        let mut helper = StdCommand::new(std::env::current_exe().expect("current test binary"))
2347            .args(["--exact", PIPE_OWNER_TEST_NAME, "--nocapture"])
2348            .env(PIPE_OWNER_HELPER_ENV, "1")
2349            .stdout(StdStdio::piped())
2350            .spawn()
2351            .expect("spawn isolated pipe owner test");
2352        let mut output = String::new();
2353        helper
2354            .stdout
2355            .take()
2356            .expect("helper stdout")
2357            .read_to_string(&mut output)
2358            .expect("read helper stdout");
2359
2360        match helper.wait() {
2361            Ok(status) => assert!(status.success(), "helper failed: {status}\n{output}"),
2362            Err(error) if error.raw_os_error() == Some(libc::ECHILD) => {}
2363            Err(error) => panic!("wait for helper: {error}"),
2364        }
2365        assert!(
2366            output.contains(PIPE_OWNER_HELPER_SENTINEL),
2367            "helper did not complete the pipe owner scenario:\n{output}"
2368        );
2369    }
2370
2371    fn run_piped_process_exit_scenario() {
2372        let process_manager = ProcessManager::get().expect("get process manager");
2373        let spawning_runtime = tokio::runtime::Builder::new_current_thread()
2374            .enable_all()
2375            .build()
2376            .expect("spawning runtime");
2377        let exit_watcher = {
2378            let _runtime_guard = spawning_runtime.enter();
2379            let mut command = Command::new("/bin/sh");
2380            command
2381                .args(["-c", "exit 63"])
2382                .stdin(Stdio::piped())
2383                .stdout(Stdio::piped())
2384                .stderr(Stdio::piped());
2385            let process =
2386                spawn_piped_process(command, &process_manager).expect("spawn piped process");
2387            let PipedProcess { exit_watcher, .. } = process;
2388            exit_watcher
2389        };
2390        drop(spawning_runtime);
2391
2392        let waiting_runtime = tokio::runtime::Builder::new_current_thread()
2393            .enable_all()
2394            .build()
2395            .expect("waiting runtime");
2396        let code = waiting_runtime.block_on(async {
2397            time::timeout(Duration::from_secs(5), exit_watcher)
2398                .await
2399                .expect("wait for piped process exit")
2400        });
2401        assert_eq!(code, 63);
2402    }
2403
2404    // Run on Linux as root with CAP_SYS_RESOURCE and a finite memlock hard limit.
2405    // The test executable must be static so nofile=0 cannot block its loader.
2406    // For Docker: --cap-add SYS_RESOURCE --ulimit memlock=8388608:8388608.
2407    #[test]
2408    #[ignore = "requires root, CAP_SYS_RESOURCE, finite memlock, and a static test binary"]
2409    fn test_piped_nonroot_rlimits() {
2410        check_nonroot_rlimits(false);
2411    }
2412
2413    #[test]
2414    #[ignore = "requires root, CAP_SYS_RESOURCE, finite memlock, and a static test binary"]
2415    fn test_pty_nonroot_rlimits() {
2416        check_nonroot_rlimits(true);
2417    }
2418
2419    fn check_nonroot_rlimits(tty: bool) {
2420        const PROBE_ENV: &str = "MSB_AGENTD_RLIMIT_PROBE";
2421        if let Ok(expected) = std::env::var(PROBE_ENV) {
2422            // This branch runs after the real session fork/exec and user switch.
2423            assert_eq!(unsafe { libc::geteuid() }, 65534);
2424            assert_eq!(unsafe { libc::getegid() }, 65534);
2425            let expected: u64 = expected.parse().unwrap();
2426            let mut limit = libc::rlimit {
2427                rlim_cur: 0,
2428                rlim_max: 0,
2429            };
2430            assert_eq!(
2431                unsafe { libc::getrlimit(libc::RLIMIT_MEMLOCK, &mut limit) },
2432                0
2433            );
2434            assert_eq!((limit.rlim_cur, limit.rlim_max), (expected, expected));
2435            assert_eq!(
2436                unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limit) },
2437                0
2438            );
2439            assert_eq!((limit.rlim_cur, limit.rlim_max), (0, 0));
2440            return;
2441        }
2442
2443        assert_eq!(unsafe { libc::geteuid() }, 0, "requires guest root");
2444        let mut baseline = libc::rlimit {
2445            rlim_cur: 0,
2446            rlim_max: 0,
2447        };
2448        assert_eq!(
2449            unsafe { libc::getrlimit(libc::RLIMIT_MEMLOCK, &mut baseline) },
2450            0
2451        );
2452        assert_ne!(
2453            baseline.rlim_max,
2454            libc::RLIM_INFINITY,
2455            "requires a finite baseline"
2456        );
2457        let raised = baseline.rlim_max.checked_add(1024 * 1024).unwrap();
2458        let user = lookup_passwd_by_uid(65534).expect("look up nobody");
2459        assert!(
2460            matches!(user, ResolvedUserLookup::Known(_)),
2461            "requires UID 65534 in /etc/passwd to exercise initgroups"
2462        );
2463        let test_name = if tty {
2464            "session::tests::test_pty_nonroot_rlimits"
2465        } else {
2466            "session::tests::test_piped_nonroot_rlimits"
2467        };
2468        let runtime = tokio::runtime::Builder::new_current_thread()
2469            .enable_all()
2470            .build()
2471            .unwrap();
2472        runtime.block_on(async {
2473            for profile in [SecurityProfile::Default, SecurityProfile::Restricted] {
2474                let (tx, mut rx) = SessionOutputSender::channel();
2475                let req = ExecRequest {
2476                    cmd: std::env::current_exe().unwrap().to_str().unwrap().into(),
2477                    args: vec![
2478                        "--exact".into(),
2479                        test_name.into(),
2480                        "--ignored".into(),
2481                        "--nocapture".into(),
2482                    ],
2483                    env: vec![format!("{PROBE_ENV}={raised}")],
2484                    cwd: None,
2485                    user: Some("65534:65534".into()),
2486                    tty,
2487                    rows: 24,
2488                    cols: 80,
2489                    rlimits: vec![
2490                        microsandbox_protocol::exec::ExecRlimit {
2491                            resource: "memlock".into(),
2492                            soft: raised,
2493                            hard: raised,
2494                        },
2495                        microsandbox_protocol::exec::ExecRlimit {
2496                            resource: "nofile".into(),
2497                            soft: 0,
2498                            hard: 0,
2499                        },
2500                    ],
2501                };
2502                let session = ExecSession::spawn(7, &req, tx, None, profile, None)
2503                    .expect("spawn non-root command with raised memlock and nofile=0");
2504                let result = time::timeout(Duration::from_secs(10), async {
2505                    let mut output = Vec::new();
2506                    while let Some(envelope) = rx.recv().await {
2507                        match envelope.output {
2508                            SessionOutput::Exited(code) => return (code, output),
2509                            SessionOutput::Stdout(data) | SessionOutput::Stderr(data) => {
2510                                output.extend(data)
2511                            }
2512                            _ => {}
2513                        }
2514                    }
2515                    panic!("session closed without an exit status");
2516                })
2517                .await;
2518                if result.is_err() {
2519                    let _ = session.send_signal(libc::SIGKILL);
2520                }
2521                let (code, output) = result.expect("non-root command timed out");
2522                assert_eq!(code, 0, "{}", String::from_utf8_lossy(&output));
2523            }
2524        });
2525
2526        let mut after = libc::rlimit {
2527            rlim_cur: 0,
2528            rlim_max: 0,
2529        };
2530        assert_eq!(
2531            unsafe { libc::getrlimit(libc::RLIMIT_MEMLOCK, &mut after) },
2532            0
2533        );
2534        assert_eq!(
2535            (after.rlim_cur, after.rlim_max),
2536            (baseline.rlim_cur, baseline.rlim_max)
2537        );
2538    }
2539
2540    #[tokio::test]
2541    async fn test_pty_reader_drains_ready_fd() {
2542        let (tx, mut rx) = SessionOutputSender::channel();
2543        let req = ExecRequest {
2544            cmd: "/bin/sh".to_string(),
2545            args: vec![
2546                "-c".to_string(),
2547                "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"
2548                    .to_string(),
2549            ],
2550            env: vec!["PATH=/usr/local/bin:/usr/bin:/bin".to_string()],
2551            cwd: None,
2552            user: None,
2553            tty: true,
2554            rows: 24,
2555            cols: 80,
2556            rlimits: Vec::new(),
2557        };
2558
2559        let session = ExecSession::spawn(7, &req, tx, None, SecurityProfile::Default, None)
2560            .expect("spawn pty session");
2561        let mut stdout = Vec::new();
2562        let mut exit = None;
2563
2564        let recv_result = time::timeout(Duration::from_secs(15), async {
2565            while let Some(envelope) = rx.recv().await {
2566                assert_eq!(envelope.id, 7);
2567                match envelope.output {
2568                    SessionOutput::Stdout(data) => stdout.extend_from_slice(&data),
2569                    SessionOutput::Exited(code) => {
2570                        exit = Some(code);
2571                        break;
2572                    }
2573                    SessionOutput::Stderr(_) | SessionOutput::Raw(_) | SessionOutput::Bulk(_) => {}
2574                }
2575            }
2576        })
2577        .await;
2578
2579        if recv_result.is_err() {
2580            let _ = session.send_signal(libc::SIGKILL);
2581            panic!("timed out waiting for PTY output");
2582        }
2583
2584        assert_eq!(exit, Some(0));
2585
2586        let second = stdout
2587            .windows(b"SECOND".len())
2588            .position(|window| window == b"SECOND");
2589        let end = stdout
2590            .windows(b"<END>".len())
2591            .position(|window| window == b"<END>");
2592
2593        assert!(
2594            matches!((second, end), (Some(second), Some(end)) if second < end),
2595            "expected immediate PTY write to arrive before later output; got {:?}",
2596            String::from_utf8_lossy(&stdout),
2597        );
2598    }
2599
2600    #[test]
2601    fn test_resolve_user_spec_for_current_uid_gid() {
2602        let uid = unsafe { libc::getuid() };
2603        let gid = unsafe { libc::getgid() };
2604        let resolved = resolve_user_spec(&format!("{uid}:{gid}")).expect("resolve numeric user");
2605        assert_eq!(resolved.uid, uid);
2606        assert_eq!(resolved.gid, gid);
2607    }
2608
2609    #[test]
2610    fn test_request_user_overrides_config_default() {
2611        let req = ExecRequest {
2612            cmd: "/bin/true".to_string(),
2613            args: Vec::new(),
2614            env: Vec::new(),
2615            cwd: None,
2616            user: Some("1:1".to_string()),
2617            tty: false,
2618            rows: 24,
2619            cols: 80,
2620            rlimits: Vec::new(),
2621        };
2622
2623        let resolved = resolve_requested_user(&req, Some("0:0")).expect("resolve requested user");
2624        assert_eq!(resolved.unwrap().uid, 1);
2625    }
2626
2627    #[test]
2628    fn test_config_default_user_used_when_request_has_none() {
2629        let req = ExecRequest {
2630            cmd: "/bin/true".to_string(),
2631            args: Vec::new(),
2632            env: Vec::new(),
2633            cwd: None,
2634            user: None,
2635            tty: false,
2636            rows: 24,
2637            cols: 80,
2638            rlimits: Vec::new(),
2639        };
2640
2641        let uid = unsafe { libc::getuid() };
2642        let gid = unsafe { libc::getgid() };
2643        let resolved = resolve_requested_user(&req, Some(&format!("{uid}:{gid}")))
2644            .expect("resolve with config default");
2645        let resolved = resolved.expect("should resolve to a user");
2646        assert_eq!(resolved.uid, uid);
2647        assert_eq!(resolved.gid, gid);
2648    }
2649
2650    #[test]
2651    fn test_request_without_user_does_not_apply_user_switch() {
2652        let req = ExecRequest {
2653            cmd: "/bin/true".to_string(),
2654            args: Vec::new(),
2655            env: Vec::new(),
2656            cwd: None,
2657            user: None,
2658            tty: false,
2659            rows: 24,
2660            cols: 80,
2661            rlimits: Vec::new(),
2662        };
2663
2664        let resolved = resolve_requested_user(&req, None).expect("resolve absent user");
2665        assert!(resolved.is_none());
2666    }
2667
2668    #[test]
2669    fn test_default_user_absent_resolves_to_root() {
2670        let resolved = resolve_default_user(None).expect("resolve absent default user");
2671        assert_eq!(resolved, (0, 0));
2672    }
2673
2674    #[test]
2675    fn test_default_home_dir_uses_resolved_user_home() {
2676        let req = ExecRequest {
2677            cmd: "/bin/true".to_string(),
2678            args: Vec::new(),
2679            env: Vec::new(),
2680            cwd: None,
2681            user: None,
2682            tty: false,
2683            rows: 24,
2684            cols: 80,
2685            rlimits: Vec::new(),
2686        };
2687        let user = ResolvedUser {
2688            uid: 1000,
2689            gid: 1000,
2690            initgroups_user: None,
2691            home_dir: Some(CString::new("/home/tester").unwrap()),
2692        };
2693
2694        assert_eq!(
2695            default_home_dir(&req, Some(&user))
2696                .expect("resolve default home")
2697                .as_deref()
2698                .map(CStr::to_string_lossy),
2699            Some("/home/tester".into()),
2700        );
2701    }
2702
2703    #[test]
2704    fn test_default_home_dir_uses_root_when_user_absent() {
2705        let req = ExecRequest {
2706            cmd: "/bin/true".to_string(),
2707            args: Vec::new(),
2708            env: Vec::new(),
2709            cwd: None,
2710            user: None,
2711            tty: false,
2712            rows: 24,
2713            cols: 80,
2714            rlimits: Vec::new(),
2715        };
2716        let root = resolve_user_spec(DEFAULT_USER_SPEC).expect("resolve implicit root");
2717
2718        assert_eq!(
2719            default_home_dir(&req, None)
2720                .expect("resolve default home")
2721                .as_deref()
2722                .map(CStr::to_string_lossy),
2723            root.home_dir.as_deref().map(CStr::to_string_lossy),
2724        );
2725    }
2726
2727    #[test]
2728    fn test_default_home_dir_respects_explicit_home_env() {
2729        let req = ExecRequest {
2730            cmd: "/bin/true".to_string(),
2731            args: Vec::new(),
2732            env: vec!["HOME=/tmp/custom".to_string()],
2733            cwd: None,
2734            user: None,
2735            tty: false,
2736            rows: 24,
2737            cols: 80,
2738            rlimits: Vec::new(),
2739        };
2740        let user = ResolvedUser {
2741            uid: 1000,
2742            gid: 1000,
2743            initgroups_user: None,
2744            home_dir: Some(CString::new("/home/tester").unwrap()),
2745        };
2746
2747        assert!(
2748            default_home_dir(&req, Some(&user))
2749                .expect("resolve default home")
2750                .is_none()
2751        );
2752    }
2753
2754    #[tokio::test]
2755    async fn test_spawn_pipe_error_does_not_include_probe_details() {
2756        let (tx, _rx) = SessionOutputSender::channel();
2757        let req = ExecRequest {
2758            cmd: "/definitely/not/a/real/binary".to_string(),
2759            args: Vec::new(),
2760            env: Vec::new(),
2761            cwd: None,
2762            user: None,
2763            tty: false,
2764            rows: 24,
2765            cols: 80,
2766            rlimits: Vec::new(),
2767        };
2768
2769        // Use the process-wide manager because other tests may have already
2770        // started its reaper thread. A private manager cannot guard this spawn
2771        // from the global `waitpid(-1, ...)` owner.
2772        let process_manager = ProcessManager::get().expect("get process manager");
2773        let err = ExecSession::spawn_pipe(
2774            9,
2775            &req,
2776            tx,
2777            None,
2778            SecurityProfile::Default,
2779            &process_manager,
2780            None,
2781        )
2782        .expect_err("spawn should fail");
2783
2784        // Spawn failures now produce the typed `ExecSpawnFailed` so
2785        // the host can render a useful message + hint. The classifier
2786        // maps ENOENT on the binary path to `NotFound`.
2787        let payload = match &err {
2788            AgentdError::ExecSpawnFailed(p) => p,
2789            other => panic!("expected ExecSpawnFailed, got: {other:?}"),
2790        };
2791        assert_eq!(payload.kind, ExecFailureKind::NotFound);
2792        assert_eq!(payload.errno, Some(libc::ENOENT));
2793        assert_eq!(payload.errno_name.as_deref(), Some("ENOENT"));
2794
2795        // The original intent of the test: probe internals leak into
2796        // the error message. The format is now
2797        // `spawn "<cmd>": <io::Error>` from
2798        // `exec_failed_from_io_error`. Verify that none of the old
2799        // probe-detail keys snuck back into the message.
2800        let message = &payload.message;
2801        assert!(message.contains("spawn"));
2802        assert!(!message.contains("symlink_metadata="));
2803        assert!(!message.contains("metadata="));
2804        assert!(!message.contains("magic="));
2805        assert!(!message.contains("path_probe="));
2806        assert!(!message.contains("cwd_probe="));
2807        assert!(!message.contains("target_probe="));
2808    }
2809}