1use std::collections::HashMap;
43use std::io::Read;
44use std::path::PathBuf;
45use std::sync::{Arc, Mutex};
46use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
47
48use portable_pty::{native_pty_system, CommandBuilder, PtySize};
49use thiserror::Error;
50use tokio::sync::{mpsc, oneshot};
51use tokio::task;
52
53use crate::beholders::{registry_with_user_beholders, BeholderSelect};
54use crate::store::{RunFilter, StoreError, TaskStore};
55use crate::types::{BeholderStatus, Initiator, OutputChunk, RunStatus, Stream, TaskRunId, TaskRunMeta};
56
57const DEFAULT_GRACE: Duration = Duration::from_secs(5);
58const READ_BUF_SIZE: usize = 4096;
59const SIGTERM: i32 = 15;
60const SIGKILL: i32 = 9;
61
62#[derive(Debug, Error)]
65pub enum DriverError {
66 #[error("store: {0}")]
67 Store(#[from] StoreError),
68 #[error("pty: {0}")]
69 Pty(String),
70 #[error("run not found: {0}")]
71 NotFound(String),
72 #[error("io: {0}")]
73 Io(#[from] std::io::Error),
74}
75
76#[derive(Debug, Clone)]
80pub struct SpawnOpts {
81 pub cwd: PathBuf,
82 pub env: Vec<(String, String)>,
84 pub label: Option<String>,
85 pub initiator: Initiator,
86 pub pty_cols: u16,
88 pub pty_rows: u16,
90 pub stdin_enabled: bool,
92 pub pin: bool,
94 pub beholder_select: BeholderSelect,
96 pub verbatim_output: bool,
105 pub log_fd_enabled: bool,
109 pub origin: Option<String>,
112 pub argv: Option<Vec<String>>,
130 pub pipe: bool,
153 pub pipefail: bool,
177}
178
179const PIPEFAIL_PRELUDE: &str = "if (set -o pipefail) 2>/dev/null; then set -o pipefail; fi\n";
190
191impl Default for SpawnOpts {
192 fn default() -> Self {
193 Self {
194 cwd: std::env::current_dir().unwrap_or_else(|_| PathBuf::from("/")),
195 env: vec![],
196 label: None,
197 initiator: Initiator::Human { camp: "local".to_string() },
198 pty_cols: 80,
199 pty_rows: 24,
200 stdin_enabled: false,
201 pin: false,
202 beholder_select: BeholderSelect::Auto,
203 verbatim_output: false,
204 log_fd_enabled: true,
205 origin: None,
206 argv: None,
207 pipe: false,
208 pipefail: false,
209 }
210 }
211}
212
213#[derive(Default)]
218pub struct DriverChannels {
219 pub completion: Option<mpsc::UnboundedSender<(TaskRunId, RunStatus)>>,
222 pub output: Option<mpsc::UnboundedSender<OutputChunk>>,
229}
230
231#[derive(Debug, Clone, Default, PartialEq, Eq)]
248pub enum StaleRunPolicy {
249 #[default]
254 LostOnDisappear,
255 AdoptLiveHosts { origins: Vec<String> },
266}
267
268impl StaleRunPolicy {
269 fn tombstones(&self, meta: &TaskRunMeta) -> bool {
271 match self {
272 StaleRunPolicy::LostOnDisappear => true,
273 StaleRunPolicy::AdoptLiveHosts { origins } => {
274 let exempt_origin = origins.is_empty()
275 || meta
276 .origin
277 .as_deref()
278 .is_some_and(|o| origins.iter().any(|want| want == o));
279 if !exempt_origin {
280 return true;
281 }
282 match meta.host_pid {
283 Some(pid) => !host_process_alive(pid),
284 None => true,
285 }
286 }
287 }
288 }
289}
290
291#[cfg(unix)]
303fn host_process_alive(pid: u32) -> bool {
304 if pid == 0 {
305 return false;
306 }
307 if pid == std::process::id() {
308 return true;
309 }
310 let rc = unsafe { libc::kill(pid as libc::pid_t, 0) };
313 rc == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
314}
315
316#[cfg(not(unix))]
319fn host_process_alive(_pid: u32) -> bool {
320 false
321}
322
323struct RunControl {
326 kill_tx: mpsc::Sender<KillRequest>,
327 stdin_tx: Option<mpsc::Sender<Vec<u8>>>,
328 master: Option<Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>>>,
335 origin: Option<String>,
339 last_attached_at: Instant,
349}
350
351struct ReaderDone(Option<oneshot::Sender<()>>);
359
360impl Drop for ReaderDone {
361 fn drop(&mut self) {
362 if let Some(tx) = self.0.take() {
363 let _ = tx.send(());
364 }
365 }
366}
367
368#[derive(Debug)]
369struct KillRequest {
370 signal: i32,
371}
372
373#[cfg(unix)]
380#[derive(serde::Deserialize)]
381struct ShimRecord {
382 level: String,
383 target: String,
384 msg: String,
385 #[serde(default)]
386 fields: serde_json::Value,
387 #[serde(rename = "_lib", default)]
390 lib: Option<String>,
391 #[serde(rename = "_lib_ver", default)]
393 lib_version: Option<String>,
394}
395
396#[cfg(unix)]
403struct FdCloser(libc::c_int);
404
405#[cfg(unix)]
406impl Drop for FdCloser {
407 fn drop(&mut self) {
408 unsafe { libc::close(self.0) };
409 }
410}
411
412#[cfg(unix)]
415unsafe impl Send for FdCloser {}
416
417pub struct TaskDriver {
423 store: Arc<TaskStore>,
424 active: Arc<Mutex<HashMap<String, RunControl>>>,
425 channels: DriverChannels,
427}
428
429impl TaskDriver {
430 pub async fn new(store: Arc<TaskStore>) -> Result<Self, DriverError> {
435 Self::with_channels(store, DriverChannels::default()).await
436 }
437
438 pub async fn with_channels(
441 store: Arc<TaskStore>,
442 channels: DriverChannels,
443 ) -> Result<Self, DriverError> {
444 Self::with_config(store, channels, StaleRunPolicy::default()).await
445 }
446
447 pub async fn with_config(
455 store: Arc<TaskStore>,
456 channels: DriverChannels,
457 stale_policy: StaleRunPolicy,
458 ) -> Result<Self, DriverError> {
459 let stale = store
460 .list_runs(&RunFilter {
461 status: Some("running".to_string()),
462 ..Default::default()
463 })
464 .await?;
465 for meta in stale {
466 if !stale_policy.tombstones(&meta) {
467 continue;
468 }
469 let status = RunStatus::Lost {
470 reason: "daemon restarted while run was in-flight".to_string(),
471 };
472 store.update_status(&meta.id, &status).await?;
473 if let Some(ref tx) = channels.completion {
480 let _ = tx.send((meta.id.clone(), status));
481 }
482 }
483 Ok(Self {
484 store,
485 active: Arc::new(Mutex::new(HashMap::new())),
486 channels,
487 })
488 }
489
490 pub async fn spawn_run(&self, cmd: &str, opts: SpawnOpts) -> Result<TaskRunId, DriverError> {
504 let id = TaskRunId::new();
505 let started_at = unix_now_secs();
506 let started_at_ms: u64 = started_at.saturating_mul(1000);
507
508 let user_dir = std::env::var_os("YAH_BEHOLDERS_DIR")
511 .map(std::path::PathBuf::from)
512 .or_else(|| {
513 std::env::var_os("HOME")
514 .map(|h| std::path::PathBuf::from(h).join(".yah/beholders"))
515 });
516 let registry = registry_with_user_beholders(user_dir.as_deref());
517 let select = if opts.argv.is_some() {
523 &BeholderSelect::None
524 } else {
525 &opts.beholder_select
526 };
527 let attach = registry.attach(cmd, select, opts.verbatim_output);
528 let effective_cmd = match &attach.status.rewrite_added {
537 Some(added) if !added.is_empty() && !attach.argv.is_empty() => attach.argv.join(" "),
538 _ => cmd.to_string(),
539 };
540
541 self.store.insert_run(&TaskRunMeta {
542 id: id.clone(),
543 command: cmd.to_string(),
544 cwd: opts.cwd.clone(),
545 env: opts.env.clone(),
546 started_at,
547 status: RunStatus::Running,
548 label: opts.label.clone(),
549 initiator: opts.initiator.clone(),
550 beholder_status: Some(attach.status),
551 pinned: opts.pin,
552 origin: opts.origin.clone(),
553 host_pid: Some(std::process::id()),
559 }).await?;
560
561 let (program, args): (String, Vec<String>) = match opts.argv.as_deref() {
573 Some([p, rest @ ..]) => (p.clone(), rest.to_vec()),
574 _ => {
575 let line = if opts.pipefail {
576 format!("{PIPEFAIL_PRELUDE}{effective_cmd}")
577 } else {
578 effective_cmd.clone()
579 };
580 ("sh".to_string(), vec!["-c".to_string(), line])
581 }
582 };
583
584 #[cfg(unix)]
597 let log_fifo: Option<(libc::c_int, FdCloser, std::path::PathBuf)> = if opts.log_fd_enabled {
598 let fifo_path = std::env::temp_dir().join(format!("yah-log-{}.fifo", id));
599 let path_cstr = match std::ffi::CString::new(fifo_path.to_string_lossy().as_bytes()) {
600 Ok(s) => s,
601 Err(_) => {
602 return Err(DriverError::Io(std::io::Error::new(
604 std::io::ErrorKind::InvalidInput,
605 "log FIFO path contained nul byte",
606 )));
607 }
608 };
609 let mkfifo_ret = unsafe { libc::mkfifo(path_cstr.as_ptr(), 0o600) };
610 if mkfifo_ret != 0 {
611 None } else {
613 let rfd = unsafe {
615 libc::open(path_cstr.as_ptr(), libc::O_RDONLY | libc::O_NONBLOCK)
616 };
617 if rfd < 0 {
618 let _ = unsafe { libc::unlink(path_cstr.as_ptr()) };
619 None
620 } else {
621 unsafe { libc::fcntl(rfd, libc::F_SETFL, 0) };
623 let wfd = unsafe {
625 libc::open(path_cstr.as_ptr(), libc::O_WRONLY)
626 };
627 if wfd < 0 {
628 unsafe { libc::close(rfd) };
629 let _ = unsafe { libc::unlink(path_cstr.as_ptr()) };
630 None
631 } else {
632 Some((rfd, FdCloser(wfd), fifo_path))
633 }
634 }
635 }
636 } else {
637 None
638 };
639
640 #[cfg(unix)]
642 let fifo_env: Option<(String, String)> = log_fifo
643 .as_ref()
644 .map(|(_, _, path)| (id.to_string(), path.to_string_lossy().into_owned()));
645 #[cfg(not(unix))]
646 let fifo_env: Option<(String, String)> = None;
647
648 let pid: u32;
654 let reap: Box<dyn FnOnce() -> Option<u32> + Send>;
655 let stdin_tx: Option<mpsc::Sender<Vec<u8>>>;
656 let master: Option<Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>>>;
657 let mut sources: Vec<(Box<dyn Read + Send>, Stream)> = Vec::new();
660
661 if opts.pipe {
662 use std::process::{Command, Stdio};
663
664 let mut cmd = Command::new(&program);
665 cmd.args(&args);
666 cmd.current_dir(&opts.cwd);
667 for (k, v) in &opts.env {
668 cmd.env(k, v);
669 }
670 if let Some((run_id_env, fifo_path)) = &fifo_env {
677 cmd.env("YAH_TASK_RUN", run_id_env);
678 cmd.env("YAH_LOG_PIPE", fifo_path);
679 }
680 cmd.stdout(Stdio::piped());
681 cmd.stderr(Stdio::piped());
682 cmd.stdin(if opts.stdin_enabled { Stdio::piped() } else { Stdio::null() });
683
684 let mut child = cmd.spawn().map_err(DriverError::Io)?;
685 pid = child.id();
686
687 if let Some(out) = child.stdout.take() {
688 sources.push((Box::new(out), Stream::Stdout));
689 }
690 if let Some(err) = child.stderr.take() {
691 sources.push((Box::new(err), Stream::Stderr));
692 }
693
694 stdin_tx = child.stdin.take().map(|mut writer| {
695 let (tx, mut rx) = mpsc::channel::<Vec<u8>>(64);
696 task::spawn(async move {
697 use std::io::Write;
698 while let Some(bytes) = rx.recv().await {
699 let _ = writer.write_all(&bytes);
700 let _ = writer.flush();
701 }
702 });
703 tx
704 });
705
706 master = None;
707 reap = Box::new(move || child.wait().ok().and_then(|s| s.code()).map(|c| c as u32));
708 } else {
709 let pty_sys = native_pty_system();
711 let pair = pty_sys
712 .openpty(PtySize {
713 rows: opts.pty_rows,
714 cols: opts.pty_cols,
715 pixel_width: 0,
716 pixel_height: 0,
717 })
718 .map_err(|e| DriverError::Pty(e.to_string()))?;
719
720 let pty_reader = pair
722 .master
723 .try_clone_reader()
724 .map_err(|e| DriverError::Pty(e.to_string()))?;
725 sources.push((Box::new(pty_reader), Stream::Stdout));
726
727 stdin_tx = if opts.stdin_enabled {
729 let mut writer = pair
730 .master
731 .take_writer()
732 .map_err(|e| DriverError::Pty(e.to_string()))?;
733 let (tx, mut rx) = mpsc::channel::<Vec<u8>>(64);
734 task::spawn(async move {
735 use std::io::Write;
736 while let Some(bytes) = rx.recv().await {
737 let _ = writer.write_all(&bytes);
738 let _ = writer.flush();
739 }
740 });
741 Some(tx)
742 } else {
743 None
744 };
745
746 let mut cb = CommandBuilder::new(&program);
747 cb.args(&args);
748 cb.cwd(&opts.cwd);
749 for (k, v) in &opts.env {
750 cb.env(k, v);
751 }
752 cb.env("TERM", "xterm-256color");
753 if let Some((run_id_env, fifo_path)) = &fifo_env {
754 cb.env("YAH_TASK_RUN", run_id_env);
755 cb.env("YAH_LOG_PIPE", fifo_path);
756 }
757
758 let child = pair
759 .slave
760 .spawn_command(cb)
761 .map_err(|e| DriverError::Pty(e.to_string()))?;
762 drop(pair.slave);
764
765 pid = child.process_id().unwrap_or(0);
766
767 let m: Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>> =
770 Arc::new(Mutex::new(pair.master));
771 master = Some(Arc::clone(&m));
772 reap = Box::new(move || {
773 let mut c = child;
774 let _m = m; c.wait().ok().map(|s| s.exit_code())
776 });
777 }
778
779 #[cfg(unix)]
787 let log_wfd_holder: Option<FdCloser> = if let Some((rfd, wfd, fifo_path)) = log_fifo {
788 let store_log = Arc::clone(&self.store);
789 let id_log = id.clone();
790 let rt = tokio::runtime::Handle::current();
791 tokio::task::spawn_blocking(move || {
794 run_log_receiver(rt, store_log, id_log, rfd, fifo_path, started_at_ms);
795 });
796 Some(wfd)
797 } else {
798 None
799 };
800
801 let (kill_tx, kill_rx) = mpsc::channel::<KillRequest>(4);
803 let (reader_done_tx, reader_done_rx) = oneshot::channel::<()>();
804
805 {
811 let done = Arc::new(ReaderDone(Some(reader_done_tx)));
812 let mut beholder = attach.beholder;
819 for (reader, stream) in sources {
820 spawn_output_pump(
821 reader,
822 stream,
823 Arc::clone(&self.store),
824 id.clone(),
825 started_at_ms,
826 self.channels.output.clone(),
827 if stream == Stream::Stdout { beholder.take() } else { None },
828 Arc::clone(&done),
829 );
830 }
831 }
832
833 {
837 let store_l = Arc::clone(&self.store);
838 let active_l = Arc::clone(&self.active);
839 let id_l = id.clone();
840 let completion_tx_l = self.channels.completion.clone();
841 #[cfg(unix)]
842 let wfd_l = log_wfd_holder;
843 task::spawn(async move {
844 run_lifecycle(
845 store_l,
846 active_l,
847 id_l,
848 pid,
849 reap,
850 kill_rx,
851 reader_done_rx,
852 completion_tx_l,
853 #[cfg(unix)]
854 wfd_l,
855 )
856 .await;
857 });
858 }
859
860 self.active
861 .lock()
862 .unwrap()
863 .insert(
864 id.to_string(),
865 RunControl {
866 kill_tx,
867 stdin_tx,
868 master,
869 origin: opts.origin.clone(),
870 last_attached_at: Instant::now(),
871 },
872 );
873
874 Ok(id)
875 }
876
877 pub async fn resize_run(
886 &self,
887 id: &TaskRunId,
888 cols: u16,
889 rows: u16,
890 ) -> Result<(), DriverError> {
891 let master = self
892 .active
893 .lock()
894 .unwrap()
895 .get(&id.to_string())
896 .and_then(|c| c.master.as_ref().map(Arc::clone));
897
898 match master {
899 Some(m) => {
900 let size = PtySize { rows, cols, pixel_width: 0, pixel_height: 0 };
901 m.lock()
902 .unwrap()
903 .resize(size)
904 .map_err(|e| DriverError::Pty(e.to_string()))
905 }
906 None => Err(DriverError::NotFound(id.to_string())),
907 }
908 }
909
910 pub fn foreground_pid(&self, id: &TaskRunId) -> Option<u32> {
924 let master = self
925 .active
926 .lock()
927 .unwrap()
928 .get(&id.to_string())
929 .and_then(|c| c.master.as_ref().map(Arc::clone))?;
930 #[cfg(unix)]
931 {
932 let pid = master.lock().unwrap().process_group_leader()?;
933 u32::try_from(pid).ok()
934 }
935 #[cfg(not(unix))]
936 {
937 let _ = master;
938 None
939 }
940 }
941
942 pub async fn kill_run(&self, id: &TaskRunId, signal: Option<i32>) -> Result<(), DriverError> {
949 let kill_tx = self
950 .active
951 .lock()
952 .unwrap()
953 .get(&id.to_string())
954 .map(|c| c.kill_tx.clone());
955
956 match kill_tx {
957 Some(tx) => tx
958 .send(KillRequest { signal: signal.unwrap_or(SIGTERM) })
959 .await
960 .map_err(|_| DriverError::NotFound(id.to_string())),
961 None => Err(DriverError::NotFound(id.to_string())),
962 }
963 }
964
965 pub async fn send_stdin(&self, id: &TaskRunId, bytes: Vec<u8>) -> Result<(), DriverError> {
967 let stdin_tx = self
968 .active
969 .lock()
970 .unwrap()
971 .get(&id.to_string())
972 .and_then(|c| c.stdin_tx.clone());
973
974 match stdin_tx {
975 Some(tx) => tx
976 .send(bytes)
977 .await
978 .map_err(|_| DriverError::NotFound(id.to_string())),
979 None => Err(DriverError::NotFound(id.to_string())),
980 }
981 }
982
983 pub fn note_attached(&self, id: &TaskRunId) {
989 if let Some(control) = self.active.lock().unwrap().get_mut(&id.to_string()) {
990 control.last_attached_at = Instant::now();
991 }
992 }
993
994 pub fn attached_age(&self, id: &TaskRunId) -> Option<Duration> {
997 self.active
998 .lock()
999 .unwrap()
1000 .get(&id.to_string())
1001 .map(|c| c.last_attached_at.elapsed())
1002 }
1003
1004 pub async fn reap_unattached(&self, idle: Duration, origins: &[String]) -> Vec<TaskRunId> {
1031 if origins.is_empty() {
1032 return Vec::new();
1033 }
1034 let candidates: Vec<TaskRunId> = {
1035 let active = self.active.lock().unwrap();
1036 active
1037 .iter()
1038 .filter(|(_, c)| {
1039 c.origin
1040 .as_deref()
1041 .is_some_and(|o| origins.iter().any(|want| want == o))
1042 && c.last_attached_at.elapsed() >= idle
1043 })
1044 .filter_map(|(id, _)| id.parse::<TaskRunId>().ok())
1045 .collect()
1046 };
1047
1048 let mut reaped = Vec::new();
1049 for id in candidates {
1050 if self.kill_run(&id, None).await.is_ok() {
1055 reaped.push(id);
1056 }
1057 }
1058 reaped
1059 }
1060}
1061
1062#[cfg(unix)]
1071fn run_log_receiver(
1072 rt: tokio::runtime::Handle,
1073 store: Arc<TaskStore>,
1074 run_id: TaskRunId,
1075 read_fd: libc::c_int,
1076 fifo_path: std::path::PathBuf,
1077 started_at_ms: u64,
1078) {
1079 use std::io::BufRead;
1080 use std::os::unix::io::FromRawFd;
1081
1082 let file = unsafe { std::fs::File::from_raw_fd(read_fd) };
1085 let reader = std::io::BufReader::new(file);
1086
1087 for line in reader.lines() {
1088 let line = match line {
1089 Ok(l) => l,
1090 Err(_) => break,
1091 };
1092 let trimmed = line.trim();
1093 if trimmed.is_empty() {
1094 continue;
1095 }
1096 let rec: ShimRecord = match serde_json::from_str(trimmed) {
1097 Ok(r) => r,
1098 Err(_) => continue, };
1100 let level = rec.level.parse::<crate::types::Level>().unwrap_or(crate::types::Level::Info);
1101 let source = crate::types::EventSource::Shim {
1102 lib: rec.lib.unwrap_or_else(|| "unknown".to_string()),
1103 version: rec.lib_version.unwrap_or_else(|| "0.0.0".to_string()),
1104 };
1105 let fields = if rec.fields.is_object() {
1106 rec.fields
1107 } else {
1108 serde_json::Value::Object(Default::default())
1109 };
1110 let offset = elapsed_ms(started_at_ms);
1111 let _ = rt.block_on(store.append_event(
1112 &run_id,
1113 offset,
1114 level,
1115 &rec.target,
1116 &rec.msg,
1117 &fields,
1118 None,
1119 &source,
1120 ));
1121 }
1122
1123 let _ = std::fs::remove_file(&fifo_path);
1125}
1126
1127#[allow(clippy::too_many_arguments)]
1138fn spawn_output_pump(
1139 reader: Box<dyn Read + Send>,
1140 stream: Stream,
1141 store: Arc<TaskStore>,
1142 id: TaskRunId,
1143 started_at_ms: u64,
1144 output_tx: Option<mpsc::UnboundedSender<OutputChunk>>,
1145 beholder: Option<Box<dyn crate::beholders::Beholder>>,
1146 done: Arc<ReaderDone>,
1147) {
1148 let rt = tokio::runtime::Handle::current();
1149 tokio::task::spawn_blocking(move || {
1150 let _done = done;
1151 let mut beholder = beholder;
1152 let mut buf = [0u8; READ_BUF_SIZE];
1153 let mut reader = reader;
1154 loop {
1155 match reader.read(&mut buf) {
1156 Ok(0) | Err(_) => break,
1157 Ok(n) => {
1158 let offset = elapsed_ms(started_at_ms);
1159 let append_res =
1160 rt.block_on(store.append_chunk(&id, offset, stream, &buf[..n]));
1161 if let Ok(seq) = append_res {
1162 let chunk = (output_tx.is_some() || beholder.is_some()).then(|| {
1166 OutputChunk {
1167 run_id: id.clone(),
1168 seq,
1169 offset_ms: offset,
1170 stream,
1171 bytes: buf[..n].to_vec(),
1172 }
1173 });
1174 if let (Some(tx), Some(c)) = (&output_tx, &chunk) {
1178 let _ = tx.send(c.clone());
1179 }
1180 let mut detach_beholder = false;
1181 if let (Some(b), Some(chunk)) = (beholder.as_mut(), &chunk) {
1182 for ev in b.parse_chunk(chunk) {
1183 let _ = rt.block_on(store.append_event(
1184 &ev.run_id,
1185 ev.offset_ms,
1186 ev.level,
1187 &ev.target,
1188 &ev.msg,
1189 &ev.fields,
1190 ev.anchor.as_ref().map(|a| a.seq),
1191 &ev.source,
1192 ));
1193 }
1194 if let Some(reason) = b.unknown_format_reason() {
1195 let new_status =
1196 BeholderStatus::unknown_format_with_reason(b.name(), reason);
1197 let _ =
1198 rt.block_on(store.update_beholder_status(&id, &new_status));
1199 detach_beholder = true;
1200 }
1201 }
1202 if detach_beholder {
1203 beholder = None;
1204 }
1205 }
1206 }
1207 }
1208 }
1209 if let Some(ref mut b) = beholder {
1210 let final_offset = elapsed_ms(started_at_ms);
1211 for ev in b.on_done(&id, final_offset) {
1212 let _ = rt.block_on(store.append_event(
1213 &ev.run_id,
1214 ev.offset_ms,
1215 ev.level,
1216 &ev.target,
1217 &ev.msg,
1218 &ev.fields,
1219 ev.anchor.as_ref().map(|a| a.seq),
1220 &ev.source,
1221 ));
1222 }
1223 if let Some(reason) = b.unknown_format_reason() {
1224 let new_status = BeholderStatus::unknown_format_with_reason(b.name(), reason);
1225 let _ = rt.block_on(store.update_beholder_status(&id, &new_status));
1226 }
1227 }
1228 });
1229}
1230
1231#[allow(clippy::too_many_arguments)]
1232async fn run_lifecycle(
1233 store: Arc<TaskStore>,
1234 active: Arc<Mutex<HashMap<String, RunControl>>>,
1235 id: TaskRunId,
1236 pid: u32,
1237 reap: Box<dyn FnOnce() -> Option<u32> + Send>,
1241 mut kill_rx: mpsc::Receiver<KillRequest>,
1242 reader_done_rx: oneshot::Receiver<()>,
1243 completion_tx: Option<tokio::sync::mpsc::UnboundedSender<(TaskRunId, RunStatus)>>,
1244 #[cfg(unix)]
1248 _log_wfd: Option<FdCloser>,
1249) {
1250 let reader_done = async { reader_done_rx.await.ok(); };
1253 tokio::pin!(reader_done);
1254
1255 let sent_signal: Option<i32>;
1256
1257 tokio::select! {
1258 req = kill_rx.recv() => {
1259 match req {
1260 Some(KillRequest { signal }) => {
1261 send_unix_signal(pid, signal);
1262 if signal == SIGKILL {
1263 sent_signal = Some(SIGKILL);
1264 } else {
1265 tokio::select! {
1267 _ = &mut reader_done => {
1268 sent_signal = Some(signal);
1270 }
1271 _ = tokio::time::sleep(DEFAULT_GRACE) => {
1272 send_unix_signal(pid, SIGKILL);
1274 sent_signal = Some(SIGKILL);
1275 }
1276 }
1277 }
1278 }
1279 None => {
1281 send_unix_signal(pid, SIGKILL);
1282 sent_signal = Some(SIGKILL);
1283 }
1284 }
1285 }
1286 _ = &mut reader_done => {
1287 sent_signal = None;
1288 }
1289 }
1290
1291 let exit_code = task::spawn_blocking(reap).await.ok().flatten();
1296
1297 let ended_at = unix_now_secs();
1298 let status = match sent_signal {
1299 Some(sig) => RunStatus::Killed { signal: sig, ended_at },
1300 None => match exit_code {
1301 Some(code) => RunStatus::Done { exit_code: code as i32, ended_at },
1302 None => RunStatus::Lost {
1303 reason: "process exited without an exit code".to_string(),
1304 },
1305 },
1306 };
1307
1308 if let Err(e) = store.update_status(&id, &status).await {
1314 eprintln!("[yah task-runs] failed to record terminal status for run {id}: {e}");
1315 }
1316 if let Some(ref tx) = completion_tx {
1317 let _ = tx.send((id.clone(), status));
1318 }
1319 active.lock().unwrap().remove(&id.to_string());
1320}
1321
1322fn send_unix_signal(pid: u32, signal: i32) {
1325 #[cfg(unix)]
1326 unsafe {
1327 libc::kill(pid as libc::pid_t, signal);
1328 }
1329 }
1331
1332fn unix_now_secs() -> u64 {
1333 SystemTime::now()
1334 .duration_since(UNIX_EPOCH)
1335 .unwrap_or_default()
1336 .as_secs()
1337}
1338
1339fn elapsed_ms(started_at_ms: u64) -> u32 {
1340 let now_ms = SystemTime::now()
1341 .duration_since(UNIX_EPOCH)
1342 .unwrap_or_default()
1343 .as_millis() as u64;
1344 now_ms.saturating_sub(started_at_ms).min(u32::MAX as u64) as u32
1345}
1346
1347#[cfg(test)]
1350mod tests {
1351 use super::*;
1352 use crate::store::ChunkFilter;
1353
1354 async fn open_store(dir: &tempfile::TempDir) -> Arc<TaskStore> {
1355 Arc::new(TaskStore::open(&dir.path().join("tr.turso")).await.unwrap())
1356 }
1357
1358 #[tokio::test]
1361 async fn lost_on_disappear_marks_stale_running_runs() {
1362 let dir = tempfile::tempdir().unwrap();
1363 let store = open_store(&dir).await;
1364
1365 let stale_id = TaskRunId::new();
1367 store
1368 .insert_run(&TaskRunMeta {
1369 id: stale_id.clone(),
1370 command: "sleep 9999".to_string(),
1371 cwd: "/tmp".into(),
1372 env: vec![],
1373 started_at: unix_now_secs() - 60,
1374 status: RunStatus::Running,
1375 label: None,
1376 initiator: Initiator::Human { camp: "test".to_string() },
1377 beholder_status: None,
1378 pinned: false,
1379 origin: None,
1380 host_pid: None,
1381 })
1382 .await
1383 .unwrap();
1384
1385 let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1387
1388 let meta = store.get_run(&stale_id).await.unwrap().unwrap();
1389 assert!(
1390 matches!(meta.status, RunStatus::Lost { .. }),
1391 "stale run should be Lost, got {:?}",
1392 meta.status
1393 );
1394 }
1395
1396 #[tokio::test]
1404 async fn a_sweep_tombstone_fires_the_completion_channel() {
1405 let dir = tempfile::tempdir().unwrap();
1406 let store = open_store(&dir).await;
1407
1408 let stale_id = TaskRunId::new();
1409 store
1410 .insert_run(&TaskRunMeta {
1411 id: stale_id.clone(),
1412 command: "sleep 9999".to_string(),
1413 cwd: "/tmp".into(),
1414 env: vec![],
1415 started_at: unix_now_secs() - 60,
1416 status: RunStatus::Running,
1417 label: None,
1418 initiator: Initiator::Human { camp: "test".to_string() },
1419 beholder_status: None,
1420 pinned: false,
1421 origin: None,
1422 host_pid: None,
1423 })
1424 .await
1425 .unwrap();
1426
1427 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1428 let _driver = TaskDriver::with_channels(
1429 Arc::clone(&store),
1430 DriverChannels { completion: Some(tx), output: None },
1431 )
1432 .await
1433 .unwrap();
1434
1435 let (id, status) = rx.try_recv().expect("the sweep must announce what it tombstoned");
1436 assert_eq!(id, stale_id);
1437 assert!(
1438 matches!(status, RunStatus::Lost { .. }),
1439 "expected Lost, got {status:?}"
1440 );
1441 }
1442
1443 async fn plant_running(
1447 store: &Arc<TaskStore>,
1448 origin: Option<&str>,
1449 host_pid: Option<u32>,
1450 ) -> TaskRunId {
1451 let id = TaskRunId::new();
1452 store
1453 .insert_run(&TaskRunMeta {
1454 id: id.clone(),
1455 command: "sleep 9999".to_string(),
1456 cwd: "/tmp".into(),
1457 env: vec![],
1458 started_at: unix_now_secs() - 60,
1459 status: RunStatus::Running,
1460 label: None,
1461 initiator: Initiator::Human {
1462 camp: "test".to_string(),
1463 },
1464 beholder_status: None,
1465 pinned: false,
1466 origin: origin.map(str::to_string),
1467 host_pid,
1468 })
1469 .await
1470 .unwrap();
1471 id
1472 }
1473
1474 async fn is_lost(store: &Arc<TaskStore>, id: &TaskRunId) -> bool {
1475 matches!(
1476 store.get_run(id).await.unwrap().unwrap().status,
1477 RunStatus::Lost { .. }
1478 )
1479 }
1480
1481 fn adopt_terminal() -> StaleRunPolicy {
1482 StaleRunPolicy::AdoptLiveHosts {
1483 origins: vec!["terminal".to_string()],
1484 }
1485 }
1486
1487 #[tokio::test]
1490 async fn a_run_owned_by_a_live_host_survives_a_new_driver() {
1491 let dir = tempfile::tempdir().unwrap();
1492 let store = open_store(&dir).await;
1493 let id = plant_running(&store, Some("terminal"), Some(std::process::id())).await;
1496
1497 let _driver = TaskDriver::with_config(
1498 Arc::clone(&store),
1499 DriverChannels::default(),
1500 adopt_terminal(),
1501 )
1502 .await
1503 .unwrap();
1504
1505 assert!(
1506 !is_lost(&store, &id).await,
1507 "a terminal run whose owner is alive must stay Running — \
1508 tombstoning it is what made a surviving shell read as dead"
1509 );
1510 }
1511
1512 #[tokio::test]
1515 async fn a_run_whose_host_is_gone_is_still_tombstoned() {
1516 let dir = tempfile::tempdir().unwrap();
1517 let store = open_store(&dir).await;
1518 let dead_pid = {
1520 let child = std::process::Command::new("true").spawn().unwrap();
1521 let pid = child.id();
1522 let mut child = child;
1523 let _ = child.wait();
1524 pid
1525 };
1526 let id = plant_running(&store, Some("terminal"), Some(dead_pid)).await;
1527
1528 let _driver = TaskDriver::with_config(
1529 Arc::clone(&store),
1530 DriverChannels::default(),
1531 adopt_terminal(),
1532 )
1533 .await
1534 .unwrap();
1535
1536 assert!(
1537 is_lost(&store, &id).await,
1538 "pid {dead_pid} was reaped; its run has no owner left and must be Lost"
1539 );
1540 }
1541
1542 #[tokio::test]
1546 async fn a_non_matching_origin_is_tombstoned_even_with_a_live_host() {
1547 let dir = tempfile::tempdir().unwrap();
1548 let store = open_store(&dir).await;
1549 let job = plant_running(&store, None, Some(std::process::id())).await;
1550 let other = plant_running(&store, Some("gnome"), Some(std::process::id())).await;
1551
1552 let _driver = TaskDriver::with_config(
1553 Arc::clone(&store),
1554 DriverChannels::default(),
1555 adopt_terminal(),
1556 )
1557 .await
1558 .unwrap();
1559
1560 assert!(is_lost(&store, &job).await, "an origin-less job is not exempt");
1561 assert!(
1562 is_lost(&store, &other).await,
1563 "an origin outside the list is not exempt"
1564 );
1565 }
1566
1567 #[tokio::test]
1571 async fn an_unattributed_run_is_tombstoned() {
1572 let dir = tempfile::tempdir().unwrap();
1573 let store = open_store(&dir).await;
1574 let id = plant_running(&store, Some("terminal"), None).await;
1575
1576 let _driver = TaskDriver::with_config(
1577 Arc::clone(&store),
1578 DriverChannels::default(),
1579 adopt_terminal(),
1580 )
1581 .await
1582 .unwrap();
1583
1584 assert!(is_lost(&store, &id).await);
1585 }
1586
1587 #[tokio::test]
1590 async fn the_default_policy_is_still_lost_on_disappear() {
1591 assert_eq!(StaleRunPolicy::default(), StaleRunPolicy::LostOnDisappear);
1592
1593 let dir = tempfile::tempdir().unwrap();
1594 let store = open_store(&dir).await;
1595 let id = plant_running(&store, Some("terminal"), Some(std::process::id())).await;
1596
1597 let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1598
1599 assert!(
1600 is_lost(&store, &id).await,
1601 "the default must tombstone regardless of origin or owner liveness"
1602 );
1603 }
1604
1605 #[tokio::test]
1608 async fn spawn_run_stamps_this_process_as_the_owner() {
1609 let dir = tempfile::tempdir().unwrap();
1610 let store = open_store(&dir).await;
1611 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1612
1613 let id = driver
1614 .spawn_run(
1615 "true",
1616 SpawnOpts {
1617 cwd: "/tmp".into(),
1618 origin: Some("terminal".to_string()),
1619 ..Default::default()
1620 },
1621 )
1622 .await
1623 .unwrap();
1624
1625 let meta = store.get_run(&id).await.unwrap().unwrap();
1626 assert_eq!(meta.host_pid, Some(std::process::id()));
1627 }
1628
1629 #[tokio::test]
1630 async fn new_driver_does_not_touch_completed_runs() {
1631 let dir = tempfile::tempdir().unwrap();
1632 let store = open_store(&dir).await;
1633
1634 let done_id = TaskRunId::new();
1635 store
1636 .insert_run(&TaskRunMeta {
1637 id: done_id.clone(),
1638 command: "true".to_string(),
1639 cwd: "/tmp".into(),
1640 env: vec![],
1641 started_at: unix_now_secs() - 10,
1642 status: RunStatus::Running,
1643 label: None,
1644 initiator: Initiator::Human { camp: "test".to_string() },
1645 beholder_status: None,
1646 pinned: false,
1647 origin: None,
1648 host_pid: None,
1649 })
1650 .await
1651 .unwrap();
1652 store
1653 .update_status(&done_id, &RunStatus::Done { exit_code: 0, ended_at: unix_now_secs() })
1654 .await
1655 .unwrap();
1656
1657 let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1658
1659 let meta = store.get_run(&done_id).await.unwrap().unwrap();
1660 assert!(
1661 matches!(meta.status, RunStatus::Done { .. }),
1662 "completed run must not be touched"
1663 );
1664 }
1665
1666 #[tokio::test]
1669 async fn spawn_echo_and_read_chunks() {
1670 let dir = tempfile::tempdir().unwrap();
1671 let store = open_store(&dir).await;
1672 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1673
1674 let id = driver
1675 .spawn_run(
1676 "echo hello_world",
1677 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1678 )
1679 .await
1680 .unwrap();
1681
1682 let deadline = std::time::Instant::now() + Duration::from_secs(5);
1684 loop {
1685 let meta = store.get_run(&id).await.unwrap().unwrap();
1686 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1687 break;
1688 }
1689 if std::time::Instant::now() > deadline {
1690 panic!("run did not complete in time, status={:?}", meta.status);
1691 }
1692 tokio::time::sleep(Duration::from_millis(50)).await;
1693 }
1694
1695 let chunks = store
1697 .get_chunks(&id, &ChunkFilter::default())
1698 .await
1699 .unwrap();
1700 let output: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
1701 let text = String::from_utf8_lossy(&output);
1702 assert!(
1703 text.contains("hello_world"),
1704 "expected 'hello_world' in output, got: {text:?}"
1705 );
1706
1707 let meta = store.get_run(&id).await.unwrap().unwrap();
1708 assert!(
1709 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
1710 "expected Done(0), got {:?}",
1711 meta.status
1712 );
1713 }
1714
1715 async fn run_to_completion(
1719 store: &Arc<TaskStore>,
1720 driver: &TaskDriver,
1721 cmd: &str,
1722 opts: SpawnOpts,
1723 ) -> Vec<OutputChunk> {
1724 let id = driver.spawn_run(cmd, opts).await.unwrap();
1725 let deadline = std::time::Instant::now() + Duration::from_secs(10);
1726 loop {
1727 let meta = store.get_run(&id).await.unwrap().unwrap();
1728 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1729 break;
1730 }
1731 if std::time::Instant::now() > deadline {
1732 panic!("run did not complete in time, status={:?}", meta.status);
1733 }
1734 tokio::time::sleep(Duration::from_millis(25)).await;
1735 }
1736 store.get_chunks(&id, &ChunkFilter::default()).await.unwrap()
1737 }
1738
1739 fn joined(chunks: &[OutputChunk]) -> Vec<u8> {
1740 chunks.iter().flat_map(|c| c.bytes.clone()).collect()
1741 }
1742
1743 fn joined_stream(chunks: &[OutputChunk], stream: Stream) -> Vec<u8> {
1744 chunks
1745 .iter()
1746 .filter(|c| c.stream == stream)
1747 .flat_map(|c| c.bytes.clone())
1748 .collect()
1749 }
1750
1751 #[tokio::test]
1754 async fn pipe_mode_child_sees_no_tty_on_stdout() {
1755 let dir = tempfile::tempdir().unwrap();
1756 let store = open_store(&dir).await;
1757 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1758 let cmd = "if [ -t 1 ]; then echo TTY; else echo PIPE; fi";
1759
1760 let piped = run_to_completion(
1761 &store,
1762 &driver,
1763 cmd,
1764 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1765 )
1766 .await;
1767 assert_eq!(joined(&piped), b"PIPE\n");
1768
1769 let ptied = run_to_completion(
1771 &store,
1772 &driver,
1773 cmd,
1774 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1775 )
1776 .await;
1777 assert_eq!(joined(&ptied), b"TTY\r\n");
1778 }
1779
1780 #[tokio::test]
1783 async fn pipe_mode_does_not_translate_newlines() {
1784 let dir = tempfile::tempdir().unwrap();
1785 let store = open_store(&dir).await;
1786 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1787
1788 let piped = run_to_completion(
1789 &store,
1790 &driver,
1791 r"printf 'a\nb\n'",
1792 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1793 )
1794 .await;
1795 assert_eq!(joined(&piped), b"a\nb\n");
1796 }
1797
1798 #[tokio::test]
1801 async fn pipe_mode_keeps_stderr_separate_from_stdout() {
1802 let dir = tempfile::tempdir().unwrap();
1803 let store = open_store(&dir).await;
1804 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1805 let cmd = "printf 'to-out\n'; printf 'to-err\n' >&2";
1806
1807 let piped = run_to_completion(
1808 &store,
1809 &driver,
1810 cmd,
1811 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1812 )
1813 .await;
1814 assert_eq!(joined_stream(&piped, Stream::Stdout), b"to-out\n");
1815 assert_eq!(joined_stream(&piped, Stream::Stderr), b"to-err\n");
1816
1817 let ptied = run_to_completion(
1820 &store,
1821 &driver,
1822 cmd,
1823 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1824 )
1825 .await;
1826 assert!(
1827 joined_stream(&ptied, Stream::Stderr).is_empty(),
1828 "PTY runs have no stderr chunks; that is the behaviour pipe mode exists to fix",
1829 );
1830 }
1831
1832 #[tokio::test]
1837 async fn pipe_mode_drains_both_streams_before_the_run_is_terminal() {
1838 let dir = tempfile::tempdir().unwrap();
1839 let store = open_store(&dir).await;
1840 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1841
1842 let piped = run_to_completion(
1843 &store,
1844 &driver,
1845 "head -c 4096 /dev/zero | tr '\\0' 'x'; head -c 4096 /dev/zero | tr '\\0' 'y' >&2",
1846 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1847 )
1848 .await;
1849 assert_eq!(joined_stream(&piped, Stream::Stdout).len(), 4096);
1850 assert_eq!(joined_stream(&piped, Stream::Stderr).len(), 4096);
1851 }
1852
1853 #[tokio::test]
1856 async fn pipe_mode_records_the_childs_exit_code() {
1857 let dir = tempfile::tempdir().unwrap();
1858 let store = open_store(&dir).await;
1859 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1860
1861 let id = driver
1862 .spawn_run(
1863 "exit 101",
1864 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1865 )
1866 .await
1867 .unwrap();
1868
1869 let deadline = std::time::Instant::now() + Duration::from_secs(10);
1870 loop {
1871 let meta = store.get_run(&id).await.unwrap().unwrap();
1872 match meta.status {
1873 RunStatus::Done { exit_code, .. } => {
1874 assert_eq!(exit_code, 101);
1875 return;
1876 }
1877 RunStatus::Lost { .. } | RunStatus::Killed { .. } => {
1878 panic!("unexpected terminal status {:?}", meta.status)
1879 }
1880 _ => {}
1881 }
1882 if std::time::Instant::now() > deadline {
1883 panic!("run did not complete in time");
1884 }
1885 tokio::time::sleep(Duration::from_millis(25)).await;
1886 }
1887 }
1888
1889 #[tokio::test]
1892 async fn pipe_mode_has_no_terminal_to_resize_or_read_a_foreground_pid_from() {
1893 let dir = tempfile::tempdir().unwrap();
1894 let store = open_store(&dir).await;
1895 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1896
1897 let id = driver
1898 .spawn_run(
1899 "sleep 2",
1900 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1901 )
1902 .await
1903 .unwrap();
1904
1905 assert!(matches!(
1906 driver.resize_run(&id, 100, 40).await,
1907 Err(DriverError::NotFound(_))
1908 ));
1909 assert_eq!(driver.foreground_pid(&id), None);
1910 let _ = driver.kill_run(&id, Some(SIGKILL)).await;
1911 }
1912
1913 async fn await_done(store: &TaskStore, id: &TaskRunId) -> TaskRunMeta {
1915 let deadline = std::time::Instant::now() + Duration::from_secs(5);
1916 loop {
1917 let meta = store.get_run(id).await.unwrap().unwrap();
1918 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1919 return meta;
1920 }
1921 if std::time::Instant::now() > deadline {
1922 panic!("run did not complete in time, status={:?}", meta.status);
1923 }
1924 tokio::time::sleep(Duration::from_millis(50)).await;
1925 }
1926 }
1927
1928 async fn output_of(store: &TaskStore, id: &TaskRunId) -> String {
1929 let chunks = store.get_chunks(id, &ChunkFilter::default()).await.unwrap();
1930 let bytes: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
1931 String::from_utf8_lossy(&bytes).into_owned()
1932 }
1933
1934 #[cfg(unix)]
1943 #[tokio::test]
1944 async fn a_multi_line_command_is_not_flattened_into_one_line() {
1945 let dir = tempfile::tempdir().unwrap();
1946 let store = open_store(&dir).await;
1947 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
1948
1949 let id = driver
1953 .spawn_run(
1954 "echo one\necho two",
1955 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1956 )
1957 .await
1958 .unwrap();
1959 await_done(&store, &id).await;
1960
1961 let out = output_of(&store, &id).await;
1962 assert!(out.contains("one"), "got: {out:?}");
1963 assert!(
1964 out.contains("two"),
1965 "the second line must have run as its own command; got: {out:?}"
1966 );
1967 assert!(
1968 !out.contains("one echo two"),
1969 "the newline was flattened into a space; got: {out:?}"
1970 );
1971 }
1972
1973 #[cfg(unix)]
1977 #[tokio::test]
1978 async fn a_wrapper_the_caller_wrote_is_not_stripped_from_the_spawned_command() {
1979 let dir = tempfile::tempdir().unwrap();
1980 let store = open_store(&dir).await;
1981 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
1982
1983 let id = driver
1987 .spawn_run(
1988 "npx r739s2-nonexistent-tool --version",
1989 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1990 )
1991 .await
1992 .unwrap();
1993 let meta = await_done(&store, &id).await;
1994 let out = output_of(&store, &id).await;
1995 assert!(
1996 !matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
1997 "expected a failure, got {:?} with output {out:?}",
1998 meta.status
1999 );
2000 assert!(
2001 !out.contains("--version: "),
2002 "the wrapper was stripped and the shell tried to run the flag; got: {out:?}"
2003 );
2004 }
2005
2006 #[tokio::test]
2009 async fn explicit_argv_execs_the_program_directly() {
2010 let dir = tempfile::tempdir().unwrap();
2011 let store = open_store(&dir).await;
2012 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2013
2014 let id = driver
2018 .spawn_run(
2019 "unused-because-argv-wins",
2020 SpawnOpts {
2021 cwd: "/tmp".into(),
2022 argv: Some(vec![
2023 "/bin/sh".into(),
2024 "-c".into(),
2025 "printf 'argv0=%s\\n' \"$0\"".into(),
2026 "direct-exec-marker".into(),
2027 ]),
2028 ..Default::default()
2029 },
2030 )
2031 .await
2032 .unwrap();
2033
2034 await_done(&store, &id).await;
2035 let text = output_of(&store, &id).await;
2036 assert!(
2037 text.contains("argv0=direct-exec-marker"),
2038 "argv should have been exec'd verbatim, got: {text:?}"
2039 );
2040 }
2041
2042 #[tokio::test]
2043 async fn explicit_argv_still_records_the_requested_command() {
2044 let dir = tempfile::tempdir().unwrap();
2045 let store = open_store(&dir).await;
2046 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2047
2048 let id = driver
2053 .spawn_run(
2054 "$SHELL",
2055 SpawnOpts {
2056 cwd: "/tmp".into(),
2057 argv: Some(vec!["/bin/sh".into(), "-c".into(), "true".into()]),
2058 ..Default::default()
2059 },
2060 )
2061 .await
2062 .unwrap();
2063
2064 let meta = await_done(&store, &id).await;
2065 assert_eq!(meta.command, "$SHELL");
2066 assert!(
2067 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2068 "expected Done(0), got {:?}",
2069 meta.status
2070 );
2071 }
2072
2073 #[tokio::test]
2079 async fn pipefail_reports_the_failing_stage_and_posix_reports_the_last_one() {
2080 let dir = tempfile::tempdir().unwrap();
2081 let store = open_store(&dir).await;
2082 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2083
2084 let line = "(exit 101) | tail -1";
2087
2088 let posix = driver
2089 .spawn_run(
2090 line,
2091 SpawnOpts { cwd: "/tmp".into(), pipefail: false, ..Default::default() },
2092 )
2093 .await
2094 .unwrap();
2095 let meta = await_done(&store, &posix).await;
2096 assert!(
2097 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2098 "POSIX pipeline status is the LAST stage's — expected Done(0), got {:?}",
2099 meta.status
2100 );
2101
2102 let failing = driver
2103 .spawn_run(
2104 line,
2105 SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2106 )
2107 .await
2108 .unwrap();
2109 let meta = await_done(&store, &failing).await;
2110 assert!(
2111 matches!(meta.status, RunStatus::Done { exit_code: 101, .. }),
2112 "pipefail must surface the producer's 101, got {:?}",
2113 meta.status
2114 );
2115 }
2116
2117 #[tokio::test]
2120 async fn pipefail_does_not_leak_into_the_recorded_command() {
2121 let dir = tempfile::tempdir().unwrap();
2122 let store = open_store(&dir).await;
2123 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2124
2125 let id = driver
2126 .spawn_run(
2127 "echo recorded-verbatim | cat",
2128 SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2129 )
2130 .await
2131 .unwrap();
2132
2133 let meta = await_done(&store, &id).await;
2134 assert_eq!(meta.command, "echo recorded-verbatim | cat");
2135 assert!(
2136 !meta.command.contains("pipefail"),
2137 "the prelude leaked into the recorded command: {:?}",
2138 meta.command
2139 );
2140 }
2141
2142 #[tokio::test]
2149 async fn the_pipefail_probe_never_costs_the_command_that_follows_it() {
2150 let dir = tempfile::tempdir().unwrap();
2151 let store = open_store(&dir).await;
2152 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2153
2154 let id = driver
2155 .spawn_run(
2156 "echo probe-survived",
2157 SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2158 )
2159 .await
2160 .unwrap();
2161
2162 let meta = await_done(&store, &id).await;
2163 let text = output_of(&store, &id).await;
2164 assert!(
2165 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2166 "expected Done(0), got {:?}",
2167 meta.status
2168 );
2169 assert!(text.contains("probe-survived"), "command did not run, got: {text:?}");
2170 assert!(
2172 !text.contains("pipefail"),
2173 "the probe printed a diagnostic into the build's own output: {text:?}"
2174 );
2175 }
2176
2177 #[tokio::test]
2178 async fn empty_argv_falls_back_to_the_shell_path() {
2179 let dir = tempfile::tempdir().unwrap();
2180 let store = open_store(&dir).await;
2181 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2182
2183 let id = driver
2184 .spawn_run(
2185 "echo empty_argv_fallback",
2186 SpawnOpts { cwd: "/tmp".into(), argv: Some(vec![]), ..Default::default() },
2187 )
2188 .await
2189 .unwrap();
2190
2191 await_done(&store, &id).await;
2192 let text = output_of(&store, &id).await;
2193 assert!(
2194 text.contains("empty_argv_fallback"),
2195 "empty argv must not spawn nothing, got: {text:?}"
2196 );
2197 }
2198
2199 #[tokio::test]
2200 async fn spawn_failing_command_records_nonzero_exit() {
2201 let dir = tempfile::tempdir().unwrap();
2202 let store = open_store(&dir).await;
2203 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2204
2205 let id = driver
2206 .spawn_run(
2207 "exit 42",
2208 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2209 )
2210 .await
2211 .unwrap();
2212
2213 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2214 loop {
2215 let meta = store.get_run(&id).await.unwrap().unwrap();
2216 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2217 match meta.status {
2218 RunStatus::Done { exit_code, .. } => {
2219 assert_ne!(exit_code, 0, "exit 42 should produce a non-zero exit code");
2220 }
2221 other => panic!("unexpected status: {other:?}"),
2222 }
2223 break;
2224 }
2225 if std::time::Instant::now() > deadline {
2226 panic!("run did not complete in time");
2227 }
2228 tokio::time::sleep(Duration::from_millis(50)).await;
2229 }
2230 }
2231
2232 #[cfg(unix)]
2235 #[tokio::test]
2236 async fn kill_with_sigterm_transitions_to_killed() {
2237 let dir = tempfile::tempdir().unwrap();
2238 let store = open_store(&dir).await;
2239 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2240
2241 let id = driver
2242 .spawn_run(
2243 "sleep 60",
2244 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2245 )
2246 .await
2247 .unwrap();
2248
2249 tokio::time::sleep(Duration::from_millis(100)).await;
2251
2252 driver.kill_run(&id, Some(SIGTERM)).await.unwrap();
2253
2254 let deadline = std::time::Instant::now() + Duration::from_secs(10);
2255 loop {
2256 let meta = store.get_run(&id).await.unwrap().unwrap();
2257 if matches!(meta.status, RunStatus::Killed { .. } | RunStatus::Lost { .. }) {
2258 assert!(
2259 matches!(meta.status, RunStatus::Killed { .. }),
2260 "expected Killed, got {:?}",
2261 meta.status
2262 );
2263 break;
2264 }
2265 if std::time::Instant::now() > deadline {
2266 panic!("run did not become Killed in time, status={:?}", meta.status);
2267 }
2268 tokio::time::sleep(Duration::from_millis(50)).await;
2269 }
2270 }
2271
2272 #[cfg(unix)]
2273 #[tokio::test]
2274 async fn kill_run_returns_not_found_after_exit() {
2275 let dir = tempfile::tempdir().unwrap();
2276 let store = open_store(&dir).await;
2277 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2278
2279 let id = driver
2280 .spawn_run(
2281 "echo done",
2282 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2283 )
2284 .await
2285 .unwrap();
2286
2287 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2289 loop {
2290 let meta = store.get_run(&id).await.unwrap().unwrap();
2291 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2292 break;
2293 }
2294 if std::time::Instant::now() > deadline {
2295 panic!("run did not complete");
2296 }
2297 tokio::time::sleep(Duration::from_millis(50)).await;
2298 }
2299
2300 let result = driver.kill_run(&id, None).await;
2302 assert!(
2303 matches!(result, Err(DriverError::NotFound(_))),
2304 "expected NotFound, got {result:?}"
2305 );
2306 }
2307
2308 #[cfg(unix)]
2311 #[tokio::test]
2312 async fn stdin_send_reaches_child() {
2313 let dir = tempfile::tempdir().unwrap();
2314 let store = open_store(&dir).await;
2315 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2316
2317 let id = driver
2319 .spawn_run(
2320 "read line && echo got_$line",
2321 SpawnOpts {
2322 cwd: "/tmp".into(),
2323 stdin_enabled: true,
2324 ..Default::default()
2325 },
2326 )
2327 .await
2328 .unwrap();
2329
2330 tokio::time::sleep(Duration::from_millis(150)).await;
2331 driver.send_stdin(&id, b"hello\n".to_vec()).await.unwrap();
2332
2333 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2334 loop {
2335 let meta = store.get_run(&id).await.unwrap().unwrap();
2336 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2337 break;
2338 }
2339 if std::time::Instant::now() > deadline {
2340 panic!("run did not complete after stdin input");
2341 }
2342 tokio::time::sleep(Duration::from_millis(50)).await;
2343 }
2344
2345 let chunks = store.get_chunks(&id, &ChunkFilter::default()).await.unwrap();
2346 let raw: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
2347 let text = String::from_utf8_lossy(&raw);
2348 assert!(
2349 text.contains("got_hello"),
2350 "expected 'got_hello' in output, got: {text:?}"
2351 );
2352 }
2353
2354 #[tokio::test]
2358 async fn resize_run_changes_geometry_the_child_sees() {
2359 let dir = tempfile::tempdir().unwrap();
2360 let store = open_store(&dir).await;
2361 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2362
2363 let id = driver
2365 .spawn_run(
2366 "read line && stty size",
2367 SpawnOpts {
2368 cwd: "/tmp".into(),
2369 stdin_enabled: true,
2370 ..Default::default()
2373 },
2374 )
2375 .await
2376 .unwrap();
2377
2378 tokio::time::sleep(Duration::from_millis(150)).await;
2379 driver.resize_run(&id, 120, 40).await.unwrap();
2380 driver.send_stdin(&id, b"go\n".to_vec()).await.unwrap();
2381
2382 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2383 loop {
2384 let meta = store.get_run(&id).await.unwrap().unwrap();
2385 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2386 break;
2387 }
2388 if std::time::Instant::now() > deadline {
2389 panic!("run did not complete after stdin input");
2390 }
2391 tokio::time::sleep(Duration::from_millis(50)).await;
2392 }
2393
2394 let chunks = store.get_chunks(&id, &ChunkFilter::default()).await.unwrap();
2395 let raw: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
2396 let text = String::from_utf8_lossy(&raw);
2397 assert!(
2398 text.contains("40 120"),
2399 "expected resized geometry '40 120' in output, got: {text:?}"
2400 );
2401 }
2402
2403 #[tokio::test]
2406 async fn resize_run_returns_not_found_after_exit() {
2407 let dir = tempfile::tempdir().unwrap();
2408 let store = open_store(&dir).await;
2409 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2410
2411 let id = driver
2412 .spawn_run("true", SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2413 .await
2414 .unwrap();
2415
2416 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2417 loop {
2418 let meta = store.get_run(&id).await.unwrap().unwrap();
2419 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2420 break;
2421 }
2422 if std::time::Instant::now() > deadline {
2423 panic!("run did not exit");
2424 }
2425 tokio::time::sleep(Duration::from_millis(50)).await;
2426 }
2427
2428 assert!(matches!(
2429 driver.resize_run(&id, 100, 30).await,
2430 Err(DriverError::NotFound(_))
2431 ));
2432 }
2433
2434 #[cfg(unix)]
2442 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
2443 async fn log_pipe_events_land_in_store() {
2444 use crate::store::EventFilter;
2445
2446 let dir = tempfile::tempdir().unwrap();
2447 let store = open_store(&dir).await;
2448 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2449
2450 let cmd = r#"printf '{"level":"warn","target":"test.shim","msg":"hello-from-pipe","fields":{"x":42},"_lib":"test-shim","_lib_ver":"0.1.0"}\n' >> "$YAH_LOG_PIPE""#;
2453
2454 let id = driver
2455 .spawn_run(cmd, SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2456 .await
2457 .unwrap();
2458
2459 let deadline = std::time::Instant::now() + Duration::from_secs(20);
2463 loop {
2464 let meta = store.get_run(&id).await.unwrap().unwrap();
2465 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
2466 break;
2467 }
2468 if std::time::Instant::now() > deadline {
2469 panic!("run did not complete in time");
2470 }
2471 tokio::time::sleep(Duration::from_millis(50)).await;
2472 }
2473
2474 tokio::time::sleep(Duration::from_millis(500)).await;
2477
2478 let events = store.query_events(&id, &EventFilter::default()).await.unwrap();
2479 assert!(
2480 !events.is_empty(),
2481 "expected at least one shim event, got none"
2482 );
2483 let ev = events.iter().find(|e| e.target == "test.shim");
2484 let ev = ev.expect("event with target 'test.shim' not found");
2485 assert_eq!(ev.msg, "hello-from-pipe");
2486 assert_eq!(ev.level, crate::types::Level::Warn);
2487 assert!(
2488 matches!(&ev.source, crate::types::EventSource::Shim { lib, .. } if lib == "test-shim"),
2489 "unexpected source: {:?}",
2490 ev.source
2491 );
2492 assert_eq!(ev.fields.get("x"), Some(&serde_json::json!(42)));
2493 }
2494
2495 #[cfg(unix)]
2498 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2499 async fn log_pipe_disabled_produces_no_events() {
2500 use crate::store::EventFilter;
2501
2502 let dir = tempfile::tempdir().unwrap();
2503 let store = open_store(&dir).await;
2504 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2505
2506 let cmd = r#"[ -n "$YAH_LOG_PIPE" ] && printf '{"level":"info","target":"t","msg":"m","fields":{}}\n' >> "$YAH_LOG_PIPE" || true"#;
2509
2510 let id = driver
2511 .spawn_run(
2512 cmd,
2513 SpawnOpts { cwd: "/tmp".into(), log_fd_enabled: false, ..Default::default() },
2514 )
2515 .await
2516 .unwrap();
2517
2518 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2519 loop {
2520 let meta = store.get_run(&id).await.unwrap().unwrap();
2521 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
2522 break;
2523 }
2524 if std::time::Instant::now() > deadline {
2525 panic!("run did not complete");
2526 }
2527 tokio::time::sleep(Duration::from_millis(50)).await;
2528 }
2529
2530 tokio::time::sleep(Duration::from_millis(100)).await;
2531
2532 let events = store.query_events(&id, &EventFilter::default()).await.unwrap();
2533 assert!(
2534 events.is_empty(),
2535 "expected no shim events when log_fd_enabled=false, got {}",
2536 events.len()
2537 );
2538 }
2539
2540 const BUILD_RUN: &str = "build-run";
2546
2547 fn opted_in() -> Vec<String> {
2548 vec![BUILD_RUN.to_string()]
2549 }
2550
2551 async fn spawn_long_run(driver: &TaskDriver, origin: &str) -> TaskRunId {
2552 driver
2553 .spawn_run(
2554 "sleep 30",
2555 SpawnOpts {
2556 cwd: "/tmp".into(),
2557 origin: Some(origin.to_string()),
2558 ..Default::default()
2559 },
2560 )
2561 .await
2562 .unwrap()
2563 }
2564
2565 async fn await_status(
2566 store: &Arc<TaskStore>,
2567 id: &TaskRunId,
2568 want: fn(&RunStatus) -> bool,
2569 ) -> RunStatus {
2570 let deadline = std::time::Instant::now() + Duration::from_secs(10);
2571 loop {
2572 let status = store.get_run(id).await.unwrap().unwrap().status;
2573 if want(&status) {
2574 return status;
2575 }
2576 if std::time::Instant::now() > deadline {
2577 panic!("run never reached the expected status, last={status:?}");
2578 }
2579 tokio::time::sleep(Duration::from_millis(25)).await;
2580 }
2581 }
2582
2583 #[tokio::test]
2586 async fn an_unpolled_run_of_an_opted_in_origin_is_reaped() {
2587 let dir = tempfile::tempdir().unwrap();
2588 let store = open_store(&dir).await;
2589 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2590
2591 let id = spawn_long_run(&driver, BUILD_RUN).await;
2592 tokio::time::sleep(Duration::from_millis(300)).await;
2593
2594 let reaped = driver
2595 .reap_unattached(Duration::from_millis(200), &opted_in())
2596 .await;
2597 assert_eq!(reaped, vec![id.clone()], "the unattached run should be reaped");
2598
2599 let status = await_status(&store, &id, |s| {
2600 matches!(s, RunStatus::Killed { .. } | RunStatus::Done { .. })
2601 })
2602 .await;
2603 assert!(
2604 matches!(status, RunStatus::Killed { .. }),
2605 "a reaped run ends Killed, got {status:?}"
2606 );
2607 }
2608
2609 #[tokio::test]
2612 async fn a_run_a_client_is_still_polling_is_never_reaped() {
2613 let dir = tempfile::tempdir().unwrap();
2614 let store = open_store(&dir).await;
2615 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2616
2617 let id = spawn_long_run(&driver, BUILD_RUN).await;
2618
2619 for _ in 0..6 {
2622 tokio::time::sleep(Duration::from_millis(100)).await;
2623 driver.note_attached(&id);
2624 let reaped = driver
2625 .reap_unattached(Duration::from_millis(200), &opted_in())
2626 .await;
2627 assert!(reaped.is_empty(), "a polled run must survive, reaped {reaped:?}");
2628 }
2629
2630 assert!(
2631 matches!(
2632 store.get_run(&id).await.unwrap().unwrap().status,
2633 RunStatus::Running
2634 ),
2635 "the polled run should still be running"
2636 );
2637
2638 tokio::time::sleep(Duration::from_millis(300)).await;
2641 let reaped = driver
2642 .reap_unattached(Duration::from_millis(200), &opted_in())
2643 .await;
2644 assert_eq!(reaped, vec![id], "a run that stopped being polled is reapable");
2645 }
2646
2647 #[tokio::test]
2651 async fn a_terminal_tile_is_never_reaped_however_long_it_idles() {
2652 let dir = tempfile::tempdir().unwrap();
2653 let store = open_store(&dir).await;
2654 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2655
2656 let id = spawn_long_run(&driver, "terminal").await;
2657 tokio::time::sleep(Duration::from_millis(300)).await;
2658
2659 for _ in 0..3 {
2660 let reaped = driver.reap_unattached(Duration::ZERO, &opted_in()).await;
2661 assert!(
2662 reaped.is_empty(),
2663 "a terminal tile is outside the opted-in origins, reaped {reaped:?}"
2664 );
2665 tokio::time::sleep(Duration::from_millis(50)).await;
2666 }
2667
2668 assert!(
2669 matches!(
2670 store.get_run(&id).await.unwrap().unwrap().status,
2671 RunStatus::Running
2672 ),
2673 "the terminal run must still be running"
2674 );
2675
2676 driver.kill_run(&id, Some(SIGKILL)).await.unwrap();
2677 }
2678
2679 #[tokio::test]
2682 async fn an_origin_less_run_and_an_empty_opt_in_list_reap_nothing() {
2683 let dir = tempfile::tempdir().unwrap();
2684 let store = open_store(&dir).await;
2685 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2686
2687 let plain = driver
2688 .spawn_run("sleep 30", SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2689 .await
2690 .unwrap();
2691 let build = spawn_long_run(&driver, BUILD_RUN).await;
2692 tokio::time::sleep(Duration::from_millis(100)).await;
2693
2694 assert!(
2695 driver.reap_unattached(Duration::ZERO, &[]).await.is_empty(),
2696 "an empty opt-in list must reap nothing, not everything"
2697 );
2698 assert_eq!(
2699 driver.reap_unattached(Duration::ZERO, &opted_in()).await,
2700 vec![build],
2701 "only the opted-in origin is reapable"
2702 );
2703
2704 assert!(
2705 matches!(
2706 store.get_run(&plain).await.unwrap().unwrap().status,
2707 RunStatus::Running
2708 ),
2709 "the origin-less run must be untouched"
2710 );
2711 driver.kill_run(&plain, Some(SIGKILL)).await.unwrap();
2712 }
2713
2714 #[tokio::test]
2717 async fn attached_age_resets_on_a_poll_and_is_none_for_a_foreign_run() {
2718 let dir = tempfile::tempdir().unwrap();
2719 let store = open_store(&dir).await;
2720 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2721
2722 let id = spawn_long_run(&driver, BUILD_RUN).await;
2723 tokio::time::sleep(Duration::from_millis(150)).await;
2724 let aged = driver.attached_age(&id).expect("driver owns this run");
2725 assert!(aged >= Duration::from_millis(100), "age should have grown, got {aged:?}");
2726
2727 driver.note_attached(&id);
2728 let fresh = driver.attached_age(&id).unwrap();
2729 assert!(fresh < aged, "a poll resets the age: {fresh:?} vs {aged:?}");
2730
2731 assert!(
2732 driver.attached_age(&TaskRunId::new()).is_none(),
2733 "a run this driver does not own has no attachment age"
2734 );
2735
2736 driver.kill_run(&id, Some(SIGKILL)).await.unwrap();
2737 }
2738}