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 store
522 .update_status(
523 &meta.id,
524 &RunStatus::Lost {
525 reason: "daemon restarted while run was in-flight".to_string(),
526 },
527 )
528 .await?;
529 }
530 Ok(Self {
531 store,
532 active: Arc::new(Mutex::new(HashMap::new())),
533 channels,
534 })
535 }
536
537 /// Spawn `cmd` in a PTY and start capturing its output. Returns immediately
538 /// with the new [`TaskRunId`].
539 ///
540 /// A beholder is selected via `opts.beholder_select` (default `Auto`). When
541 /// a `Rewriter` beholder matches, its `adjust_argv` is applied to the
542 /// command before spawning and the diff is recorded on `beholder_status`.
543 /// When `opts.verbatim_output` is `true`, `Rewriter` beholders decline in
544 /// `Auto` mode, because something downstream renders these bytes and a
545 /// rewrite would change them.
546 ///
547 /// Output is written to the store as `Stream::Stdout` chunks (the PTY
548 /// kernel merges stdout and stderr). Signal handling and status updates
549 /// run in background tasks.
550 pub async fn spawn_run(&self, cmd: &str, opts: SpawnOpts) -> Result<TaskRunId, DriverError> {
551 let id = TaskRunId::new();
552 let started_at = unix_now_secs();
553 let started_at_ms: u64 = started_at.saturating_mul(1000);
554
555 // Attach a beholder (may rewrite argv and produce structured events).
556 // Resolve user drop-in directory: $YAH_BEHOLDERS_DIR or $HOME/.yah/beholders.
557 let user_dir = std::env::var_os("YAH_BEHOLDERS_DIR")
558 .map(std::path::PathBuf::from)
559 .or_else(|| {
560 std::env::var_os("HOME")
561 .map(|h| std::path::PathBuf::from(h).join(".yah/beholders"))
562 });
563 let registry = registry_with_user_beholders(user_dir.as_deref());
564 /* An explicit argv means the caller already chose the exact process
565 (an interactive login shell, say). Selecting a beholder there would
566 either do nothing — the rewritten argv is discarded on that path —
567 or record a rewrite that never happened, so we opt out honestly
568 instead. */
569 let select = if opts.argv.is_some() {
570 &BeholderSelect::None
571 } else {
572 &opts.beholder_select
573 };
574 let attach = registry.attach(cmd, select, opts.verbatim_output);
575 // Reconstruct the command from argv ONLY when a beholder actually
576 // rewrote it. `AttachResult.argv` is always populated — it is
577 // `resolve_argv(cmd)` even when nothing attached — so joining it
578 // unconditionally ran every run's command through a whitespace
579 // normalization nobody asked for: runs of spaces collapse and embedded
580 // newlines become spaces, which is silent corruption for a heredoc or
581 // any multi-line line. The caller's bytes go to the shell untouched
582 // unless a rewrite is the whole point.
583 let effective_cmd = match &attach.status.rewrite_added {
584 Some(added) if !added.is_empty() && !attach.argv.is_empty() => attach.argv.join(" "),
585 _ => cmd.to_string(),
586 };
587
588 self.store.insert_run(&TaskRunMeta {
589 id: id.clone(),
590 command: cmd.to_string(),
591 cwd: opts.cwd.clone(),
592 env: opts.env.clone(),
593 started_at,
594 status: RunStatus::Running,
595 label: opts.label.clone(),
596 initiator: opts.initiator.clone(),
597 beholder_status: Some(attach.status),
598 pinned: opts.pin,
599 origin: opts.origin.clone(),
600 /* R617-F6: stamp the OWNER, before the child exists. Written at
601 insert rather than after spawn so a crash between the two still
602 leaves the row attributable — an unattributed `Running` row is
603 exactly what the conservative arm of `StaleRunPolicy` has to
604 tombstone. */
605 host_pid: Some(std::process::id()),
606 }).await?;
607
608 // The program and argv both spawn modes exec. An explicit argv execs
609 // that program directly; otherwise the command line goes through `sh`
610 // so the caller's quoting, pipes and redirections mean what they say.
611 // An empty argv is a caller bug, not a request for an empty exec — fall
612 // back to the shell path rather than spawning nothing.
613 // R901-B2: `pipefail` is prepended HERE and not folded into
614 // `effective_cmd`, so `TaskRunMeta.command` keeps reading as the line
615 // the caller actually wrote. A run's recorded command is re-run by
616 // history and audited by agents against the relocation note; a prelude
617 // nobody asked for showing up in it would be the same class of lie as
618 // recording a beholder `rewrite=…` that never happened.
619 let (program, args): (String, Vec<String>) = match opts.argv.as_deref() {
620 Some([p, rest @ ..]) => (p.clone(), rest.to_vec()),
621 _ => {
622 let line = if opts.pipefail {
623 format!("{PIPEFAIL_PRELUDE}{effective_cmd}")
624 } else {
625 effective_cmd.clone()
626 };
627 ("sh".to_string(), vec!["-c".to_string(), line])
628 }
629 };
630
631 // ── Side-channel log FIFO (Tier 2 / yah-log shims) ──────────────────
632 //
633 // Create a named pipe (FIFO) so child processes can write structured
634 // events without touching stdout/stderr. We export its path via
635 // YAH_LOG_PIPE; no fd inheritance is involved, so portable-pty's
636 // close_random_fds() pre_exec hook doesn't interfere.
637 //
638 // The parent opens the FIFO twice:
639 // rfd — O_RDONLY|O_NONBLOCK, then cleared to blocking → read events
640 // wfd — O_WRONLY (wrapped in FdCloser) → keeps the FIFO alive until
641 // the lifecycle task drops it (after run completion), producing
642 // EOF for the receiver thread.
643 #[cfg(unix)]
644 let log_fifo: Option<(libc::c_int, FdCloser, std::path::PathBuf)> = if opts.log_fd_enabled {
645 let fifo_path = std::env::temp_dir().join(format!("yah-log-{}.fifo", id));
646 let path_cstr = match std::ffi::CString::new(fifo_path.to_string_lossy().as_bytes()) {
647 Ok(s) => s,
648 Err(_) => {
649 // Path contained a nul byte — extremely unlikely; skip FIFO.
650 return Err(DriverError::Io(std::io::Error::new(
651 std::io::ErrorKind::InvalidInput,
652 "log FIFO path contained nul byte",
653 )));
654 }
655 };
656 let mkfifo_ret = unsafe { libc::mkfifo(path_cstr.as_ptr(), 0o600) };
657 if mkfifo_ret != 0 {
658 None // FIFO creation failed; continue without side-channel
659 } else {
660 // Open read end without blocking (no writer yet).
661 let rfd = unsafe {
662 libc::open(path_cstr.as_ptr(), libc::O_RDONLY | libc::O_NONBLOCK)
663 };
664 if rfd < 0 {
665 let _ = unsafe { libc::unlink(path_cstr.as_ptr()) };
666 None
667 } else {
668 // Switch read end to blocking so reads yield proper data.
669 unsafe { libc::fcntl(rfd, libc::F_SETFL, 0) };
670 // Open write end — this succeeds immediately because rfd is open.
671 let wfd = unsafe {
672 libc::open(path_cstr.as_ptr(), libc::O_WRONLY)
673 };
674 if wfd < 0 {
675 unsafe { libc::close(rfd) };
676 let _ = unsafe { libc::unlink(path_cstr.as_ptr()) };
677 None
678 } else {
679 Some((rfd, FdCloser(wfd), fifo_path))
680 }
681 }
682 }
683 } else {
684 None
685 };
686
687 // The FIFO env, applied identically by both spawn modes.
688 #[cfg(unix)]
689 let fifo_env: Option<(String, String)> = log_fifo
690 .as_ref()
691 .map(|(_, _, path)| (id.to_string(), path.to_string_lossy().into_owned()));
692 #[cfg(not(unix))]
693 let fifo_env: Option<(String, String)> = None;
694
695 /* Spawn. The two modes differ only in what the child's stdio is
696 attached to, and everything downstream — reader pumps, lifecycle,
697 kill — is written against the uniform handles produced here:
698 `pid`, a `reap` closure that blocks until the child exits, and an
699 optional PTY master for resize / foreground-pid. */
700 let pid: u32;
701 let reap: Box<dyn FnOnce() -> Option<u32> + Send>;
702 let stdin_tx: Option<mpsc::Sender<Vec<u8>>>;
703 let master: Option<Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>>>;
704 // Each entry is one blocking source to pump into the store. The PTY
705 // yields a single merged stream; pipes yield stdout and stderr apart.
706 let mut sources: Vec<(Box<dyn Read + Send>, Stream)> = Vec::new();
707
708 if opts.pipe {
709 use std::process::{Command, Stdio};
710
711 let mut cmd = Command::new(&program);
712 cmd.args(&args);
713 cmd.current_dir(&opts.cwd);
714 for (k, v) in &opts.env {
715 cmd.env(k, v);
716 }
717 /* Deliberately NOT setting TERM. The PTY path forces
718 `xterm-256color` because a child on a terminal that claims no
719 terminal type degrades badly; a child on a pipe should see
720 whatever the daemon's own environment says, exactly as it would
721 under a non-interactive shell. Forcing a terminal type here is
722 how a pipe-mode run would talk itself back into colorizing. */
723 if let Some((run_id_env, fifo_path)) = &fifo_env {
724 cmd.env("YAH_TASK_RUN", run_id_env);
725 cmd.env("YAH_LOG_PIPE", fifo_path);
726 }
727 cmd.stdout(Stdio::piped());
728 cmd.stderr(Stdio::piped());
729 cmd.stdin(if opts.stdin_enabled { Stdio::piped() } else { Stdio::null() });
730
731 let mut child = cmd.spawn().map_err(DriverError::Io)?;
732 pid = child.id();
733
734 if let Some(out) = child.stdout.take() {
735 sources.push((Box::new(out), Stream::Stdout));
736 }
737 if let Some(err) = child.stderr.take() {
738 sources.push((Box::new(err), Stream::Stderr));
739 }
740
741 stdin_tx = child.stdin.take().map(|mut writer| {
742 let (tx, mut rx) = mpsc::channel::<Vec<u8>>(64);
743 task::spawn(async move {
744 use std::io::Write;
745 while let Some(bytes) = rx.recv().await {
746 let _ = writer.write_all(&bytes);
747 let _ = writer.flush();
748 }
749 });
750 tx
751 });
752
753 master = None;
754 reap = Box::new(move || child.wait().ok().and_then(|s| s.code()).map(|c| c as u32));
755 } else {
756 // Open PTY pair.
757 let pty_sys = native_pty_system();
758 let pair = pty_sys
759 .openpty(PtySize {
760 rows: opts.pty_rows,
761 cols: opts.pty_cols,
762 pixel_width: 0,
763 pixel_height: 0,
764 })
765 .map_err(|e| DriverError::Pty(e.to_string()))?;
766
767 // Clone reader before spawning so the fd is ready immediately.
768 let pty_reader = pair
769 .master
770 .try_clone_reader()
771 .map_err(|e| DriverError::Pty(e.to_string()))?;
772 sources.push((Box::new(pty_reader), Stream::Stdout));
773
774 // Optional stdin relay: take the writer before spawning the child.
775 stdin_tx = if opts.stdin_enabled {
776 let mut writer = pair
777 .master
778 .take_writer()
779 .map_err(|e| DriverError::Pty(e.to_string()))?;
780 let (tx, mut rx) = mpsc::channel::<Vec<u8>>(64);
781 task::spawn(async move {
782 use std::io::Write;
783 while let Some(bytes) = rx.recv().await {
784 let _ = writer.write_all(&bytes);
785 let _ = writer.flush();
786 }
787 });
788 Some(tx)
789 } else {
790 None
791 };
792
793 let mut cb = CommandBuilder::new(&program);
794 cb.args(&args);
795 cb.cwd(&opts.cwd);
796 for (k, v) in &opts.env {
797 cb.env(k, v);
798 }
799 cb.env("TERM", "xterm-256color");
800 if let Some((run_id_env, fifo_path)) = &fifo_env {
801 cb.env("YAH_TASK_RUN", run_id_env);
802 cb.env("YAH_LOG_PIPE", fifo_path);
803 }
804
805 let child = pair
806 .slave
807 .spawn_command(cb)
808 .map_err(|e| DriverError::Pty(e.to_string()))?;
809 // Drop the parent's slave handle so EOF propagates once the child exits.
810 drop(pair.slave);
811
812 pid = child.process_id().unwrap_or(0);
813
814 // Share the master between the lifecycle task (which must outlive
815 // `child.wait()` so the fd stays open) and `resize_run`.
816 let m: Arc<Mutex<Box<dyn portable_pty::MasterPty + Send>>> =
817 Arc::new(Mutex::new(pair.master));
818 master = Some(Arc::clone(&m));
819 reap = Box::new(move || {
820 let mut c = child;
821 let _m = m; // dropped after wait() returns, closing the PTY fd
822 c.wait().ok().map(|s| s.exit_code())
823 });
824 }
825
826 // ── FIFO: launch receiver thread; pass write-end holder to lifecycle ──
827 //
828 // The receiver thread reads until EOF. EOF arrives when ALL write-end
829 // holders close: the child's own writers (when it exits) plus the
830 // FdCloser we hand to the lifecycle task (which drops it after writing
831 // the terminal RunStatus). Events written before the last close are
832 // still drained by the receiver thread before it exits.
833 #[cfg(unix)]
834 let log_wfd_holder: Option<FdCloser> = if let Some((rfd, wfd, fifo_path)) = log_fifo {
835 let store_log = Arc::clone(&self.store);
836 let id_log = id.clone();
837 let rt = tokio::runtime::Handle::current();
838 // spawn_blocking: lets the runtime track this thread so the
839 // Handle::block_on calls inside have a worker to drive futures.
840 tokio::task::spawn_blocking(move || {
841 run_log_receiver(rt, store_log, id_log, rfd, fifo_path, started_at_ms);
842 });
843 Some(wfd)
844 } else {
845 None
846 };
847
848 // Channels.
849 let (kill_tx, kill_rx) = mpsc::channel::<KillRequest>(4);
850 let (reader_done_tx, reader_done_rx) = oneshot::channel::<()>();
851
852 /* Reader threads: child output → store chunks → beholder events. Each
853 runs on a dedicated OS thread because the reads are blocking. The
854 `ReaderDone` guard is shared across them, so the lifecycle's
855 reader-done signal fires only once every source has hit EOF — which
856 is what makes the two-pipe case correct without a reader count. */
857 {
858 let done = Arc::new(ReaderDone(Some(reader_done_tx)));
859 /* The beholder goes to stdout only. It parses a structured
860 protocol (cargo's JSON, say) that the child writes to stdout by
861 definition, and there is exactly one of it — handing the same
862 instance to two threads would need a lock for no gain, and
863 feeding it stderr would make `unknown_format_reason` fire on
864 human-readable diagnostics it was never meant to see. */
865 let mut beholder = attach.beholder;
866 for (reader, stream) in sources {
867 spawn_output_pump(
868 reader,
869 stream,
870 Arc::clone(&self.store),
871 id.clone(),
872 started_at_ms,
873 self.channels.output.clone(),
874 if stream == Stream::Stdout { beholder.take() } else { None },
875 Arc::clone(&done),
876 );
877 }
878 }
879
880 // Lifecycle task: monitor kill requests, wait for exit, update status.
881 // The task also holds the log FIFO write-end closer (if any) so that
882 // EOF propagates to the receiver thread after RunStatus is written.
883 {
884 let store_l = Arc::clone(&self.store);
885 let active_l = Arc::clone(&self.active);
886 let id_l = id.clone();
887 let completion_tx_l = self.channels.completion.clone();
888 #[cfg(unix)]
889 let wfd_l = log_wfd_holder;
890 task::spawn(async move {
891 run_lifecycle(
892 store_l,
893 active_l,
894 id_l,
895 pid,
896 reap,
897 kill_rx,
898 reader_done_rx,
899 completion_tx_l,
900 #[cfg(unix)]
901 wfd_l,
902 )
903 .await;
904 });
905 }
906
907 self.active
908 .lock()
909 .unwrap()
910 .insert(
911 id.to_string(),
912 RunControl {
913 kill_tx,
914 stdin_tx,
915 master,
916 origin: opts.origin.clone(),
917 last_attached_at: Instant::now(),
918 },
919 );
920
921 Ok(id)
922 }
923
924 /// Resize a running task's PTY and deliver `SIGWINCH` to the foreground
925 /// process group (portable-pty's `resize` does the ioctl, which is what
926 /// signals the child).
927 ///
928 /// Returns `DriverError::NotFound` when the run is not active on this
929 /// driver instance — the same contract as [`TaskDriver::send_stdin`] — and
930 /// also when it is active but was spawned in [`SpawnOpts::pipe`] mode, which
931 /// has no terminal to resize.
932 pub async fn resize_run(
933 &self,
934 id: &TaskRunId,
935 cols: u16,
936 rows: u16,
937 ) -> Result<(), DriverError> {
938 let master = self
939 .active
940 .lock()
941 .unwrap()
942 .get(&id.to_string())
943 .and_then(|c| c.master.as_ref().map(Arc::clone));
944
945 match master {
946 Some(m) => {
947 let size = PtySize { rows, cols, pixel_width: 0, pixel_height: 0 };
948 m.lock()
949 .unwrap()
950 .resize(size)
951 .map_err(|e| DriverError::Pty(e.to_string()))
952 }
953 None => Err(DriverError::NotFound(id.to_string())),
954 }
955 }
956
957 /// The pid of the run's *foreground* process — the leader of the process
958 /// group the PTY currently gives the keyboard to.
959 ///
960 /// For a shell tile that is the shell itself while it sits at a prompt,
961 /// and the command the operator is running while one is in flight. That
962 /// distinction is the whole point: asking the spawned child would report
963 /// the shell forever, so anything derived from this pid (a live cwd probe,
964 /// a "what is this pane doing" label) would answer for the wrong process.
965 ///
966 /// `None` when the run is not active on this driver instance, when it was
967 /// spawned in [`SpawnOpts::pipe`] mode (no controlling terminal, so no
968 /// foreground process group to read), or when the platform has no notion of
969 /// a foreground process group.
970 pub fn foreground_pid(&self, id: &TaskRunId) -> Option<u32> {
971 let master = self
972 .active
973 .lock()
974 .unwrap()
975 .get(&id.to_string())
976 .and_then(|c| c.master.as_ref().map(Arc::clone))?;
977 #[cfg(unix)]
978 {
979 let pid = master.lock().unwrap().process_group_leader()?;
980 u32::try_from(pid).ok()
981 }
982 #[cfg(not(unix))]
983 {
984 let _ = master;
985 None
986 }
987 }
988
989 /// Send `signal` to a running task. Defaults to SIGTERM (15).
990 ///
991 /// For SIGTERM, the driver waits up to 5 seconds for the process to exit
992 /// before escalating to SIGKILL. Returns `DriverError::NotFound` if the
993 /// run is not active (already exited or launched on a different driver
994 /// instance).
995 pub async fn kill_run(&self, id: &TaskRunId, signal: Option<i32>) -> Result<(), DriverError> {
996 let kill_tx = self
997 .active
998 .lock()
999 .unwrap()
1000 .get(&id.to_string())
1001 .map(|c| c.kill_tx.clone());
1002
1003 match kill_tx {
1004 Some(tx) => tx
1005 .send(KillRequest { signal: signal.unwrap_or(SIGTERM) })
1006 .await
1007 .map_err(|_| DriverError::NotFound(id.to_string())),
1008 None => Err(DriverError::NotFound(id.to_string())),
1009 }
1010 }
1011
1012 /// Write bytes to the stdin of a running task (requires `stdin_enabled`).
1013 pub async fn send_stdin(&self, id: &TaskRunId, bytes: Vec<u8>) -> Result<(), DriverError> {
1014 let stdin_tx = self
1015 .active
1016 .lock()
1017 .unwrap()
1018 .get(&id.to_string())
1019 .and_then(|c| c.stdin_tx.clone());
1020
1021 match stdin_tx {
1022 Some(tx) => tx
1023 .send(bytes)
1024 .await
1025 .map_err(|_| DriverError::NotFound(id.to_string())),
1026 None => Err(DriverError::NotFound(id.to_string())),
1027 }
1028 }
1029
1030 /// R739-B12 — record that a client just looked at this run.
1031 ///
1032 /// A no-op for a run this driver does not own (already finished, or
1033 /// spawned by another process against the same store): attachment only
1034 /// means anything for a run something here could still signal.
1035 pub fn note_attached(&self, id: &TaskRunId) {
1036 if let Some(control) = self.active.lock().unwrap().get_mut(&id.to_string()) {
1037 control.last_attached_at = Instant::now();
1038 }
1039 }
1040
1041 /// How long ago a client last looked at `id`, or `None` when this driver
1042 /// does not own the run. The observable half of [`Self::note_attached`].
1043 pub fn attached_age(&self, id: &TaskRunId) -> Option<Duration> {
1044 self.active
1045 .lock()
1046 .unwrap()
1047 .get(&id.to_string())
1048 .map(|c| c.last_attached_at.elapsed())
1049 }
1050
1051 /// R739-B12 — SIGTERM every run of an opted-in origin that no client has
1052 /// looked at for `idle`. Returns the runs it signalled.
1053 ///
1054 /// This exists because a run outlives the client that asked for it. When
1055 /// `yah build run` is SIGKILLed — its harness dies, the terminal goes away
1056 /// — the cargo it relocated into the daemon keeps compiling with nobody
1057 /// attached, holding the build-directory lock until a human finds the pid.
1058 /// That happened on 2026-08-28 and stalled a whole camp for ~30 minutes.
1059 /// R739-B9 closed every give-up the client is *alive* to make; this closes
1060 /// the one it is not.
1061 ///
1062 /// **Not [`StaleRunPolicy`], and not that policy on a timer.** The policy
1063 /// is a construction-time reconciliation of rows a *previous process*
1064 /// left behind: it decides on `host_pid`, only ever calls
1065 /// `store.update_status`, and tombstones any run outside its origin list
1066 /// outright — so running it periodically would mark every in-flight run of
1067 /// an un-adopted origin `Lost` while it compiles perfectly well, and would
1068 /// still never signal the process that is the actual problem. This is the
1069 /// opposite shape: it decides on *attachment*, it signals, and it touches
1070 /// nothing outside `origins`.
1071 ///
1072 /// `origins` is an opt-in list precisely because most runs must never be
1073 /// reaped on this rule. An interactive terminal tile is legitimately
1074 /// unpolled for hours, and killing one would be a far worse bug than the
1075 /// orphan this prevents — so an empty list reaps nothing at all, rather
1076 /// than meaning "every origin" the way [`StaleRunPolicy`]'s list does.
1077 pub async fn reap_unattached(&self, idle: Duration, origins: &[String]) -> Vec<TaskRunId> {
1078 if origins.is_empty() {
1079 return Vec::new();
1080 }
1081 let candidates: Vec<TaskRunId> = {
1082 let active = self.active.lock().unwrap();
1083 active
1084 .iter()
1085 .filter(|(_, c)| {
1086 c.origin
1087 .as_deref()
1088 .is_some_and(|o| origins.iter().any(|want| want == o))
1089 && c.last_attached_at.elapsed() >= idle
1090 })
1091 .filter_map(|(id, _)| id.parse::<TaskRunId>().ok())
1092 .collect()
1093 };
1094
1095 let mut reaped = Vec::new();
1096 for id in candidates {
1097 // SIGTERM, not SIGKILL: `kill_run` gives the child the same 5s
1098 // grace a `task.kill` from a live client would, then escalates.
1099 // A terminal status is also what releases the run's admission
1100 // enrollment (R739-F7), so a reaped run frees the build key.
1101 if self.kill_run(&id, None).await.is_ok() {
1102 reaped.push(id);
1103 }
1104 }
1105 reaped
1106 }
1107}
1108
1109// ─── Log fd receiver ─────────────────────────────────────────────────────────
1110
1111/// Read JSON-lines from the side-channel FIFO read end and store them as
1112/// [`EventSource::Shim`] events.
1113///
1114/// Runs on a dedicated OS thread; exits when the read end sees EOF. EOF
1115/// arrives after both the child process AND the lifecycle task have closed
1116/// their write ends of the FIFO. The FIFO file is deleted on exit.
1117#[cfg(unix)]
1118fn run_log_receiver(
1119 rt: tokio::runtime::Handle,
1120 store: Arc<TaskStore>,
1121 run_id: TaskRunId,
1122 read_fd: libc::c_int,
1123 fifo_path: std::path::PathBuf,
1124 started_at_ms: u64,
1125) {
1126 use std::io::BufRead;
1127 use std::os::unix::io::FromRawFd;
1128
1129 // SAFETY: `read_fd` is a valid, open FIFO fd handed exclusively to this
1130 // thread. `File` takes ownership and closes the fd on drop.
1131 let file = unsafe { std::fs::File::from_raw_fd(read_fd) };
1132 let reader = std::io::BufReader::new(file);
1133
1134 for line in reader.lines() {
1135 let line = match line {
1136 Ok(l) => l,
1137 Err(_) => break,
1138 };
1139 let trimmed = line.trim();
1140 if trimmed.is_empty() {
1141 continue;
1142 }
1143 let rec: ShimRecord = match serde_json::from_str(trimmed) {
1144 Ok(r) => r,
1145 Err(_) => continue, // skip malformed lines silently
1146 };
1147 let level = rec.level.parse::<crate::types::Level>().unwrap_or(crate::types::Level::Info);
1148 let source = crate::types::EventSource::Shim {
1149 lib: rec.lib.unwrap_or_else(|| "unknown".to_string()),
1150 version: rec.lib_version.unwrap_or_else(|| "0.0.0".to_string()),
1151 };
1152 let fields = if rec.fields.is_object() {
1153 rec.fields
1154 } else {
1155 serde_json::Value::Object(Default::default())
1156 };
1157 let offset = elapsed_ms(started_at_ms);
1158 let _ = rt.block_on(store.append_event(
1159 &run_id,
1160 offset,
1161 level,
1162 &rec.target,
1163 &rec.msg,
1164 &fields,
1165 None,
1166 &source,
1167 ));
1168 }
1169
1170 // Clean up the FIFO file now that the receiver has drained.
1171 let _ = std::fs::remove_file(&fifo_path);
1172}
1173
1174// ─── Lifecycle task ───────────────────────────────────────────────────────────
1175
1176/// Pump one blocking output source into the store, tapping and beholding on the
1177/// way past.
1178///
1179/// Split out of `spawn_run` for R739-F6: a PTY run has one source and a piped
1180/// run has two, and the only thing that differs between them is which [`Stream`]
1181/// the chunks are stored under. `done` is the shared [`ReaderDone`] guard —
1182/// dropping it here, after `on_done`, is what tells the lifecycle this source is
1183/// finished.
1184#[allow(clippy::too_many_arguments)]
1185fn spawn_output_pump(
1186 reader: Box<dyn Read + Send>,
1187 stream: Stream,
1188 store: Arc<TaskStore>,
1189 id: TaskRunId,
1190 started_at_ms: u64,
1191 output_tx: Option<mpsc::UnboundedSender<OutputChunk>>,
1192 beholder: Option<Box<dyn crate::beholders::Beholder>>,
1193 done: Arc<ReaderDone>,
1194) {
1195 let rt = tokio::runtime::Handle::current();
1196 tokio::task::spawn_blocking(move || {
1197 let _done = done;
1198 let mut beholder = beholder;
1199 let mut buf = [0u8; READ_BUF_SIZE];
1200 let mut reader = reader;
1201 loop {
1202 match reader.read(&mut buf) {
1203 Ok(0) | Err(_) => break,
1204 Ok(n) => {
1205 let offset = elapsed_ms(started_at_ms);
1206 let append_res =
1207 rt.block_on(store.append_chunk(&id, offset, stream, &buf[..n]));
1208 if let Ok(seq) = append_res {
1209 /* Both the tap and the beholder want the same owned
1210 chunk; build it once, and only when someone is
1211 listening. */
1212 let chunk = (output_tx.is_some() || beholder.is_some()).then(|| {
1213 OutputChunk {
1214 run_id: id.clone(),
1215 seq,
1216 offset_ms: offset,
1217 stream,
1218 bytes: buf[..n].to_vec(),
1219 }
1220 });
1221 /* Tap first: it feeds live views, where latency is
1222 visible to a human. Send failure means the host
1223 dropped its receiver — never fatal. */
1224 if let (Some(tx), Some(c)) = (&output_tx, &chunk) {
1225 let _ = tx.send(c.clone());
1226 }
1227 let mut detach_beholder = false;
1228 if let (Some(b), Some(chunk)) = (beholder.as_mut(), &chunk) {
1229 for ev in b.parse_chunk(chunk) {
1230 let _ = rt.block_on(store.append_event(
1231 &ev.run_id,
1232 ev.offset_ms,
1233 ev.level,
1234 &ev.target,
1235 &ev.msg,
1236 &ev.fields,
1237 ev.anchor.as_ref().map(|a| a.seq),
1238 &ev.source,
1239 ));
1240 }
1241 if let Some(reason) = b.unknown_format_reason() {
1242 let new_status =
1243 BeholderStatus::unknown_format_with_reason(b.name(), reason);
1244 let _ =
1245 rt.block_on(store.update_beholder_status(&id, &new_status));
1246 detach_beholder = true;
1247 }
1248 }
1249 if detach_beholder {
1250 beholder = None;
1251 }
1252 }
1253 }
1254 }
1255 }
1256 if let Some(ref mut b) = beholder {
1257 let final_offset = elapsed_ms(started_at_ms);
1258 for ev in b.on_done(&id, final_offset) {
1259 let _ = rt.block_on(store.append_event(
1260 &ev.run_id,
1261 ev.offset_ms,
1262 ev.level,
1263 &ev.target,
1264 &ev.msg,
1265 &ev.fields,
1266 ev.anchor.as_ref().map(|a| a.seq),
1267 &ev.source,
1268 ));
1269 }
1270 if let Some(reason) = b.unknown_format_reason() {
1271 let new_status = BeholderStatus::unknown_format_with_reason(b.name(), reason);
1272 let _ = rt.block_on(store.update_beholder_status(&id, &new_status));
1273 }
1274 }
1275 });
1276}
1277
1278#[allow(clippy::too_many_arguments)]
1279async fn run_lifecycle(
1280 store: Arc<TaskStore>,
1281 active: Arc<Mutex<HashMap<String, RunControl>>>,
1282 id: TaskRunId,
1283 pid: u32,
1284 // `reap` blocks until the child exits and yields its exit code. It owns
1285 // whatever the spawn mode has to keep alive across the wait — for a PTY run
1286 // that includes the master fd, which must outlive `wait()`.
1287 reap: Box<dyn FnOnce() -> Option<u32> + Send>,
1288 mut kill_rx: mpsc::Receiver<KillRequest>,
1289 reader_done_rx: oneshot::Receiver<()>,
1290 completion_tx: Option<tokio::sync::mpsc::UnboundedSender<(TaskRunId, RunStatus)>>,
1291 // Holds the write end of the log FIFO open until this task completes.
1292 // Dropping it produces EOF for the receiver thread, which happens after
1293 // the terminal RunStatus is written below.
1294 #[cfg(unix)]
1295 _log_wfd: Option<FdCloser>,
1296) {
1297 // Pin the reader-done future so it can be polled by reference in
1298 // nested select! arms without consuming ownership.
1299 let reader_done = async { reader_done_rx.await.ok(); };
1300 tokio::pin!(reader_done);
1301
1302 let sent_signal: Option<i32>;
1303
1304 tokio::select! {
1305 req = kill_rx.recv() => {
1306 match req {
1307 Some(KillRequest { signal }) => {
1308 send_unix_signal(pid, signal);
1309 if signal == SIGKILL {
1310 sent_signal = Some(SIGKILL);
1311 } else {
1312 // Grace period: give the process a chance to exit cleanly.
1313 tokio::select! {
1314 _ = &mut reader_done => {
1315 // Exited within grace — no SIGKILL needed.
1316 sent_signal = Some(signal);
1317 }
1318 _ = tokio::time::sleep(DEFAULT_GRACE) => {
1319 // Grace expired — escalate.
1320 send_unix_signal(pid, SIGKILL);
1321 sent_signal = Some(SIGKILL);
1322 }
1323 }
1324 }
1325 }
1326 // kill_tx dropped (driver shutting down) — force kill.
1327 None => {
1328 send_unix_signal(pid, SIGKILL);
1329 sent_signal = Some(SIGKILL);
1330 }
1331 }
1332 }
1333 _ = &mut reader_done => {
1334 sent_signal = None;
1335 }
1336 }
1337
1338 // Reap the child (blocking) on a dedicated thread-pool slot. For a PTY run
1339 // the closure also owns our master handle, so the fd outlives the wait; the
1340 // matching `RunControl` (removed from `active` below) holds the other `Arc`,
1341 // so the fd actually closes once both are gone.
1342 let exit_code = task::spawn_blocking(reap).await.ok().flatten();
1343
1344 let ended_at = unix_now_secs();
1345 let status = match sent_signal {
1346 Some(sig) => RunStatus::Killed { signal: sig, ended_at },
1347 None => match exit_code {
1348 Some(code) => RunStatus::Done { exit_code: code as i32, ended_at },
1349 None => RunStatus::Lost {
1350 reason: "process exited without an exit code".to_string(),
1351 },
1352 },
1353 };
1354
1355 /* Losing this write is not cosmetic: the run stays `Running` in the store
1356 forever and every reader — tail loops, the terminal UI, the next
1357 daemon's Lost-on-disappear sweep — believes a dead process is alive.
1358 `update_status` already retries through lock contention, so a failure
1359 here is terminal and worth saying out loud. */
1360 if let Err(e) = store.update_status(&id, &status).await {
1361 eprintln!("[yah task-runs] failed to record terminal status for run {id}: {e}");
1362 }
1363 if let Some(ref tx) = completion_tx {
1364 let _ = tx.send((id.clone(), status));
1365 }
1366 active.lock().unwrap().remove(&id.to_string());
1367}
1368
1369// ─── Helpers ──────────────────────────────────────────────────────────────────
1370
1371fn send_unix_signal(pid: u32, signal: i32) {
1372 #[cfg(unix)]
1373 unsafe {
1374 libc::kill(pid as libc::pid_t, signal);
1375 }
1376 // On non-Unix platforms signal delivery is not implemented here.
1377}
1378
1379fn unix_now_secs() -> u64 {
1380 SystemTime::now()
1381 .duration_since(UNIX_EPOCH)
1382 .unwrap_or_default()
1383 .as_secs()
1384}
1385
1386fn elapsed_ms(started_at_ms: u64) -> u32 {
1387 let now_ms = SystemTime::now()
1388 .duration_since(UNIX_EPOCH)
1389 .unwrap_or_default()
1390 .as_millis() as u64;
1391 now_ms.saturating_sub(started_at_ms).min(u32::MAX as u64) as u32
1392}
1393
1394// ─── Tests ────────────────────────────────────────────────────────────────────
1395
1396#[cfg(test)]
1397mod tests {
1398 use super::*;
1399 use crate::store::ChunkFilter;
1400
1401 async fn open_store(dir: &tempfile::TempDir) -> Arc<TaskStore> {
1402 Arc::new(TaskStore::open(&dir.path().join("tr.turso")).await.unwrap())
1403 }
1404
1405 // ── Lost-on-disappear (pure store, no PTY) ────────────────────────────────
1406
1407 #[tokio::test]
1408 async fn lost_on_disappear_marks_stale_running_runs() {
1409 let dir = tempfile::tempdir().unwrap();
1410 let store = open_store(&dir).await;
1411
1412 // Simulate a run left in "Running" state by a prior daemon.
1413 let stale_id = TaskRunId::new();
1414 store
1415 .insert_run(&TaskRunMeta {
1416 id: stale_id.clone(),
1417 command: "sleep 9999".to_string(),
1418 cwd: "/tmp".into(),
1419 env: vec![],
1420 started_at: unix_now_secs() - 60,
1421 status: RunStatus::Running,
1422 label: None,
1423 initiator: Initiator::Human { camp: "test".to_string() },
1424 beholder_status: None,
1425 pinned: false,
1426 origin: None,
1427 host_pid: None,
1428 })
1429 .await
1430 .unwrap();
1431
1432 // Creating a new driver must mark stale runs Lost.
1433 let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1434
1435 let meta = store.get_run(&stale_id).await.unwrap().unwrap();
1436 assert!(
1437 matches!(meta.status, RunStatus::Lost { .. }),
1438 "stale run should be Lost, got {:?}",
1439 meta.status
1440 );
1441 }
1442
1443 // ── Stale-run policy (R617-F6) ───────────────────────────────────────────
1444
1445 /// Plant a `Running` row as if some other process had spawned it.
1446 async fn plant_running(
1447 store: &Arc<TaskStore>,
1448 origin: Option<&str>,
1449 host_pid: Option<u32>,
1450 ) -> TaskRunId {
1451 let id = TaskRunId::new();
1452 store
1453 .insert_run(&TaskRunMeta {
1454 id: id.clone(),
1455 command: "sleep 9999".to_string(),
1456 cwd: "/tmp".into(),
1457 env: vec![],
1458 started_at: unix_now_secs() - 60,
1459 status: RunStatus::Running,
1460 label: None,
1461 initiator: Initiator::Human {
1462 camp: "test".to_string(),
1463 },
1464 beholder_status: None,
1465 pinned: false,
1466 origin: origin.map(str::to_string),
1467 host_pid,
1468 })
1469 .await
1470 .unwrap();
1471 id
1472 }
1473
1474 async fn is_lost(store: &Arc<TaskStore>, id: &TaskRunId) -> bool {
1475 matches!(
1476 store.get_run(id).await.unwrap().unwrap().status,
1477 RunStatus::Lost { .. }
1478 )
1479 }
1480
1481 fn adopt_terminal() -> StaleRunPolicy {
1482 StaleRunPolicy::AdoptLiveHosts {
1483 origins: vec!["terminal".to_string()],
1484 }
1485 }
1486
1487 /// The property the whole ticket exists for: attaching to a store must not
1488 /// declare another live process's shell dead.
1489 #[tokio::test]
1490 async fn a_run_owned_by_a_live_host_survives_a_new_driver() {
1491 let dir = tempfile::tempdir().unwrap();
1492 let store = open_store(&dir).await;
1493 // Our own pid is by definition a live process, and is the cheapest
1494 // honest stand-in for "a peer that is still running".
1495 let id = plant_running(&store, Some("terminal"), Some(std::process::id())).await;
1496
1497 let _driver = TaskDriver::with_config(
1498 Arc::clone(&store),
1499 DriverChannels::default(),
1500 adopt_terminal(),
1501 )
1502 .await
1503 .unwrap();
1504
1505 assert!(
1506 !is_lost(&store, &id).await,
1507 "a terminal run whose owner is alive must stay Running — \
1508 tombstoning it is what made a surviving shell read as dead"
1509 );
1510 }
1511
1512 /// The other half: a genuinely abandoned shell must still be tombstoned,
1513 /// or a crashed host leaves permanent zombie tiles.
1514 #[tokio::test]
1515 async fn a_run_whose_host_is_gone_is_still_tombstoned() {
1516 let dir = tempfile::tempdir().unwrap();
1517 let store = open_store(&dir).await;
1518 // Reaped in-test, so the pid is real-but-dead rather than guessed.
1519 let dead_pid = {
1520 let child = std::process::Command::new("true").spawn().unwrap();
1521 let pid = child.id();
1522 let mut child = child;
1523 let _ = child.wait();
1524 pid
1525 };
1526 let id = plant_running(&store, Some("terminal"), Some(dead_pid)).await;
1527
1528 let _driver = TaskDriver::with_config(
1529 Arc::clone(&store),
1530 DriverChannels::default(),
1531 adopt_terminal(),
1532 )
1533 .await
1534 .unwrap();
1535
1536 assert!(
1537 is_lost(&store, &id).await,
1538 "pid {dead_pid} was reaped; its run has no owner left and must be Lost"
1539 );
1540 }
1541
1542 /// The exemption is narrowed by origin, so ordinary jobs keep the old rule
1543 /// even when their owner happens to still be alive — an in-flight `cargo
1544 /// build` whose driver is gone has nobody left to record its exit.
1545 #[tokio::test]
1546 async fn a_non_matching_origin_is_tombstoned_even_with_a_live_host() {
1547 let dir = tempfile::tempdir().unwrap();
1548 let store = open_store(&dir).await;
1549 let job = plant_running(&store, None, Some(std::process::id())).await;
1550 let other = plant_running(&store, Some("gnome"), Some(std::process::id())).await;
1551
1552 let _driver = TaskDriver::with_config(
1553 Arc::clone(&store),
1554 DriverChannels::default(),
1555 adopt_terminal(),
1556 )
1557 .await
1558 .unwrap();
1559
1560 assert!(is_lost(&store, &job).await, "an origin-less job is not exempt");
1561 assert!(
1562 is_lost(&store, &other).await,
1563 "an origin outside the list is not exempt"
1564 );
1565 }
1566
1567 /// A row written before `host_pid` existed reads back `None`. Unknown
1568 /// ownership must fall back to the old behaviour rather than stranding the
1569 /// run `Running` forever.
1570 #[tokio::test]
1571 async fn an_unattributed_run_is_tombstoned() {
1572 let dir = tempfile::tempdir().unwrap();
1573 let store = open_store(&dir).await;
1574 let id = plant_running(&store, Some("terminal"), None).await;
1575
1576 let _driver = TaskDriver::with_config(
1577 Arc::clone(&store),
1578 DriverChannels::default(),
1579 adopt_terminal(),
1580 )
1581 .await
1582 .unwrap();
1583
1584 assert!(is_lost(&store, &id).await);
1585 }
1586
1587 /// `TaskDriver::new` must not have quietly changed behaviour — every
1588 /// existing embedder still gets Lost-on-disappear.
1589 #[tokio::test]
1590 async fn the_default_policy_is_still_lost_on_disappear() {
1591 assert_eq!(StaleRunPolicy::default(), StaleRunPolicy::LostOnDisappear);
1592
1593 let dir = tempfile::tempdir().unwrap();
1594 let store = open_store(&dir).await;
1595 let id = plant_running(&store, Some("terminal"), Some(std::process::id())).await;
1596
1597 let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1598
1599 assert!(
1600 is_lost(&store, &id).await,
1601 "the default must tombstone regardless of origin or owner liveness"
1602 );
1603 }
1604
1605 /// The owner is recorded by `spawn_run` itself, not by the caller — the
1606 /// policy is worthless if rows arrive unattributed.
1607 #[tokio::test]
1608 async fn spawn_run_stamps_this_process_as_the_owner() {
1609 let dir = tempfile::tempdir().unwrap();
1610 let store = open_store(&dir).await;
1611 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1612
1613 let id = driver
1614 .spawn_run(
1615 "true",
1616 SpawnOpts {
1617 cwd: "/tmp".into(),
1618 origin: Some("terminal".to_string()),
1619 ..Default::default()
1620 },
1621 )
1622 .await
1623 .unwrap();
1624
1625 let meta = store.get_run(&id).await.unwrap().unwrap();
1626 assert_eq!(meta.host_pid, Some(std::process::id()));
1627 }
1628
1629 #[tokio::test]
1630 async fn new_driver_does_not_touch_completed_runs() {
1631 let dir = tempfile::tempdir().unwrap();
1632 let store = open_store(&dir).await;
1633
1634 let done_id = TaskRunId::new();
1635 store
1636 .insert_run(&TaskRunMeta {
1637 id: done_id.clone(),
1638 command: "true".to_string(),
1639 cwd: "/tmp".into(),
1640 env: vec![],
1641 started_at: unix_now_secs() - 10,
1642 status: RunStatus::Running,
1643 label: None,
1644 initiator: Initiator::Human { camp: "test".to_string() },
1645 beholder_status: None,
1646 pinned: false,
1647 origin: None,
1648 host_pid: None,
1649 })
1650 .await
1651 .unwrap();
1652 store
1653 .update_status(&done_id, &RunStatus::Done { exit_code: 0, ended_at: unix_now_secs() })
1654 .await
1655 .unwrap();
1656
1657 let _driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1658
1659 let meta = store.get_run(&done_id).await.unwrap().unwrap();
1660 assert!(
1661 matches!(meta.status, RunStatus::Done { .. }),
1662 "completed run must not be touched"
1663 );
1664 }
1665
1666 // ── PTY spawn + capture ───────────────────────────────────────────────────
1667
1668 #[tokio::test]
1669 async fn spawn_echo_and_read_chunks() {
1670 let dir = tempfile::tempdir().unwrap();
1671 let store = open_store(&dir).await;
1672 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1673
1674 let id = driver
1675 .spawn_run(
1676 "echo hello_world",
1677 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1678 )
1679 .await
1680 .unwrap();
1681
1682 // Wait for the run to complete (poll status up to 5 s).
1683 let deadline = std::time::Instant::now() + Duration::from_secs(5);
1684 loop {
1685 let meta = store.get_run(&id).await.unwrap().unwrap();
1686 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1687 break;
1688 }
1689 if std::time::Instant::now() > deadline {
1690 panic!("run did not complete in time, status={:?}", meta.status);
1691 }
1692 tokio::time::sleep(Duration::from_millis(50)).await;
1693 }
1694
1695 // Chunks must contain "hello_world".
1696 let chunks = store
1697 .get_chunks(&id, &ChunkFilter::default())
1698 .await
1699 .unwrap();
1700 let output: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
1701 let text = String::from_utf8_lossy(&output);
1702 assert!(
1703 text.contains("hello_world"),
1704 "expected 'hello_world' in output, got: {text:?}"
1705 );
1706
1707 let meta = store.get_run(&id).await.unwrap().unwrap();
1708 assert!(
1709 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
1710 "expected Done(0), got {:?}",
1711 meta.status
1712 );
1713 }
1714
1715 // ── Pipe mode (R739-F6) ───────────────────────────────────────────────────
1716
1717 /// Run `cmd` to completion and return its stored chunks.
1718 async fn run_to_completion(
1719 store: &Arc<TaskStore>,
1720 driver: &TaskDriver,
1721 cmd: &str,
1722 opts: SpawnOpts,
1723 ) -> Vec<OutputChunk> {
1724 let id = driver.spawn_run(cmd, opts).await.unwrap();
1725 let deadline = std::time::Instant::now() + Duration::from_secs(10);
1726 loop {
1727 let meta = store.get_run(&id).await.unwrap().unwrap();
1728 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1729 break;
1730 }
1731 if std::time::Instant::now() > deadline {
1732 panic!("run did not complete in time, status={:?}", meta.status);
1733 }
1734 tokio::time::sleep(Duration::from_millis(25)).await;
1735 }
1736 store.get_chunks(&id, &ChunkFilter::default()).await.unwrap()
1737 }
1738
1739 fn joined(chunks: &[OutputChunk]) -> Vec<u8> {
1740 chunks.iter().flat_map(|c| c.bytes.clone()).collect()
1741 }
1742
1743 fn joined_stream(chunks: &[OutputChunk], stream: Stream) -> Vec<u8> {
1744 chunks
1745 .iter()
1746 .filter(|c| c.stream == stream)
1747 .flat_map(|c| c.bytes.clone())
1748 .collect()
1749 }
1750
1751 /// Divergence 1 of 3 (R739-F4): the child must not think it is on a
1752 /// terminal. This is the one that makes cargo colorize.
1753 #[tokio::test]
1754 async fn pipe_mode_child_sees_no_tty_on_stdout() {
1755 let dir = tempfile::tempdir().unwrap();
1756 let store = open_store(&dir).await;
1757 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1758 let cmd = "if [ -t 1 ]; then echo TTY; else echo PIPE; fi";
1759
1760 let piped = run_to_completion(
1761 &store,
1762 &driver,
1763 cmd,
1764 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1765 )
1766 .await;
1767 assert_eq!(joined(&piped), b"PIPE\n");
1768
1769 // The PTY default is unchanged — the terminal tiles depend on it.
1770 let ptied = run_to_completion(
1771 &store,
1772 &driver,
1773 cmd,
1774 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1775 )
1776 .await;
1777 assert_eq!(joined(&ptied), b"TTY\r\n");
1778 }
1779
1780 /// Divergence 2 of 3: no `ONLCR`, so a `\n` the child wrote stays a `\n`.
1781 /// This is what `build_run.rs::undo_onlcr` used to compensate for.
1782 #[tokio::test]
1783 async fn pipe_mode_does_not_translate_newlines() {
1784 let dir = tempfile::tempdir().unwrap();
1785 let store = open_store(&dir).await;
1786 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1787
1788 let piped = run_to_completion(
1789 &store,
1790 &driver,
1791 r"printf 'a\nb\n'",
1792 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1793 )
1794 .await;
1795 assert_eq!(joined(&piped), b"a\nb\n");
1796 }
1797
1798 /// Divergence 3 of 3: stdout and stderr stay apart, under their true
1799 /// [`Stream`], instead of being merged by the terminal.
1800 #[tokio::test]
1801 async fn pipe_mode_keeps_stderr_separate_from_stdout() {
1802 let dir = tempfile::tempdir().unwrap();
1803 let store = open_store(&dir).await;
1804 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1805 let cmd = "printf 'to-out\n'; printf 'to-err\n' >&2";
1806
1807 let piped = run_to_completion(
1808 &store,
1809 &driver,
1810 cmd,
1811 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1812 )
1813 .await;
1814 assert_eq!(joined_stream(&piped, Stream::Stdout), b"to-out\n");
1815 assert_eq!(joined_stream(&piped, Stream::Stderr), b"to-err\n");
1816
1817 // Under a PTY the kernel merges them and everything lands on stdout —
1818 // the property that made stream separation unrecoverable downstream.
1819 let ptied = run_to_completion(
1820 &store,
1821 &driver,
1822 cmd,
1823 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1824 )
1825 .await;
1826 assert!(
1827 joined_stream(&ptied, Stream::Stderr).is_empty(),
1828 "PTY runs have no stderr chunks; that is the behaviour pipe mode exists to fix",
1829 );
1830 }
1831
1832 /// Both pipes must reach EOF before the child is reaped, or a run whose
1833 /// last bytes went to stderr would be marked terminal with output still
1834 /// unread. The 4 KiB write is larger than a pipe's atomic-write buffer, so
1835 /// this fails if either pump is dropped rather than awaited.
1836 #[tokio::test]
1837 async fn pipe_mode_drains_both_streams_before_the_run_is_terminal() {
1838 let dir = tempfile::tempdir().unwrap();
1839 let store = open_store(&dir).await;
1840 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1841
1842 let piped = run_to_completion(
1843 &store,
1844 &driver,
1845 "head -c 4096 /dev/zero | tr '\\0' 'x'; head -c 4096 /dev/zero | tr '\\0' 'y' >&2",
1846 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1847 )
1848 .await;
1849 assert_eq!(joined_stream(&piped, Stream::Stdout).len(), 4096);
1850 assert_eq!(joined_stream(&piped, Stream::Stderr).len(), 4096);
1851 }
1852
1853 /// Exit codes have to survive the move to `std::process::Child`, which
1854 /// reports them through a different type than `portable_pty::Child`.
1855 #[tokio::test]
1856 async fn pipe_mode_records_the_childs_exit_code() {
1857 let dir = tempfile::tempdir().unwrap();
1858 let store = open_store(&dir).await;
1859 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1860
1861 let id = driver
1862 .spawn_run(
1863 "exit 101",
1864 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1865 )
1866 .await
1867 .unwrap();
1868
1869 let deadline = std::time::Instant::now() + Duration::from_secs(10);
1870 loop {
1871 let meta = store.get_run(&id).await.unwrap().unwrap();
1872 match meta.status {
1873 RunStatus::Done { exit_code, .. } => {
1874 assert_eq!(exit_code, 101);
1875 return;
1876 }
1877 RunStatus::Lost { .. } | RunStatus::Killed { .. } => {
1878 panic!("unexpected terminal status {:?}", meta.status)
1879 }
1880 _ => {}
1881 }
1882 if std::time::Instant::now() > deadline {
1883 panic!("run did not complete in time");
1884 }
1885 tokio::time::sleep(Duration::from_millis(25)).await;
1886 }
1887 }
1888
1889 /// A pipe run has no terminal, and the two PTY-only verbs must say so
1890 /// rather than reaching into a `None` master.
1891 #[tokio::test]
1892 async fn pipe_mode_has_no_terminal_to_resize_or_read_a_foreground_pid_from() {
1893 let dir = tempfile::tempdir().unwrap();
1894 let store = open_store(&dir).await;
1895 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
1896
1897 let id = driver
1898 .spawn_run(
1899 "sleep 2",
1900 SpawnOpts { cwd: "/tmp".into(), pipe: true, ..Default::default() },
1901 )
1902 .await
1903 .unwrap();
1904
1905 assert!(matches!(
1906 driver.resize_run(&id, 100, 40).await,
1907 Err(DriverError::NotFound(_))
1908 ));
1909 assert_eq!(driver.foreground_pid(&id), None);
1910 let _ = driver.kill_run(&id, Some(SIGKILL)).await;
1911 }
1912
1913 /// Wait for a run to reach a terminal status, or panic.
1914 async fn await_done(store: &TaskStore, id: &TaskRunId) -> TaskRunMeta {
1915 let deadline = std::time::Instant::now() + Duration::from_secs(5);
1916 loop {
1917 let meta = store.get_run(id).await.unwrap().unwrap();
1918 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
1919 return meta;
1920 }
1921 if std::time::Instant::now() > deadline {
1922 panic!("run did not complete in time, status={:?}", meta.status);
1923 }
1924 tokio::time::sleep(Duration::from_millis(50)).await;
1925 }
1926 }
1927
1928 async fn output_of(store: &TaskStore, id: &TaskRunId) -> String {
1929 let chunks = store.get_chunks(id, &ChunkFilter::default()).await.unwrap();
1930 let bytes: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
1931 String::from_utf8_lossy(&bytes).into_owned()
1932 }
1933
1934 // ── The caller's bytes reach the shell unchanged (R739-S2) ───────────────
1935
1936 /// `AttachResult.argv` is populated on every run, rewrite or not, so
1937 /// `spawn_run` used to join it back into the command line unconditionally.
1938 /// That put every `task.run` command through a whitespace normalization
1939 /// nobody asked for. A multi-line command is the case where that is not
1940 /// cosmetic: the newline the caller wrote becomes a space, and two
1941 /// commands become one nonsense command.
1942 #[cfg(unix)]
1943 #[tokio::test]
1944 async fn a_multi_line_command_is_not_flattened_into_one_line() {
1945 let dir = tempfile::tempdir().unwrap();
1946 let store = open_store(&dir).await;
1947 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
1948
1949 // Flattened to one line this is `echo one echo two`, which prints
1950 // "one echo two" — a different answer, not a failure, which is what
1951 // makes the old behaviour dangerous rather than merely wrong.
1952 let id = driver
1953 .spawn_run(
1954 "echo one\necho two",
1955 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1956 )
1957 .await
1958 .unwrap();
1959 await_done(&store, &id).await;
1960
1961 let out = output_of(&store, &id).await;
1962 assert!(out.contains("one"), "got: {out:?}");
1963 assert!(
1964 out.contains("two"),
1965 "the second line must have run as its own command; got: {out:?}"
1966 );
1967 assert!(
1968 !out.contains("one echo two"),
1969 "the newline was flattened into a space; got: {out:?}"
1970 );
1971 }
1972
1973 /// `resolve_argv` strips `bunx`/`npx`/`pnpm` so a beholder's `matches` sees
1974 /// the bare tool. That is a *matching* concern; it must never reach the
1975 /// spawn, or the wrapper the caller needed is gone from the command.
1976 #[cfg(unix)]
1977 #[tokio::test]
1978 async fn a_wrapper_the_caller_wrote_is_not_stripped_from_the_spawned_command() {
1979 let dir = tempfile::tempdir().unwrap();
1980 let store = open_store(&dir).await;
1981 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
1982
1983 // `npx` is almost certainly absent in test environments, and that is
1984 // the point: if the wrapper survived, the shell reports it missing. If
1985 // it were stripped we would be running bare `--version`.
1986 let id = driver
1987 .spawn_run(
1988 "npx r739s2-nonexistent-tool --version",
1989 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
1990 )
1991 .await
1992 .unwrap();
1993 let meta = await_done(&store, &id).await;
1994 let out = output_of(&store, &id).await;
1995 assert!(
1996 !matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
1997 "expected a failure, got {:?} with output {out:?}",
1998 meta.status
1999 );
2000 assert!(
2001 !out.contains("--version: "),
2002 "the wrapper was stripped and the shell tried to run the flag; got: {out:?}"
2003 );
2004 }
2005
2006 // ── Direct argv (R652-T6) ────────────────────────────────────────────────
2007
2008 #[tokio::test]
2009 async fn explicit_argv_execs_the_program_directly() {
2010 let dir = tempfile::tempdir().unwrap();
2011 let store = open_store(&dir).await;
2012 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2013
2014 /* The distinguishing observation: under `sh -c` the child is `sh` and
2015 `$0` is `sh`; exec'd directly it is the program itself. Printing
2016 `$0` is the cheapest way to see which of the two happened. */
2017 let id = driver
2018 .spawn_run(
2019 "unused-because-argv-wins",
2020 SpawnOpts {
2021 cwd: "/tmp".into(),
2022 argv: Some(vec![
2023 "/bin/sh".into(),
2024 "-c".into(),
2025 "printf 'argv0=%s\\n' \"$0\"".into(),
2026 "direct-exec-marker".into(),
2027 ]),
2028 ..Default::default()
2029 },
2030 )
2031 .await
2032 .unwrap();
2033
2034 await_done(&store, &id).await;
2035 let text = output_of(&store, &id).await;
2036 assert!(
2037 text.contains("argv0=direct-exec-marker"),
2038 "argv should have been exec'd verbatim, got: {text:?}"
2039 );
2040 }
2041
2042 #[tokio::test]
2043 async fn explicit_argv_still_records_the_requested_command() {
2044 let dir = tempfile::tempdir().unwrap();
2045 let store = open_store(&dir).await;
2046 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2047
2048 /* A shell tile asks for "$SHELL" and the daemon resolves it to a real
2049 argv. The run must still read back as what was asked for, or the
2050 rail row and the history re-run both show an implementation
2051 detail. */
2052 let id = driver
2053 .spawn_run(
2054 "$SHELL",
2055 SpawnOpts {
2056 cwd: "/tmp".into(),
2057 argv: Some(vec!["/bin/sh".into(), "-c".into(), "true".into()]),
2058 ..Default::default()
2059 },
2060 )
2061 .await
2062 .unwrap();
2063
2064 let meta = await_done(&store, &id).await;
2065 assert_eq!(meta.command, "$SHELL");
2066 assert!(
2067 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2068 "expected Done(0), got {:?}",
2069 meta.status
2070 );
2071 }
2072
2073 /// R901-B2. The control is the whole point: without the `pipefail: false`
2074 /// half this would pass if `pipefail` stopped existing, and with only the
2075 /// `true` half it would pass if every pipeline had always reported its
2076 /// leftmost failure. The pair pins the *difference*, which is the thing
2077 /// that cost this camp ~50 minutes of red tree.
2078 #[tokio::test]
2079 async fn pipefail_reports_the_failing_stage_and_posix_reports_the_last_one() {
2080 let dir = tempfile::tempdir().unwrap();
2081 let store = open_store(&dir).await;
2082 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2083
2084 // `(exit 101) | tail -1` is `cargo check 2>&1 | tail -40` with the
2085 // compile stripped out: a failing producer feeding a succeeding tail.
2086 let line = "(exit 101) | tail -1";
2087
2088 let posix = driver
2089 .spawn_run(
2090 line,
2091 SpawnOpts { cwd: "/tmp".into(), pipefail: false, ..Default::default() },
2092 )
2093 .await
2094 .unwrap();
2095 let meta = await_done(&store, &posix).await;
2096 assert!(
2097 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2098 "POSIX pipeline status is the LAST stage's — expected Done(0), got {:?}",
2099 meta.status
2100 );
2101
2102 let failing = driver
2103 .spawn_run(
2104 line,
2105 SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2106 )
2107 .await
2108 .unwrap();
2109 let meta = await_done(&store, &failing).await;
2110 assert!(
2111 matches!(meta.status, RunStatus::Done { exit_code: 101, .. }),
2112 "pipefail must surface the producer's 101, got {:?}",
2113 meta.status
2114 );
2115 }
2116
2117 /// The prelude must not reach [`TaskRunMeta::command`]: that string is what
2118 /// history re-runs and what an agent audits the relocation note against.
2119 #[tokio::test]
2120 async fn pipefail_does_not_leak_into_the_recorded_command() {
2121 let dir = tempfile::tempdir().unwrap();
2122 let store = open_store(&dir).await;
2123 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2124
2125 let id = driver
2126 .spawn_run(
2127 "echo recorded-verbatim | cat",
2128 SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2129 )
2130 .await
2131 .unwrap();
2132
2133 let meta = await_done(&store, &id).await;
2134 assert_eq!(meta.command, "echo recorded-verbatim | cat");
2135 assert!(
2136 !meta.command.contains("pipefail"),
2137 "the prelude leaked into the recorded command: {:?}",
2138 meta.command
2139 );
2140 }
2141
2142 /// The portability guard. On a `/bin/sh` that rejects `pipefail` (dash, i.e.
2143 /// most Linux camps) the probe must degrade to plain POSIX semantics — it
2144 /// must NOT take the shell down with it, because `set` is a special builtin
2145 /// and a bare `set -o pipefail` there is entitled to exit before the
2146 /// caller's command runs at all. Asserting the command still produces its
2147 /// output is asserting exactly that.
2148 #[tokio::test]
2149 async fn the_pipefail_probe_never_costs_the_command_that_follows_it() {
2150 let dir = tempfile::tempdir().unwrap();
2151 let store = open_store(&dir).await;
2152 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2153
2154 let id = driver
2155 .spawn_run(
2156 "echo probe-survived",
2157 SpawnOpts { cwd: "/tmp".into(), pipefail: true, ..Default::default() },
2158 )
2159 .await
2160 .unwrap();
2161
2162 let meta = await_done(&store, &id).await;
2163 let text = output_of(&store, &id).await;
2164 assert!(
2165 matches!(meta.status, RunStatus::Done { exit_code: 0, .. }),
2166 "expected Done(0), got {:?}",
2167 meta.status
2168 );
2169 assert!(text.contains("probe-survived"), "command did not run, got: {text:?}");
2170 // The probe itself must be silent — it runs on every relocated build.
2171 assert!(
2172 !text.contains("pipefail"),
2173 "the probe printed a diagnostic into the build's own output: {text:?}"
2174 );
2175 }
2176
2177 #[tokio::test]
2178 async fn empty_argv_falls_back_to_the_shell_path() {
2179 let dir = tempfile::tempdir().unwrap();
2180 let store = open_store(&dir).await;
2181 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2182
2183 let id = driver
2184 .spawn_run(
2185 "echo empty_argv_fallback",
2186 SpawnOpts { cwd: "/tmp".into(), argv: Some(vec![]), ..Default::default() },
2187 )
2188 .await
2189 .unwrap();
2190
2191 await_done(&store, &id).await;
2192 let text = output_of(&store, &id).await;
2193 assert!(
2194 text.contains("empty_argv_fallback"),
2195 "empty argv must not spawn nothing, got: {text:?}"
2196 );
2197 }
2198
2199 #[tokio::test]
2200 async fn spawn_failing_command_records_nonzero_exit() {
2201 let dir = tempfile::tempdir().unwrap();
2202 let store = open_store(&dir).await;
2203 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2204
2205 let id = driver
2206 .spawn_run(
2207 "exit 42",
2208 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2209 )
2210 .await
2211 .unwrap();
2212
2213 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2214 loop {
2215 let meta = store.get_run(&id).await.unwrap().unwrap();
2216 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2217 match meta.status {
2218 RunStatus::Done { exit_code, .. } => {
2219 assert_ne!(exit_code, 0, "exit 42 should produce a non-zero exit code");
2220 }
2221 other => panic!("unexpected status: {other:?}"),
2222 }
2223 break;
2224 }
2225 if std::time::Instant::now() > deadline {
2226 panic!("run did not complete in time");
2227 }
2228 tokio::time::sleep(Duration::from_millis(50)).await;
2229 }
2230 }
2231
2232 // ── Signal handling ───────────────────────────────────────────────────────
2233
2234 #[cfg(unix)]
2235 #[tokio::test]
2236 async fn kill_with_sigterm_transitions_to_killed() {
2237 let dir = tempfile::tempdir().unwrap();
2238 let store = open_store(&dir).await;
2239 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2240
2241 let id = driver
2242 .spawn_run(
2243 "sleep 60",
2244 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2245 )
2246 .await
2247 .unwrap();
2248
2249 // Give the process a moment to start.
2250 tokio::time::sleep(Duration::from_millis(100)).await;
2251
2252 driver.kill_run(&id, Some(SIGTERM)).await.unwrap();
2253
2254 let deadline = std::time::Instant::now() + Duration::from_secs(10);
2255 loop {
2256 let meta = store.get_run(&id).await.unwrap().unwrap();
2257 if matches!(meta.status, RunStatus::Killed { .. } | RunStatus::Lost { .. }) {
2258 assert!(
2259 matches!(meta.status, RunStatus::Killed { .. }),
2260 "expected Killed, got {:?}",
2261 meta.status
2262 );
2263 break;
2264 }
2265 if std::time::Instant::now() > deadline {
2266 panic!("run did not become Killed in time, status={:?}", meta.status);
2267 }
2268 tokio::time::sleep(Duration::from_millis(50)).await;
2269 }
2270 }
2271
2272 #[cfg(unix)]
2273 #[tokio::test]
2274 async fn kill_run_returns_not_found_after_exit() {
2275 let dir = tempfile::tempdir().unwrap();
2276 let store = open_store(&dir).await;
2277 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2278
2279 let id = driver
2280 .spawn_run(
2281 "echo done",
2282 SpawnOpts { cwd: "/tmp".into(), ..Default::default() },
2283 )
2284 .await
2285 .unwrap();
2286
2287 // Wait for natural exit.
2288 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2289 loop {
2290 let meta = store.get_run(&id).await.unwrap().unwrap();
2291 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2292 break;
2293 }
2294 if std::time::Instant::now() > deadline {
2295 panic!("run did not complete");
2296 }
2297 tokio::time::sleep(Duration::from_millis(50)).await;
2298 }
2299
2300 // Kill on a completed run should return NotFound.
2301 let result = driver.kill_run(&id, None).await;
2302 assert!(
2303 matches!(result, Err(DriverError::NotFound(_))),
2304 "expected NotFound, got {result:?}"
2305 );
2306 }
2307
2308 // ── Stdin relay ───────────────────────────────────────────────────────────
2309
2310 #[cfg(unix)]
2311 #[tokio::test]
2312 async fn stdin_send_reaches_child() {
2313 let dir = tempfile::tempdir().unwrap();
2314 let store = open_store(&dir).await;
2315 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2316
2317 // Shell that reads a line from stdin and echoes it back.
2318 let id = driver
2319 .spawn_run(
2320 "read line && echo got_$line",
2321 SpawnOpts {
2322 cwd: "/tmp".into(),
2323 stdin_enabled: true,
2324 ..Default::default()
2325 },
2326 )
2327 .await
2328 .unwrap();
2329
2330 tokio::time::sleep(Duration::from_millis(150)).await;
2331 driver.send_stdin(&id, b"hello\n".to_vec()).await.unwrap();
2332
2333 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2334 loop {
2335 let meta = store.get_run(&id).await.unwrap().unwrap();
2336 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2337 break;
2338 }
2339 if std::time::Instant::now() > deadline {
2340 panic!("run did not complete after stdin input");
2341 }
2342 tokio::time::sleep(Duration::from_millis(50)).await;
2343 }
2344
2345 let chunks = store.get_chunks(&id, &ChunkFilter::default()).await.unwrap();
2346 let raw: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
2347 let text = String::from_utf8_lossy(&raw);
2348 assert!(
2349 text.contains("got_hello"),
2350 "expected 'got_hello' in output, got: {text:?}"
2351 );
2352 }
2353
2354 /// `resize_run` must change the geometry the *child* sees, not just the
2355 /// master fd — so the assertion reads `stty size` from inside the PTY
2356 /// after the resize rather than inspecting the driver's own state.
2357 #[tokio::test]
2358 async fn resize_run_changes_geometry_the_child_sees() {
2359 let dir = tempfile::tempdir().unwrap();
2360 let store = open_store(&dir).await;
2361 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2362
2363 // Wait for a line on stdin, then report the geometry as of that moment.
2364 let id = driver
2365 .spawn_run(
2366 "read line && stty size",
2367 SpawnOpts {
2368 cwd: "/tmp".into(),
2369 stdin_enabled: true,
2370 // Spawn at the default 80x24 so the assertion can't pass by
2371 // accident if the resize is a no-op.
2372 ..Default::default()
2373 },
2374 )
2375 .await
2376 .unwrap();
2377
2378 tokio::time::sleep(Duration::from_millis(150)).await;
2379 driver.resize_run(&id, 120, 40).await.unwrap();
2380 driver.send_stdin(&id, b"go\n".to_vec()).await.unwrap();
2381
2382 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2383 loop {
2384 let meta = store.get_run(&id).await.unwrap().unwrap();
2385 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2386 break;
2387 }
2388 if std::time::Instant::now() > deadline {
2389 panic!("run did not complete after stdin input");
2390 }
2391 tokio::time::sleep(Duration::from_millis(50)).await;
2392 }
2393
2394 let chunks = store.get_chunks(&id, &ChunkFilter::default()).await.unwrap();
2395 let raw: Vec<u8> = chunks.into_iter().flat_map(|c| c.bytes).collect();
2396 let text = String::from_utf8_lossy(&raw);
2397 assert!(
2398 text.contains("40 120"),
2399 "expected resized geometry '40 120' in output, got: {text:?}"
2400 );
2401 }
2402
2403 /// A run that is not active on this driver (finished, or never existed) is
2404 /// `NotFound` rather than a panic — same contract as `send_stdin`.
2405 #[tokio::test]
2406 async fn resize_run_returns_not_found_after_exit() {
2407 let dir = tempfile::tempdir().unwrap();
2408 let store = open_store(&dir).await;
2409 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2410
2411 let id = driver
2412 .spawn_run("true", SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2413 .await
2414 .unwrap();
2415
2416 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2417 loop {
2418 let meta = store.get_run(&id).await.unwrap().unwrap();
2419 if !matches!(meta.status, RunStatus::Running | RunStatus::Pending) {
2420 break;
2421 }
2422 if std::time::Instant::now() > deadline {
2423 panic!("run did not exit");
2424 }
2425 tokio::time::sleep(Duration::from_millis(50)).await;
2426 }
2427
2428 assert!(matches!(
2429 driver.resize_run(&id, 100, 30).await,
2430 Err(DriverError::NotFound(_))
2431 ));
2432 }
2433
2434 // ── Tier-2 side-channel log fd ────────────────────────────────────────────
2435
2436 /// Verify that a child writing a JSON-line to `YAH_LOG_PIPE` (via
2437 /// `printf ... >> $YAH_LOG_PIPE`) produces a shim event with the correct
2438 /// fields in the store.
2439 ///
2440 /// The child opens the FIFO path for writing — no fd inheritance needed.
2441 #[cfg(unix)]
2442 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
2443 async fn log_pipe_events_land_in_store() {
2444 use crate::store::EventFilter;
2445
2446 let dir = tempfile::tempdir().unwrap();
2447 let store = open_store(&dir).await;
2448 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2449
2450 // The shell writes one JSON-line to the FIFO by redirecting printf
2451 // output to the path stored in YAH_LOG_PIPE.
2452 let cmd = r#"printf '{"level":"warn","target":"test.shim","msg":"hello-from-pipe","fields":{"x":42},"_lib":"test-shim","_lib_ver":"0.1.0"}\n' >> "$YAH_LOG_PIPE""#;
2453
2454 let id = driver
2455 .spawn_run(cmd, SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2456 .await
2457 .unwrap();
2458
2459 // Wait for run completion. Deadline is generous because parallel-test
2460 // load + the rt.block_on hops from the reader/log threads can slow
2461 // child-process scheduling.
2462 let deadline = std::time::Instant::now() + Duration::from_secs(20);
2463 loop {
2464 let meta = store.get_run(&id).await.unwrap().unwrap();
2465 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
2466 break;
2467 }
2468 if std::time::Instant::now() > deadline {
2469 panic!("run did not complete in time");
2470 }
2471 tokio::time::sleep(Duration::from_millis(50)).await;
2472 }
2473
2474 // The log receiver thread drains after the lifecycle task drops the
2475 // write-end FdCloser; give it a brief moment.
2476 tokio::time::sleep(Duration::from_millis(500)).await;
2477
2478 let events = store.query_events(&id, &EventFilter::default()).await.unwrap();
2479 assert!(
2480 !events.is_empty(),
2481 "expected at least one shim event, got none"
2482 );
2483 let ev = events.iter().find(|e| e.target == "test.shim");
2484 let ev = ev.expect("event with target 'test.shim' not found");
2485 assert_eq!(ev.msg, "hello-from-pipe");
2486 assert_eq!(ev.level, crate::types::Level::Warn);
2487 assert!(
2488 matches!(&ev.source, crate::types::EventSource::Shim { lib, .. } if lib == "test-shim"),
2489 "unexpected source: {:?}",
2490 ev.source
2491 );
2492 assert_eq!(ev.fields.get("x"), Some(&serde_json::json!(42)));
2493 }
2494
2495 /// When `log_fd_enabled` is false, neither `YAH_TASK_RUN` nor
2496 /// `YAH_LOG_PIPE` are exported, and no shim events are written.
2497 #[cfg(unix)]
2498 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2499 async fn log_pipe_disabled_produces_no_events() {
2500 use crate::store::EventFilter;
2501
2502 let dir = tempfile::tempdir().unwrap();
2503 let store = open_store(&dir).await;
2504 let driver = Arc::new(TaskDriver::new(Arc::clone(&store)).await.unwrap());
2505
2506 // Try to write to YAH_LOG_PIPE; the conditional guards against
2507 // the variable being absent, so the command always exits 0.
2508 let cmd = r#"[ -n "$YAH_LOG_PIPE" ] && printf '{"level":"info","target":"t","msg":"m","fields":{}}\n' >> "$YAH_LOG_PIPE" || true"#;
2509
2510 let id = driver
2511 .spawn_run(
2512 cmd,
2513 SpawnOpts { cwd: "/tmp".into(), log_fd_enabled: false, ..Default::default() },
2514 )
2515 .await
2516 .unwrap();
2517
2518 let deadline = std::time::Instant::now() + Duration::from_secs(5);
2519 loop {
2520 let meta = store.get_run(&id).await.unwrap().unwrap();
2521 if matches!(meta.status, RunStatus::Done { .. } | RunStatus::Lost { .. }) {
2522 break;
2523 }
2524 if std::time::Instant::now() > deadline {
2525 panic!("run did not complete");
2526 }
2527 tokio::time::sleep(Duration::from_millis(50)).await;
2528 }
2529
2530 tokio::time::sleep(Duration::from_millis(100)).await;
2531
2532 let events = store.query_events(&id, &EventFilter::default()).await.unwrap();
2533 assert!(
2534 events.is_empty(),
2535 "expected no shim events when log_fd_enabled=false, got {}",
2536 events.len()
2537 );
2538 }
2539
2540 // ── Unattached-run reaper (R739-B12) ─────────────────────────────────────
2541
2542 /// The origin `yah build run` tags its relocated builds with. Spelled out
2543 /// here rather than imported: what the reaper must do is defined by the
2544 /// string on the wire, not by any constant this crate owns.
2545 const BUILD_RUN: &str = "build-run";
2546
2547 fn opted_in() -> Vec<String> {
2548 vec![BUILD_RUN.to_string()]
2549 }
2550
2551 async fn spawn_long_run(driver: &TaskDriver, origin: &str) -> TaskRunId {
2552 driver
2553 .spawn_run(
2554 "sleep 30",
2555 SpawnOpts {
2556 cwd: "/tmp".into(),
2557 origin: Some(origin.to_string()),
2558 ..Default::default()
2559 },
2560 )
2561 .await
2562 .unwrap()
2563 }
2564
2565 async fn await_status(
2566 store: &Arc<TaskStore>,
2567 id: &TaskRunId,
2568 want: fn(&RunStatus) -> bool,
2569 ) -> RunStatus {
2570 let deadline = std::time::Instant::now() + Duration::from_secs(10);
2571 loop {
2572 let status = store.get_run(id).await.unwrap().unwrap().status;
2573 if want(&status) {
2574 return status;
2575 }
2576 if std::time::Instant::now() > deadline {
2577 panic!("run never reached the expected status, last={status:?}");
2578 }
2579 tokio::time::sleep(Duration::from_millis(25)).await;
2580 }
2581 }
2582
2583 /// The orphan this ticket exists for: the client is gone, so nothing polls
2584 /// the run, so the daemon must end it.
2585 #[tokio::test]
2586 async fn an_unpolled_run_of_an_opted_in_origin_is_reaped() {
2587 let dir = tempfile::tempdir().unwrap();
2588 let store = open_store(&dir).await;
2589 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2590
2591 let id = spawn_long_run(&driver, BUILD_RUN).await;
2592 tokio::time::sleep(Duration::from_millis(300)).await;
2593
2594 let reaped = driver
2595 .reap_unattached(Duration::from_millis(200), &opted_in())
2596 .await;
2597 assert_eq!(reaped, vec![id.clone()], "the unattached run should be reaped");
2598
2599 let status = await_status(&store, &id, |s| {
2600 matches!(s, RunStatus::Killed { .. } | RunStatus::Done { .. })
2601 })
2602 .await;
2603 assert!(
2604 matches!(status, RunStatus::Killed { .. }),
2605 "a reaped run ends Killed, got {status:?}"
2606 );
2607 }
2608
2609 /// The regression that protects a healthy long build: a client that is
2610 /// still polling keeps its run alive however long the build takes.
2611 #[tokio::test]
2612 async fn a_run_a_client_is_still_polling_is_never_reaped() {
2613 let dir = tempfile::tempdir().unwrap();
2614 let store = open_store(&dir).await;
2615 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2616
2617 let id = spawn_long_run(&driver, BUILD_RUN).await;
2618
2619 // Six polls across three idle windows — what `yah build run`'s tail
2620 // loop does, slowed down.
2621 for _ in 0..6 {
2622 tokio::time::sleep(Duration::from_millis(100)).await;
2623 driver.note_attached(&id);
2624 let reaped = driver
2625 .reap_unattached(Duration::from_millis(200), &opted_in())
2626 .await;
2627 assert!(reaped.is_empty(), "a polled run must survive, reaped {reaped:?}");
2628 }
2629
2630 assert!(
2631 matches!(
2632 store.get_run(&id).await.unwrap().unwrap().status,
2633 RunStatus::Running
2634 ),
2635 "the polled run should still be running"
2636 );
2637
2638 // And the moment the polling stops, it becomes reapable — same run,
2639 // so this pins the refresh rather than a missing origin match.
2640 tokio::time::sleep(Duration::from_millis(300)).await;
2641 let reaped = driver
2642 .reap_unattached(Duration::from_millis(200), &opted_in())
2643 .await;
2644 assert_eq!(reaped, vec![id], "a run that stopped being polled is reapable");
2645 }
2646
2647 /// The regression that protects real users' terminals. An interactive tile
2648 /// sits unpolled for hours by design and must never be touched, however
2649 /// long the reaper runs.
2650 #[tokio::test]
2651 async fn a_terminal_tile_is_never_reaped_however_long_it_idles() {
2652 let dir = tempfile::tempdir().unwrap();
2653 let store = open_store(&dir).await;
2654 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2655
2656 let id = spawn_long_run(&driver, "terminal").await;
2657 tokio::time::sleep(Duration::from_millis(300)).await;
2658
2659 for _ in 0..3 {
2660 let reaped = driver.reap_unattached(Duration::ZERO, &opted_in()).await;
2661 assert!(
2662 reaped.is_empty(),
2663 "a terminal tile is outside the opted-in origins, reaped {reaped:?}"
2664 );
2665 tokio::time::sleep(Duration::from_millis(50)).await;
2666 }
2667
2668 assert!(
2669 matches!(
2670 store.get_run(&id).await.unwrap().unwrap().status,
2671 RunStatus::Running
2672 ),
2673 "the terminal run must still be running"
2674 );
2675
2676 driver.kill_run(&id, Some(SIGKILL)).await.unwrap();
2677 }
2678
2679 /// A run with no origin at all — an ordinary `task.run` job — is outside
2680 /// every opt-in list, and an empty list reaps nothing.
2681 #[tokio::test]
2682 async fn an_origin_less_run_and_an_empty_opt_in_list_reap_nothing() {
2683 let dir = tempfile::tempdir().unwrap();
2684 let store = open_store(&dir).await;
2685 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2686
2687 let plain = driver
2688 .spawn_run("sleep 30", SpawnOpts { cwd: "/tmp".into(), ..Default::default() })
2689 .await
2690 .unwrap();
2691 let build = spawn_long_run(&driver, BUILD_RUN).await;
2692 tokio::time::sleep(Duration::from_millis(100)).await;
2693
2694 assert!(
2695 driver.reap_unattached(Duration::ZERO, &[]).await.is_empty(),
2696 "an empty opt-in list must reap nothing, not everything"
2697 );
2698 assert_eq!(
2699 driver.reap_unattached(Duration::ZERO, &opted_in()).await,
2700 vec![build],
2701 "only the opted-in origin is reapable"
2702 );
2703
2704 assert!(
2705 matches!(
2706 store.get_run(&plain).await.unwrap().unwrap().status,
2707 RunStatus::Running
2708 ),
2709 "the origin-less run must be untouched"
2710 );
2711 driver.kill_run(&plain, Some(SIGKILL)).await.unwrap();
2712 }
2713
2714 /// `note_attached` is observable, and the age it resets is what the sweep
2715 /// reads.
2716 #[tokio::test]
2717 async fn attached_age_resets_on_a_poll_and_is_none_for_a_foreign_run() {
2718 let dir = tempfile::tempdir().unwrap();
2719 let store = open_store(&dir).await;
2720 let driver = TaskDriver::new(Arc::clone(&store)).await.unwrap();
2721
2722 let id = spawn_long_run(&driver, BUILD_RUN).await;
2723 tokio::time::sleep(Duration::from_millis(150)).await;
2724 let aged = driver.attached_age(&id).expect("driver owns this run");
2725 assert!(aged >= Duration::from_millis(100), "age should have grown, got {aged:?}");
2726
2727 driver.note_attached(&id);
2728 let fresh = driver.attached_age(&id).unwrap();
2729 assert!(fresh < aged, "a poll resets the age: {fresh:?} vs {aged:?}");
2730
2731 assert!(
2732 driver.attached_age(&TaskRunId::new()).is_none(),
2733 "a run this driver does not own has no attachment age"
2734 );
2735
2736 driver.kill_run(&id, Some(SIGKILL)).await.unwrap();
2737 }
2738}