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}