Skip to main content

task_runs/
driver.rs

1//! @arch:layer(kg_store)
2//! @arch:role(substrate)
3//! @arch:see(.yah/docs/working/yah-task-runs.md)
4//!
5//! PTY subprocess driver — spawn commands, capture output as append-only
6//! chunks, handle SIGTERM/SIGKILL with a grace period, and mark stale
7//! `Running` runs as `Lost` when the daemon restarts.
8//!
9//! ## Tier 2 side-channel (yah-log shims)
10//!
11//! When `SpawnOpts::log_fd_enabled` is true (the default), the driver creates
12//! a named pipe (FIFO) and exports two env vars into the child:
13//!
14//! - `YAH_TASK_RUN`  — the `TaskRunId` as a hyphenated UUID string.
15//! - `YAH_LOG_PIPE`  — absolute path to the FIFO.
16//!
17//! The child opens `YAH_LOG_PIPE` for writing and emits JSON-lines. The
18//! driver reads those lines in a background thread and stores them as
19//! [`EventSource::Shim`] events.
20//!
21//! **Why FIFO instead of a raw fd?** `portable-pty` calls `close_random_fds()`
22//! in its `pre_exec` hook, closing every fd ≥ 3 before exec. A raw-pipe write
23//! fd is always ≥ 3 and would be closed before the child could use it. Opening
24//! a FIFO by path requires no fd inheritance.
25//!
26//! Wire format — one JSON object per line:
27//! ```json
28//! {"level":"info","target":"myapp::module","msg":"text","fields":{"key":"val"}}
29//! ```
30//! Optional shim-identity keys: `"_lib"` (string), `"_lib_ver"` (string).
31//! Unknown keys in `fields` pass through as freeform JSON.
32//!
33//! The driver holds the write end of the FIFO open until the run lifecycle
34//! task completes, which triggers EOF for the receiver thread. The FIFO file
35//! is deleted after the receiver thread drains the last line.
36//!
37//! On non-Unix platforms `YAH_TASK_RUN` and `YAH_LOG_PIPE` are not exported.
38//! Shim libraries must treat absent `YAH_TASK_RUN` as "not inside a TaskRun".
39//!
40//! @arch:see(.yah/docs/working/W280-durable-terminal-sessions.md)
41
42use 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// ─── Error ────────────────────────────────────────────────────────────────────
63
64#[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// ─── SpawnOpts ────────────────────────────────────────────────────────────────
77
78/// Options for [`TaskDriver::spawn_run`].
79#[derive(Debug, Clone)]
80pub struct SpawnOpts {
81    pub cwd: PathBuf,
82    /// Env vars set on the child process (merged on top of the current env).
83    pub env: Vec<(String, String)>,
84    pub label: Option<String>,
85    pub initiator: Initiator,
86    /// PTY column count. Defaults to 80.
87    pub pty_cols: u16,
88    /// PTY row count. Defaults to 24.
89    pub pty_rows: u16,
90    /// Enable stdin relay via [`TaskDriver::send_stdin`].
91    pub stdin_enabled: bool,
92    /// Pin the run so the GC sweep does not drop its output during warm rolloff.
93    pub pin: bool,
94    /// Beholder attachment policy. Defaults to [`BeholderSelect::Auto`].
95    pub beholder_select: BeholderSelect,
96    /// `true` when this run's output is read as text by somebody downstream —
97    /// a human watching a terminal tile, or a client that promised its caller
98    /// byte-identical passthrough. Causes `Rewriter` beholders to decline in
99    /// `Auto` mode, since a rewrite changes what that reader gets.
100    ///
101    /// R739-B10 renamed this from `tty_attached`: a PTY was only ever a proxy
102    /// for "someone is reading this", and the proxy broke the moment
103    /// `yah build run` moved onto pipes (R739-F6).
104    pub verbatim_output: bool,
105    /// Create a side-channel FIFO and export `YAH_TASK_RUN` / `YAH_LOG_PIPE`
106    /// so Tier-2 shim libraries (yah-log-rust, @yah/log) can emit structured
107    /// events. Has no effect on non-Unix platforms. Defaults to `true`.
108    pub log_fd_enabled: bool,
109    /// Provenance tag stored on the run's `TaskRunMeta.origin` (e.g.
110    /// `Some("terminal")` for an interactive shell). `None` is an ordinary job.
111    pub origin: Option<String>,
112    /// Exec this argv directly instead of wrapping `cmd` in `sh -c`.
113    ///
114    /// The default `sh -c <cmd>` is right for a job — the caller wrote a
115    /// command line and expects a shell to parse it. It is wrong for an
116    /// *interactive shell*: `sh -c "zsh -l"` leaves an inert `sh` as the PTY's
117    /// foreground process group leader, so job control misbehaves, signals go
118    /// to the wrong process, and anything that reads the foreground pid (a
119    /// live-cwd probe, say) sees `sh` instead of the shell the operator is
120    /// typing into. Handing the exact argv here makes the shell itself the
121    /// child, which is also the only way to pass `-l` as a real argv element
122    /// so `.zprofile` / `.profile` actually run.
123    ///
124    /// `cmd` is still what gets recorded on `TaskRunMeta.command`, so the run
125    /// reads the way the caller asked for it. Beholder argv rewriting is
126    /// bypassed when this is set: the caller has already decided the exact
127    /// process to exec, and a recorded `rewrite=…` that didn't happen would be
128    /// a lie in the run metadata.
129    pub argv: Option<Vec<String>>,
130    /// R739-F6 — spawn on **pipes** instead of a PTY. Defaults to `false`,
131    /// which is the PTY behaviour every existing caller already has.
132    ///
133    /// A PTY is right for an interactive terminal tile: the child gets a
134    /// controlling terminal, job control works, and `isatty` says yes, which is
135    /// what a human sitting in front of it expects. It is wrong for *emulating
136    /// a non-interactive shell invocation*, where three PTY properties show up
137    /// as divergence from running the same command directly (all three measured
138    /// in R739-F4 against `cargo check`):
139    ///
140    /// 1. `isatty(1)` is true, so tools colorize — plain `error: …` arrives as
141    ///    `\x1b[1m\x1b[91merror\x1b[0m: …`, which also defeats `| rg "^error"`.
142    /// 2. The line discipline's `ONLCR` rewrites every `\n` the child wrote
143    ///    into `\r\n`.
144    /// 3. The terminal merges stderr into stdout, so stream separation is gone
145    ///    by the time anything reads the capture.
146    ///
147    /// In pipe mode the child gets `pipe(2)` for stdout and stderr, chunks are
148    /// stored under their true [`Stream`], and `TERM` is left alone rather than
149    /// forced to `xterm-256color`. [`TaskDriver::resize_run`] and
150    /// [`TaskDriver::foreground_pid`] have no PTY to answer for and report
151    /// `NotFound` / `None`.
152    pub pipe: bool,
153    /// R901-B2 — run the `sh -c` line with `pipefail`, so a pipeline reports
154    /// the **leftmost** failing stage instead of its last one. Defaults to
155    /// `false`, i.e. POSIX behaviour, which is what every existing caller has.
156    ///
157    /// Without it a pipeline's status is the last stage's and nothing else:
158    /// `cargo check 2>&1 | tail -40` exits **0** on a build with 101 errors,
159    /// because `tail` succeeded. That is not a wrapper lying — the wrapper is
160    /// faithful, and the shell is answering the question it was actually
161    /// asked — but it is indistinguishable from a green build to everything
162    /// downstream, including the harness task notification an agent reads to
163    /// decide whether it is done. On 2026-09-13 two sessions read that 0 as a
164    /// pass and left `cargo check -p yah` red camp-wide for ~50 minutes.
165    ///
166    /// Only meaningful when [`SpawnOpts::argv`] is `None`; an explicit argv is
167    /// not a shell line and has no pipeline to take a status from.
168    ///
169    /// # Known cost, accepted deliberately
170    ///
171    /// `pipefail` also surfaces a producer killed by `SIGPIPE`, so
172    /// `cargo check 2>&1 | head -40` can now report failure once `head` closes
173    /// the pipe early on a build that was fine. That is a false RED, and it is
174    /// the right trade against the false GREEN above: a red is investigated,
175    /// a green ends the turn. Prefer `| tail` over `| head` on a build line.
176    pub pipefail: bool,
177}
178
179/// Prefix that turns `pipefail` on for the rest of a `sh -c` line.
180///
181/// Probing in a subshell rather than running `set -o pipefail` directly is
182/// load-bearing for portability, not caution. `pipefail` is a bash/ksh/zsh
183/// option; `/bin/sh` is bash on macOS but **dash** on most Linux distros, and
184/// dash rejects it. `set` is a POSIX *special* builtin, so a failure in one is
185/// entitled to terminate a non-interactive shell — which would turn "your
186/// pipeline now reports the truth" into "your command never ran at all" on
187/// every Linux camp. The subshell absorbs that exit; the outer shell only ever
188/// runs `set -o pipefail` on a shell that has already proved it accepts it.
189const 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// ─── Driver channels ─────────────────────────────────────────────────────────
214
215/// Optional side-channels a driver can publish to. Both are fire-and-forget:
216/// a closed receiver never stalls or fails a run.
217#[derive(Default)]
218pub struct DriverChannels {
219    /// Fires `(run_id, status)` after each run's lifecycle task writes the
220    /// terminal status. Drives completion listeners (e.g. a triage worker).
221    pub completion: Option<mpsc::UnboundedSender<(TaskRunId, RunStatus)>>,
222    /// Mirrors every PTY output chunk as it is captured, *before* any consumer
223    /// polls the store. Lets a host attach a live view (VT parser, log
224    /// forwarder) to a run without a read-back loop over the store.
225    ///
226    /// The driver deliberately stays ignorant of what the tap is for — the
227    /// chunk carries `run_id`, so the host decides which runs it cares about.
228    pub output: Option<mpsc::UnboundedSender<OutputChunk>>,
229}
230
231// ─── Stale-run policy ────────────────────────────────────────────────────────
232
233/// What a freshly-constructed [`TaskDriver`] does with `Running` rows it finds
234/// already in the store.
235///
236/// The historical rule — tombstone every one of them — bakes in an assumption
237/// that stops being true the moment a second process attaches to the same
238/// store: that any `Running` row must be a corpse from *this* process's
239/// predecessor. When two processes share a store, a driver starting up in one
240/// will happily mark the other's live runs `Lost`, and the run keeps producing
241/// output under a status that says it is dead.
242///
243/// The policy is deliberately origin-agnostic in its mechanism — it decides on
244/// **who owns the run** ([`TaskRunMeta::host_pid`]) — and takes the origin list
245/// as data, so an embedder names the runs it wants exempted without this crate
246/// knowing what any of them mean.
247#[derive(Debug, Clone, Default, PartialEq, Eq)]
248pub enum StaleRunPolicy {
249    /// Tombstone every leftover `Running` run as `Lost`.
250    ///
251    /// Correct, and the default, whenever this process is the only writer:
252    /// a run whose driver is gone has no one left to notice it exit.
253    #[default]
254    LostOnDisappear,
255    /// Spare runs whose recorded owner process is still alive.
256    ///
257    /// A leftover run is tombstoned only when its `host_pid` is absent (owner
258    /// unknown — a row from before the column existed) or names a process that
259    /// no longer exists. Anything else belongs to a live peer and is left
260    /// `Running` for that peer to finish.
261    ///
262    /// `origins` narrows the exemption to runs whose
263    /// [`TaskRunMeta::origin`] is in the list; empty means every origin
264    /// qualifies. A run with no origin never matches a non-empty list.
265    AdoptLiveHosts { origins: Vec<String> },
266}
267
268impl StaleRunPolicy {
269    /// Whether `meta` should be tombstoned `Lost` at driver construction.
270    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/// Is a process with this pid still around?
292///
293/// `kill(pid, 0)` is the portable liveness probe: it performs the permission
294/// check and existence lookup without delivering anything. `EPERM` counts as
295/// alive — the process exists, it just is not ours to signal.
296///
297/// Pid reuse can make a dead owner read as alive. That is why
298/// [`StaleRunPolicy::AdoptLiveHosts`] is opt-in and origin-narrowed: the cost
299/// of a false "alive" is one run left `Running` until something closes it,
300/// which is strictly better for an interactive session than the false "dead"
301/// this replaces — which kills a *live* session's status.
302#[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    // SAFETY: `kill` with signal 0 delivers nothing; it only reports whether
311    // the pid exists and is signallable.
312    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/// No `kill(2)` off Unix. Reporting every owner dead keeps the historical
317/// Lost-on-disappear behaviour rather than stranding runs `Running` forever.
318#[cfg(not(unix))]
319fn host_process_alive(_pid: u32) -> bool {
320    false
321}
322
323// ─── Internal run-control handle ─────────────────────────────────────────────
324
325struct RunControl {
326    kill_tx: mpsc::Sender<KillRequest>,
327    stdin_tx: Option<mpsc::Sender<Vec<u8>>>,
328    /// Shared with the lifecycle task, which holds the same `Arc` so the PTY fd
329    /// outlives `child.wait()`. `MasterPty::resize` takes `&self`, so a mutex is
330    /// enough to make the `Box<dyn MasterPty + Send>` `Sync` across the two.
331    ///
332    /// `None` for a [`SpawnOpts::pipe`] run, which has no terminal to resize or
333    /// to ask for a foreground process group.
334    master: Option<Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>>>,
335    /// R739-B12 — the run's [`SpawnOpts::origin`], copied here so
336    /// [`TaskDriver::reap_unattached`] can narrow to an opted-in origin set
337    /// without a store round-trip per candidate.
338    origin: Option<String>,
339    /// R739-B12 — when a client last looked at this run.
340    ///
341    /// Set at spawn (the caller that asked for the run is attached to it by
342    /// definition) and refreshed by [`TaskDriver::note_attached`], which the
343    /// embedder calls from whatever its "a client is watching" surface is —
344    /// for the camp daemon, `task.tail` and `task.status`.
345    ///
346    /// Monotonic rather than a wall clock: a clock step must not be able to
347    /// make a healthy build look abandoned.
348    last_attached_at: Instant,
349}
350
351/// Fires the reader-done signal when the LAST holder drops.
352///
353/// A PTY run has one reader; a piped run has two (stdout and stderr) and the
354/// lifecycle must not reap the child until both have hit EOF. Making this a
355/// drop guard behind an `Arc` means neither path has to count readers: the
356/// signal goes out when the refcount reaches zero, after each pump has finished
357/// its own `on_done` work.
358struct 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// ─── ShimRecord ───────────────────────────────────────────────────────────────
374
375/// One JSON-line record emitted by a Tier-2 shim to the side-channel FIFO.
376///
377/// The shim (Rust `yah-log` layer or TS `@yah/log` pino transport) writes one
378/// of these per log call. Unknown keys inside `fields` pass through unchanged.
379#[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    /// Shim library name, e.g. `"yah-log-rust"`. Populates
388    /// [`EventSource::Shim::lib`].
389    #[serde(rename = "_lib", default)]
390    lib: Option<String>,
391    /// Shim library version string.
392    #[serde(rename = "_lib_ver", default)]
393    lib_version: Option<String>,
394}
395
396// ─── FdCloser ─────────────────────────────────────────────────────────────────
397
398/// RAII wrapper that closes a raw fd on drop.
399///
400/// Used to hold the write end of the log FIFO open until the lifecycle task
401/// completes. Dropping it signals EOF to the receiver thread.
402#[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// SAFETY: a raw fd number is an integer; closing it from any thread is safe
413// provided we never duplicate ownership (enforced by move semantics here).
414#[cfg(unix)]
415unsafe impl Send for FdCloser {}
416
417// ─── TaskDriver ───────────────────────────────────────────────────────────────
418
419/// Manages in-flight task runs for a single camp.
420///
421/// Wrap in `Arc` to share across tasks; internal state is mutex-protected.
422pub struct TaskDriver {
423    store: Arc<TaskStore>,
424    active: Arc<Mutex<HashMap<String, RunControl>>>,
425    /// Side-channels published to by every run this driver owns.
426    channels: DriverChannels,
427}
428
429impl TaskDriver {
430    /// Create a driver backed by `store`, with no side-channels.
431    ///
432    /// Immediately scans the store for `Running` runs left over from a prior
433    /// daemon process and marks them `Lost` ("Lost-on-disappear").
434    pub async fn new(store: Arc<TaskStore>) -> Result<Self, DriverError> {
435        Self::with_channels(store, DriverChannels::default()).await
436    }
437
438    /// Like `new` but wires the optional [`DriverChannels`] side-channels
439    /// (completion notifications, live output tap).
440    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    /// Full constructor: side-channels plus the [`StaleRunPolicy`] applied to
448    /// `Running` rows already in the store.
449    ///
450    /// R617-F6 — annotation in this file's header. Splitting the sweep out of
451    /// the constructor's fixed behaviour is what lets a store be shared: a
452    /// process that is not the run's owner can now attach without declaring
453    /// the owner's live work dead.
454    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            /* A sweep tombstone is a terminal status like any other, and a
474               client watching that run needs it as much as it needs an exit
475               code. Until R267-T12 this was the one status transition in the
476               driver that wrote the store and told nobody — which is precisely
477               the transition a client cannot discover on its own, because the
478               process that would have reported it is the one that died. */
479            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    /// Spawn `cmd` in a PTY and start capturing its output. Returns immediately
491    /// with the new [`TaskRunId`].
492    ///
493    /// A beholder is selected via `opts.beholder_select` (default `Auto`). When
494    /// a `Rewriter` beholder matches, its `adjust_argv` is applied to the
495    /// command before spawning and the diff is recorded on `beholder_status`.
496    /// When `opts.verbatim_output` is `true`, `Rewriter` beholders decline in
497    /// `Auto` mode, because something downstream renders these bytes and a
498    /// rewrite would change them.
499    ///
500    /// Output is written to the store as `Stream::Stdout` chunks (the PTY
501    /// kernel merges stdout and stderr). Signal handling and status updates
502    /// run in background tasks.
503    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        // Attach a beholder (may rewrite argv and produce structured events).
509        // Resolve user drop-in directory: $YAH_BEHOLDERS_DIR or $HOME/.yah/beholders.
510        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        /* An explicit argv means the caller already chose the exact process
518           (an interactive login shell, say). Selecting a beholder there would
519           either do nothing — the rewritten argv is discarded on that path —
520           or record a rewrite that never happened, so we opt out honestly
521           instead. */
522        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        // Reconstruct the command from argv ONLY when a beholder actually
529        // rewrote it. `AttachResult.argv` is always populated — it is
530        // `resolve_argv(cmd)` even when nothing attached — so joining it
531        // unconditionally ran every run's command through a whitespace
532        // normalization nobody asked for: runs of spaces collapse and embedded
533        // newlines become spaces, which is silent corruption for a heredoc or
534        // any multi-line line. The caller's bytes go to the shell untouched
535        // unless a rewrite is the whole point.
536        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            /* R617-F6: stamp the OWNER, before the child exists. Written at
554               insert rather than after spawn so a crash between the two still
555               leaves the row attributable — an unattributed `Running` row is
556               exactly what the conservative arm of `StaleRunPolicy` has to
557               tombstone. */
558            host_pid: Some(std::process::id()),
559        }).await?;
560
561        // The program and argv both spawn modes exec. An explicit argv execs
562        // that program directly; otherwise the command line goes through `sh`
563        // so the caller's quoting, pipes and redirections mean what they say.
564        // An empty argv is a caller bug, not a request for an empty exec — fall
565        // back to the shell path rather than spawning nothing.
566        // R901-B2: `pipefail` is prepended HERE and not folded into
567        // `effective_cmd`, so `TaskRunMeta.command` keeps reading as the line
568        // the caller actually wrote. A run's recorded command is re-run by
569        // history and audited by agents against the relocation note; a prelude
570        // nobody asked for showing up in it would be the same class of lie as
571        // recording a beholder `rewrite=…` that never happened.
572        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        // ── Side-channel log FIFO (Tier 2 / yah-log shims) ──────────────────
585        //
586        // Create a named pipe (FIFO) so child processes can write structured
587        // events without touching stdout/stderr. We export its path via
588        // YAH_LOG_PIPE; no fd inheritance is involved, so portable-pty's
589        // close_random_fds() pre_exec hook doesn't interfere.
590        //
591        // The parent opens the FIFO twice:
592        //   rfd — O_RDONLY|O_NONBLOCK, then cleared to blocking → read events
593        //   wfd — O_WRONLY (wrapped in FdCloser) → keeps the FIFO alive until
594        //          the lifecycle task drops it (after run completion), producing
595        //          EOF for the receiver thread.
596        #[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                    // Path contained a nul byte — extremely unlikely; skip FIFO.
603                    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 // FIFO creation failed; continue without side-channel
612            } else {
613                // Open read end without blocking (no writer yet).
614                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                    // Switch read end to blocking so reads yield proper data.
622                    unsafe { libc::fcntl(rfd, libc::F_SETFL, 0) };
623                    // Open write end — this succeeds immediately because rfd is open.
624                    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        // The FIFO env, applied identically by both spawn modes.
641        #[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        /* Spawn. The two modes differ only in what the child's stdio is
649           attached to, and everything downstream — reader pumps, lifecycle,
650           kill — is written against the uniform handles produced here:
651           `pid`, a `reap` closure that blocks until the child exits, and an
652           optional PTY master for resize / foreground-pid. */
653        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        // Each entry is one blocking source to pump into the store. The PTY
658        // yields a single merged stream; pipes yield stdout and stderr apart.
659        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            /* Deliberately NOT setting TERM. The PTY path forces
671               `xterm-256color` because a child on a terminal that claims no
672               terminal type degrades badly; a child on a pipe should see
673               whatever the daemon's own environment says, exactly as it would
674               under a non-interactive shell. Forcing a terminal type here is
675               how a pipe-mode run would talk itself back into colorizing. */
676            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            // Open PTY pair.
710            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            // Clone reader before spawning so the fd is ready immediately.
721            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            // Optional stdin relay: take the writer before spawning the child.
728            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 the parent's slave handle so EOF propagates once the child exits.
763            drop(pair.slave);
764
765            pid = child.process_id().unwrap_or(0);
766
767            // Share the master between the lifecycle task (which must outlive
768            // `child.wait()` so the fd stays open) and `resize_run`.
769            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; // dropped after wait() returns, closing the PTY fd
775                c.wait().ok().map(|s| s.exit_code())
776            });
777        }
778
779        // ── FIFO: launch receiver thread; pass write-end holder to lifecycle ──
780        //
781        // The receiver thread reads until EOF. EOF arrives when ALL write-end
782        // holders close: the child's own writers (when it exits) plus the
783        // FdCloser we hand to the lifecycle task (which drops it after writing
784        // the terminal RunStatus). Events written before the last close are
785        // still drained by the receiver thread before it exits.
786        #[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            // spawn_blocking: lets the runtime track this thread so the
792            // Handle::block_on calls inside have a worker to drive futures.
793            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        // Channels.
802        let (kill_tx, kill_rx) = mpsc::channel::<KillRequest>(4);
803        let (reader_done_tx, reader_done_rx) = oneshot::channel::<()>();
804
805        /* Reader threads: child output → store chunks → beholder events. Each
806           runs on a dedicated OS thread because the reads are blocking. The
807           `ReaderDone` guard is shared across them, so the lifecycle's
808           reader-done signal fires only once every source has hit EOF — which
809           is what makes the two-pipe case correct without a reader count. */
810        {
811            let done = Arc::new(ReaderDone(Some(reader_done_tx)));
812            /* The beholder goes to stdout only. It parses a structured
813               protocol (cargo's JSON, say) that the child writes to stdout by
814               definition, and there is exactly one of it — handing the same
815               instance to two threads would need a lock for no gain, and
816               feeding it stderr would make `unknown_format_reason` fire on
817               human-readable diagnostics it was never meant to see. */
818            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        // Lifecycle task: monitor kill requests, wait for exit, update status.
834        // The task also holds the log FIFO write-end closer (if any) so that
835        // EOF propagates to the receiver thread after RunStatus is written.
836        {
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    /// Resize a running task's PTY and deliver `SIGWINCH` to the foreground
878    /// process group (portable-pty's `resize` does the ioctl, which is what
879    /// signals the child).
880    ///
881    /// Returns `DriverError::NotFound` when the run is not active on this
882    /// driver instance — the same contract as [`TaskDriver::send_stdin`] — and
883    /// also when it is active but was spawned in [`SpawnOpts::pipe`] mode, which
884    /// has no terminal to resize.
885    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    /// The pid of the run's *foreground* process — the leader of the process
911    /// group the PTY currently gives the keyboard to.
912    ///
913    /// For a shell tile that is the shell itself while it sits at a prompt,
914    /// and the command the operator is running while one is in flight. That
915    /// distinction is the whole point: asking the spawned child would report
916    /// the shell forever, so anything derived from this pid (a live cwd probe,
917    /// a "what is this pane doing" label) would answer for the wrong process.
918    ///
919    /// `None` when the run is not active on this driver instance, when it was
920    /// spawned in [`SpawnOpts::pipe`] mode (no controlling terminal, so no
921    /// foreground process group to read), or when the platform has no notion of
922    /// a foreground process group.
923    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    /// Send `signal` to a running task. Defaults to SIGTERM (15).
943    ///
944    /// For SIGTERM, the driver waits up to 5 seconds for the process to exit
945    /// before escalating to SIGKILL. Returns `DriverError::NotFound` if the
946    /// run is not active (already exited or launched on a different driver
947    /// instance).
948    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    /// Write bytes to the stdin of a running task (requires `stdin_enabled`).
966    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    /// R739-B12 — record that a client just looked at this run.
984    ///
985    /// A no-op for a run this driver does not own (already finished, or
986    /// spawned by another process against the same store): attachment only
987    /// means anything for a run something here could still signal.
988    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    /// How long ago a client last looked at `id`, or `None` when this driver
995    /// does not own the run. The observable half of [`Self::note_attached`].
996    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    /// R739-B12 — SIGTERM every run of an opted-in origin that no client has
1005    /// looked at for `idle`. Returns the runs it signalled.
1006    ///
1007    /// This exists because a run outlives the client that asked for it. When
1008    /// `yah build run` is SIGKILLed — its harness dies, the terminal goes away
1009    /// — the cargo it relocated into the daemon keeps compiling with nobody
1010    /// attached, holding the build-directory lock until a human finds the pid.
1011    /// That happened on 2026-08-28 and stalled a whole camp for ~30 minutes.
1012    /// R739-B9 closed every give-up the client is *alive* to make; this closes
1013    /// the one it is not.
1014    ///
1015    /// **Not [`StaleRunPolicy`], and not that policy on a timer.** The policy
1016    /// is a construction-time reconciliation of rows a *previous process*
1017    /// left behind: it decides on `host_pid`, only ever calls
1018    /// `store.update_status`, and tombstones any run outside its origin list
1019    /// outright — so running it periodically would mark every in-flight run of
1020    /// an un-adopted origin `Lost` while it compiles perfectly well, and would
1021    /// still never signal the process that is the actual problem. This is the
1022    /// opposite shape: it decides on *attachment*, it signals, and it touches
1023    /// nothing outside `origins`.
1024    ///
1025    /// `origins` is an opt-in list precisely because most runs must never be
1026    /// reaped on this rule. An interactive terminal tile is legitimately
1027    /// unpolled for hours, and killing one would be a far worse bug than the
1028    /// orphan this prevents — so an empty list reaps nothing at all, rather
1029    /// than meaning "every origin" the way [`StaleRunPolicy`]'s list does.
1030    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            // SIGTERM, not SIGKILL: `kill_run` gives the child the same 5s
1051            // grace a `task.kill` from a live client would, then escalates.
1052            // A terminal status is also what releases the run's admission
1053            // enrollment (R739-F7), so a reaped run frees the build key.
1054            if self.kill_run(&id, None).await.is_ok() {
1055                reaped.push(id);
1056            }
1057        }
1058        reaped
1059    }
1060}
1061
1062// ─── Log fd receiver ─────────────────────────────────────────────────────────
1063
1064/// Read JSON-lines from the side-channel FIFO read end and store them as
1065/// [`EventSource::Shim`] events.
1066///
1067/// Runs on a dedicated OS thread; exits when the read end sees EOF. EOF
1068/// arrives after both the child process AND the lifecycle task have closed
1069/// their write ends of the FIFO. The FIFO file is deleted on exit.
1070#[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    // SAFETY: `read_fd` is a valid, open FIFO fd handed exclusively to this
1083    // thread. `File` takes ownership and closes the fd on drop.
1084    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, // skip malformed lines silently
1099        };
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    // Clean up the FIFO file now that the receiver has drained.
1124    let _ = std::fs::remove_file(&fifo_path);
1125}
1126
1127// ─── Lifecycle task ───────────────────────────────────────────────────────────
1128
1129/// Pump one blocking output source into the store, tapping and beholding on the
1130/// way past.
1131///
1132/// Split out of `spawn_run` for R739-F6: a PTY run has one source and a piped
1133/// run has two, and the only thing that differs between them is which [`Stream`]
1134/// the chunks are stored under. `done` is the shared [`ReaderDone`] guard —
1135/// dropping it here, after `on_done`, is what tells the lifecycle this source is
1136/// finished.
1137#[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                        /* Both the tap and the beholder want the same owned
1163                           chunk; build it once, and only when someone is
1164                           listening. */
1165                        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                        /* Tap first: it feeds live views, where latency is
1175                           visible to a human. Send failure means the host
1176                           dropped its receiver — never fatal. */
1177                        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` blocks until the child exits and yields its exit code. It owns
1238    // whatever the spawn mode has to keep alive across the wait — for a PTY run
1239    // that includes the master fd, which must outlive `wait()`.
1240    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    // Holds the write end of the log FIFO open until this task completes.
1245    // Dropping it produces EOF for the receiver thread, which happens after
1246    // the terminal RunStatus is written below.
1247    #[cfg(unix)]
1248    _log_wfd: Option<FdCloser>,
1249) {
1250    // Pin the reader-done future so it can be polled by reference in
1251    // nested select! arms without consuming ownership.
1252    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                        // Grace period: give the process a chance to exit cleanly.
1266                        tokio::select! {
1267                            _ = &mut reader_done => {
1268                                // Exited within grace — no SIGKILL needed.
1269                                sent_signal = Some(signal);
1270                            }
1271                            _ = tokio::time::sleep(DEFAULT_GRACE) => {
1272                                // Grace expired — escalate.
1273                                send_unix_signal(pid, SIGKILL);
1274                                sent_signal = Some(SIGKILL);
1275                            }
1276                        }
1277                    }
1278                }
1279                // kill_tx dropped (driver shutting down) — force kill.
1280                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    // Reap the child (blocking) on a dedicated thread-pool slot. For a PTY run
1292    // the closure also owns our master handle, so the fd outlives the wait; the
1293    // matching `RunControl` (removed from `active` below) holds the other `Arc`,
1294    // so the fd actually closes once both are gone.
1295    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    /* Losing this write is not cosmetic: the run stays `Running` in the store
1309       forever and every reader — tail loops, the terminal UI, the next
1310       daemon's Lost-on-disappear sweep — believes a dead process is alive.
1311       `update_status` already retries through lock contention, so a failure
1312       here is terminal and worth saying out loud. */
1313    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
1322// ─── Helpers ──────────────────────────────────────────────────────────────────
1323
1324fn send_unix_signal(pid: u32, signal: i32) {
1325    #[cfg(unix)]
1326    unsafe {
1327        libc::kill(pid as libc::pid_t, signal);
1328    }
1329    // On non-Unix platforms signal delivery is not implemented here.
1330}
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// ─── Tests ────────────────────────────────────────────────────────────────────
1348
1349#[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    // ── Lost-on-disappear (pure store, no PTY) ────────────────────────────────
1359
1360    #[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        // Simulate a run left in "Running" state by a prior daemon.
1366        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        // Creating a new driver must mark stale runs Lost.
1386        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    /// R267-T12: a sweep tombstone must reach the completion channel.
1397    ///
1398    /// This was the one status transition that wrote the store and told
1399    /// nobody — and it is the transition a client can least discover on its
1400    /// own, because the process that would have reported the exit is the one
1401    /// that died. A renderer that trusts pushes would leave the pane on
1402    /// "running" for as long as it stayed open.
1403    #[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    // ── Stale-run policy (R617-F6) ───────────────────────────────────────────
1444
1445    /// Plant a `Running` row as if some other process had spawned it.
1446    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    /// The property the whole ticket exists for: attaching to a store must not
1488    /// declare another live process's shell dead.
1489    #[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        // Our own pid is by definition a live process, and is the cheapest
1494        // honest stand-in for "a peer that is still running".
1495        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    /// The other half: a genuinely abandoned shell must still be tombstoned,
1513    /// or a crashed host leaves permanent zombie tiles.
1514    #[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        // Reaped in-test, so the pid is real-but-dead rather than guessed.
1519        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    /// The exemption is narrowed by origin, so ordinary jobs keep the old rule
1543    /// even when their owner happens to still be alive — an in-flight `cargo
1544    /// build` whose driver is gone has nobody left to record its exit.
1545    #[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    /// A row written before `host_pid` existed reads back `None`. Unknown
1568    /// ownership must fall back to the old behaviour rather than stranding the
1569    /// run `Running` forever.
1570    #[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    /// `TaskDriver::new` must not have quietly changed behaviour — every
1588    /// existing embedder still gets Lost-on-disappear.
1589    #[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    /// The owner is recorded by `spawn_run` itself, not by the caller — the
1606    /// policy is worthless if rows arrive unattributed.
1607    #[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    // ── PTY spawn + capture ───────────────────────────────────────────────────
1667
1668    #[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        // Wait for the run to complete (poll status up to 5 s).
1683        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        // Chunks must contain "hello_world".
1696        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    // ── Pipe mode (R739-F6) ───────────────────────────────────────────────────
1716
1717    /// Run `cmd` to completion and return its stored chunks.
1718    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    /// Divergence 1 of 3 (R739-F4): the child must not think it is on a
1752    /// terminal. This is the one that makes cargo colorize.
1753    #[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        // The PTY default is unchanged — the terminal tiles depend on it.
1770        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    /// Divergence 2 of 3: no `ONLCR`, so a `\n` the child wrote stays a `\n`.
1781    /// This is what `build_run.rs::undo_onlcr` used to compensate for.
1782    #[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    /// Divergence 3 of 3: stdout and stderr stay apart, under their true
1799    /// [`Stream`], instead of being merged by the terminal.
1800    #[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        // Under a PTY the kernel merges them and everything lands on stdout —
1818        // the property that made stream separation unrecoverable downstream.
1819        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    /// Both pipes must reach EOF before the child is reaped, or a run whose
1833    /// last bytes went to stderr would be marked terminal with output still
1834    /// unread. The 4 KiB write is larger than a pipe's atomic-write buffer, so
1835    /// this fails if either pump is dropped rather than awaited.
1836    #[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    /// Exit codes have to survive the move to `std::process::Child`, which
1854    /// reports them through a different type than `portable_pty::Child`.
1855    #[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    /// A pipe run has no terminal, and the two PTY-only verbs must say so
1890    /// rather than reaching into a `None` master.
1891    #[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    /// Wait for a run to reach a terminal status, or panic.
1914    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    // ── The caller's bytes reach the shell unchanged (R739-S2) ───────────────
1935
1936    /// `AttachResult.argv` is populated on every run, rewrite or not, so
1937    /// `spawn_run` used to join it back into the command line unconditionally.
1938    /// That put every `task.run` command through a whitespace normalization
1939    /// nobody asked for. A multi-line command is the case where that is not
1940    /// cosmetic: the newline the caller wrote becomes a space, and two
1941    /// commands become one nonsense command.
1942    #[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        // Flattened to one line this is `echo one echo two`, which prints
1950        // "one echo two" — a different answer, not a failure, which is what
1951        // makes the old behaviour dangerous rather than merely wrong.
1952        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    /// `resolve_argv` strips `bunx`/`npx`/`pnpm` so a beholder's `matches` sees
1974    /// the bare tool. That is a *matching* concern; it must never reach the
1975    /// spawn, or the wrapper the caller needed is gone from the command.
1976    #[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        // `npx` is almost certainly absent in test environments, and that is
1984        // the point: if the wrapper survived, the shell reports it missing. If
1985        // it were stripped we would be running bare `--version`.
1986        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    // ── Direct argv (R652-T6) ────────────────────────────────────────────────
2007
2008    #[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        /* The distinguishing observation: under `sh -c` the child is `sh` and
2015           `$0` is `sh`; exec'd directly it is the program itself. Printing
2016           `$0` is the cheapest way to see which of the two happened. */
2017        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        /* A shell tile asks for "$SHELL" and the daemon resolves it to a real
2049           argv. The run must still read back as what was asked for, or the
2050           rail row and the history re-run both show an implementation
2051           detail. */
2052        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    /// R901-B2. The control is the whole point: without the `pipefail: false`
2074    /// half this would pass if `pipefail` stopped existing, and with only the
2075    /// `true` half it would pass if every pipeline had always reported its
2076    /// leftmost failure. The pair pins the *difference*, which is the thing
2077    /// that cost this camp ~50 minutes of red tree.
2078    #[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        // `(exit 101) | tail -1` is `cargo check 2>&1 | tail -40` with the
2085        // compile stripped out: a failing producer feeding a succeeding tail.
2086        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    /// The prelude must not reach [`TaskRunMeta::command`]: that string is what
2118    /// history re-runs and what an agent audits the relocation note against.
2119    #[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    /// The portability guard. On a `/bin/sh` that rejects `pipefail` (dash, i.e.
2143    /// most Linux camps) the probe must degrade to plain POSIX semantics — it
2144    /// must NOT take the shell down with it, because `set` is a special builtin
2145    /// and a bare `set -o pipefail` there is entitled to exit before the
2146    /// caller's command runs at all. Asserting the command still produces its
2147    /// output is asserting exactly that.
2148    #[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        // The probe itself must be silent — it runs on every relocated build.
2171        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    // ── Signal handling ───────────────────────────────────────────────────────
2233
2234    #[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        // Give the process a moment to start.
2250        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        // Wait for natural exit.
2288        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        // Kill on a completed run should return NotFound.
2301        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    // ── Stdin relay ───────────────────────────────────────────────────────────
2309
2310    #[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        // Shell that reads a line from stdin and echoes it back.
2318        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    /// `resize_run` must change the geometry the *child* sees, not just the
2355    /// master fd — so the assertion reads `stty size` from inside the PTY
2356    /// after the resize rather than inspecting the driver's own state.
2357    #[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        // Wait for a line on stdin, then report the geometry as of that moment.
2364        let id = driver
2365            .spawn_run(
2366                "read line && stty size",
2367                SpawnOpts {
2368                    cwd: "/tmp".into(),
2369                    stdin_enabled: true,
2370                    // Spawn at the default 80x24 so the assertion can't pass by
2371                    // accident if the resize is a no-op.
2372                    ..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    /// A run that is not active on this driver (finished, or never existed) is
2404    /// `NotFound` rather than a panic — same contract as `send_stdin`.
2405    #[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    // ── Tier-2 side-channel log fd ────────────────────────────────────────────
2435
2436    /// Verify that a child writing a JSON-line to `YAH_LOG_PIPE` (via
2437    /// `printf ... >> $YAH_LOG_PIPE`) produces a shim event with the correct
2438    /// fields in the store.
2439    ///
2440    /// The child opens the FIFO path for writing — no fd inheritance needed.
2441    #[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        // The shell writes one JSON-line to the FIFO by redirecting printf
2451        // output to the path stored in YAH_LOG_PIPE.
2452        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        // Wait for run completion. Deadline is generous because parallel-test
2460        // load + the rt.block_on hops from the reader/log threads can slow
2461        // child-process scheduling.
2462        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        // The log receiver thread drains after the lifecycle task drops the
2475        // write-end FdCloser; give it a brief moment.
2476        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    /// When `log_fd_enabled` is false, neither `YAH_TASK_RUN` nor
2496    /// `YAH_LOG_PIPE` are exported, and no shim events are written.
2497    #[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        // Try to write to YAH_LOG_PIPE; the conditional guards against
2507        // the variable being absent, so the command always exits 0.
2508        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    // ── Unattached-run reaper (R739-B12) ─────────────────────────────────────
2541
2542    /// The origin `yah build run` tags its relocated builds with. Spelled out
2543    /// here rather than imported: what the reaper must do is defined by the
2544    /// string on the wire, not by any constant this crate owns.
2545    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    /// The orphan this ticket exists for: the client is gone, so nothing polls
2584    /// the run, so the daemon must end it.
2585    #[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    /// The regression that protects a healthy long build: a client that is
2610    /// still polling keeps its run alive however long the build takes.
2611    #[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        // Six polls across three idle windows — what `yah build run`'s tail
2620        // loop does, slowed down.
2621        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        // And the moment the polling stops, it becomes reapable — same run,
2639        // so this pins the refresh rather than a missing origin match.
2640        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    /// The regression that protects real users' terminals. An interactive tile
2648    /// sits unpolled for hours by design and must never be touched, however
2649    /// long the reaper runs.
2650    #[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    /// A run with no origin at all — an ordinary `task.run` job — is outside
2680    /// every opt-in list, and an empty list reaps nothing.
2681    #[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    /// `note_attached` is observable, and the age it resets is what the sweep
2715    /// reads.
2716    #[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}