1use 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
30const 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
42const SESSION_OUTPUT_BYTE_CAPACITY: usize = 32 * 1024 * 1024;
44
45const SESSION_OUTPUT_BUDGET_GRANULE: usize = 4096;
47
48const SESSION_OUTPUT_ITEM_CAPACITY: usize = 1024;
50
51const SESSION_BULK_OUTPUT_ITEM_CAPACITY: usize = 256;
53
54const SESSION_BULK_COMMAND_CAPACITY: usize = 128;
56
57fn 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
90fn 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
118fn 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#[derive(Debug)]
144pub struct ExecSession {
145 process_identity: ProcessIdentity,
147
148 process_manager: Arc<ProcessManager>,
150
151 pty_master: Option<AsyncFd<OwnedFd>>,
153
154 stdin: Option<AsyncFd<OwnedFd>>,
156
157 pending_stdin: VecDeque<PendingStdin>,
159}
160
161#[derive(Debug)]
162struct PendingStdin {
163 data: Vec<u8>,
164 written: usize,
165 _charge: Option<InputCharge>,
166}
167
168pub enum SessionOutput {
170 Stdout(Vec<u8>),
172
173 Stderr(Vec<u8>),
175
176 Exited(i32),
178
179 Raw(RawSessionOutput),
181
182 Bulk(BulkSessionOutput),
184}
185
186pub struct SessionOutputEnvelope {
188 pub generation: u64,
190 pub id: u32,
192
193 pub incarnation: Option<ClientIncarnation>,
195
196 pub output: SessionOutput,
198
199 _permit: Option<tokio::sync::OwnedSemaphorePermit>,
201}
202
203pub enum BulkOutputCommand {
205 Park {
207 completion: oneshot::Sender<u64>,
209 },
210 Resume {
212 completion: oneshot::Sender<()>,
214 },
215 Restore {
217 generation: u64,
219 completion: oneshot::Sender<()>,
221 },
222 DropFlow {
224 incarnation: ClientIncarnation,
226 id: u32,
228 completion: oneshot::Sender<()>,
230 },
231
232 DropIncarnation {
234 incarnation: ClientIncarnation,
236 completion: oneshot::Sender<()>,
238 },
239}
240
241pub struct SessionOutputPermit(tokio::sync::OwnedSemaphorePermit);
243
244#[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
256pub struct RawSessionOutput {
258 pub frame: Vec<u8>,
260
261 pub activity: RawActivity,
263
264 pub completion: Option<RawSessionCompletion>,
266}
267
268pub struct BulkSessionOutput {
270 pub record: BulkRecord,
272
273 pub activity: RawActivity,
275}
276
277#[derive(Debug, Clone, Copy, Default)]
279pub struct RawActivity {
280 pub guest_messages: usize,
282
283 pub fs_bytes: usize,
285
286 pub tcp_bytes: usize,
288}
289
290#[derive(Debug, Clone, Copy)]
292pub enum RawSessionCompletion {
293 FsRead,
295
296 FsWrite,
298
299 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
326struct 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
349impl RawSessionOutput {
354 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 pub fn new(record: BulkRecord, activity: RawActivity) -> Self {
371 Self { record, activity }
372 }
373}
374
375impl SessionOutput {
376 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 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 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 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 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 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 pub fn disable_bulk_scheduler(&mut self) {
508 self.bulk_tx = None;
511 self.bulk_command_tx = None;
512 }
513
514 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 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 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 pub async fn reserve(&self, max_bytes: usize) -> Option<SessionOutputPermit> {
569 self.reserve_from(&self.control_budget, max_bytes).await
570 }
571
572 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 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 pub fn guest_message() -> Self {
632 Self {
633 guest_messages: 1,
634 ..Self::default()
635 }
636 }
637
638 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 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 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 pub fn pid(&self) -> u32 {
696 self.process_identity.pid() as u32
697 }
698
699 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 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 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 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 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 self.pending_stdin.push_back(PendingStdin {
788 data: Vec::new(),
789 written: 0,
790 _charge: None,
791 });
792 }
793 }
794
795 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 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 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 pub fn close_stdin(&mut self) {
875 self.stdin.take();
876 }
877}
878
879impl ExecSession {
880 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 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 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 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 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 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 let parsed_rlimits = rlimit::to_libc(&req.rlimits);
957
958 let spawn_guard = process_manager.spawn_guard()?;
961
962 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 drop(pty.master);
975 drop(err_pipe.read_end);
976
977 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 if unsafe { libc::setsid() } < 0 {
987 unsafe { libc::_exit(1) };
988 }
989
990 if unsafe { libc::ioctl(slave_fd, libc::TIOCSCTTY, 0) } < 0 {
992 unsafe { libc::_exit(1) };
993 }
994
995 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 for (key, val) in &c_env {
1013 unsafe {
1014 libc::setenv(key.as_ptr(), val.as_ptr(), 1);
1015 }
1016 }
1017
1018 if let Some(ref dir) = c_cwd {
1020 unsafe {
1021 libc::chdir(dir.as_ptr());
1022 }
1023 }
1024
1025 if apply_exec_security_profile(security_profile).is_err() {
1026 unsafe { libc::_exit(1) };
1027 }
1028
1029 if let Some(ref user) = resolved_user
1031 && apply_resolved_groups(user).is_err()
1032 {
1033 unsafe { libc::_exit(1) };
1034 }
1035
1036 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 unsafe {
1057 libc::execvp(argv_ptrs[0], argv_ptrs.as_ptr());
1058 }
1059
1060 write_exec_error_and_exit(err_pipe.write_end.as_raw_fd());
1062 }
1063
1064 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 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 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 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 let parsed_rlimits = rlimit::to_libc(&req.rlimits);
1146 unsafe {
1147 cmd.pre_exec(move || {
1148 if let Some(ref placement) = workload_placement {
1151 placement.place_current()?;
1152 }
1153 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 if let Some(ref user) = resolved_user {
1164 apply_resolved_groups(user).map_err(agentd_to_io_error)?;
1165 }
1166
1167 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 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
1210impl Drop for ExecSession {
1215 fn drop(&mut self) {
1216 self.process_manager.release(self.process_identity);
1219 }
1220}
1221
1222fn 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 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 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 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 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
1404pub(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
1683fn 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 return Ok(());
1729 }
1730}
1731
1732async 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 let raw = master_fd.as_raw_fd();
1746 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
1794async 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#[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 #[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 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 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 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 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}