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//! @yah:ticket(R617-F6, "Reattach-by-run_id replaces Lost-on-disappear for origin=terminal shells")
41//! @yah:status(review)
42//! @yah:assignee(agent:bundle-anthropic-ashguard)
43//! @yah:at(2026-07-24T01:26:41Z)
44//! @yah:phase(P3)
45//! @yah:parent(R617)
46//! @arch:see(.yah/docs/working/W280-durable-terminal-sessions.md)
47//! @yah:depends_on(R617-F13)
48//! @yah:handoff("DELIVERED. Verified: `cd oss/qed && cargo test -p task-runs --lib` 243/243 (was 237 — 6 new); `cargo test -p kg-daemon --lib shell_vt` 9/9; `cargo test -p yah --lib r617` 9/9; `cargo test -p desktop --lib` 357 pass / 2 fail, both pre-existing and in files this ticket does not touch (agent.rs rules-view expects 12 rows and a peer's approval-rule change makes 19; agent_process reader-finished is a known timing flake).")
49//! @yah:handoff("THE TICKET'S OWN FRAMING WAS WRONG ABOUT THE MECHANISM, and the correction is the design. @yah:next said to 're-adopt' a live shell by 'control channel rebuilt, reader thread restarted against the surviving PTY'. That is not possible and never was: you cannot re-open another process's PTY master fd. The real defect is narrower and worse — a driver was tombstoning runs IT DID NOT OWN. `.yah/db/task-runs.turso` has several writers (desktop, the R617-F13 shell host, one CampService per MCP sidecar), and `TaskDriver::new` assumed any leftover `Running` row must be its own predecessor's corpse. So every attach marked some other LIVE process's shell `Lost`, and that shell kept producing output under a status saying it was dead. The fix is therefore 'do not tombstone what you do not own', not 'reattach'. Actual PTY reattach is unnecessary once F13 puts the PTY in a process that outlives the desktop.")
50//! @yah:handoff("HOW OWNERSHIP IS KNOWN: new `TaskRunMeta::host_pid` — the pid of the process whose driver spawned the run, NOT the child's. Stamped by `spawn_run` at INSERT, before the child exists, so a crash between insert and spawn still leaves the row attributable. Store column added by the same idempotent `ALTER TABLE ... ADD COLUMN` pattern `origin` used, and `row_to_meta` reads index 15 with `.ok().flatten()` so a DB with no such column reads `None` rather than erroring.")
51//! @yah:handoff("THE SEAM IS ORIGIN-AGNOSTIC, per this ticket's gotcha. New `task_runs::StaleRunPolicy` in oss/qed/crates/task-runs/src/driver.rs: `LostOnDisappear` (the default — `TaskDriver::new` and `with_channels` behave exactly as before, so no existing embedder changed) and `AdoptLiveHosts { origins: Vec<String> }`, which spares a leftover run only when its `host_pid` names a process that still exists. The crate decides on OWNERSHIP and takes the origin list as data — it never learns what 'terminal' means. New `TaskDriver::with_config` is the constructor that takes it.")
52//! @yah:handoff("yah side: `crates/yah/kg-daemon/src/service.rs::open_task_store` now passes `AdoptLiveHosts { origins: [ORIGIN_TERMINAL] }`. Also replaced the magic string — new `kg_daemon::shell_vt::ORIGIN_TERMINAL` now backs the two live `origin == \"terminal\"` gates in shell_vt.rs plus the policy, so the VT-parsing gate and the tombstone-exemption gate cannot drift apart by a typo. The constant lives on the yah side, NOT in task-runs, precisely to keep the crate generic.")
53//! @yah:handoff("Also stamped at app/yah/desktop/src/terminal.rs:519 — the desktop-local PTY path (terminal_open_local's scrollback mint) owns its own PTYs, so those rows carry the desktop's pid. Without it the shell host's driver would tombstone a live desktop-local session on attach, which is the same bug pointing the other way.")
54//! @yah:handoff("PID REUSE is the honest weakness and is why the policy is opt-in and origin-narrowed. `kill(pid, 0)` (EPERM counts as alive — the process exists, it is just not ours to signal) can read a recycled pid as the original owner. The failure mode of a false 'alive' is one run left `Running` until something closes it; the false 'dead' this replaces kills a live session's status. Strictly the better direction for an interactive shell, and the exposure is bounded to origins the embedder opted in. Non-unix has no kill(2), so `host_process_alive` reports false there and the platform keeps the old behaviour rather than stranding runs forever.")
55//! @yah:handoff("SIX NEW TESTS, each pinned to a failure rather than a code path: a live-owner terminal run survives a new driver (the ticket's whole point); a run whose owner pid was spawned and reaped in-test IS tombstoned (a crashed host must not leave zombie tiles); origin-less and non-matching origins are tombstoned even with a live owner (an in-flight `cargo build` whose driver is gone has nobody left to record its exit); an unattributed row (pre-migration) is tombstoned; `TaskDriver::new` still tombstones unconditionally (no silent behaviour change for existing embedders); and `spawn_run` stamps this process — the policy is worthless if rows arrive unattributed.")
56//! @yah:verify("cd oss/qed && cargo test -p task-runs --lib  # 243/243, 6 new under driver::tests")
57//! @yah:verify("cargo test -p kg-daemon --lib shell_vt  # 9/9")
58//! @yah:verify("cargo test -p yah --lib r617  # 9/9")
59//! @yah:verify("Manual (needs a desktop rebuild): open a shell, run `sleep 300`, quit and relaunch the desktop — the run is still Running, not Lost")
60//! @yah:verify("sqlite3 .yah/db/task-runs.turso \"select id, origin, host_pid, status from runs where status='running';\"  # every live row names a pid that ps shows")
61//! @yah:gotcha("This is an oss/qed crate — changes land in-tree under oss/qed/crates/task-runs and flow outward via scripts/export-oss.sh. The seam was kept origin-agnostic (StaleRunPolicy decides on host_pid, takes origins as data); the one yah-ism, ORIGIN_TERMINAL, lives in crates/yah/kg-daemon/src/shell_vt.rs instead.")
62//! @yah:gotcha("`host_pid` is NOT on the wire. rpc::WireRunMeta does not carry it, so a client cannot ask 'is this run's owner alive'. Nothing needs it today — the policy runs entirely daemon-side — but R617-F7 should check whether reattaching tiles want it before adding a second liveness notion of their own.")
63//! @yah:gotcha("pid reuse can make a dead owner read alive, leaving a run `Running` with nobody driving it. Bounded on purpose (opt-in + origin-narrowed) and strictly safer than the false-dead it replaces, but it is a real edge: if zombie terminal rows ever accumulate, this is why.")
64//! @yah:gotcha("TaskRunMeta gained a required field, so every struct-literal construction site had to be updated (velveteen-exec x4, scryer, task-runs fixtures, kg-daemon fixtures, desktop/terminal.rs x2). A new construction site added by anyone else will fail to compile until they pick a value — which is the intended forcing function: a run with no recorded owner is a run the policy has to tombstone.")
65//!
66//! @yah:ticket(R617-B9, "Pre-existing: task-runs log_pipe_events_land_in_store never completes (233 pass / 1 fail)")
67//! @yah:status(review)
68//! @yah:assignee(agent:bundle-anthropic-ashguard)
69//! @yah:at(2026-07-22T19:50:25Z)
70//! @yah:phase(P1)
71//! @yah:parent(R617)
72//! @yah:handoff("Root cause: not the FIFO, not the PTY. The whole pipeline completed correctly every time (child wrote the JSON line, receiver drained it, reader hit EOF, child.wait returned 0) — but the lifecycle's terminal `store.update_status` returned `Sql(Busy(\"database is locked\"))` and run_lifecycle swallowed it with `let _ =`, so the run stayed Running forever and the 20s poll deadline blew. A live run has three concurrent turso writers (PTY chunk appends, shim-FIFO event appends, lifecycle status) on independent connections with no busy handling at all.")
73//! @yah:handoff("Fix in oss/qed/crates/task-runs/src/store.rs: (1) `conn()` now sets `busy_timeout(5s)` on every connection; (2) new `exec_retry()` wraps writes in an outer exponential-backoff retry on the `Busy`/`BusySnapshot` class, because turso caps its internal backoff and then hands `Busy` back; (3) insert_run / update_status / update_beholder_status / append_chunk / append_event all routed through it.")
74//! @yah:handoff("driver.rs run_lifecycle no longer swallows the terminal status write — a genuine failure after retries now prints `[yah task-runs] failed to record terminal status for run <id>`, matching the crate's existing eprintln convention.")
75//! @yah:handoff("New regression test store.rs::concurrent_writers_do_not_lose_the_terminal_status — two background tasks hammer append_chunk/append_event while update_status lands. Verified it has teeth: with busy_timeout and the retry disabled it fails 3/3 with the exact `Busy(\"database is locked\")`; with them it passes 5/5.")
76//! @yah:verify("cd oss/qed && cargo test -p task-runs --lib — 237 passed / 0 failed (was 235 pass / 1 fail)")
77//! @yah:verify("log_pipe_events_land_in_store run 8x sequentially: 8/8 green in ~0.58s each. Before the fix the same loop was 11/12 red at the 20s timeout.")
78//!
79//! @yah:ticket(R652-T6, "Login shell: when cmd is the resolved shell, exec it directly (not sh -c) with -l")
80//! @yah:at(2026-08-02T00:03:08Z)
81//! @yah:status(review)
82//! @yah:assignee(agent:bundle-ollama-cloud-boulder)
83//! @yah:phase(P1)
84//! @yah:parent(R652)
85//! @yah:handoff("Login shells now exec directly with -l instead of going through sh -c. SpawnOpts (oss/qed/crates/task-runs/src/driver.rs) gained `argv: Option<Vec<String>>`: when set, spawn_run builds the CommandBuilder from that argv verbatim instead of wrapping `cmd` in `sh -c`. camp-service task_run sets it to [resolved_shell, \"-l\"] whenever the request is a shell request.")
86//! @yah:handoff("Why an argv escape hatch rather than a `login_shell: bool` flag in the driver: task-runs is an oss/qed crate and has no business knowing what a login shell is. The caller names the exact process; the driver just execs it. This also made R652-T4 a two-line addition rather than a second flag.")
87//! @yah:handoff("Three things this fixes beyond .zprofile finally running. (1) `sh -c \"zsh -l\"` left an inert `sh` as the PTY's foreground process group leader, so job control misbehaved and signals went to the wrong process. (2) That same inert sh is what the foreground-pid cwd probe (R652-T2) would have reported for, so T2 could not have worked without this. (3) -l is now a real argv element instead of text inside a shell string, so no quoting layer can eat it.")
88//! @yah:handoff("`cmd` is still what lands on TaskRunMeta.command, so a shell run reads back as \"$SHELL\" -- the rail label and the history re-run path both keep working. Beholder argv rewriting is bypassed when argv is set (the attach runs with BeholderSelect::None): the rewritten argv would be discarded on that path, so recording a `rewrite=...` that never happened would be a lie in the run metadata.")
89//! @yah:handoff("An empty argv falls back to the sh -c path rather than spawning nothing -- a caller bug should not become an exec of the empty string.")
90//! @yah:verify("cd oss/qed && cargo test -p task-runs --lib  # 246/246 green (3 new: explicit_argv_execs_the_program_directly, explicit_argv_still_records_the_requested_command, empty_argv_falls_back_to_the_shell_path)")
91//! @yah:verify("Manual (needs desktop rebuild): add `echo W289-login-test >> /tmp/w289.log` to ~/.zprofile, open a shell tile, confirm the file gets a line")
92//! @yah:gotcha("driver.rs is an oss/qed crate -- this lands in-tree under oss/qed/crates/task-runs and flows outward via scripts/export-oss.sh on the next release. SpawnOpts gained a field, but every in-tree construction site uses ..Default::default(), so nothing else needed touching.")
93
94use std::collections::HashMap;
95use std::io::Read;
96use std::path::PathBuf;
97use std::sync::{Arc, Mutex};
98use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
99
100use portable_pty::{native_pty_system, CommandBuilder, PtySize};
101use thiserror::Error;
102use tokio::sync::{mpsc, oneshot};
103use tokio::task;
104
105use crate::beholders::{registry_with_user_beholders, BeholderSelect};
106use crate::store::{RunFilter, StoreError, TaskStore};
107use crate::types::{BeholderStatus, Initiator, OutputChunk, RunStatus, Stream, TaskRunId, TaskRunMeta};
108
109const DEFAULT_GRACE: Duration = Duration::from_secs(5);
110const READ_BUF_SIZE: usize = 4096;
111const SIGTERM: i32 = 15;
112const SIGKILL: i32 = 9;
113
114// ─── Error ────────────────────────────────────────────────────────────────────
115
116#[derive(Debug, Error)]
117pub enum DriverError {
118    #[error("store: {0}")]
119    Store(#[from] StoreError),
120    #[error("pty: {0}")]
121    Pty(String),
122    #[error("run not found: {0}")]
123    NotFound(String),
124    #[error("io: {0}")]
125    Io(#[from] std::io::Error),
126}
127
128// ─── SpawnOpts ────────────────────────────────────────────────────────────────
129
130/// Options for [`TaskDriver::spawn_run`].
131#[derive(Debug, Clone)]
132pub struct SpawnOpts {
133    pub cwd: PathBuf,
134    /// Env vars set on the child process (merged on top of the current env).
135    pub env: Vec<(String, String)>,
136    pub label: Option<String>,
137    pub initiator: Initiator,
138    /// PTY column count. Defaults to 80.
139    pub pty_cols: u16,
140    /// PTY row count. Defaults to 24.
141    pub pty_rows: u16,
142    /// Enable stdin relay via [`TaskDriver::send_stdin`].
143    pub stdin_enabled: bool,
144    /// Pin the run so the GC sweep does not drop its output during warm rolloff.
145    pub pin: bool,
146    /// Beholder attachment policy. Defaults to [`BeholderSelect::Auto`].
147    pub beholder_select: BeholderSelect,
148    /// `true` when this run's output is read as text by somebody downstream —
149    /// a human watching a terminal tile, or a client that promised its caller
150    /// byte-identical passthrough. Causes `Rewriter` beholders to decline in
151    /// `Auto` mode, since a rewrite changes what that reader gets.
152    ///
153    /// R739-B10 renamed this from `tty_attached`: a PTY was only ever a proxy
154    /// for "someone is reading this", and the proxy broke the moment
155    /// `yah build run` moved onto pipes (R739-F6).
156    pub verbatim_output: bool,
157    /// Create a side-channel FIFO and export `YAH_TASK_RUN` / `YAH_LOG_PIPE`
158    /// so Tier-2 shim libraries (yah-log-rust, @yah/log) can emit structured
159    /// events. Has no effect on non-Unix platforms. Defaults to `true`.
160    pub log_fd_enabled: bool,
161    /// Provenance tag stored on the run's `TaskRunMeta.origin` (e.g.
162    /// `Some("terminal")` for an interactive shell). `None` is an ordinary job.
163    pub origin: Option<String>,
164    /// Exec this argv directly instead of wrapping `cmd` in `sh -c`.
165    ///
166    /// The default `sh -c <cmd>` is right for a job — the caller wrote a
167    /// command line and expects a shell to parse it. It is wrong for an
168    /// *interactive shell*: `sh -c "zsh -l"` leaves an inert `sh` as the PTY's
169    /// foreground process group leader, so job control misbehaves, signals go
170    /// to the wrong process, and anything that reads the foreground pid (a
171    /// live-cwd probe, say) sees `sh` instead of the shell the operator is
172    /// typing into. Handing the exact argv here makes the shell itself the
173    /// child, which is also the only way to pass `-l` as a real argv element
174    /// so `.zprofile` / `.profile` actually run.
175    ///
176    /// `cmd` is still what gets recorded on `TaskRunMeta.command`, so the run
177    /// reads the way the caller asked for it. Beholder argv rewriting is
178    /// bypassed when this is set: the caller has already decided the exact
179    /// process to exec, and a recorded `rewrite=…` that didn't happen would be
180    /// a lie in the run metadata.
181    pub argv: Option<Vec<String>>,
182    /// R739-F6 — spawn on **pipes** instead of a PTY. Defaults to `false`,
183    /// which is the PTY behaviour every existing caller already has.
184    ///
185    /// A PTY is right for an interactive terminal tile: the child gets a
186    /// controlling terminal, job control works, and `isatty` says yes, which is
187    /// what a human sitting in front of it expects. It is wrong for *emulating
188    /// a non-interactive shell invocation*, where three PTY properties show up
189    /// as divergence from running the same command directly (all three measured
190    /// in R739-F4 against `cargo check`):
191    ///
192    /// 1. `isatty(1)` is true, so tools colorize — plain `error: …` arrives as
193    ///    `\x1b[1m\x1b[91merror\x1b[0m: …`, which also defeats `| rg "^error"`.
194    /// 2. The line discipline's `ONLCR` rewrites every `\n` the child wrote
195    ///    into `\r\n`.
196    /// 3. The terminal merges stderr into stdout, so stream separation is gone
197    ///    by the time anything reads the capture.
198    ///
199    /// In pipe mode the child gets `pipe(2)` for stdout and stderr, chunks are
200    /// stored under their true [`Stream`], and `TERM` is left alone rather than
201    /// forced to `xterm-256color`. [`TaskDriver::resize_run`] and
202    /// [`TaskDriver::foreground_pid`] have no PTY to answer for and report
203    /// `NotFound` / `None`.
204    pub pipe: bool,
205    /// R901-B2 — run the `sh -c` line with `pipefail`, so a pipeline reports
206    /// the **leftmost** failing stage instead of its last one. Defaults to
207    /// `false`, i.e. POSIX behaviour, which is what every existing caller has.
208    ///
209    /// Without it a pipeline's status is the last stage's and nothing else:
210    /// `cargo check 2>&1 | tail -40` exits **0** on a build with 101 errors,
211    /// because `tail` succeeded. That is not a wrapper lying — the wrapper is
212    /// faithful, and the shell is answering the question it was actually
213    /// asked — but it is indistinguishable from a green build to everything
214    /// downstream, including the harness task notification an agent reads to
215    /// decide whether it is done. On 2026-09-13 two sessions read that 0 as a
216    /// pass and left `cargo check -p yah` red camp-wide for ~50 minutes.
217    ///
218    /// Only meaningful when [`SpawnOpts::argv`] is `None`; an explicit argv is
219    /// not a shell line and has no pipeline to take a status from.
220    ///
221    /// # Known cost, accepted deliberately
222    ///
223    /// `pipefail` also surfaces a producer killed by `SIGPIPE`, so
224    /// `cargo check 2>&1 | head -40` can now report failure once `head` closes
225    /// the pipe early on a build that was fine. That is a false RED, and it is
226    /// the right trade against the false GREEN above: a red is investigated,
227    /// a green ends the turn. Prefer `| tail` over `| head` on a build line.
228    pub pipefail: bool,
229}
230
231/// Prefix that turns `pipefail` on for the rest of a `sh -c` line.
232///
233/// Probing in a subshell rather than running `set -o pipefail` directly is
234/// load-bearing for portability, not caution. `pipefail` is a bash/ksh/zsh
235/// option; `/bin/sh` is bash on macOS but **dash** on most Linux distros, and
236/// dash rejects it. `set` is a POSIX *special* builtin, so a failure in one is
237/// entitled to terminate a non-interactive shell — which would turn "your
238/// pipeline now reports the truth" into "your command never ran at all" on
239/// every Linux camp. The subshell absorbs that exit; the outer shell only ever
240/// runs `set -o pipefail` on a shell that has already proved it accepts it.
241const PIPEFAIL_PRELUDE: &str = "if (set -o pipefail) 2>/dev/null; then set -o pipefail; fi\n";
242
243impl Default for SpawnOpts {
244    fn default() -> Self {
245        Self {
246            cwd: std::env::current_dir().unwrap_or_else(|_| PathBuf::from("/")),
247            env: vec![],
248            label: None,
249            initiator: Initiator::Human { camp: "local".to_string() },
250            pty_cols: 80,
251            pty_rows: 24,
252            stdin_enabled: false,
253            pin: false,
254            beholder_select: BeholderSelect::Auto,
255            verbatim_output: false,
256            log_fd_enabled: true,
257            origin: None,
258            argv: None,
259            pipe: false,
260            pipefail: false,
261        }
262    }
263}
264
265// ─── Driver channels ─────────────────────────────────────────────────────────
266
267/// Optional side-channels a driver can publish to. Both are fire-and-forget:
268/// a closed receiver never stalls or fails a run.
269#[derive(Default)]
270pub struct DriverChannels {
271    /// Fires `(run_id, status)` after each run's lifecycle task writes the
272    /// terminal status. Drives completion listeners (e.g. a triage worker).
273    pub completion: Option<mpsc::UnboundedSender<(TaskRunId, RunStatus)>>,
274    /// Mirrors every PTY output chunk as it is captured, *before* any consumer
275    /// polls the store. Lets a host attach a live view (VT parser, log
276    /// forwarder) to a run without a read-back loop over the store.
277    ///
278    /// The driver deliberately stays ignorant of what the tap is for — the
279    /// chunk carries `run_id`, so the host decides which runs it cares about.
280    pub output: Option<mpsc::UnboundedSender<OutputChunk>>,
281}
282
283// ─── Stale-run policy ────────────────────────────────────────────────────────
284
285/// What a freshly-constructed [`TaskDriver`] does with `Running` rows it finds
286/// already in the store.
287///
288/// The historical rule — tombstone every one of them — bakes in an assumption
289/// that stops being true the moment a second process attaches to the same
290/// store: that any `Running` row must be a corpse from *this* process's
291/// predecessor. When two processes share a store, a driver starting up in one
292/// will happily mark the other's live runs `Lost`, and the run keeps producing
293/// output under a status that says it is dead.
294///
295/// The policy is deliberately origin-agnostic in its mechanism — it decides on
296/// **who owns the run** ([`TaskRunMeta::host_pid`]) — and takes the origin list
297/// as data, so an embedder names the runs it wants exempted without this crate
298/// knowing what any of them mean.
299#[derive(Debug, Clone, Default, PartialEq, Eq)]
300pub enum StaleRunPolicy {
301    /// Tombstone every leftover `Running` run as `Lost`.
302    ///
303    /// Correct, and the default, whenever this process is the only writer:
304    /// a run whose driver is gone has no one left to notice it exit.
305    #[default]
306    LostOnDisappear,
307    /// Spare runs whose recorded owner process is still alive.
308    ///
309    /// A leftover run is tombstoned only when its `host_pid` is absent (owner
310    /// unknown — a row from before the column existed) or names a process that
311    /// no longer exists. Anything else belongs to a live peer and is left
312    /// `Running` for that peer to finish.
313    ///
314    /// `origins` narrows the exemption to runs whose
315    /// [`TaskRunMeta::origin`] is in the list; empty means every origin
316    /// qualifies. A run with no origin never matches a non-empty list.
317    AdoptLiveHosts { origins: Vec<String> },
318}
319
320impl StaleRunPolicy {
321    /// Whether `meta` should be tombstoned `Lost` at driver construction.
322    fn tombstones(&self, meta: &TaskRunMeta) -> bool {
323        match self {
324            StaleRunPolicy::LostOnDisappear => true,
325            StaleRunPolicy::AdoptLiveHosts { origins } => {
326                let exempt_origin = origins.is_empty()
327                    || meta
328                        .origin
329                        .as_deref()
330                        .is_some_and(|o| origins.iter().any(|want| want == o));
331                if !exempt_origin {
332                    return true;
333                }
334                match meta.host_pid {
335                    Some(pid) => !host_process_alive(pid),
336                    None => true,
337                }
338            }
339        }
340    }
341}
342
343/// Is a process with this pid still around?
344///
345/// `kill(pid, 0)` is the portable liveness probe: it performs the permission
346/// check and existence lookup without delivering anything. `EPERM` counts as
347/// alive — the process exists, it just is not ours to signal.
348///
349/// Pid reuse can make a dead owner read as alive. That is why
350/// [`StaleRunPolicy::AdoptLiveHosts`] is opt-in and origin-narrowed: the cost
351/// of a false "alive" is one run left `Running` until something closes it,
352/// which is strictly better for an interactive session than the false "dead"
353/// this replaces — which kills a *live* session's status.
354#[cfg(unix)]
355fn host_process_alive(pid: u32) -> bool {
356    if pid == 0 {
357        return false;
358    }
359    if pid == std::process::id() {
360        return true;
361    }
362    // SAFETY: `kill` with signal 0 delivers nothing; it only reports whether
363    // the pid exists and is signallable.
364    let rc = unsafe { libc::kill(pid as libc::pid_t, 0) };
365    rc == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
366}
367
368/// No `kill(2)` off Unix. Reporting every owner dead keeps the historical
369/// Lost-on-disappear behaviour rather than stranding runs `Running` forever.
370#[cfg(not(unix))]
371fn host_process_alive(_pid: u32) -> bool {
372    false
373}
374
375// ─── Internal run-control handle ─────────────────────────────────────────────
376
377struct RunControl {
378    kill_tx: mpsc::Sender<KillRequest>,
379    stdin_tx: Option<mpsc::Sender<Vec<u8>>>,
380    /// Shared with the lifecycle task, which holds the same `Arc` so the PTY fd
381    /// outlives `child.wait()`. `MasterPty::resize` takes `&self`, so a mutex is
382    /// enough to make the `Box<dyn MasterPty + Send>` `Sync` across the two.
383    ///
384    /// `None` for a [`SpawnOpts::pipe`] run, which has no terminal to resize or
385    /// to ask for a foreground process group.
386    master: Option<Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>>>,
387    /// R739-B12 — the run's [`SpawnOpts::origin`], copied here so
388    /// [`TaskDriver::reap_unattached`] can narrow to an opted-in origin set
389    /// without a store round-trip per candidate.
390    origin: Option<String>,
391    /// R739-B12 — when a client last looked at this run.
392    ///
393    /// Set at spawn (the caller that asked for the run is attached to it by
394    /// definition) and refreshed by [`TaskDriver::note_attached`], which the
395    /// embedder calls from whatever its "a client is watching" surface is —
396    /// for the camp daemon, `task.tail` and `task.status`.
397    ///
398    /// Monotonic rather than a wall clock: a clock step must not be able to
399    /// make a healthy build look abandoned.
400    last_attached_at: Instant,
401}
402
403/// Fires the reader-done signal when the LAST holder drops.
404///
405/// A PTY run has one reader; a piped run has two (stdout and stderr) and the
406/// lifecycle must not reap the child until both have hit EOF. Making this a
407/// drop guard behind an `Arc` means neither path has to count readers: the
408/// signal goes out when the refcount reaches zero, after each pump has finished
409/// its own `on_done` work.
410struct ReaderDone(Option<oneshot::Sender<()>>);
411
412impl Drop for ReaderDone {
413    fn drop(&mut self) {
414        if let Some(tx) = self.0.take() {
415            let _ = tx.send(());
416        }
417    }
418}
419
420#[derive(Debug)]
421struct KillRequest {
422    signal: i32,
423}
424
425// ─── ShimRecord ───────────────────────────────────────────────────────────────
426
427/// One JSON-line record emitted by a Tier-2 shim to the side-channel FIFO.
428///
429/// The shim (Rust `yah-log` layer or TS `@yah/log` pino transport) writes one
430/// of these per log call. Unknown keys inside `fields` pass through unchanged.
431#[cfg(unix)]
432#[derive(serde::Deserialize)]
433struct ShimRecord {
434    level: String,
435    target: String,
436    msg: String,
437    #[serde(default)]
438    fields: serde_json::Value,
439    /// Shim library name, e.g. `"yah-log-rust"`. Populates
440    /// [`EventSource::Shim::lib`].
441    #[serde(rename = "_lib", default)]
442    lib: Option<String>,
443    /// Shim library version string.
444    #[serde(rename = "_lib_ver", default)]
445    lib_version: Option<String>,
446}
447
448// ─── FdCloser ─────────────────────────────────────────────────────────────────
449
450/// RAII wrapper that closes a raw fd on drop.
451///
452/// Used to hold the write end of the log FIFO open until the lifecycle task
453/// completes. Dropping it signals EOF to the receiver thread.
454#[cfg(unix)]
455struct FdCloser(libc::c_int);
456
457#[cfg(unix)]
458impl Drop for FdCloser {
459    fn drop(&mut self) {
460        unsafe { libc::close(self.0) };
461    }
462}
463
464// SAFETY: a raw fd number is an integer; closing it from any thread is safe
465// provided we never duplicate ownership (enforced by move semantics here).
466#[cfg(unix)]
467unsafe impl Send for FdCloser {}
468
469// ─── TaskDriver ───────────────────────────────────────────────────────────────
470
471/// Manages in-flight task runs for a single camp.
472///
473/// Wrap in `Arc` to share across tasks; internal state is mutex-protected.
474pub struct TaskDriver {
475    store: Arc<TaskStore>,
476    active: Arc<Mutex<HashMap<String, RunControl>>>,
477    /// Side-channels published to by every run this driver owns.
478    channels: DriverChannels,
479}
480
481impl TaskDriver {
482    /// Create a driver backed by `store`, with no side-channels.
483    ///
484    /// Immediately scans the store for `Running` runs left over from a prior
485    /// daemon process and marks them `Lost` ("Lost-on-disappear").
486    pub async fn new(store: Arc<TaskStore>) -> Result<Self, DriverError> {
487        Self::with_channels(store, DriverChannels::default()).await
488    }
489
490    /// Like `new` but wires the optional [`DriverChannels`] side-channels
491    /// (completion notifications, live output tap).
492    pub async fn with_channels(
493        store: Arc<TaskStore>,
494        channels: DriverChannels,
495    ) -> Result<Self, DriverError> {
496        Self::with_config(store, channels, StaleRunPolicy::default()).await
497    }
498
499    /// Full constructor: side-channels plus the [`StaleRunPolicy`] applied to
500    /// `Running` rows already in the store.
501    ///
502    /// R617-F6 — annotation in this file's header. Splitting the sweep out of
503    /// the constructor's fixed behaviour is what lets a store be shared: a
504    /// process that is not the run's owner can now attach without declaring
505    /// the owner's live work dead.
506    pub async fn with_config(
507        store: Arc<TaskStore>,
508        channels: DriverChannels,
509        stale_policy: StaleRunPolicy,
510    ) -> Result<Self, DriverError> {
511        let stale = store
512            .list_runs(&RunFilter {
513                status: Some("running".to_string()),
514                ..Default::default()
515            })
516            .await?;
517        for meta in stale {
518            if !stale_policy.tombstones(&meta) {
519                continue;
520            }
521            let status = RunStatus::Lost {
522                reason: "daemon restarted while run was in-flight".to_string(),
523            };
524            store.update_status(&meta.id, &status).await?;
525            /* A sweep tombstone is a terminal status like any other, and a
526               client watching that run needs it as much as it needs an exit
527               code. Until R267-T12 this was the one status transition in the
528               driver that wrote the store and told nobody — which is precisely
529               the transition a client cannot discover on its own, because the
530               process that would have reported it is the one that died. */
531            if let Some(ref tx) = channels.completion {
532                let _ = tx.send((meta.id.clone(), status));
533            }
534        }
535        Ok(Self {
536            store,
537            active: Arc::new(Mutex::new(HashMap::new())),
538            channels,
539        })
540    }
541
542    /// Spawn `cmd` in a PTY and start capturing its output. Returns immediately
543    /// with the new [`TaskRunId`].
544    ///
545    /// A beholder is selected via `opts.beholder_select` (default `Auto`). When
546    /// a `Rewriter` beholder matches, its `adjust_argv` is applied to the
547    /// command before spawning and the diff is recorded on `beholder_status`.
548    /// When `opts.verbatim_output` is `true`, `Rewriter` beholders decline in
549    /// `Auto` mode, because something downstream renders these bytes and a
550    /// rewrite would change them.
551    ///
552    /// Output is written to the store as `Stream::Stdout` chunks (the PTY
553    /// kernel merges stdout and stderr). Signal handling and status updates
554    /// run in background tasks.
555    pub async fn spawn_run(&self, cmd: &str, opts: SpawnOpts) -> Result<TaskRunId, DriverError> {
556        let id = TaskRunId::new();
557        let started_at = unix_now_secs();
558        let started_at_ms: u64 = started_at.saturating_mul(1000);
559
560        // Attach a beholder (may rewrite argv and produce structured events).
561        // Resolve user drop-in directory: $YAH_BEHOLDERS_DIR or $HOME/.yah/beholders.
562        let user_dir = std::env::var_os("YAH_BEHOLDERS_DIR")
563            .map(std::path::PathBuf::from)
564            .or_else(|| {
565                std::env::var_os("HOME")
566                    .map(|h| std::path::PathBuf::from(h).join(".yah/beholders"))
567            });
568        let registry = registry_with_user_beholders(user_dir.as_deref());
569        /* An explicit argv means the caller already chose the exact process
570           (an interactive login shell, say). Selecting a beholder there would
571           either do nothing — the rewritten argv is discarded on that path —
572           or record a rewrite that never happened, so we opt out honestly
573           instead. */
574        let select = if opts.argv.is_some() {
575            &BeholderSelect::None
576        } else {
577            &opts.beholder_select
578        };
579        let attach = registry.attach(cmd, select, opts.verbatim_output);
580        // Reconstruct the command from argv ONLY when a beholder actually
581        // rewrote it. `AttachResult.argv` is always populated — it is
582        // `resolve_argv(cmd)` even when nothing attached — so joining it
583        // unconditionally ran every run's command through a whitespace
584        // normalization nobody asked for: runs of spaces collapse and embedded
585        // newlines become spaces, which is silent corruption for a heredoc or
586        // any multi-line line. The caller's bytes go to the shell untouched
587        // unless a rewrite is the whole point.
588        let effective_cmd = match &attach.status.rewrite_added {
589            Some(added) if !added.is_empty() && !attach.argv.is_empty() => attach.argv.join(" "),
590            _ => cmd.to_string(),
591        };
592
593        self.store.insert_run(&TaskRunMeta {
594            id: id.clone(),
595            command: cmd.to_string(),
596            cwd: opts.cwd.clone(),
597            env: opts.env.clone(),
598            started_at,
599            status: RunStatus::Running,
600            label: opts.label.clone(),
601            initiator: opts.initiator.clone(),
602            beholder_status: Some(attach.status),
603            pinned: opts.pin,
604            origin: opts.origin.clone(),
605            /* R617-F6: stamp the OWNER, before the child exists. Written at
606               insert rather than after spawn so a crash between the two still
607               leaves the row attributable — an unattributed `Running` row is
608               exactly what the conservative arm of `StaleRunPolicy` has to
609               tombstone. */
610            host_pid: Some(std::process::id()),
611        }).await?;
612
613        // The program and argv both spawn modes exec. An explicit argv execs
614        // that program directly; otherwise the command line goes through `sh`
615        // so the caller's quoting, pipes and redirections mean what they say.
616        // An empty argv is a caller bug, not a request for an empty exec — fall
617        // back to the shell path rather than spawning nothing.
618        // R901-B2: `pipefail` is prepended HERE and not folded into
619        // `effective_cmd`, so `TaskRunMeta.command` keeps reading as the line
620        // the caller actually wrote. A run's recorded command is re-run by
621        // history and audited by agents against the relocation note; a prelude
622        // nobody asked for showing up in it would be the same class of lie as
623        // recording a beholder `rewrite=…` that never happened.
624        let (program, args): (String, Vec<String>) = match opts.argv.as_deref() {
625            Some([p, rest @ ..]) => (p.clone(), rest.to_vec()),
626            _ => {
627                let line = if opts.pipefail {
628                    format!("{PIPEFAIL_PRELUDE}{effective_cmd}")
629                } else {
630                    effective_cmd.clone()
631                };
632                ("sh".to_string(), vec!["-c".to_string(), line])
633            }
634        };
635
636        // ── Side-channel log FIFO (Tier 2 / yah-log shims) ──────────────────
637        //
638        // Create a named pipe (FIFO) so child processes can write structured
639        // events without touching stdout/stderr. We export its path via
640        // YAH_LOG_PIPE; no fd inheritance is involved, so portable-pty's
641        // close_random_fds() pre_exec hook doesn't interfere.
642        //
643        // The parent opens the FIFO twice:
644        //   rfd — O_RDONLY|O_NONBLOCK, then cleared to blocking → read events
645        //   wfd — O_WRONLY (wrapped in FdCloser) → keeps the FIFO alive until
646        //          the lifecycle task drops it (after run completion), producing
647        //          EOF for the receiver thread.
648        #[cfg(unix)]
649        let log_fifo: Option<(libc::c_int, FdCloser, std::path::PathBuf)> = if opts.log_fd_enabled {
650            let fifo_path = std::env::temp_dir().join(format!("yah-log-{}.fifo", id));
651            let path_cstr = match std::ffi::CString::new(fifo_path.to_string_lossy().as_bytes()) {
652                Ok(s) => s,
653                Err(_) => {
654                    // Path contained a nul byte — extremely unlikely; skip FIFO.
655                    return Err(DriverError::Io(std::io::Error::new(
656                        std::io::ErrorKind::InvalidInput,
657                        "log FIFO path contained nul byte",
658                    )));
659                }
660            };
661            let mkfifo_ret = unsafe { libc::mkfifo(path_cstr.as_ptr(), 0o600) };
662            if mkfifo_ret != 0 {
663                None // FIFO creation failed; continue without side-channel
664            } else {
665                // Open read end without blocking (no writer yet).
666                let rfd = unsafe {
667                    libc::open(path_cstr.as_ptr(), libc::O_RDONLY | libc::O_NONBLOCK)
668                };
669                if rfd < 0 {
670                    let _ = unsafe { libc::unlink(path_cstr.as_ptr()) };
671                    None
672                } else {
673                    // Switch read end to blocking so reads yield proper data.
674                    unsafe { libc::fcntl(rfd, libc::F_SETFL, 0) };
675                    // Open write end — this succeeds immediately because rfd is open.
676                    let wfd = unsafe {
677                        libc::open(path_cstr.as_ptr(), libc::O_WRONLY)
678                    };
679                    if wfd < 0 {
680                        unsafe { libc::close(rfd) };
681                        let _ = unsafe { libc::unlink(path_cstr.as_ptr()) };
682                        None
683                    } else {
684                        Some((rfd, FdCloser(wfd), fifo_path))
685                    }
686                }
687            }
688        } else {
689            None
690        };
691
692        // The FIFO env, applied identically by both spawn modes.
693        #[cfg(unix)]
694        let fifo_env: Option<(String, String)> = log_fifo
695            .as_ref()
696            .map(|(_, _, path)| (id.to_string(), path.to_string_lossy().into_owned()));
697        #[cfg(not(unix))]
698        let fifo_env: Option<(String, String)> = None;
699
700        /* Spawn. The two modes differ only in what the child's stdio is
701           attached to, and everything downstream — reader pumps, lifecycle,
702           kill — is written against the uniform handles produced here:
703           `pid`, a `reap` closure that blocks until the child exits, and an
704           optional PTY master for resize / foreground-pid. */
705        let pid: u32;
706        let reap: Box<dyn FnOnce() -> Option<u32> + Send>;
707        let stdin_tx: Option<mpsc::Sender<Vec<u8>>>;
708        let master: Option<Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>>>;
709        // Each entry is one blocking source to pump into the store. The PTY
710        // yields a single merged stream; pipes yield stdout and stderr apart.
711        let mut sources: Vec<(Box<dyn Read + Send>, Stream)> = Vec::new();
712
713        if opts.pipe {
714            use std::process::{Command, Stdio};
715
716            let mut cmd = Command::new(&program);
717            cmd.args(&args);
718            cmd.current_dir(&opts.cwd);
719            for (k, v) in &opts.env {
720                cmd.env(k, v);
721            }
722            /* Deliberately NOT setting TERM. The PTY path forces
723               `xterm-256color` because a child on a terminal that claims no
724               terminal type degrades badly; a child on a pipe should see
725               whatever the daemon's own environment says, exactly as it would
726               under a non-interactive shell. Forcing a terminal type here is
727               how a pipe-mode run would talk itself back into colorizing. */
728            if let Some((run_id_env, fifo_path)) = &fifo_env {
729                cmd.env("YAH_TASK_RUN", run_id_env);
730                cmd.env("YAH_LOG_PIPE", fifo_path);
731            }
732            cmd.stdout(Stdio::piped());
733            cmd.stderr(Stdio::piped());
734            cmd.stdin(if opts.stdin_enabled { Stdio::piped() } else { Stdio::null() });
735
736            let mut child = cmd.spawn().map_err(DriverError::Io)?;
737            pid = child.id();
738
739            if let Some(out) = child.stdout.take() {
740                sources.push((Box::new(out), Stream::Stdout));
741            }
742            if let Some(err) = child.stderr.take() {
743                sources.push((Box::new(err), Stream::Stderr));
744            }
745
746            stdin_tx = child.stdin.take().map(|mut writer| {
747                let (tx, mut rx) = mpsc::channel::<Vec<u8>>(64);
748                task::spawn(async move {
749                    use std::io::Write;
750                    while let Some(bytes) = rx.recv().await {
751                        let _ = writer.write_all(&bytes);
752                        let _ = writer.flush();
753                    }
754                });
755                tx
756            });
757
758            master = None;
759            reap = Box::new(move || child.wait().ok().and_then(|s| s.code()).map(|c| c as u32));
760        } else {
761            // Open PTY pair.
762            let pty_sys = native_pty_system();
763            let pair = pty_sys
764                .openpty(PtySize {
765                    rows: opts.pty_rows,
766                    cols: opts.pty_cols,
767                    pixel_width: 0,
768                    pixel_height: 0,
769                })
770                .map_err(|e| DriverError::Pty(e.to_string()))?;
771
772            // Clone reader before spawning so the fd is ready immediately.
773            let pty_reader = pair
774                .master
775                .try_clone_reader()
776                .map_err(|e| DriverError::Pty(e.to_string()))?;
777            sources.push((Box::new(pty_reader), Stream::Stdout));
778
779            // Optional stdin relay: take the writer before spawning the child.
780            stdin_tx = if opts.stdin_enabled {
781                let mut writer = pair
782                    .master
783                    .take_writer()
784                    .map_err(|e| DriverError::Pty(e.to_string()))?;
785                let (tx, mut rx) = mpsc::channel::<Vec<u8>>(64);
786                task::spawn(async move {
787                    use std::io::Write;
788                    while let Some(bytes) = rx.recv().await {
789                        let _ = writer.write_all(&bytes);
790                        let _ = writer.flush();
791                    }
792                });
793                Some(tx)
794            } else {
795                None
796            };
797
798            let mut cb = CommandBuilder::new(&program);
799            cb.args(&args);
800            cb.cwd(&opts.cwd);
801            for (k, v) in &opts.env {
802                cb.env(k, v);
803            }
804            cb.env("TERM", "xterm-256color");
805            if let Some((run_id_env, fifo_path)) = &fifo_env {
806                cb.env("YAH_TASK_RUN", run_id_env);
807                cb.env("YAH_LOG_PIPE", fifo_path);
808            }
809
810            let child = pair
811                .slave
812                .spawn_command(cb)
813                .map_err(|e| DriverError::Pty(e.to_string()))?;
814            // Drop the parent's slave handle so EOF propagates once the child exits.
815            drop(pair.slave);
816
817            pid = child.process_id().unwrap_or(0);
818
819            // Share the master between the lifecycle task (which must outlive
820            // `child.wait()` so the fd stays open) and `resize_run`.
821            let m: Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>> =
822                Arc::new(Mutex::new(pair.master));
823            master = Some(Arc::clone(&m));
824            reap = Box::new(move || {
825                let mut c = child;
826                let _m = m; // dropped after wait() returns, closing the PTY fd
827                c.wait().ok().map(|s| s.exit_code())
828            });
829        }
830
831        // ── FIFO: launch receiver thread; pass write-end holder to lifecycle ──
832        //
833        // The receiver thread reads until EOF. EOF arrives when ALL write-end
834        // holders close: the child's own writers (when it exits) plus the
835        // FdCloser we hand to the lifecycle task (which drops it after writing
836        // the terminal RunStatus). Events written before the last close are
837        // still drained by the receiver thread before it exits.
838        #[cfg(unix)]
839        let log_wfd_holder: Option<FdCloser> = if let Some((rfd, wfd, fifo_path)) = log_fifo {
840            let store_log = Arc::clone(&self.store);
841            let id_log = id.clone();
842            let rt = tokio::runtime::Handle::current();
843            // spawn_blocking: lets the runtime track this thread so the
844            // Handle::block_on calls inside have a worker to drive futures.
845            tokio::task::spawn_blocking(move || {
846                run_log_receiver(rt, store_log, id_log, rfd, fifo_path, started_at_ms);
847            });
848            Some(wfd)
849        } else {
850            None
851        };
852
853        // Channels.
854        let (kill_tx, kill_rx) = mpsc::channel::<KillRequest>(4);
855        let (reader_done_tx, reader_done_rx) = oneshot::channel::<()>();
856
857        /* Reader threads: child output → store chunks → beholder events. Each
858           runs on a dedicated OS thread because the reads are blocking. The
859           `ReaderDone` guard is shared across them, so the lifecycle's
860           reader-done signal fires only once every source has hit EOF — which
861           is what makes the two-pipe case correct without a reader count. */
862        {
863            let done = Arc::new(ReaderDone(Some(reader_done_tx)));
864            /* The beholder goes to stdout only. It parses a structured
865               protocol (cargo's JSON, say) that the child writes to stdout by
866               definition, and there is exactly one of it — handing the same
867               instance to two threads would need a lock for no gain, and
868               feeding it stderr would make `unknown_format_reason` fire on
869               human-readable diagnostics it was never meant to see. */
870            let mut beholder = attach.beholder;
871            for (reader, stream) in sources {
872                spawn_output_pump(
873                    reader,
874                    stream,
875                    Arc::clone(&self.store),
876                    id.clone(),
877                    started_at_ms,
878                    self.channels.output.clone(),
879                    if stream == Stream::Stdout { beholder.take() } else { None },
880                    Arc::clone(&done),
881                );
882            }
883        }
884
885        // Lifecycle task: monitor kill requests, wait for exit, update status.
886        // The task also holds the log FIFO write-end closer (if any) so that
887        // EOF propagates to the receiver thread after RunStatus is written.
888        {
889            let store_l = Arc::clone(&self.store);
890            let active_l = Arc::clone(&self.active);
891            let id_l = id.clone();
892            let completion_tx_l = self.channels.completion.clone();
893            #[cfg(unix)]
894            let wfd_l = log_wfd_holder;
895            task::spawn(async move {
896                run_lifecycle(
897                    store_l,
898                    active_l,
899                    id_l,
900                    pid,
901                    reap,
902                    kill_rx,
903                    reader_done_rx,
904                    completion_tx_l,
905                    #[cfg(unix)]
906                    wfd_l,
907                )
908                .await;
909            });
910        }
911
912        self.active
913            .lock()
914            .unwrap()
915            .insert(
916                id.to_string(),
917                RunControl {
918                    kill_tx,
919                    stdin_tx,
920                    master,
921                    origin: opts.origin.clone(),
922                    last_attached_at: Instant::now(),
923                },
924            );
925
926        Ok(id)
927    }
928
929    /// Resize a running task's PTY and deliver `SIGWINCH` to the foreground
930    /// process group (portable-pty's `resize` does the ioctl, which is what
931    /// signals the child).
932    ///
933    /// Returns `DriverError::NotFound` when the run is not active on this
934    /// driver instance — the same contract as [`TaskDriver::send_stdin`] — and
935    /// also when it is active but was spawned in [`SpawnOpts::pipe`] mode, which
936    /// has no terminal to resize.
937    pub async fn resize_run(
938        &self,
939        id: &TaskRunId,
940        cols: u16,
941        rows: u16,
942    ) -> Result<(), DriverError> {
943        let master = self
944            .active
945            .lock()
946            .unwrap()
947            .get(&id.to_string())
948            .and_then(|c| c.master.as_ref().map(Arc::clone));
949
950        match master {
951            Some(m) => {
952                let size = PtySize { rows, cols, pixel_width: 0, pixel_height: 0 };
953                m.lock()
954                    .unwrap()
955                    .resize(size)
956                    .map_err(|e| DriverError::Pty(e.to_string()))
957            }
958            None => Err(DriverError::NotFound(id.to_string())),
959        }
960    }
961
962    /// The pid of the run's *foreground* process — the leader of the process
963    /// group the PTY currently gives the keyboard to.
964    ///
965    /// For a shell tile that is the shell itself while it sits at a prompt,
966    /// and the command the operator is running while one is in flight. That
967    /// distinction is the whole point: asking the spawned child would report
968    /// the shell forever, so anything derived from this pid (a live cwd probe,
969    /// a "what is this pane doing" label) would answer for the wrong process.
970    ///
971    /// `None` when the run is not active on this driver instance, when it was
972    /// spawned in [`SpawnOpts::pipe`] mode (no controlling terminal, so no
973    /// foreground process group to read), or when the platform has no notion of
974    /// a foreground process group.
975    pub fn foreground_pid(&self, id: &TaskRunId) -> Option<u32> {
976        let master = self
977            .active
978            .lock()
979            .unwrap()
980            .get(&id.to_string())
981            .and_then(|c| c.master.as_ref().map(Arc::clone))?;
982        #[cfg(unix)]
983        {
984            let pid = master.lock().unwrap().process_group_leader()?;
985            u32::try_from(pid).ok()
986        }
987        #[cfg(not(unix))]
988        {
989            let _ = master;
990            None
991        }
992    }
993
994    /// Send `signal` to a running task. Defaults to SIGTERM (15).
995    ///
996    /// For SIGTERM, the driver waits up to 5 seconds for the process to exit
997    /// before escalating to SIGKILL. Returns `DriverError::NotFound` if the
998    /// run is not active (already exited or launched on a different driver
999    /// instance).
1000    pub async fn kill_run(&self, id: &TaskRunId, signal: Option<i32>) -> Result<(), DriverError> {
1001        let kill_tx = self
1002            .active
1003            .lock()
1004            .unwrap()
1005            .get(&id.to_string())
1006            .map(|c| c.kill_tx.clone());
1007
1008        match kill_tx {
1009            Some(tx) => tx
1010                .send(KillRequest { signal: signal.unwrap_or(SIGTERM) })
1011                .await
1012                .map_err(|_| DriverError::NotFound(id.to_string())),
1013            None => Err(DriverError::NotFound(id.to_string())),
1014        }
1015    }
1016
1017    /// Write bytes to the stdin of a running task (requires `stdin_enabled`).
1018    pub async fn send_stdin(&self, id: &TaskRunId, bytes: Vec<u8>) -> Result<(), DriverError> {
1019        let stdin_tx = self
1020            .active
1021            .lock()
1022            .unwrap()
1023            .get(&id.to_string())
1024            .and_then(|c| c.stdin_tx.clone());
1025
1026        match stdin_tx {
1027            Some(tx) => tx
1028                .send(bytes)
1029                .await
1030                .map_err(|_| DriverError::NotFound(id.to_string())),
1031            None => Err(DriverError::NotFound(id.to_string())),
1032        }
1033    }
1034
1035    /// R739-B12 — record that a client just looked at this run.
1036    ///
1037    /// A no-op for a run this driver does not own (already finished, or
1038    /// spawned by another process against the same store): attachment only
1039    /// means anything for a run something here could still signal.
1040    pub fn note_attached(&self, id: &TaskRunId) {
1041        if let Some(control) = self.active.lock().unwrap().get_mut(&id.to_string()) {
1042            control.last_attached_at = Instant::now();
1043        }
1044    }
1045
1046    /// How long ago a client last looked at `id`, or `None` when this driver
1047    /// does not own the run. The observable half of [`Self::note_attached`].
1048    pub fn attached_age(&self, id: &TaskRunId) -> Option<Duration> {
1049        self.active
1050            .lock()
1051            .unwrap()
1052            .get(&id.to_string())
1053            .map(|c| c.last_attached_at.elapsed())
1054    }
1055
1056    /// R739-B12 — SIGTERM every run of an opted-in origin that no client has
1057    /// looked at for `idle`. Returns the runs it signalled.
1058    ///
1059    /// This exists because a run outlives the client that asked for it. When
1060    /// `yah build run` is SIGKILLed — its harness dies, the terminal goes away
1061    /// — the cargo it relocated into the daemon keeps compiling with nobody
1062    /// attached, holding the build-directory lock until a human finds the pid.
1063    /// That happened on 2026-08-28 and stalled a whole camp for ~30 minutes.
1064    /// R739-B9 closed every give-up the client is *alive* to make; this closes
1065    /// the one it is not.
1066    ///
1067    /// **Not [`StaleRunPolicy`], and not that policy on a timer.** The policy
1068    /// is a construction-time reconciliation of rows a *previous process*
1069    /// left behind: it decides on `host_pid`, only ever calls
1070    /// `store.update_status`, and tombstones any run outside its origin list
1071    /// outright — so running it periodically would mark every in-flight run of
1072    /// an un-adopted origin `Lost` while it compiles perfectly well, and would
1073    /// still never signal the process that is the actual problem. This is the
1074    /// opposite shape: it decides on *attachment*, it signals, and it touches
1075    /// nothing outside `origins`.
1076    ///
1077    /// `origins` is an opt-in list precisely because most runs must never be
1078    /// reaped on this rule. An interactive terminal tile is legitimately
1079    /// unpolled for hours, and killing one would be a far worse bug than the
1080    /// orphan this prevents — so an empty list reaps nothing at all, rather
1081    /// than meaning "every origin" the way [`StaleRunPolicy`]'s list does.
1082    pub async fn reap_unattached(&self, idle: Duration, origins: &[String]) -> Vec<TaskRunId> {
1083        if origins.is_empty() {
1084            return Vec::new();
1085        }
1086        let candidates: Vec<TaskRunId> = {
1087            let active = self.active.lock().unwrap();
1088            active
1089                .iter()
1090                .filter(|(_, c)| {
1091                    c.origin
1092                        .as_deref()
1093                        .is_some_and(|o| origins.iter().any(|want| want == o))
1094                        && c.last_attached_at.elapsed() >= idle
1095                })
1096                .filter_map(|(id, _)| id.parse::<TaskRunId>().ok())
1097                .collect()
1098        };
1099
1100        let mut reaped = Vec::new();
1101        for id in candidates {
1102            // SIGTERM, not SIGKILL: `kill_run` gives the child the same 5s
1103            // grace a `task.kill` from a live client would, then escalates.
1104            // A terminal status is also what releases the run's admission
1105            // enrollment (R739-F7), so a reaped run frees the build key.
1106            if self.kill_run(&id, None).await.is_ok() {
1107                reaped.push(id);
1108            }
1109        }
1110        reaped
1111    }
1112}
1113
1114// ─── Log fd receiver ─────────────────────────────────────────────────────────
1115
1116/// Read JSON-lines from the side-channel FIFO read end and store them as
1117/// [`EventSource::Shim`] events.
1118///
1119/// Runs on a dedicated OS thread; exits when the read end sees EOF. EOF
1120/// arrives after both the child process AND the lifecycle task have closed
1121/// their write ends of the FIFO. The FIFO file is deleted on exit.
1122#[cfg(unix)]
1123fn run_log_receiver(
1124    rt: tokio::runtime::Handle,
1125    store: Arc<TaskStore>,
1126    run_id: TaskRunId,
1127    read_fd: libc::c_int,
1128    fifo_path: std::path::PathBuf,
1129    started_at_ms: u64,
1130) {
1131    use std::io::BufRead;
1132    use std::os::unix::io::FromRawFd;
1133
1134    // SAFETY: `read_fd` is a valid, open FIFO fd handed exclusively to this
1135    // thread. `File` takes ownership and closes the fd on drop.
1136    let file = unsafe { std::fs::File::from_raw_fd(read_fd) };
1137    let reader = std::io::BufReader::new(file);
1138
1139    for line in reader.lines() {
1140        let line = match line {
1141            Ok(l) => l,
1142            Err(_) => break,
1143        };
1144        let trimmed = line.trim();
1145        if trimmed.is_empty() {
1146            continue;
1147        }
1148        let rec: ShimRecord = match serde_json::from_str(trimmed) {
1149            Ok(r) => r,
1150            Err(_) => continue, // skip malformed lines silently
1151        };
1152        let level = rec.level.parse::<crate::types::Level>().unwrap_or(crate::types::Level::Info);
1153        let source = crate::types::EventSource::Shim {
1154            lib: rec.lib.unwrap_or_else(|| "unknown".to_string()),
1155            version: rec.lib_version.unwrap_or_else(|| "0.0.0".to_string()),
1156        };
1157        let fields = if rec.fields.is_object() {
1158            rec.fields
1159        } else {
1160            serde_json::Value::Object(Default::default())
1161        };
1162        let offset = elapsed_ms(started_at_ms);
1163        let _ = rt.block_on(store.append_event(
1164            &run_id,
1165            offset,
1166            level,
1167            &rec.target,
1168            &rec.msg,
1169            &fields,
1170            None,
1171            &source,
1172        ));
1173    }
1174
1175    // Clean up the FIFO file now that the receiver has drained.
1176    let _ = std::fs::remove_file(&fifo_path);
1177}
1178
1179// ─── Lifecycle task ───────────────────────────────────────────────────────────
1180
1181/// Pump one blocking output source into the store, tapping and beholding on the
1182/// way past.
1183///
1184/// Split out of `spawn_run` for R739-F6: a PTY run has one source and a piped
1185/// run has two, and the only thing that differs between them is which [`Stream`]
1186/// the chunks are stored under. `done` is the shared [`ReaderDone`] guard —
1187/// dropping it here, after `on_done`, is what tells the lifecycle this source is
1188/// finished.
1189#[allow(clippy::too_many_arguments)]
1190fn spawn_output_pump(
1191    reader: Box<dyn Read + Send>,
1192    stream: Stream,
1193    store: Arc<TaskStore>,
1194    id: TaskRunId,
1195    started_at_ms: u64,
1196    output_tx: Option<mpsc::UnboundedSender<OutputChunk>>,
1197    beholder: Option<Box<dyn crate::beholders::Beholder>>,
1198    done: Arc<ReaderDone>,
1199) {
1200    let rt = tokio::runtime::Handle::current();
1201    tokio::task::spawn_blocking(move || {
1202        let _done = done;
1203        let mut beholder = beholder;
1204        let mut buf = [0u8; READ_BUF_SIZE];
1205        let mut reader = reader;
1206        loop {
1207            match reader.read(&mut buf) {
1208                Ok(0) | Err(_) => break,
1209                Ok(n) => {
1210                    let offset = elapsed_ms(started_at_ms);
1211                    let append_res =
1212                        rt.block_on(store.append_chunk(&id, offset, stream, &buf[..n]));
1213                    if let Ok(seq) = append_res {
1214                        /* Both the tap and the beholder want the same owned
1215                           chunk; build it once, and only when someone is
1216                           listening. */
1217                        let chunk = (output_tx.is_some() || beholder.is_some()).then(|| {
1218                            OutputChunk {
1219                                run_id: id.clone(),
1220                                seq,
1221                                offset_ms: offset,
1222                                stream,
1223                                bytes: buf[..n].to_vec(),
1224                            }
1225                        });
1226                        /* Tap first: it feeds live views, where latency is
1227                           visible to a human. Send failure means the host
1228                           dropped its receiver — never fatal. */
1229                        if let (Some(tx), Some(c)) = (&output_tx, &chunk) {
1230                            let _ = tx.send(c.clone());
1231                        }
1232                        let mut detach_beholder = false;
1233                        if let (Some(b), Some(chunk)) = (beholder.as_mut(), &chunk) {
1234                            for ev in b.parse_chunk(chunk) {
1235                                let _ = rt.block_on(store.append_event(
1236                                    &ev.run_id,
1237                                    ev.offset_ms,
1238                                    ev.level,
1239                                    &ev.target,
1240                                    &ev.msg,
1241                                    &ev.fields,
1242                                    ev.anchor.as_ref().map(|a| a.seq),
1243                                    &ev.source,
1244                                ));
1245                            }
1246                            if let Some(reason) = b.unknown_format_reason() {
1247                                let new_status =
1248                                    BeholderStatus::unknown_format_with_reason(b.name(), reason);
1249                                let _ =
1250                                    rt.block_on(store.update_beholder_status(&id, &new_status));
1251                                detach_beholder = true;
1252                            }
1253                        }
1254                        if detach_beholder {
1255                            beholder = None;
1256                        }
1257                    }
1258                }
1259            }
1260        }
1261        if let Some(ref mut b) = beholder {
1262            let final_offset = elapsed_ms(started_at_ms);
1263            for ev in b.on_done(&id, final_offset) {
1264                let _ = rt.block_on(store.append_event(
1265                    &ev.run_id,
1266                    ev.offset_ms,
1267                    ev.level,
1268                    &ev.target,
1269                    &ev.msg,
1270                    &ev.fields,
1271                    ev.anchor.as_ref().map(|a| a.seq),
1272                    &ev.source,
1273                ));
1274            }
1275            if let Some(reason) = b.unknown_format_reason() {
1276                let new_status = BeholderStatus::unknown_format_with_reason(b.name(), reason);
1277                let _ = rt.block_on(store.update_beholder_status(&id, &new_status));
1278            }
1279        }
1280    });
1281}
1282
1283#[allow(clippy::too_many_arguments)]
1284async fn run_lifecycle(
1285    store: Arc<TaskStore>,
1286    active: Arc<Mutex<HashMap<String, RunControl>>>,
1287    id: TaskRunId,
1288    pid: u32,
1289    // `reap` blocks until the child exits and yields its exit code. It owns
1290    // whatever the spawn mode has to keep alive across the wait — for a PTY run
1291    // that includes the master fd, which must outlive `wait()`.
1292    reap: Box<dyn FnOnce() -> Option<u32> + Send>,
1293    mut kill_rx: mpsc::Receiver<KillRequest>,
1294    reader_done_rx: oneshot::Receiver<()>,
1295    completion_tx: Option<tokio::sync::mpsc::UnboundedSender<(TaskRunId, RunStatus)>>,
1296    // Holds the write end of the log FIFO open until this task completes.
1297    // Dropping it produces EOF for the receiver thread, which happens after
1298    // the terminal RunStatus is written below.
1299    #[cfg(unix)]
1300    _log_wfd: Option<FdCloser>,
1301) {
1302    // Pin the reader-done future so it can be polled by reference in
1303    // nested select! arms without consuming ownership.
1304    let reader_done = async { reader_done_rx.await.ok(); };
1305    tokio::pin!(reader_done);
1306
1307    let sent_signal: Option<i32>;
1308
1309    tokio::select! {
1310        req = kill_rx.recv() => {
1311            match req {
1312                Some(KillRequest { signal }) => {
1313                    send_unix_signal(pid, signal);
1314                    if signal == SIGKILL {
1315                        sent_signal = Some(SIGKILL);
1316                    } else {
1317                        // Grace period: give the process a chance to exit cleanly.
1318                        tokio::select! {
1319                            _ = &mut reader_done => {
1320                                // Exited within grace — no SIGKILL needed.
1321                                sent_signal = Some(signal);
1322                            }
1323                            _ = tokio::time::sleep(DEFAULT_GRACE) => {
1324                                // Grace expired — escalate.
1325                                send_unix_signal(pid, SIGKILL);
1326                                sent_signal = Some(SIGKILL);
1327                            }
1328                        }
1329                    }
1330                }
1331                // kill_tx dropped (driver shutting down) — force kill.
1332                None => {
1333                    send_unix_signal(pid, SIGKILL);
1334                    sent_signal = Some(SIGKILL);
1335                }
1336            }
1337        }
1338        _ = &mut reader_done => {
1339            sent_signal = None;
1340        }
1341    }
1342
1343    // Reap the child (blocking) on a dedicated thread-pool slot. For a PTY run
1344    // the closure also owns our master handle, so the fd outlives the wait; the
1345    // matching `RunControl` (removed from `active` below) holds the other `Arc`,
1346    // so the fd actually closes once both are gone.
1347    let exit_code = task::spawn_blocking(reap).await.ok().flatten();
1348
1349    let ended_at = unix_now_secs();
1350    let status = match sent_signal {
1351        Some(sig) => RunStatus::Killed { signal: sig, ended_at },
1352        None => match exit_code {
1353            Some(code) => RunStatus::Done { exit_code: code as i32, ended_at },
1354            None => RunStatus::Lost {
1355                reason: "process exited without an exit code".to_string(),
1356            },
1357        },
1358    };
1359
1360    /* Losing this write is not cosmetic: the run stays `Running` in the store
1361       forever and every reader — tail loops, the terminal UI, the next
1362       daemon's Lost-on-disappear sweep — believes a dead process is alive.
1363       `update_status` already retries through lock contention, so a failure
1364       here is terminal and worth saying out loud. */
1365    if let Err(e) = store.update_status(&id, &status).await {
1366        eprintln!("[yah task-runs] failed to record terminal status for run {id}: {e}");
1367    }
1368    if let Some(ref tx) = completion_tx {
1369        let _ = tx.send((id.clone(), status));
1370    }
1371    active.lock().unwrap().remove(&id.to_string());
1372}
1373
1374// ─── Helpers ──────────────────────────────────────────────────────────────────
1375
1376fn send_unix_signal(pid: u32, signal: i32) {
1377    #[cfg(unix)]
1378    unsafe {
1379        libc::kill(pid as libc::pid_t, signal);
1380    }
1381    // On non-Unix platforms signal delivery is not implemented here.
1382}
1383
1384fn unix_now_secs() -> u64 {
1385    SystemTime::now()
1386        .duration_since(UNIX_EPOCH)
1387        .unwrap_or_default()
1388        .as_secs()
1389}
1390
1391fn elapsed_ms(started_at_ms: u64) -> u32 {
1392    let now_ms = SystemTime::now()
1393        .duration_since(UNIX_EPOCH)
1394        .unwrap_or_default()
1395        .as_millis() as u64;
1396    now_ms.saturating_sub(started_at_ms).min(u32::MAX as u64) as u32
1397}
1398
1399// ─── Tests ────────────────────────────────────────────────────────────────────
1400
1401#[cfg(test)]
1402mod tests {
1403    use super::*;
1404    use crate::store::ChunkFilter;
1405
1406    async fn open_store(dir: &tempfile::TempDir) -> Arc<TaskStore> {
1407        Arc::new(TaskStore::open(&dir.path().join("tr.turso")).await.unwrap())
1408    }
1409
1410    // ── Lost-on-disappear (pure store, no PTY) ────────────────────────────────
1411
1412    #[tokio::test]
1413    async fn lost_on_disappear_marks_stale_running_runs() {
1414        let dir = tempfile::tempdir().unwrap();
1415        let store = open_store(&dir).await;
1416
1417        // Simulate a run left in "Running" state by a prior daemon.
1418        let stale_id = TaskRunId::new();
1419        store
1420            .insert_run(&TaskRunMeta {
1421                id: stale_id.clone(),
1422                command: "sleep 9999".to_string(),
1423                cwd: "/tmp".into(),
1424                env: vec![],
1425                started_at: unix_now_secs() - 60,
1426                status: RunStatus::Running,
1427                label: None,
1428                initiator: Initiator::Human { camp: "test".to_string() },
1429                beholder_status: None,
1430                pinned: false,
1431                origin: None,
1432                host_pid: None,
1433            })
1434            .await
1435            .unwrap();
1436
1437        // Creating a new driver must mark stale runs Lost.
1438        let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1439
1440        let meta = store.get_run(&stale_id).await.unwrap().unwrap();
1441        assert!(
1442            matches!(meta.status, RunStatus::Lost { .. }),
1443            "stale run should be Lost, got {:?}",
1444            meta.status
1445        );
1446    }
1447
1448    /// R267-T12: a sweep tombstone must reach the completion channel.
1449    ///
1450    /// This was the one status transition that wrote the store and told
1451    /// nobody — and it is the transition a client can least discover on its
1452    /// own, because the process that would have reported the exit is the one
1453    /// that died. A renderer that trusts pushes would leave the pane on
1454    /// "running" for as long as it stayed open.
1455    #[tokio::test]
1456    async fn a_sweep_tombstone_fires_the_completion_channel() {
1457        let dir = tempfile::tempdir().unwrap();
1458        let store = open_store(&dir).await;
1459
1460        let stale_id = TaskRunId::new();
1461        store
1462            .insert_run(&TaskRunMeta {
1463                id: stale_id.clone(),
1464                command: "sleep 9999".to_string(),
1465                cwd: "/tmp".into(),
1466                env: vec![],
1467                started_at: unix_now_secs() - 60,
1468                status: RunStatus::Running,
1469                label: None,
1470                initiator: Initiator::Human { camp: "test".to_string() },
1471                beholder_status: None,
1472                pinned: false,
1473                origin: None,
1474                host_pid: None,
1475            })
1476            .await
1477            .unwrap();
1478
1479        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1480        let _driver = TaskDriver::with_channels(
1481            Arc::clone(&store),
1482            DriverChannels { completion: Some(tx), output: None },
1483        )
1484        .await
1485        .unwrap();
1486
1487        let (id, status) = rx.try_recv().expect("the sweep must announce what it tombstoned");
1488        assert_eq!(id, stale_id);
1489        assert!(
1490            matches!(status, RunStatus::Lost { .. }),
1491            "expected Lost, got {status:?}"
1492        );
1493    }
1494
1495    // ── Stale-run policy (R617-F6) ───────────────────────────────────────────
1496
1497    /// Plant a `Running` row as if some other process had spawned it.
1498    async fn plant_running(
1499        store: &Arc<TaskStore>,
1500        origin: Option<&str>,
1501        host_pid: Option<u32>,
1502    ) -> TaskRunId {
1503        let id = TaskRunId::new();
1504        store
1505            .insert_run(&TaskRunMeta {
1506                id: id.clone(),
1507                command: "sleep 9999".to_string(),
1508                cwd: "/tmp".into(),
1509                env: vec![],
1510                started_at: unix_now_secs() - 60,
1511                status: RunStatus::Running,
1512                label: None,
1513                initiator: Initiator::Human {
1514                    camp: "test".to_string(),
1515                },
1516                beholder_status: None,
1517                pinned: false,
1518                origin: origin.map(str::to_string),
1519                host_pid,
1520            })
1521            .await
1522            .unwrap();
1523        id
1524    }
1525
1526    async fn is_lost(store: &Arc<TaskStore>, id: &TaskRunId) -> bool {
1527        matches!(
1528            store.get_run(id).await.unwrap().unwrap().status,
1529            RunStatus::Lost { .. }
1530        )
1531    }
1532
1533    fn adopt_terminal() -> StaleRunPolicy {
1534        StaleRunPolicy::AdoptLiveHosts {
1535            origins: vec!["terminal".to_string()],
1536        }
1537    }
1538
1539    /// The property the whole ticket exists for: attaching to a store must not
1540    /// declare another live process's shell dead.
1541    #[tokio::test]
1542    async fn a_run_owned_by_a_live_host_survives_a_new_driver() {
1543        let dir = tempfile::tempdir().unwrap();
1544        let store = open_store(&dir).await;
1545        // Our own pid is by definition a live process, and is the cheapest
1546        // honest stand-in for "a peer that is still running".
1547        let id = plant_running(&store, Some("terminal"), Some(std::process::id())).await;
1548
1549        let _driver = TaskDriver::with_config(
1550            Arc::clone(&store),
1551            DriverChannels::default(),
1552            adopt_terminal(),
1553        )
1554        .await
1555        .unwrap();
1556
1557        assert!(
1558            !is_lost(&store, &id).await,
1559            "a terminal run whose owner is alive must stay Running — \
1560             tombstoning it is what made a surviving shell read as dead"
1561        );
1562    }
1563
1564    /// The other half: a genuinely abandoned shell must still be tombstoned,
1565    /// or a crashed host leaves permanent zombie tiles.
1566    #[tokio::test]
1567    async fn a_run_whose_host_is_gone_is_still_tombstoned() {
1568        let dir = tempfile::tempdir().unwrap();
1569        let store = open_store(&dir).await;
1570        // Reaped in-test, so the pid is real-but-dead rather than guessed.
1571        let dead_pid = {
1572            let child = std::process::Command::new("true").spawn().unwrap();
1573            let pid = child.id();
1574            let mut child = child;
1575            let _ = child.wait();
1576            pid
1577        };
1578        let id = plant_running(&store, Some("terminal"), Some(dead_pid)).await;
1579
1580        let _driver = TaskDriver::with_config(
1581            Arc::clone(&store),
1582            DriverChannels::default(),
1583            adopt_terminal(),
1584        )
1585        .await
1586        .unwrap();
1587
1588        assert!(
1589            is_lost(&store, &id).await,
1590            "pid {dead_pid} was reaped; its run has no owner left and must be Lost"
1591        );
1592    }
1593
1594    /// The exemption is narrowed by origin, so ordinary jobs keep the old rule
1595    /// even when their owner happens to still be alive — an in-flight `cargo
1596    /// build` whose driver is gone has nobody left to record its exit.
1597    #[tokio::test]
1598    async fn a_non_matching_origin_is_tombstoned_even_with_a_live_host() {
1599        let dir = tempfile::tempdir().unwrap();
1600        let store = open_store(&dir).await;
1601        let job = plant_running(&store, None, Some(std::process::id())).await;
1602        let other = plant_running(&store, Some("gnome"), Some(std::process::id())).await;
1603
1604        let _driver = TaskDriver::with_config(
1605            Arc::clone(&store),
1606            DriverChannels::default(),
1607            adopt_terminal(),
1608        )
1609        .await
1610        .unwrap();
1611
1612        assert!(is_lost(&store, &job).await, "an origin-less job is not exempt");
1613        assert!(
1614            is_lost(&store, &other).await,
1615            "an origin outside the list is not exempt"
1616        );
1617    }
1618
1619    /// A row written before `host_pid` existed reads back `None`. Unknown
1620    /// ownership must fall back to the old behaviour rather than stranding the
1621    /// run `Running` forever.
1622    #[tokio::test]
1623    async fn an_unattributed_run_is_tombstoned() {
1624        let dir = tempfile::tempdir().unwrap();
1625        let store = open_store(&dir).await;
1626        let id = plant_running(&store, Some("terminal"), None).await;
1627
1628        let _driver = TaskDriver::with_config(
1629            Arc::clone(&store),
1630            DriverChannels::default(),
1631            adopt_terminal(),
1632        )
1633        .await
1634        .unwrap();
1635
1636        assert!(is_lost(&store, &id).await);
1637    }
1638
1639    /// `TaskDriver::new` must not have quietly changed behaviour — every
1640    /// existing embedder still gets Lost-on-disappear.
1641    #[tokio::test]
1642    async fn the_default_policy_is_still_lost_on_disappear() {
1643        assert_eq!(StaleRunPolicy::default(), StaleRunPolicy::LostOnDisappear);
1644
1645        let dir = tempfile::tempdir().unwrap();
1646        let store = open_store(&dir).await;
1647        let id = plant_running(&store, Some("terminal"), Some(std::process::id())).await;
1648
1649        let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1650
1651        assert!(
1652            is_lost(&store, &id).await,
1653            "the default must tombstone regardless of origin or owner liveness"
1654        );
1655    }
1656
1657    /// The owner is recorded by `spawn_run` itself, not by the caller — the
1658    /// policy is worthless if rows arrive unattributed.
1659    #[tokio::test]
1660    async fn spawn_run_stamps_this_process_as_the_owner() {
1661        let dir = tempfile::tempdir().unwrap();
1662        let store = open_store(&dir).await;
1663        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1664
1665        let id = driver
1666            .spawn_run(
1667                "true",
1668                SpawnOpts {
1669                    cwd: "/tmp".into(),
1670                    origin: Some("terminal".to_string()),
1671                    ..Default::default()
1672                },
1673            )
1674            .await
1675            .unwrap();
1676
1677        let meta = store.get_run(&id).await.unwrap().unwrap();
1678        assert_eq!(meta.host_pid, Some(std::process::id()));
1679    }
1680
1681    #[tokio::test]
1682    async fn new_driver_does_not_touch_completed_runs() {
1683        let dir = tempfile::tempdir().unwrap();
1684        let store = open_store(&dir).await;
1685
1686        let done_id = TaskRunId::new();
1687        store
1688            .insert_run(&TaskRunMeta {
1689                id: done_id.clone(),
1690                command: "true".to_string(),
1691                cwd: "/tmp".into(),
1692                env: vec![],
1693                started_at: unix_now_secs() - 10,
1694                status: RunStatus::Running,
1695                label: None,
1696                initiator: Initiator::Human { camp: "test".to_string() },
1697                beholder_status: None,
1698                pinned: false,
1699                origin: None,
1700                host_pid: None,
1701            })
1702            .await
1703            .unwrap();
1704        store
1705            .update_status(&done_id, &RunStatus::Done { exit_code: 0, ended_at: unix_now_secs() })
1706            .await
1707            .unwrap();
1708
1709        let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1710
1711        let meta = store.get_run(&done_id).await.unwrap().unwrap();
1712        assert!(
1713            matches!(meta.status, RunStatus::Done { .. }),
1714            "completed run must not be touched"
1715        );
1716    }
1717
1718    // ── PTY spawn + capture ───────────────────────────────────────────────────
1719
1720    #[tokio::test]
1721    async fn spawn_echo_and_read_chunks() {
1722        let dir = tempfile::tempdir().unwrap();
1723        let store = open_store(&dir).await;
1724        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1725
1726        let id = driver
1727            .spawn_run(
1728                "echo hello_world",
1729                SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1730            )
1731            .await
1732            .unwrap();
1733
1734        // Wait for the run to complete (poll status up to 5 s).
1735        let deadline = std::time::Instant::now() + Duration::from_secs(5);
1736        loop {
1737            let meta = store.get_run(&id).await.unwrap().unwrap();
1738            if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1739                break;
1740            }
1741            if std::time::Instant::now() > deadline {
1742                panic!("run did not complete in time, status={:?}", meta.status);
1743            }
1744            tokio::time::sleep(Duration::from_millis(50)).await;
1745        }
1746
1747        // Chunks must contain "hello_world".
1748        let chunks = store
1749            .get_chunks(&id, &ChunkFilter::default())
1750            .await
1751            .unwrap();
1752        let output: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
1753        let text = String::from_utf8_lossy(&output);
1754        assert!(
1755            text.contains("hello_world"),
1756            "expected 'hello_world' in output, got: {text:?}"
1757        );
1758
1759        let meta = store.get_run(&id).await.unwrap().unwrap();
1760        assert!(
1761            matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
1762            "expected Done(0), got {:?}",
1763            meta.status
1764        );
1765    }
1766
1767    // ── Pipe mode (R739-F6) ───────────────────────────────────────────────────
1768
1769    /// Run `cmd` to completion and return its stored chunks.
1770    async fn run_to_completion(
1771        store: &Arc<TaskStore>,
1772        driver: &TaskDriver,
1773        cmd: &str,
1774        opts: SpawnOpts,
1775    ) -> Vec<OutputChunk> {
1776        let id = driver.spawn_run(cmd, opts).await.unwrap();
1777        let deadline = std::time::Instant::now() + Duration::from_secs(10);
1778        loop {
1779            let meta = store.get_run(&id).await.unwrap().unwrap();
1780            if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1781                break;
1782            }
1783            if std::time::Instant::now() > deadline {
1784                panic!("run did not complete in time, status={:?}", meta.status);
1785            }
1786            tokio::time::sleep(Duration::from_millis(25)).await;
1787        }
1788        store.get_chunks(&id, &ChunkFilter::default()).await.unwrap()
1789    }
1790
1791    fn joined(chunks: &[OutputChunk]) -> Vec<u8> {
1792        chunks.iter().flat_map(|c| c.bytes.clone()).collect()
1793    }
1794
1795    fn joined_stream(chunks: &[OutputChunk], stream: Stream) -> Vec<u8> {
1796        chunks
1797            .iter()
1798            .filter(|c| c.stream == stream)
1799            .flat_map(|c| c.bytes.clone())
1800            .collect()
1801    }
1802
1803    /// Divergence 1 of 3 (R739-F4): the child must not think it is on a
1804    /// terminal. This is the one that makes cargo colorize.
1805    #[tokio::test]
1806    async fn pipe_mode_child_sees_no_tty_on_stdout() {
1807        let dir = tempfile::tempdir().unwrap();
1808        let store = open_store(&dir).await;
1809        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1810        let cmd = "if [ -t 1 ]; then echo TTY; else echo PIPE; fi";
1811
1812        let piped = run_to_completion(
1813            &store,
1814            &driver,
1815            cmd,
1816            SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1817        )
1818        .await;
1819        assert_eq!(joined(&piped), b"PIPE\n");
1820
1821        // The PTY default is unchanged — the terminal tiles depend on it.
1822        let ptied = run_to_completion(
1823            &store,
1824            &driver,
1825            cmd,
1826            SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1827        )
1828        .await;
1829        assert_eq!(joined(&ptied), b"TTY\r\n");
1830    }
1831
1832    /// Divergence 2 of 3: no `ONLCR`, so a `\n` the child wrote stays a `\n`.
1833    /// This is what `build_run.rs::undo_onlcr` used to compensate for.
1834    #[tokio::test]
1835    async fn pipe_mode_does_not_translate_newlines() {
1836        let dir = tempfile::tempdir().unwrap();
1837        let store = open_store(&dir).await;
1838        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1839
1840        let piped = run_to_completion(
1841            &store,
1842            &driver,
1843            r"printf 'a\nb\n'",
1844            SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1845        )
1846        .await;
1847        assert_eq!(joined(&piped), b"a\nb\n");
1848    }
1849
1850    /// Divergence 3 of 3: stdout and stderr stay apart, under their true
1851    /// [`Stream`], instead of being merged by the terminal.
1852    #[tokio::test]
1853    async fn pipe_mode_keeps_stderr_separate_from_stdout() {
1854        let dir = tempfile::tempdir().unwrap();
1855        let store = open_store(&dir).await;
1856        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1857        let cmd = "printf 'to-out\n'; printf 'to-err\n' >&2";
1858
1859        let piped = run_to_completion(
1860            &store,
1861            &driver,
1862            cmd,
1863            SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1864        )
1865        .await;
1866        assert_eq!(joined_stream(&piped, Stream::Stdout), b"to-out\n");
1867        assert_eq!(joined_stream(&piped, Stream::Stderr), b"to-err\n");
1868
1869        // Under a PTY the kernel merges them and everything lands on stdout —
1870        // the property that made stream separation unrecoverable downstream.
1871        let ptied = run_to_completion(
1872            &store,
1873            &driver,
1874            cmd,
1875            SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1876        )
1877        .await;
1878        assert!(
1879            joined_stream(&ptied, Stream::Stderr).is_empty(),
1880            "PTY runs have no stderr chunks; that is the behaviour pipe mode exists to fix",
1881        );
1882    }
1883
1884    /// Both pipes must reach EOF before the child is reaped, or a run whose
1885    /// last bytes went to stderr would be marked terminal with output still
1886    /// unread. The 4 KiB write is larger than a pipe's atomic-write buffer, so
1887    /// this fails if either pump is dropped rather than awaited.
1888    #[tokio::test]
1889    async fn pipe_mode_drains_both_streams_before_the_run_is_terminal() {
1890        let dir = tempfile::tempdir().unwrap();
1891        let store = open_store(&dir).await;
1892        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1893
1894        let piped = run_to_completion(
1895            &store,
1896            &driver,
1897            "head -c 4096 /dev/zero | tr '\\0' 'x'; head -c 4096 /dev/zero | tr '\\0' 'y' >&2",
1898            SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1899        )
1900        .await;
1901        assert_eq!(joined_stream(&piped, Stream::Stdout).len(), 4096);
1902        assert_eq!(joined_stream(&piped, Stream::Stderr).len(), 4096);
1903    }
1904
1905    /// Exit codes have to survive the move to `std::process::Child`, which
1906    /// reports them through a different type than `portable_pty::Child`.
1907    #[tokio::test]
1908    async fn pipe_mode_records_the_childs_exit_code() {
1909        let dir = tempfile::tempdir().unwrap();
1910        let store = open_store(&dir).await;
1911        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1912
1913        let id = driver
1914            .spawn_run(
1915                "exit 101",
1916                SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1917            )
1918            .await
1919            .unwrap();
1920
1921        let deadline = std::time::Instant::now() + Duration::from_secs(10);
1922        loop {
1923            let meta = store.get_run(&id).await.unwrap().unwrap();
1924            match meta.status {
1925                RunStatus::Done { exit_code, .. } => {
1926                    assert_eq!(exit_code, 101);
1927                    return;
1928                }
1929                RunStatus::Lost { .. } | RunStatus::Killed { .. } => {
1930                    panic!("unexpected terminal status {:?}", meta.status)
1931                }
1932                _ => {}
1933            }
1934            if std::time::Instant::now() > deadline {
1935                panic!("run did not complete in time");
1936            }
1937            tokio::time::sleep(Duration::from_millis(25)).await;
1938        }
1939    }
1940
1941    /// A pipe run has no terminal, and the two PTY-only verbs must say so
1942    /// rather than reaching into a `None` master.
1943    #[tokio::test]
1944    async fn pipe_mode_has_no_terminal_to_resize_or_read_a_foreground_pid_from() {
1945        let dir = tempfile::tempdir().unwrap();
1946        let store = open_store(&dir).await;
1947        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1948
1949        let id = driver
1950            .spawn_run(
1951                "sleep 2",
1952                SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1953            )
1954            .await
1955            .unwrap();
1956
1957        assert!(matches!(
1958            driver.resize_run(&id, 100, 40).await,
1959            Err(DriverError::NotFound(_))
1960        ));
1961        assert_eq!(driver.foreground_pid(&id), None);
1962        let _ = driver.kill_run(&id, Some(SIGKILL)).await;
1963    }
1964
1965    /// Wait for a run to reach a terminal status, or panic.
1966    async fn await_done(store: &TaskStore, id: &TaskRunId) -> TaskRunMeta {
1967        let deadline = std::time::Instant::now() + Duration::from_secs(5);
1968        loop {
1969            let meta = store.get_run(id).await.unwrap().unwrap();
1970            if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1971                return meta;
1972            }
1973            if std::time::Instant::now() > deadline {
1974                panic!("run did not complete in time, status={:?}", meta.status);
1975            }
1976            tokio::time::sleep(Duration::from_millis(50)).await;
1977        }
1978    }
1979
1980    async fn output_of(store: &TaskStore, id: &TaskRunId) -> String {
1981        let chunks = store.get_chunks(id, &ChunkFilter::default()).await.unwrap();
1982        let bytes: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
1983        String::from_utf8_lossy(&bytes).into_owned()
1984    }
1985
1986    // ── The caller's bytes reach the shell unchanged (R739-S2) ───────────────
1987
1988    /// `AttachResult.argv` is populated on every run, rewrite or not, so
1989    /// `spawn_run` used to join it back into the command line unconditionally.
1990    /// That put every `task.run` command through a whitespace normalization
1991    /// nobody asked for. A multi-line command is the case where that is not
1992    /// cosmetic: the newline the caller wrote becomes a space, and two
1993    /// commands become one nonsense command.
1994    #[cfg(unix)]
1995    #[tokio::test]
1996    async fn a_multi_line_command_is_not_flattened_into_one_line() {
1997        let dir = tempfile::tempdir().unwrap();
1998        let store = open_store(&dir).await;
1999        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2000
2001        // Flattened to one line this is `echo one echo two`, which prints
2002        // "one echo two" — a different answer, not a failure, which is what
2003        // makes the old behaviour dangerous rather than merely wrong.
2004        let id = driver
2005            .spawn_run(
2006                "echo one\necho two",
2007                SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2008            )
2009            .await
2010            .unwrap();
2011        await_done(&store, &id).await;
2012
2013        let out = output_of(&store, &id).await;
2014        assert!(out.contains("one"), "got: {out:?}");
2015        assert!(
2016            out.contains("two"),
2017            "the second line must have run as its own command; got: {out:?}"
2018        );
2019        assert!(
2020            !out.contains("one echo two"),
2021            "the newline was flattened into a space; got: {out:?}"
2022        );
2023    }
2024
2025    /// `resolve_argv` strips `bunx`/`npx`/`pnpm` so a beholder's `matches` sees
2026    /// the bare tool. That is a *matching* concern; it must never reach the
2027    /// spawn, or the wrapper the caller needed is gone from the command.
2028    #[cfg(unix)]
2029    #[tokio::test]
2030    async fn a_wrapper_the_caller_wrote_is_not_stripped_from_the_spawned_command() {
2031        let dir = tempfile::tempdir().unwrap();
2032        let store = open_store(&dir).await;
2033        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2034
2035        // `npx` is almost certainly absent in test environments, and that is
2036        // the point: if the wrapper survived, the shell reports it missing. If
2037        // it were stripped we would be running bare `--version`.
2038        let id = driver
2039            .spawn_run(
2040                "npx r739s2-nonexistent-tool --version",
2041                SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2042            )
2043            .await
2044            .unwrap();
2045        let meta = await_done(&store, &id).await;
2046        let out = output_of(&store, &id).await;
2047        assert!(
2048            !matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2049            "expected a failure, got {:?} with output {out:?}",
2050            meta.status
2051        );
2052        assert!(
2053            !out.contains("--version: "),
2054            "the wrapper was stripped and the shell tried to run the flag; got: {out:?}"
2055        );
2056    }
2057
2058    // ── Direct argv (R652-T6) ────────────────────────────────────────────────
2059
2060    #[tokio::test]
2061    async fn explicit_argv_execs_the_program_directly() {
2062        let dir = tempfile::tempdir().unwrap();
2063        let store = open_store(&dir).await;
2064        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2065
2066        /* The distinguishing observation: under `sh -c` the child is `sh` and
2067           `$0` is `sh`; exec'd directly it is the program itself. Printing
2068           `$0` is the cheapest way to see which of the two happened. */
2069        let id = driver
2070            .spawn_run(
2071                "unused-because-argv-wins",
2072                SpawnOpts {
2073                    cwd: "/tmp".into(),
2074                    argv: Some(vec![
2075                        "/bin/sh".into(),
2076                        "-c".into(),
2077                        "printf 'argv0=%s\\n' \"$0\"".into(),
2078                        "direct-exec-marker".into(),
2079                    ]),
2080                    ..Default::default()
2081                },
2082            )
2083            .await
2084            .unwrap();
2085
2086        await_done(&store, &id).await;
2087        let text = output_of(&store, &id).await;
2088        assert!(
2089            text.contains("argv0=direct-exec-marker"),
2090            "argv should have been exec'd verbatim, got: {text:?}"
2091        );
2092    }
2093
2094    #[tokio::test]
2095    async fn explicit_argv_still_records_the_requested_command() {
2096        let dir = tempfile::tempdir().unwrap();
2097        let store = open_store(&dir).await;
2098        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2099
2100        /* A shell tile asks for "$SHELL" and the daemon resolves it to a real
2101           argv. The run must still read back as what was asked for, or the
2102           rail row and the history re-run both show an implementation
2103           detail. */
2104        let id = driver
2105            .spawn_run(
2106                "$SHELL",
2107                SpawnOpts {
2108                    cwd: "/tmp".into(),
2109                    argv: Some(vec!["/bin/sh".into(), "-c".into(), "true".into()]),
2110                    ..Default::default()
2111                },
2112            )
2113            .await
2114            .unwrap();
2115
2116        let meta = await_done(&store, &id).await;
2117        assert_eq!(meta.command, "$SHELL");
2118        assert!(
2119            matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2120            "expected Done(0), got {:?}",
2121            meta.status
2122        );
2123    }
2124
2125    /// R901-B2. The control is the whole point: without the `pipefail: false`
2126    /// half this would pass if `pipefail` stopped existing, and with only the
2127    /// `true` half it would pass if every pipeline had always reported its
2128    /// leftmost failure. The pair pins the *difference*, which is the thing
2129    /// that cost this camp ~50 minutes of red tree.
2130    #[tokio::test]
2131    async fn pipefail_reports_the_failing_stage_and_posix_reports_the_last_one() {
2132        let dir = tempfile::tempdir().unwrap();
2133        let store = open_store(&dir).await;
2134        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2135
2136        // `(exit 101) | tail -1` is `cargo check 2>&1 | tail -40` with the
2137        // compile stripped out: a failing producer feeding a succeeding tail.
2138        let line = "(exit 101) | tail -1";
2139
2140        let posix = driver
2141            .spawn_run(
2142                line,
2143                SpawnOpts { cwd: "/tmp".into(), pipefail: false, ..Default::default() },
2144            )
2145            .await
2146            .unwrap();
2147        let meta = await_done(&store, &posix).await;
2148        assert!(
2149            matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2150            "POSIX pipeline status is the LAST stage's — expected Done(0), got {:?}",
2151            meta.status
2152        );
2153
2154        let failing = driver
2155            .spawn_run(
2156                line,
2157                SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2158            )
2159            .await
2160            .unwrap();
2161        let meta = await_done(&store, &failing).await;
2162        assert!(
2163            matches!(meta.status, RunStatus::Done { exit_code: 101, .. }),
2164            "pipefail must surface the producer's 101, got {:?}",
2165            meta.status
2166        );
2167    }
2168
2169    /// The prelude must not reach [`TaskRunMeta::command`]: that string is what
2170    /// history re-runs and what an agent audits the relocation note against.
2171    #[tokio::test]
2172    async fn pipefail_does_not_leak_into_the_recorded_command() {
2173        let dir = tempfile::tempdir().unwrap();
2174        let store = open_store(&dir).await;
2175        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2176
2177        let id = driver
2178            .spawn_run(
2179                "echo recorded-verbatim | cat",
2180                SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2181            )
2182            .await
2183            .unwrap();
2184
2185        let meta = await_done(&store, &id).await;
2186        assert_eq!(meta.command, "echo recorded-verbatim | cat");
2187        assert!(
2188            !meta.command.contains("pipefail"),
2189            "the prelude leaked into the recorded command: {:?}",
2190            meta.command
2191        );
2192    }
2193
2194    /// The portability guard. On a `/bin/sh` that rejects `pipefail` (dash, i.e.
2195    /// most Linux camps) the probe must degrade to plain POSIX semantics — it
2196    /// must NOT take the shell down with it, because `set` is a special builtin
2197    /// and a bare `set -o pipefail` there is entitled to exit before the
2198    /// caller's command runs at all. Asserting the command still produces its
2199    /// output is asserting exactly that.
2200    #[tokio::test]
2201    async fn the_pipefail_probe_never_costs_the_command_that_follows_it() {
2202        let dir = tempfile::tempdir().unwrap();
2203        let store = open_store(&dir).await;
2204        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2205
2206        let id = driver
2207            .spawn_run(
2208                "echo probe-survived",
2209                SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2210            )
2211            .await
2212            .unwrap();
2213
2214        let meta = await_done(&store, &id).await;
2215        let text = output_of(&store, &id).await;
2216        assert!(
2217            matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2218            "expected Done(0), got {:?}",
2219            meta.status
2220        );
2221        assert!(text.contains("probe-survived"), "command did not run, got: {text:?}");
2222        // The probe itself must be silent — it runs on every relocated build.
2223        assert!(
2224            !text.contains("pipefail"),
2225            "the probe printed a diagnostic into the build's own output: {text:?}"
2226        );
2227    }
2228
2229    #[tokio::test]
2230    async fn empty_argv_falls_back_to_the_shell_path() {
2231        let dir = tempfile::tempdir().unwrap();
2232        let store = open_store(&dir).await;
2233        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2234
2235        let id = driver
2236            .spawn_run(
2237                "echo empty_argv_fallback",
2238                SpawnOpts { cwd: "/tmp".into(), argv: Some(vec![]), ..Default::default() },
2239            )
2240            .await
2241            .unwrap();
2242
2243        await_done(&store, &id).await;
2244        let text = output_of(&store, &id).await;
2245        assert!(
2246            text.contains("empty_argv_fallback"),
2247            "empty argv must not spawn nothing, got: {text:?}"
2248        );
2249    }
2250
2251    #[tokio::test]
2252    async fn spawn_failing_command_records_nonzero_exit() {
2253        let dir = tempfile::tempdir().unwrap();
2254        let store = open_store(&dir).await;
2255        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2256
2257        let id = driver
2258            .spawn_run(
2259                "exit 42",
2260                SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2261            )
2262            .await
2263            .unwrap();
2264
2265        let deadline = std::time::Instant::now() + Duration::from_secs(5);
2266        loop {
2267            let meta = store.get_run(&id).await.unwrap().unwrap();
2268            if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2269                match meta.status {
2270                    RunStatus::Done { exit_code, .. } => {
2271                        assert_ne!(exit_code, 0, "exit 42 should produce a non-zero exit code");
2272                    }
2273                    other => panic!("unexpected status: {other:?}"),
2274                }
2275                break;
2276            }
2277            if std::time::Instant::now() > deadline {
2278                panic!("run did not complete in time");
2279            }
2280            tokio::time::sleep(Duration::from_millis(50)).await;
2281        }
2282    }
2283
2284    // ── Signal handling ───────────────────────────────────────────────────────
2285
2286    #[cfg(unix)]
2287    #[tokio::test]
2288    async fn kill_with_sigterm_transitions_to_killed() {
2289        let dir = tempfile::tempdir().unwrap();
2290        let store = open_store(&dir).await;
2291        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2292
2293        let id = driver
2294            .spawn_run(
2295                "sleep 60",
2296                SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2297            )
2298            .await
2299            .unwrap();
2300
2301        // Give the process a moment to start.
2302        tokio::time::sleep(Duration::from_millis(100)).await;
2303
2304        driver.kill_run(&id, Some(SIGTERM)).await.unwrap();
2305
2306        let deadline = std::time::Instant::now() + Duration::from_secs(10);
2307        loop {
2308            let meta = store.get_run(&id).await.unwrap().unwrap();
2309            if matches!(meta.status, RunStatus::Killed { .. } | RunStatus::Lost { .. }) {
2310                assert!(
2311                    matches!(meta.status, RunStatus::Killed { .. }),
2312                    "expected Killed, got {:?}",
2313                    meta.status
2314                );
2315                break;
2316            }
2317            if std::time::Instant::now() > deadline {
2318                panic!("run did not become Killed in time, status={:?}", meta.status);
2319            }
2320            tokio::time::sleep(Duration::from_millis(50)).await;
2321        }
2322    }
2323
2324    #[cfg(unix)]
2325    #[tokio::test]
2326    async fn kill_run_returns_not_found_after_exit() {
2327        let dir = tempfile::tempdir().unwrap();
2328        let store = open_store(&dir).await;
2329        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2330
2331        let id = driver
2332            .spawn_run(
2333                "echo done",
2334                SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2335            )
2336            .await
2337            .unwrap();
2338
2339        // Wait for natural exit.
2340        let deadline = std::time::Instant::now() + Duration::from_secs(5);
2341        loop {
2342            let meta = store.get_run(&id).await.unwrap().unwrap();
2343            if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2344                break;
2345            }
2346            if std::time::Instant::now() > deadline {
2347                panic!("run did not complete");
2348            }
2349            tokio::time::sleep(Duration::from_millis(50)).await;
2350        }
2351
2352        // Kill on a completed run should return NotFound.
2353        let result = driver.kill_run(&id, None).await;
2354        assert!(
2355            matches!(result, Err(DriverError::NotFound(_))),
2356            "expected NotFound, got {result:?}"
2357        );
2358    }
2359
2360    // ── Stdin relay ───────────────────────────────────────────────────────────
2361
2362    #[cfg(unix)]
2363    #[tokio::test]
2364    async fn stdin_send_reaches_child() {
2365        let dir = tempfile::tempdir().unwrap();
2366        let store = open_store(&dir).await;
2367        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2368
2369        // Shell that reads a line from stdin and echoes it back.
2370        let id = driver
2371            .spawn_run(
2372                "read line && echo got_$line",
2373                SpawnOpts {
2374                    cwd: "/tmp".into(),
2375                    stdin_enabled: true,
2376                    ..Default::default()
2377                },
2378            )
2379            .await
2380            .unwrap();
2381
2382        tokio::time::sleep(Duration::from_millis(150)).await;
2383        driver.send_stdin(&id, b"hello\n".to_vec()).await.unwrap();
2384
2385        let deadline = std::time::Instant::now() + Duration::from_secs(5);
2386        loop {
2387            let meta = store.get_run(&id).await.unwrap().unwrap();
2388            if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2389                break;
2390            }
2391            if std::time::Instant::now() > deadline {
2392                panic!("run did not complete after stdin input");
2393            }
2394            tokio::time::sleep(Duration::from_millis(50)).await;
2395        }
2396
2397        let chunks = store.get_chunks(&id, &ChunkFilter::default()).await.unwrap();
2398        let raw: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
2399        let text = String::from_utf8_lossy(&raw);
2400        assert!(
2401            text.contains("got_hello"),
2402            "expected 'got_hello' in output, got: {text:?}"
2403        );
2404    }
2405
2406    /// `resize_run` must change the geometry the *child* sees, not just the
2407    /// master fd — so the assertion reads `stty size` from inside the PTY
2408    /// after the resize rather than inspecting the driver's own state.
2409    #[tokio::test]
2410    async fn resize_run_changes_geometry_the_child_sees() {
2411        let dir = tempfile::tempdir().unwrap();
2412        let store = open_store(&dir).await;
2413        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2414
2415        // Wait for a line on stdin, then report the geometry as of that moment.
2416        let id = driver
2417            .spawn_run(
2418                "read line && stty size",
2419                SpawnOpts {
2420                    cwd: "/tmp".into(),
2421                    stdin_enabled: true,
2422                    // Spawn at the default 80x24 so the assertion can't pass by
2423                    // accident if the resize is a no-op.
2424                    ..Default::default()
2425                },
2426            )
2427            .await
2428            .unwrap();
2429
2430        tokio::time::sleep(Duration::from_millis(150)).await;
2431        driver.resize_run(&id, 120, 40).await.unwrap();
2432        driver.send_stdin(&id, b"go\n".to_vec()).await.unwrap();
2433
2434        let deadline = std::time::Instant::now() + Duration::from_secs(5);
2435        loop {
2436            let meta = store.get_run(&id).await.unwrap().unwrap();
2437            if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2438                break;
2439            }
2440            if std::time::Instant::now() > deadline {
2441                panic!("run did not complete after stdin input");
2442            }
2443            tokio::time::sleep(Duration::from_millis(50)).await;
2444        }
2445
2446        let chunks = store.get_chunks(&id, &ChunkFilter::default()).await.unwrap();
2447        let raw: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
2448        let text = String::from_utf8_lossy(&raw);
2449        assert!(
2450            text.contains("40 120"),
2451            "expected resized geometry '40 120' in output, got: {text:?}"
2452        );
2453    }
2454
2455    /// A run that is not active on this driver (finished, or never existed) is
2456    /// `NotFound` rather than a panic — same contract as `send_stdin`.
2457    #[tokio::test]
2458    async fn resize_run_returns_not_found_after_exit() {
2459        let dir = tempfile::tempdir().unwrap();
2460        let store = open_store(&dir).await;
2461        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2462
2463        let id = driver
2464            .spawn_run("true", SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2465            .await
2466            .unwrap();
2467
2468        let deadline = std::time::Instant::now() + Duration::from_secs(5);
2469        loop {
2470            let meta = store.get_run(&id).await.unwrap().unwrap();
2471            if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2472                break;
2473            }
2474            if std::time::Instant::now() > deadline {
2475                panic!("run did not exit");
2476            }
2477            tokio::time::sleep(Duration::from_millis(50)).await;
2478        }
2479
2480        assert!(matches!(
2481            driver.resize_run(&id, 100, 30).await,
2482            Err(DriverError::NotFound(_))
2483        ));
2484    }
2485
2486    // ── Tier-2 side-channel log fd ────────────────────────────────────────────
2487
2488    /// Verify that a child writing a JSON-line to `YAH_LOG_PIPE` (via
2489    /// `printf ... >> $YAH_LOG_PIPE`) produces a shim event with the correct
2490    /// fields in the store.
2491    ///
2492    /// The child opens the FIFO path for writing — no fd inheritance needed.
2493    #[cfg(unix)]
2494    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
2495    async fn log_pipe_events_land_in_store() {
2496        use crate::store::EventFilter;
2497
2498        let dir = tempfile::tempdir().unwrap();
2499        let store = open_store(&dir).await;
2500        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2501
2502        // The shell writes one JSON-line to the FIFO by redirecting printf
2503        // output to the path stored in YAH_LOG_PIPE.
2504        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""#;
2505
2506        let id = driver
2507            .spawn_run(cmd, SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2508            .await
2509            .unwrap();
2510
2511        // Wait for run completion. Deadline is generous because parallel-test
2512        // load + the rt.block_on hops from the reader/log threads can slow
2513        // child-process scheduling.
2514        let deadline = std::time::Instant::now() + Duration::from_secs(20);
2515        loop {
2516            let meta = store.get_run(&id).await.unwrap().unwrap();
2517            if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
2518                break;
2519            }
2520            if std::time::Instant::now() > deadline {
2521                panic!("run did not complete in time");
2522            }
2523            tokio::time::sleep(Duration::from_millis(50)).await;
2524        }
2525
2526        // The log receiver thread drains after the lifecycle task drops the
2527        // write-end FdCloser; give it a brief moment.
2528        tokio::time::sleep(Duration::from_millis(500)).await;
2529
2530        let events = store.query_events(&id, &EventFilter::default()).await.unwrap();
2531        assert!(
2532            !events.is_empty(),
2533            "expected at least one shim event, got none"
2534        );
2535        let ev = events.iter().find(|e| e.target == "test.shim");
2536        let ev = ev.expect("event with target 'test.shim' not found");
2537        assert_eq!(ev.msg, "hello-from-pipe");
2538        assert_eq!(ev.level, crate::types::Level::Warn);
2539        assert!(
2540            matches!(&ev.source, crate::types::EventSource::Shim { lib, .. } if lib == "test-shim"),
2541            "unexpected source: {:?}",
2542            ev.source
2543        );
2544        assert_eq!(ev.fields.get("x"), Some(&serde_json::json!(42)));
2545    }
2546
2547    /// When `log_fd_enabled` is false, neither `YAH_TASK_RUN` nor
2548    /// `YAH_LOG_PIPE` are exported, and no shim events are written.
2549    #[cfg(unix)]
2550    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2551    async fn log_pipe_disabled_produces_no_events() {
2552        use crate::store::EventFilter;
2553
2554        let dir = tempfile::tempdir().unwrap();
2555        let store = open_store(&dir).await;
2556        let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2557
2558        // Try to write to YAH_LOG_PIPE; the conditional guards against
2559        // the variable being absent, so the command always exits 0.
2560        let cmd = r#"[ -n "$YAH_LOG_PIPE" ] && printf '{"level":"info","target":"t","msg":"m","fields":{}}\n' >> "$YAH_LOG_PIPE" || true"#;
2561
2562        let id = driver
2563            .spawn_run(
2564                cmd,
2565                SpawnOpts { cwd: "/tmp".into(), log_fd_enabled: false, ..Default::default() },
2566            )
2567            .await
2568            .unwrap();
2569
2570        let deadline = std::time::Instant::now() + Duration::from_secs(5);
2571        loop {
2572            let meta = store.get_run(&id).await.unwrap().unwrap();
2573            if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
2574                break;
2575            }
2576            if std::time::Instant::now() > deadline {
2577                panic!("run did not complete");
2578            }
2579            tokio::time::sleep(Duration::from_millis(50)).await;
2580        }
2581
2582        tokio::time::sleep(Duration::from_millis(100)).await;
2583
2584        let events = store.query_events(&id, &EventFilter::default()).await.unwrap();
2585        assert!(
2586            events.is_empty(),
2587            "expected no shim events when log_fd_enabled=false, got {}",
2588            events.len()
2589        );
2590    }
2591
2592    // ── Unattached-run reaper (R739-B12) ─────────────────────────────────────
2593
2594    /// The origin `yah build run` tags its relocated builds with. Spelled out
2595    /// here rather than imported: what the reaper must do is defined by the
2596    /// string on the wire, not by any constant this crate owns.
2597    const BUILD_RUN: &str = "build-run";
2598
2599    fn opted_in() -> Vec<String> {
2600        vec![BUILD_RUN.to_string()]
2601    }
2602
2603    async fn spawn_long_run(driver: &TaskDriver, origin: &str) -> TaskRunId {
2604        driver
2605            .spawn_run(
2606                "sleep 30",
2607                SpawnOpts {
2608                    cwd: "/tmp".into(),
2609                    origin: Some(origin.to_string()),
2610                    ..Default::default()
2611                },
2612            )
2613            .await
2614            .unwrap()
2615    }
2616
2617    async fn await_status(
2618        store: &Arc<TaskStore>,
2619        id: &TaskRunId,
2620        want: fn(&RunStatus) -> bool,
2621    ) -> RunStatus {
2622        let deadline = std::time::Instant::now() + Duration::from_secs(10);
2623        loop {
2624            let status = store.get_run(id).await.unwrap().unwrap().status;
2625            if want(&status) {
2626                return status;
2627            }
2628            if std::time::Instant::now() > deadline {
2629                panic!("run never reached the expected status, last={status:?}");
2630            }
2631            tokio::time::sleep(Duration::from_millis(25)).await;
2632        }
2633    }
2634
2635    /// The orphan this ticket exists for: the client is gone, so nothing polls
2636    /// the run, so the daemon must end it.
2637    #[tokio::test]
2638    async fn an_unpolled_run_of_an_opted_in_origin_is_reaped() {
2639        let dir = tempfile::tempdir().unwrap();
2640        let store = open_store(&dir).await;
2641        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2642
2643        let id = spawn_long_run(&driver, BUILD_RUN).await;
2644        tokio::time::sleep(Duration::from_millis(300)).await;
2645
2646        let reaped = driver
2647            .reap_unattached(Duration::from_millis(200), &opted_in())
2648            .await;
2649        assert_eq!(reaped, vec![id.clone()], "the unattached run should be reaped");
2650
2651        let status = await_status(&store, &id, |s| {
2652            matches!(s, RunStatus::Killed { .. } | RunStatus::Done { .. })
2653        })
2654        .await;
2655        assert!(
2656            matches!(status, RunStatus::Killed { .. }),
2657            "a reaped run ends Killed, got {status:?}"
2658        );
2659    }
2660
2661    /// The regression that protects a healthy long build: a client that is
2662    /// still polling keeps its run alive however long the build takes.
2663    #[tokio::test]
2664    async fn a_run_a_client_is_still_polling_is_never_reaped() {
2665        let dir = tempfile::tempdir().unwrap();
2666        let store = open_store(&dir).await;
2667        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2668
2669        let id = spawn_long_run(&driver, BUILD_RUN).await;
2670
2671        // Six polls across three idle windows — what `yah build run`'s tail
2672        // loop does, slowed down.
2673        for _ in 0..6 {
2674            tokio::time::sleep(Duration::from_millis(100)).await;
2675            driver.note_attached(&id);
2676            let reaped = driver
2677                .reap_unattached(Duration::from_millis(200), &opted_in())
2678                .await;
2679            assert!(reaped.is_empty(), "a polled run must survive, reaped {reaped:?}");
2680        }
2681
2682        assert!(
2683            matches!(
2684                store.get_run(&id).await.unwrap().unwrap().status,
2685                RunStatus::Running
2686            ),
2687            "the polled run should still be running"
2688        );
2689
2690        // And the moment the polling stops, it becomes reapable — same run,
2691        // so this pins the refresh rather than a missing origin match.
2692        tokio::time::sleep(Duration::from_millis(300)).await;
2693        let reaped = driver
2694            .reap_unattached(Duration::from_millis(200), &opted_in())
2695            .await;
2696        assert_eq!(reaped, vec![id], "a run that stopped being polled is reapable");
2697    }
2698
2699    /// The regression that protects real users' terminals. An interactive tile
2700    /// sits unpolled for hours by design and must never be touched, however
2701    /// long the reaper runs.
2702    #[tokio::test]
2703    async fn a_terminal_tile_is_never_reaped_however_long_it_idles() {
2704        let dir = tempfile::tempdir().unwrap();
2705        let store = open_store(&dir).await;
2706        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2707
2708        let id = spawn_long_run(&driver, "terminal").await;
2709        tokio::time::sleep(Duration::from_millis(300)).await;
2710
2711        for _ in 0..3 {
2712            let reaped = driver.reap_unattached(Duration::ZERO, &opted_in()).await;
2713            assert!(
2714                reaped.is_empty(),
2715                "a terminal tile is outside the opted-in origins, reaped {reaped:?}"
2716            );
2717            tokio::time::sleep(Duration::from_millis(50)).await;
2718        }
2719
2720        assert!(
2721            matches!(
2722                store.get_run(&id).await.unwrap().unwrap().status,
2723                RunStatus::Running
2724            ),
2725            "the terminal run must still be running"
2726        );
2727
2728        driver.kill_run(&id, Some(SIGKILL)).await.unwrap();
2729    }
2730
2731    /// A run with no origin at all — an ordinary `task.run` job — is outside
2732    /// every opt-in list, and an empty list reaps nothing.
2733    #[tokio::test]
2734    async fn an_origin_less_run_and_an_empty_opt_in_list_reap_nothing() {
2735        let dir = tempfile::tempdir().unwrap();
2736        let store = open_store(&dir).await;
2737        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2738
2739        let plain = driver
2740            .spawn_run("sleep 30", SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2741            .await
2742            .unwrap();
2743        let build = spawn_long_run(&driver, BUILD_RUN).await;
2744        tokio::time::sleep(Duration::from_millis(100)).await;
2745
2746        assert!(
2747            driver.reap_unattached(Duration::ZERO, &[]).await.is_empty(),
2748            "an empty opt-in list must reap nothing, not everything"
2749        );
2750        assert_eq!(
2751            driver.reap_unattached(Duration::ZERO, &opted_in()).await,
2752            vec![build],
2753            "only the opted-in origin is reapable"
2754        );
2755
2756        assert!(
2757            matches!(
2758                store.get_run(&plain).await.unwrap().unwrap().status,
2759                RunStatus::Running
2760            ),
2761            "the origin-less run must be untouched"
2762        );
2763        driver.kill_run(&plain, Some(SIGKILL)).await.unwrap();
2764    }
2765
2766    /// `note_attached` is observable, and the age it resets is what the sweep
2767    /// reads.
2768    #[tokio::test]
2769    async fn attached_age_resets_on_a_poll_and_is_none_for_a_foreign_run() {
2770        let dir = tempfile::tempdir().unwrap();
2771        let store = open_store(&dir).await;
2772        let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2773
2774        let id = spawn_long_run(&driver, BUILD_RUN).await;
2775        tokio::time::sleep(Duration::from_millis(150)).await;
2776        let aged = driver.attached_age(&id).expect("driver owns this run");
2777        assert!(aged >= Duration::from_millis(100), "age should have grown, got {aged:?}");
2778
2779        driver.note_attached(&id);
2780        let fresh = driver.attached_age(&id).unwrap();
2781        assert!(fresh < aged, "a poll resets the age: {fresh:?} vs {aged:?}");
2782
2783        assert!(
2784            driver.attached_age(&TaskRunId::new()).is_none(),
2785            "a run this driver does not own has no attachment age"
2786        );
2787
2788        driver.kill_run(&id, Some(SIGKILL)).await.unwrap();
2789    }
2790}