Skip to main content

yah_qed/
events.rs

1//! Live per-step events emitted by [`crate::runner::PipelineRunner`] as a
2//! pipeline executes (R325-F2).
3//!
4//! A runner constructed with [`crate::runner::PipelineRunner::with_events`]
5//! pushes a [`QedEvent`] onto an unbounded channel at each lifecycle boundary:
6//! the run starts, each step starts, every stdout/stderr line, each step
7//! finishes, the run finishes. The camp daemon drains these into a per-run
8//! buffer that `qed.tail` serves as a cursor-tailable feed; the CLI prints
9//! them to the console as they arrive.
10//!
11//! Without a sink the runner is silent — `run()` still returns the terminal
12//! [`crate::types::QedRunMeta`], so existing callers are unaffected.
13//!
14//! @yah:ticket(R488-F5, "Event-stream wiring + desktop QED-pane nested-tree render for sub-pipelines")
15//! @yah:assignee(agent:claude)
16//! @yah:at(2026-06-08T02:54:33Z)
17//! @yah:status(review)
18//! @yah:phase(P5)
19//! @yah:parent(R488)
20//! @yah:next("New event variants: sub_pipeline_started / sub_pipeline_finished emit on the parent's run with the child run_id")
21//! @yah:next("DB: parent_run_id foreign key on the runs table; read paths join when surfacing the tree")
22//! @yah:next("app/yah/desktop: QED pane renders nested sub-pipelines as expandable subtrees with own status pills")
23//! @yah:verify("Run a 2-child composite in dev; desktop pane shows both children as expandable rows under the parent")
24//! @arch:see(.yah/docs/working/W201-qed-pipeline-composition.md)
25//! @yah:depends_on(R488-F2)
26//! @yah:tier(Cleric)
27//! @yah:handoff("F5 shipped the structural surface for nested sub-pipeline event tracking. (1) qed::QedEvent gained SubPipelineStarted{index,name,target,child_run_id,at} and SubPipelineFinished{index,name,child_run_id,status,at} variants. (2) runner.rs::execute_step_sub_pipeline now emits these bookends around child recursion; new helper sub_pipeline_target_label() renders the resolver-token string (builtin:<n> | path:<p> | gha:<p>). The child runner's `events` sink is now decoupled (None) — previously F2 shared the parent's sink, which caused apply_qed_event_to_meta to push child steps onto the parent's StepStatus list, corrupting the parent's run snapshot. Bookends restore that information at the right altitude. (3) qed::QedRunMeta gained parent_run_id: Option<QedRunId> with serde(default, skip_serializing_if). Runner threads it into child runs via a new field on PipelineRunner; top-level runs leave it None. Persistence (JSON shards under .yah/jit/qed/) is automatic — no DB schema change needed because run history is JSON, not SQL. (4) rpc::QedRunWire grew parent_run_id; QedEventWire grew SubPipelineStarted/Finished mirrors (kebab-case, RFC3339 ts). camp.rs qed_meta_to_wire + qed_event_to_wire converters updated; apply_qed_event_to_meta ignores the bookends (parent's step list stays untouched — the StepStarted/StepFinished around the SubPipeline parent step already track it). (5) Two new runner tests: sub_pipeline_emits_started_finished_bookends_with_child_run_id (asserts pairing by child_run_id, distinct from parent run_id, child events DO NOT leak onto parent stream) and sub_pipeline_finished_emits_failed_status_when_child_fails. qed --lib: 196 pass + 1 pre-existing unrelated failure (test_builtin_release_build_pipeline 4-vs-6, documented across R380-T3/R381-T2/R407 handoffs). cargo check --workspace clean.")
28//! @yah:next("Daemon-side child run registration: today the child's QedRun is run inline by the parent runner inside qed_run_handler; the child's events vanish (sink decoupled to keep parent's meta clean) and the child isn't in the qed_runs HashMap, so qed.tail{run_id=child_run_id} returns null. To make child runs tail-able as their own stream, qed_run_handler needs to (a) intercept the new SubPipelineStarted to register a fresh QedRunState with parent_run_id set, (b) inject a per-child sink so the child runner's events drain into THAT buffer (requires runner-side hook — perhaps OutcomeDispatcher-style 'spawn_child_sink' or a new SubPipelineSinkProvider trait), (c) flush + persist on SubPipelineFinished. Bigger change than F5 — file as a followup ticket under R488 or a fresh relay.")
29//! @yah:next("Desktop QED pane render: no TS pane exists yet in app/yah/web or the desktop frontend (R325 was 'QED desktop UI blank slate'). The structural surface is in place; render lands when the pane is built. Treat the F5 verify ('desktop pane shows both children as expandable rows') as deferred until that pane exists.")
30//! @yah:next("yah qed run CLI: the live tail loop in app/yah/cli/src/qed.rs prints StepStarted + Output but does not yet handle SubPipelineStarted/Finished. Small followup to print a nested indent (e.g. '↳ sub-pipeline started: builtin:child-a (run_id=...)' / '↲ finished: success').")
31//! @yah:verify("cargo test -p qed --lib sub_pipeline  # 17 pass incl. 2 new F5 tests")
32//! @yah:verify("cargo test -p qed --lib  # 196 pass + 1 pre-existing unrelated failure")
33//! @yah:verify("cargo check --workspace  # clean")
34//! @yah:gotcha("Child runner's events sink is intentionally decoupled (None) — previously F2 shared parent's sink and child Step events corrupted the parent's QedRunMeta.steps via apply_qed_event_to_meta. Don't re-couple without first adding per-event run_id discrimination or a per-child sink provider.")
35//! @yah:gotcha("QedRunMeta.parent_run_id is serde(default, skip_serializing_if=Option::is_none) so existing on-disk shards in .yah/jit/qed/ load fine; no migration.")
36//!
37//! @yah:ticket(R365-F20, "session_started: record dispatchKind (assist|dispatch); mirror into meta.json")
38//! @yah:status(review)
39//! @yah:assignee(agent:claude)
40//! @yah:at(2026-06-10T01:58:23Z)
41//! @yah:parent(R365)
42//! @yah:next("Add dispatchKind: 'assist' | 'dispatch' field to session_started event in crates/yah/qed/src/events.rs and emit it from the party.assist / party.dispatch / subagent.spawn paths in crates/yah/agent-tools/src/tools.rs.")
43//! @yah:next("Mirror parentSessionId + dispatchDepth + dispatchKind from session_started into the meta.json sidecar so board.session and sidecar-readers see parentage without parsing the jsonl head.")
44//! @yah:next("Backfill: existing sessions have parentSessionId in jsonl but no dispatchKind — default unknown sessions to 'dispatch' on read (the common case) so old sessions don't crash the renderer.")
45//! @yah:gotcha("Observed 2026-06-09 on R365-T14: party.assist and subagent.spawn produce byte-identical session_started events for the same Yamli character (parentSessionId, dispatchDepth, agentId all match). UI can't tell which back-link should produce nesting in non-compact Party Column rendering. meta.json sidecar carries no parentage info at all — anything reading sessions via sidecar (board.session lookups, future tooling) is parent-blind.")
46//! @arch:see(.yah/docs/working/W167-smoke-matrix-plan.md)
47//! @yah:handoff("Added dispatchKind (Assist|Dispatch) to AgentEvent::SessionStarted and mirrored parentSessionId/dispatchDepth/dispatchKind into the meta.json sidecar.\n\nPath correction: the ticket's source pointer (crates/yah/qed/src/events.rs, and agent-tools/src/tools.rs for the tool paths) was wrong on both counts -- those are unrelated types (qed's pipeline-run QedEvent; tools.rs is the read-only KG tool registry). The real session_started type is AgentEvent::SessionStarted in crates/yah/party/src/agent.rs; the real party.assist/party.dispatch/subagent.spawn implementations are PartyAssist/PartyDispatch/AgentSpawn in crates/yah/agent-tools/src/agent_dispatch_tools.rs, which all funnel through one daemon-side fn (subagent_spawn in app/yah/cli/src/camp.rs) via party_dispatch.\n\nSemantics: party.assist, party.dispatch, and subagent.spawn ALL mint the identical sub-relationship child today (parent_session_id + dispatch_depth+1, nest-worthy) -- confirmed by the PartyDispatch tool docstring ('sub-relationship... assistant') and by the FE's own isAssistRelationship fallback comment in nestedAssistChildren.test.ts ('gate on parentInfo.dispatchKind === \"assist\"'). So all three stamp DispatchKind::Assist; DispatchKind::Dispatch is reserved for the not-yet-built peer-booking party.book verb.\n\nPlumbing (mirrors the existing parent_session_id/dispatch_depth rails exactly): RPC method match in camp.rs -> subagent_spawn/party_dispatch (new dispatch_kind param) -> LaunchCharacterParams.dispatch_kind -> launch_character_session's 4 engine forks -> start_claude_session/start_runner_session/start_codex_oauth_session/agent_process::start_process_session (new trailing param) -> ClaudeSessionInit/RunnerSessionInit.dispatch_kind -> register_claude_session/register_runner_session stamp AgentEvent::SessionStarted.dispatch_kind AND SessionPromptMeta.dispatch_kind. Also extended recover_dispatch_lineage (desktop/agent.rs, R534-B5 helper) to a 3-tuple so all 6 resume/fork/rewind paths re-emit dispatch_kind, not just parent_session_id/dispatch_depth.\n\nBackfill: dispatch_kind is Option<DispatchKind> with #[serde(default)] end to end -- old JSONL/meta.json missing the field deserialize as None, never crash. Per the ticket's literal instruction, documented on the field that a reader needing a concrete value for a parented session should default a missing value to Dispatch (guidance for the consumer, e.g. R365-F21's renderer).\n\nEvery AgentEvent::SessionStarted construction site in the repo was updated and verified programmatically: party, runner, agent-tools, hub, hub-tauri crates + the yah CLI daemon (camp.rs, 4 RunnerSessionInit sites) + desktop (agent.rs, agent_process.rs, agent_eval.rs, keepalive_subsystem.rs, camp_socket.rs).")
48//! @yah:verify("cargo test -p party --lib -- agent:: (5/5 pass)")
49//! @yah:verify("cargo test -p runner --lib -- sessions:: session:: sink:: (56/56 pass, incl. list_summaries_preserves_claude_path_dispatch_lineage)")
50//! @yah:verify("cargo test -p hub --lib (43/43 pass)")
51//! @yah:verify("cargo check clean on party, hub-tauri, hub, agent-tools, yah (lib), desktop (lib)")
52//! @yah:gotcha("Two pre-existing/unrelated blockers on this shared tree (confirmed via git diff showing Cargo.toml churn from a concurrent peer, R409-T9): (1) turso-vs-yah #[global_allocator] conflict blocks any test-profile build of yah/desktop (lib-only cargo check unaffected); (2) agent-tools' envoy_tools.rs test module references an undeclared anyhow dev-dependency, blocking cargo test -p agent-tools --lib (verified by hand that all 7 test-fixture edits there set dispatch_kind; cargo check -p agent-tools non-test is clean).")
53
54use chrono::{DateTime, Utc};
55
56use crate::types::RunStatus;
57
58/// Return the sorted set of environment variable names from `env_iter`
59/// whose names look like credentials. Names only — values are never read,
60/// so this is safe to log to the event stream and the operator UI.
61///
62/// A name is considered credential-shaped when it either:
63/// - ends in `TOKEN`, `KEY`, `SECRET`, `PASSWORD`, or `CREDENTIAL`
64///   (case-insensitive, after splitting on `_`); or
65/// - starts with a known credential prefix: `HETZNER_`, `CLOUDFLARE_`,
66///   `CF_`, `AWS_`, `R2_`, `GITHUB_`, `GH_`, `HUGGINGFACE_`, `HF_`,
67///   `ANTHROPIC_`, `OPENAI_`.
68///
69/// Anything else (PATH, HOME, USER, etc.) is filtered out — the goal is a
70/// readable chip row, not an `env` dump.
71pub fn credential_env_keys<I, K>(env_iter: I) -> Vec<String>
72where
73    I: IntoIterator<Item = (K, K)>,
74    K: AsRef<str>,
75{
76    const SUFFIXES: &[&str] = &["TOKEN", "KEY", "SECRET", "PASSWORD", "CREDENTIAL"];
77    const PREFIXES: &[&str] = &[
78        "HETZNER_",
79        "CLOUDFLARE_",
80        "CF_",
81        "AWS_",
82        "R2_",
83        "GITHUB_",
84        "GH_",
85        "HUGGINGFACE_",
86        "HF_",
87        "ANTHROPIC_",
88        "OPENAI_",
89    ];
90    let mut keys: Vec<String> = env_iter
91        .into_iter()
92        .map(|(k, _)| k.as_ref().to_string())
93        .filter(|k| {
94            let upper = k.to_ascii_uppercase();
95            if PREFIXES.iter().any(|p| upper.starts_with(p)) {
96                return true;
97            }
98            // Split on '_' so e.g. `MY_API_KEY` matches but `KEYBOARD` doesn't.
99            upper
100                .rsplit('_')
101                .next()
102                .map(|tail| SUFFIXES.contains(&tail))
103                .unwrap_or(false)
104        })
105        .collect();
106    keys.sort();
107    keys.dedup();
108    keys
109}
110
111#[cfg(test)]
112mod tests {
113    use super::*;
114
115    #[test]
116    fn credential_env_keys_filters_by_suffix_and_prefix() {
117        let env = [
118            ("PATH", "/usr/bin"),
119            ("HOME", "/home/u"),
120            ("KEYBOARD", "us"),               // suffix-like but no underscore
121            ("HETZNER_S3_ACCESS_KEY", "xxx"), // prefix + suffix
122            ("CF_API_TOKEN", "xxx"),          // prefix + suffix
123            ("MY_API_KEY", "xxx"),            // suffix only
124            ("HF_TOKEN", "xxx"),              // prefix
125            ("DATABASE_PASSWORD", "xxx"),     // suffix
126            ("RANDOM_VAR", "xxx"),
127        ];
128        let keys = credential_env_keys(env.iter().map(|(k, v)| (*k, *v)));
129        assert_eq!(
130            keys,
131            vec![
132                "CF_API_TOKEN".to_string(),
133                "DATABASE_PASSWORD".to_string(),
134                "HETZNER_S3_ACCESS_KEY".to_string(),
135                "HF_TOKEN".to_string(),
136                "MY_API_KEY".to_string(),
137            ]
138        );
139    }
140}
141
142/// Which standard stream a line of step output came from.
143#[derive(Debug, Clone, Copy, PartialEq, Eq)]
144pub enum OutputStream {
145    Stdout,
146    Stderr,
147}
148
149/// A live event emitted while a pipeline executes.
150///
151/// Step `index` is 0-based and aligns with `Pipeline::steps`. Steps run
152/// strictly in sequence, so a consumer sees `StepStarted { index: i }` before
153/// any `StepOutput { index: i, .. }` and before `StepFinished { index: i }`,
154/// and indices arrive monotonically.
155#[derive(Debug, Clone)]
156pub enum QedEvent {
157    /// Emitted once, immediately after registration, when the run is
158    /// holding for its `concurrency_key` lock. For pipelines that
159    /// opt out (`concurrency_key = "@parallel"`), this event still
160    /// fires but is immediately followed by `RunStarted` — the queue
161    /// hop is just instantaneous.
162    RunQueued { key: String, at: DateTime<Utc> },
163    /// Emitted once, before the first step. For queued runs this fires
164    /// when the key lock is acquired; for parallel pipelines it fires
165    /// right after `RunQueued`.
166    RunStarted {
167        total_steps: usize,
168        at: DateTime<Utc>,
169    },
170    /// A step is about to execute.
171    ///
172    /// `argv` is the substituted command line for Subprocess-kind steps
173    /// (empty for BuildImage / PackageNativeTarball / etc. — those have
174    /// their own command shape). `env_keys` is the set of credential-shaped
175    /// environment variable names that were present in the runner's
176    /// environment at spawn time (see [`credential_env_keys`]) — KEYS
177    /// ONLY, never values. Surfaced in the QED detail pane so an operator
178    /// can confirm at a glance which secrets the step inherited without
179    /// shell-pasting or re-running.
180    StepStarted {
181        index: usize,
182        name: String,
183        argv: Vec<String>,
184        env_keys: Vec<String>,
185        at: DateTime<Utc>,
186    },
187    /// A step was dispatched to a remote build-worker and now has a durable
188    /// yubaba workload identity (`forge_id`), emitted BEFORE the runner blocks
189    /// on `handle.wait()`. This is the reattach anchor (R603-T1): the camp
190    /// daemon stamps `forge_id` onto the running step's `task_run_id` and
191    /// persists a non-terminal `<run_id>.json` so that, if the daemon restarts
192    /// mid-build, boot reconcile (R603-T2) can re-poll the workload by this id
193    /// instead of orphaning it. Local steps never emit this.
194    StepRemoteDispatched {
195        index: usize,
196        name: String,
197        forge_id: String,
198        at: DateTime<Utc>,
199    },
200    /// One line of stdout/stderr captured from the executing step (local runs).
201    StepOutput {
202        index: usize,
203        name: String,
204        stream: OutputStream,
205        line: String,
206    },
207    /// A step reached a terminal status (`Success` or `Failed`).
208    ///
209    /// `msg` carries the failure tail (e.g. the last lines of cargo stderr)
210    /// when `status == Failed`; `None` on success. Consumers render this as a
211    /// red banner above the per-step log so the operator sees *why* without
212    /// having to scroll through the streamed output. Empty string is treated
213    /// the same as `None` by the UI.
214    StepFinished {
215        index: usize,
216        name: String,
217        status: RunStatus,
218        msg: Option<String>,
219        at: DateTime<Utc>,
220    },
221    /// Emitted once, after the last executed step (or after an aborting failure).
222    RunFinished {
223        status: RunStatus,
224        at: DateTime<Utc>,
225    },
226    /// A `kind = "sub-pipeline"` step (W201) began recursion into a child
227    /// pipeline. Carries the child's run id so a consumer can pivot to
228    /// `qed.tail { run_id = child_run_id }` for the child's step-level
229    /// detail; the parent's own stream does NOT carry the child's events
230    /// (decoupled to keep parent's `apply_qed_event_to_meta` step list
231    /// uncorrupted). `target` is the resolver token (`builtin:<name>`,
232    /// `path:<path>`, `gha:<path>`) — same discipline as the F1 walker.
233    SubPipelineStarted {
234        index: usize,
235        name: String,
236        target: String,
237        child_run_id: String,
238        at: DateTime<Utc>,
239    },
240    /// A `kind = "sub-pipeline"` child run reached a terminal status. Pairs
241    /// with the prior `SubPipelineStarted` event by `child_run_id`.
242    SubPipelineFinished {
243        index: usize,
244        name: String,
245        child_run_id: String,
246        status: RunStatus,
247        at: DateTime<Utc>,
248    },
249    /// A job instance inside a `kind = "gha-workflow"` step began executing
250    /// (W200 R487 follow-up). `index` is the qed step index of the enclosing
251    /// gha-workflow step; `job_key` is `yah_qed_gha::JobInstance::key()`
252    /// (`"<job>"` for non-matrix, `"<job>#<row>"` for matrix). The receiver
253    /// uses `(index, job_key)` to scope the per-job subtree under the
254    /// parent step.
255    GhaJobStarted {
256        index: usize,
257        name: String,
258        job_id: String,
259        matrix_index: Option<usize>,
260        job_key: String,
261        total_steps: usize,
262        at: DateTime<Utc>,
263    },
264    /// One step inside a gha-workflow job is about to run. `action_kind` is
265    /// `"run"` for bash steps and `"uses:<slug>"` for action invocations.
266    GhaStepStarted {
267        index: usize,
268        name: String,
269        job_key: String,
270        step_index: usize,
271        step_id: Option<String>,
272        step_name: Option<String>,
273        action_kind: String,
274        at: DateTime<Utc>,
275    },
276    /// One line of stdout/stderr captured from a gha-workflow bash step.
277    GhaStepOutput {
278        index: usize,
279        name: String,
280        job_key: String,
281        step_index: usize,
282        stream: OutputStream,
283        line: String,
284    },
285    /// A gha-workflow step reached a terminal conclusion. `msg` carries a
286    /// stderr tail on failure (mirroring the qed-runner StepFinished.msg
287    /// convention so the receiver can render a red banner uniformly).
288    /// `conclusion` is `"success" | "failure" | "skipped"`.
289    GhaStepFinished {
290        index: usize,
291        name: String,
292        job_key: String,
293        step_index: usize,
294        conclusion: String,
295        msg: Option<String>,
296        at: DateTime<Utc>,
297    },
298    /// A gha-workflow job instance reached a terminal result. `result` is
299    /// `"success" | "failure" | "cancelled" | "skipped"`. Pairs with the
300    /// prior `GhaJobStarted` by `(index, job_key)`.
301    GhaJobFinished {
302        index: usize,
303        name: String,
304        job_key: String,
305        result: String,
306        at: DateTime<Utc>,
307    },
308}