diff --git a/docs/internal/events.md b/docs/internal/events.md index af5566435..be498955e 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -200,6 +200,90 @@ Informational, warning, or error notice emitted during the run. | `code` | string | Machine-readable notice code | | `message` | string | Human-readable message | +### `metadata.snapshot.started` + +Emitted when Fabro begins a durable metadata snapshot operation. These are product events for Fabro metadata snapshots, not tracing spans for the underlying git or filesystem work. + +Init and finalize metadata snapshots are unscoped. Checkpoint metadata snapshots use the checkpoint stage scope, so they include the checkpoint `node_id`, `node_label`, and `stage_id`. + +```json +{ + "id": "...", "ts": "...", "run_id": "...", + "event": "metadata.snapshot.started", + "properties": { + "phase": "checkpoint", + "branch": "fabro/meta" + } +} +``` + +| Property | Type | Description | +|----------|------|-------------| +| `phase` | string | Logical metadata operation: `"init"`, `"checkpoint"`, or `"finalize"` | +| `branch` | string | Metadata branch/ref being written | + +### `metadata.snapshot.completed` + +Emitted when Fabro commits and pushes a metadata snapshot successfully. + +```json +{ + "id": "...", "ts": "...", "run_id": "...", + "event": "metadata.snapshot.completed", + "properties": { + "phase": "checkpoint", + "branch": "fabro/meta", + "duration_ms": 2800, + "entry_count": 12, + "bytes": 18432, + "commit_sha": "def456..." + } +} +``` + +| Property | Type | Description | +|----------|------|-------------| +| `phase` | string | Logical metadata operation: `"init"`, `"checkpoint"`, or `"finalize"` | +| `branch` | string | Metadata branch/ref that was written | +| `duration_ms` | number | End-to-end duration of the metadata snapshot operation | +| `entry_count` | number | Number of metadata files written into the snapshot commit | +| `bytes` | number | Sum of serialized metadata entry byte lengths | +| `commit_sha` | string | Metadata snapshot commit SHA | + +### `metadata.snapshot.failed` + +Emitted when a real metadata snapshot attempt fails. It is emitted before the matching compatibility `run.notice`, allowing human-facing consumers to suppress duplicate warning text. Compatibility notices with codes `checkpoint_metadata_write_failed` and `checkpoint_metadata_push_failed` may still appear in raw event streams. The `checkpoint_metadata_degraded` notice is a separate summary signal and should not be treated as a duplicate of this event. + +```json +{ + "id": "...", "ts": "...", "run_id": "...", + "event": "metadata.snapshot.failed", + "properties": { + "phase": "checkpoint", + "branch": "fabro/meta", + "duration_ms": 900, + "failure_kind": "push", + "error": "failed to push metadata snapshot", + "causes": ["remote rejected the push"], + "commit_sha": "def456...", + "entry_count": 12, + "bytes": 18432 + } +} +``` + +| Property | Type | Description | +|----------|------|-------------| +| `phase` | string | Logical metadata operation: `"init"`, `"checkpoint"`, or `"finalize"` | +| `branch` | string | Metadata branch/ref being written | +| `duration_ms` | number | End-to-end duration before failure | +| `failure_kind` | string | Failure phase: `"load_state"`, `"write"`, or `"push"` | +| `error` | string | Primary error summary | +| `causes` | string[] | Error cause chain; omitted when empty | +| `commit_sha` | string? | Local metadata commit SHA for push failures; omitted for load-state and write failures | +| `entry_count` | number? | Metadata entry count for push failures; omitted for load-state and write failures | +| `bytes` | number? | Serialized metadata byte count for push failures; omitted for load-state and write failures | + --- ## Stage events diff --git a/docs/superpowers/plans/2026-04-29-metadata-snapshot-events.md b/docs/superpowers/plans/2026-04-29-metadata-snapshot-events.md new file mode 100644 index 000000000..d12c9a77d --- /dev/null +++ b/docs/superpowers/plans/2026-04-29-metadata-snapshot-events.md @@ -0,0 +1,196 @@ +# Metadata Snapshot Events Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [x]`) syntax for tracking. + +**Goal:** Add first-class Fabro run events for durable metadata snapshot writes so run timing gaps are visible without emitting low-level git span events. + +**Architecture:** Metadata snapshot events are product-domain workflow events emitted around each real metadata archive attempt. Callers own event emission around the whole metadata operation, including run-store state loading; `SandboxMetadataWriter` remains responsible for creating and pushing snapshots and returns snapshot accounting data. + +**Tech Stack:** Rust, serde, strum, Fabro typed `Event`/`RunEvent` pipeline, existing CLI log/progress renderers, `cargo nextest`. + +--- + +## Summary + +Add these event names: + +- `metadata.snapshot.started` +- `metadata.snapshot.completed` +- `metadata.snapshot.failed` + +The events cover Fabro metadata snapshots only, not every underlying git or filesystem operation. Emit them for `init`, `checkpoint`, and `finalize` metadata attempts so timing gaps become visible without rebuilding tracing as events. + +## Event Contract + +- `MetadataSnapshotStartedProps { phase: MetadataSnapshotPhase, branch: String }` +- `MetadataSnapshotCompletedProps { phase: MetadataSnapshotPhase, branch: String, duration_ms: u64, entry_count: usize, bytes: u64, commit_sha: String }` +- `MetadataSnapshotFailedProps { phase: MetadataSnapshotPhase, branch: String, duration_ms: u64, failure_kind: MetadataSnapshotFailureKind, error: String, causes: Vec, commit_sha: Option, entry_count: Option, bytes: Option }` + +Enums: + +- `MetadataSnapshotPhase = init | checkpoint | finalize` +- `MetadataSnapshotFailureKind = load_state | write | push` +- Both enums must pair `#[serde(rename_all = "snake_case")]` with `#[strum(serialize_all = "snake_case")]` so serde and strum stay aligned with the project enum convention. + +Rules: + +- `metadata.snapshot.completed` means the metadata snapshot was committed and pushed successfully. There is no `pushed` field because it would always be true. +- `commit_sha` is intentionally asymmetric: completed snapshots always include `commit_sha: String`; failed snapshots include `commit_sha: Option` because push failures can have a local commit while load-state and write failures cannot. +- Failed accounting fields are optional: `entry_count: Option` and `bytes: Option` are `Some` for push failures and `None` for load-state/write failures. +- Optional fields use `#[serde(default, skip_serializing_if = "Option::is_none")]`, matching the convention in `infra.rs`. +- Failed props follow the existing failure-event convention in `infra.rs`: `error: String` contains the primary error summary and `causes: Vec` contains the cause chain with `#[serde(default, skip_serializing_if = "Vec::is_empty")]`. +- A writer `push_error` maps to `metadata.snapshot.failed { failure_kind: "push", commit_sha: Some(...), entry_count: Some(...), bytes: Some(...) }`. +- A run-store `state()` failure maps to `metadata.snapshot.failed { failure_kind: "load_state", commit_sha: None, entry_count: None, bytes: None }`. +- If metadata is already degraded and a later snapshot would currently return early, emit no metadata snapshot event for that skipped attempt. Do not emit `started`; skipped attempts are not real attempts and should not count as failures. +- Writer errors before a local commit map to `metadata.snapshot.failed { failure_kind: "write", commit_sha: None, entry_count: None, bytes: None }`. +- Emit `metadata.snapshot.failed` before the compatibility `run.notice` for the same failure so human-facing consumers can deterministically suppress duplicate warning text. +- Typed `metadata.snapshot.*` events are not deduplicated for real attempts. Existing metadata `run.notice` deduping remains compatibility-only. +- Checkpoint metadata events use the existing stage scope so they include `node_id`, `node_label`, and `stage_id`. Init/finalize metadata events are unscoped and must not set `node_id`, `node_label`, or `stage_id`. +- Keep `branch` because the exact metadata ref is useful in raw logs and for push failure diagnostics. Do not include `message`; `phase` fully identifies the logical metadata operation. + +## Implementation Tasks + +### Task 1: Add Typed Event Bodies + +**Files:** +- Modify: `lib/crates/fabro-types/src/run_event/mod.rs` +- Modify: `lib/crates/fabro-types/src/run_event/infra.rs` + +- [x] Add `MetadataSnapshotPhase` and `MetadataSnapshotFailureKind` enums in `infra.rs`. +- [x] Derive `Serialize`, `Deserialize`, `strum::Display`, `strum::EnumString`, and `strum::IntoStaticStr`. +- [x] Add both `#[serde(rename_all = "snake_case")]` and `#[strum(serialize_all = "snake_case")]` to each enum. +- [x] Add the three metadata snapshot props structs in `infra.rs` with the exact fields from the Event Contract section. +- [x] Add serde attributes for optional failed fields and empty `causes` exactly as specified in the Event Contract. +- [x] Add three `EventBody` variants in `mod.rs` with exact serde names: + - `metadata.snapshot.started` + - `metadata.snapshot.completed` + - `metadata.snapshot.failed` +- [x] Extend `EventBody::event_name()` and known-event handling for all three names. + +### Task 2: Add Workflow Event Variants And Mapping + +**Files:** +- Modify: `lib/crates/fabro-workflow/src/event.rs` + +- [x] Add internal `Event` variants matching the three new event bodies. +- [x] Extend `Event::trace()` with concise tracing fields: `phase`, `branch`, `duration_ms`, and `failure_kind`. +- [x] Extend `event_name()` with the three exact event names. +- [x] Extend `event_body_from_event()` to construct the matching `EventBody` variants. +- [x] Ensure unscoped metadata snapshot events do not set envelope `node_id`, `node_label`, or `stage_id`; checkpoint callers will use `emit_scoped()`. + +### Task 3: Return Snapshot Accounting From The Writer + +**Files:** +- Modify: `lib/crates/fabro-workflow/src/sandbox_metadata.rs` + +- [x] Extend `MetadataSnapshot` to include `entry_count: usize` and `bytes: u64`. +- [x] Compute `entry_count` and `bytes` inside `SandboxMetadataWriter::write_snapshot()` from the single `dump.git_entries()` allocation that the writer already needs. +- [x] Return those accounting values on successful local metadata commit, including the case where `push_error` is present. +- [x] Do not add per-command events or expose writer-internal steps on the wire. + +### Task 4a: Emit Init Metadata Events + +**Files:** +- Modify: `lib/crates/fabro-workflow/src/lifecycle/git.rs` + +- [x] If metadata is already degraded, return from the init metadata path without emitting `metadata.snapshot.*`. +- [x] Move the degraded check above the init `metadata.snapshot.started` emission point; the existing check inside `write_metadata_snapshot()` is not enough because skipped attempts must not leave a dangling `started`. +- [x] Keep the existing inner `metadata_degraded()` guard in `GitLifecycle::write_metadata_snapshot()` as defense-in-depth for future callers, but do not rely on it for init/checkpoint skip semantics. +- [x] Wrap the full init metadata operation in `GitLifecycle::on_run_start`, including `run_store.state()`. +- [x] Emit `metadata.snapshot.started { phase: "init" }` before loading run state for the init operation. +- [x] Emit `metadata.snapshot.completed` only when the init metadata commit and push both succeed. +- [x] Emit `metadata.snapshot.failed` for init load-state, write, and push failures using the Event Contract mapping. +- [x] Emit `metadata.snapshot.failed` before calling `emit_metadata_warning()` for the same init failure. + +### Task 4b: Emit Checkpoint Metadata Events + +**Files:** +- Modify: `lib/crates/fabro-workflow/src/lifecycle/git.rs` + +- [x] If metadata is already degraded, return from the checkpoint metadata path without emitting `metadata.snapshot.*`. +- [x] Move the degraded check above the checkpoint `metadata.snapshot.started` emission point; the existing check inside `write_metadata_snapshot()` is not enough because skipped attempts must not leave a dangling `started`. +- [x] Keep the existing inner `metadata_degraded()` guard in `GitLifecycle::write_metadata_snapshot()` as defense-in-depth for future callers, but do not rely on it for init/checkpoint skip semantics. +- [x] Wrap the full checkpoint metadata operation in `GitLifecycle::on_checkpoint`, including `run_store.state()`. +- [x] Emit scoped `metadata.snapshot.started { phase: "checkpoint" }` before loading run state for the checkpoint operation. +- [x] Emit scoped `metadata.snapshot.completed` or `metadata.snapshot.failed` before `checkpoint.completed`. +- [x] Emit scoped `metadata.snapshot.failed` before calling `emit_metadata_warning()` for the same checkpoint failure. +- [x] Preserve existing metadata-degraded `run.notice` emission for compatibility, but treat the new typed event as the primary human-facing signal. + +### Task 5: Emit Finalize Metadata Events + +**Files:** +- Modify: `lib/crates/fabro-workflow/src/pipeline/finalize.rs` + +- [x] If metadata is already degraded, return from the final metadata path without emitting `metadata.snapshot.*`. +- [x] Keep the degraded check above the final `metadata.snapshot.started` emission point; skipped attempts must not leave a dangling `started`. +- [x] Wrap the full final metadata operation in `write_finalize_commit`, including `run_store.state()`. +- [x] Emit `metadata.snapshot.started { phase: "finalize" }` before loading run state for the final metadata operation. +- [x] Emit `metadata.snapshot.completed` only when the final metadata commit and push both succeed. +- [x] Emit `metadata.snapshot.failed` for finalize load-state, write, and push failures using the Event Contract mapping. +- [x] Emit `metadata.snapshot.failed` before calling `emit_metadata_warning()` for the same final metadata failure. +- [x] Ensure final metadata events are emitted before `run.completed`. +- [x] Preserve existing `checkpoint_metadata_write_failed`, `checkpoint_metadata_push_failed`, and `checkpoint_metadata_degraded` notices for compatibility. + +### Task 6: CLI, Consumers, And Documentation + +**Files:** +- Modify: `lib/crates/fabro-cli/src/commands/run/logs.rs` +- Modify: `lib/crates/fabro-cli/src/commands/run/run_progress/event.rs` +- Modify: `lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs` +- Modify: `docs/internal/events.md` + +- [x] Audit existing consumers with `rg -n "event_name|EventBody|metadata.snapshot" lib/crates apps lib/packages docs/public/api-reference/fabro-api.yaml docs/internal --glob '!docs/superpowers/**'` and update any event-name filters that should recognize metadata snapshot events. Ignore matches in implementation-plan docs. +- [x] Render completed metadata snapshots compactly in pretty/progress output, for example `Metadata checkpoint 2.8s`. +- [x] Render `metadata.snapshot.failed` as the primary user-visible metadata warning/error. +- [x] Suppress duplicate CLI display of the compatibility `checkpoint_metadata_*` notice when the same stream already contains a matching earlier `metadata.snapshot.failed` event. +- [x] Limit suppression to per-failure compatibility notices: `checkpoint_metadata_write_failed` and `checkpoint_metadata_push_failed`. Do not suppress the `checkpoint_metadata_degraded` end-of-run summary notice; it is a distinct summary signal. +- [x] Keep `fabro logs --json` unchanged except for the new serialized event records. +- [x] Document the three event definitions in `docs/internal/events.md`. +- [x] State in the docs that these are product events for durable metadata snapshots, not tracing spans. + +## API And Client Compatibility + +Checked current API/client shape: + +- `docs/public/api-reference/fabro-api.yaml` models `RunEvent` as a generic object with `event: string` and `properties: object`. +- `lib/packages/fabro-api-client/src/models/run-event.ts` includes `[key: string]: any` and `properties?: { [key: string]: any }`. + +No OpenAPI or TypeScript client schema changes are required for this event-only addition unless implementation discovers a stricter consumer outside this model. + +## Test Plan + +- [x] Add `fabro-types` serialization/deserialization tests proving the three event names are known and props serialize to the agreed JSON shape. +- [x] Add `fabro-workflow` event conversion tests for all three variants, including scoped checkpoint metadata events. +- [x] Add lifecycle tests covering successful metadata snapshot emission: started then completed. +- [x] Add lifecycle tests covering `push_error` mapping to `metadata.snapshot.failed { failure_kind: "push", commit_sha: Some(...), entry_count: Some(...), bytes: Some(...) }`. +- [x] Add lifecycle tests covering degraded short-circuit behavior: no `metadata.snapshot.*` event is emitted and no `metadata.snapshot.started` event is left unterminated. +- [x] Add a cross-phase degraded test: after an init `metadata.snapshot.failed` marks metadata degraded, subsequent checkpoint and finalize attempts emit no `metadata.snapshot.*` events. +- [x] Add lifecycle tests covering pre-writer `run_store.state()` failure emission for init: started then failed with `failure_kind: "load_state"`. +- [x] Add lifecycle tests covering pre-writer `run_store.state()` failure emission for checkpoint: started then failed with `failure_kind: "load_state"`. +- [x] Add lifecycle tests covering pre-writer `run_store.state()` failure emission for finalize: started then failed with `failure_kind: "load_state"`. +- [x] Add tests proving `completed.entry_count` and `completed.bytes` equal the writer's `MetadataSnapshot` values. +- [x] Add tests proving push-failure `failed.entry_count` and `failed.bytes` equal the writer's `MetadataSnapshot` values. +- [x] Add ordering tests proving checkpoint metadata events occur before `checkpoint.completed`. +- [x] Add ordering tests proving finalize metadata events occur before `run.completed`. +- [x] Add tests proving `metadata.snapshot.failed` is emitted before the matching compatibility `run.notice`. +- [x] Add tests proving compatibility `run.notice` still fires for metadata degradation while CLI display avoids duplicate warnings. +- [x] Add CLI rendering tests for pretty/progress output so metadata events display compactly and do not break generic log output. +- [x] Run: + +```bash +cargo nextest run -p fabro-types +cargo nextest run -p fabro-workflow metadata +cargo nextest run -p fabro-cli logs +cargo +nightly-2026-04-14 fmt --check --all +``` + +## Assumptions + +- The public wire shape uses the exact event names in this plan. +- `entry_count` is the number of metadata files in the snapshot. +- `bytes` is the sum of serialized metadata entry byte lengths. +- Existing metadata-degraded notices stay for compatibility, but typed metadata snapshot events become the preferred signal for humans and new consumers. +- Already-degraded skipped attempts are intentionally silent; the first real failure event and compatibility notice explain why later metadata work is skipped. +- Runs with no configured metadata branch stay silent for metadata snapshot events because metadata snapshots are out of scope for those runs. +- A panic between `metadata.snapshot.started` and `metadata.snapshot.completed`/`metadata.snapshot.failed` may leave a dangling started event; this feature treats that as a run-level crash case rather than adding panic recovery around metadata event emission. +- The implementation should not introduce metadata writer sub-step events or expose low-level git command boundaries on the event stream. diff --git a/lib/crates/fabro-cli/src/commands/run/logs.rs b/lib/crates/fabro-cli/src/commands/run/logs.rs index e9c4c99ca..b5cde4456 100644 --- a/lib/crates/fabro-cli/src/commands/run/logs.rs +++ b/lib/crates/fabro-cli/src/commands/run/logs.rs @@ -14,6 +14,7 @@ use std::time::Duration; use anyhow::{Context, Result, bail}; use chrono::{DateTime, Utc}; use fabro_redact::redact_jsonl_line; +use fabro_types::run_event::is_metadata_snapshot_compat_notice_code; use fabro_util::json::normalize_json_value; use fabro_util::terminal::Styles; use tokio::time; @@ -52,10 +53,11 @@ pub(crate) async fn run(args: &LogsArgs, styles: &Styles, base_ctx: &CommandCont let is_tty = stdout.is_terminal(); let mut out = stdout.lock(); let pretty = args.pretty && !ctx.json_output(); + let mut pretty_state = PrettyEventState::default(); for line in &filtered { if pretty { - if let Some(formatted) = format_event_pretty(line, styles) { + if let Some(formatted) = format_event_pretty_streamed(line, styles, &mut pretty_state) { writeln!(out, "{formatted}")?; } } else { @@ -71,6 +73,7 @@ pub(crate) async fn run(args: &LogsArgs, styles: &Styles, base_ctx: &CommandCont pretty, styles, is_tty, + pretty_state, ) .await?; } @@ -149,6 +152,7 @@ async fn follow_store_logs( pretty: bool, styles: &Styles, _is_tty: bool, + mut pretty_state: PrettyEventState, ) -> Result<()> { let stdout = io::stdout(); let mut out = stdout.lock(); @@ -170,7 +174,9 @@ async fn follow_store_logs( for event in events { let line = event_payload_line(&event)?; if pretty { - if let Some(formatted) = format_event_pretty(&line, styles) { + if let Some(formatted) = + format_event_pretty_streamed(&line, styles, &mut pretty_state) + { writeln!(out, "{formatted}")?; } } else { @@ -199,9 +205,16 @@ async fn follow_store_logs( continue; } - let flushed_next_seq = - flush_remaining_store_events(client, run_id, next_seq, pretty, styles, &mut out) - .await?; + let flushed_next_seq = flush_remaining_store_events( + client, + run_id, + next_seq, + pretty, + styles, + &mut pretty_state, + &mut out, + ) + .await?; if flushed_next_seq > next_seq { next_seq = flushed_next_seq; terminal_deadline = Some(time::Instant::now() + FOLLOW_TERMINAL_GRACE); @@ -235,6 +248,7 @@ async fn flush_remaining_store_events( next_seq: u32, pretty: bool, styles: &Styles, + pretty_state: &mut PrettyEventState, out: &mut dyn Write, ) -> Result { let events = client @@ -246,7 +260,7 @@ async fn flush_remaining_store_events( for event in events { let line = event_payload_line(&event)?; if pretty { - if let Some(formatted) = format_event_pretty(&line, styles) { + if let Some(formatted) = format_event_pretty_streamed(&line, styles, pretty_state) { writeln!(out, "{formatted}")?; } } else { @@ -296,15 +310,51 @@ fn render_indented_markdown(styles: &Styles, text: &str, indent: &str) -> String .join("\n") } +#[derive(Debug, Default)] +struct PrettyEventState { + saw_metadata_snapshot_failure: bool, +} + +fn format_event_pretty_streamed( + line: &str, + styles: &Styles, + state: &mut PrettyEventState, +) -> Option { + let envelope: serde_json::Value = serde_json::from_str(line).ok()?; + let event = envelope.get("event")?.as_str()?; + if event == "run.notice" + && state.saw_metadata_snapshot_failure + && is_metadata_snapshot_compat_notice(&envelope) + { + return None; + } + let formatted = format_event_pretty_value(&envelope, styles); + if event == "metadata.snapshot.failed" { + state.saw_metadata_snapshot_failure = true; + } + formatted +} + +#[cfg_attr( + not(test), + allow( + dead_code, + reason = "Production pretty logs use the stateful stream formatter; unit tests exercise this single-line helper." + ) +)] pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option { let envelope: serde_json::Value = serde_json::from_str(line).ok()?; + format_event_pretty_value(&envelope, styles) +} + +fn format_event_pretty_value(envelope: &serde_json::Value, styles: &Styles) -> Option { let event = envelope.get("event")?.as_str()?; let ts = format_timestamp(envelope.get("ts")?.as_str()?); match event { "run.started" => { - let name = prop_str_field(&envelope, "name").unwrap_or("?"); - let run_id = str_field(&envelope, "run_id").unwrap_or("?"); + let name = prop_str_field(envelope, "name").unwrap_or("?"); + let run_id = str_field(envelope, "run_id").unwrap_or("?"); let header = format!( "{} {} {} {}", styles.dim.apply_to(&ts), @@ -312,7 +362,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option styles.bold.apply_to(name), styles.dim.apply_to(run_id), ); - match prop_str_field(&envelope, "goal") { + match prop_str_field(envelope, "goal") { Some(goal) if !goal.is_empty() => { let body = render_indented_markdown(styles, goal, " "); Some(format!("{header}\n{body}\n")) @@ -321,8 +371,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option } } "run.completed" => { - let duration = format_duration_ms(prop_field(&envelope, "duration_ms")); - let status_str = match prop_str_field(&envelope, "status") { + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); + let status_str = match prop_str_field(envelope, "status") { Some(status) if !status.is_empty() => status, _ => "success", }; @@ -332,8 +382,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option _ => &styles.bold_red, }; let cost = format_cost( - prop_field(&envelope, "total_usd_micros") - .or_else(|| prop_field(&envelope, "total_cost")), + prop_field(envelope, "total_usd_micros") + .or_else(|| prop_field(envelope, "total_cost")), ); let mut lines = vec![format!( @@ -345,7 +395,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )]; if let Some(billing) = - prop_field(&envelope, "billing").or_else(|| prop_field(&envelope, "usage")) + prop_field(envelope, "billing").or_else(|| prop_field(envelope, "usage")) { let total = billing .get("total_tokens") @@ -399,7 +449,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option Some(lines.join("\n")) } "run.failed" => { - let error = prop_str_field(&envelope, "error").unwrap_or("unknown error"); + let error = prop_str_field(envelope, "error").unwrap_or("unknown error"); Some(format!( "{} {} {}", styles.dim.apply_to(&ts), @@ -408,9 +458,9 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "run.notice" => { - let level = prop_str_field(&envelope, "level").unwrap_or("info"); - let code = prop_str_field(&envelope, "code").unwrap_or(""); - let message = prop_str_field(&envelope, "message").unwrap_or(""); + let level = prop_str_field(envelope, "level").unwrap_or("info"); + let code = prop_str_field(envelope, "code").unwrap_or(""); + let message = prop_str_field(envelope, "message").unwrap_or(""); let label = match level { "warn" => styles.yellow.apply_to("Warning:").to_string(), "error" => styles.bold_red.apply_to("Error:").to_string(), @@ -429,8 +479,36 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option code_suffix, )) } + "metadata.snapshot.completed" => { + let phase = prop_str_field(envelope, "phase").unwrap_or("?"); + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); + Some(format!( + "{} Metadata {} {}", + styles.dim.apply_to(&ts), + phase, + styles.dim.apply_to(&duration), + )) + } + "metadata.snapshot.failed" => { + let phase = prop_str_field(envelope, "phase").unwrap_or("?"); + let failure_kind = prop_str_field(envelope, "failure_kind").unwrap_or(""); + let error = prop_str_field(envelope, "error").unwrap_or("unknown error"); + let kind_suffix = if failure_kind.is_empty() { + String::new() + } else { + format!(" {}", styles.dim.apply_to(format!("[{failure_kind}]"))) + }; + Some(format!( + "{} {} Metadata {} failed: {}{}", + styles.dim.apply_to(&ts), + styles.yellow.apply_to("Warning:"), + phase, + error, + kind_suffix, + )) + } "stage.started" => { - let label = str_field(&envelope, "node_label").unwrap_or("?"); + let label = str_field(envelope, "node_label").unwrap_or("?"); Some(format!( "{} {} {}", styles.dim.apply_to(&ts), @@ -439,10 +517,9 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "stage.completed" => { - let label = str_field(&envelope, "node_label").unwrap_or("?"); - let duration = format_duration_ms(prop_field(&envelope, "duration_ms")); - let billing = - prop_field(&envelope, "billing").or_else(|| prop_field(&envelope, "usage")); + let label = str_field(envelope, "node_label").unwrap_or("?"); + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); + let billing = prop_field(envelope, "billing").or_else(|| prop_field(envelope, "usage")); let cost = format_cost( billing .and_then(|value| value.get("total_usd_micros")) @@ -475,8 +552,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option Some(line) } "stage.failed" => { - let label = str_field(&envelope, "node_label").unwrap_or("?"); - let error = prop_str_field(&envelope, "error").unwrap_or("unknown error"); + let label = str_field(envelope, "node_label").unwrap_or("?"); + let error = prop_str_field(envelope, "error").unwrap_or("unknown error"); Some(format!( "{} {} {} {}", styles.dim.apply_to(&ts), @@ -486,9 +563,9 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "agent.message" => { - let stage = str_field(&envelope, "node_id").unwrap_or("?"); - let model = prop_str_field(&envelope, "model").unwrap_or("?"); - let text = prop_str_field(&envelope, "text").unwrap_or(""); + let stage = str_field(envelope, "node_id").unwrap_or("?"); + let model = prop_str_field(envelope, "model").unwrap_or("?"); + let text = prop_str_field(envelope, "text").unwrap_or(""); let header = format!( "{} {} {} {}{}{}", styles.dim.apply_to(&ts), @@ -502,8 +579,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option Some(format!("{header}\n{body}\n")) } "agent.tool.started" => { - let tool = prop_str_field(&envelope, "tool_name").unwrap_or("?"); - let detail = tool_detail(&envelope); + let tool = prop_str_field(envelope, "tool_name").unwrap_or("?"); + let detail = tool_detail(envelope); let display = match detail { Some(value) => format!("{tool}({value})"), None => tool.to_string(), @@ -516,11 +593,11 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "agent.tool.completed" => { - let tool = prop_str_field(&envelope, "tool_name").unwrap_or("?"); - let is_error = prop_field(&envelope, "is_error") + let tool = prop_str_field(envelope, "tool_name").unwrap_or("?"); + let is_error = prop_field(envelope, "is_error") .and_then(serde_json::Value::as_bool) .unwrap_or(false); - let detail = tool_detail(&envelope); + let detail = tool_detail(envelope); let display = match detail { Some(value) => format!("{tool}({value})"), None => tool.to_string(), @@ -535,9 +612,9 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "edge.selected" => { - let to = prop_str_field(&envelope, "to_node").unwrap_or("?"); - let reason = prop_str_field(&envelope, "reason").unwrap_or("?"); - let condition = prop_str_field(&envelope, "condition"); + let to = prop_str_field(envelope, "to_node").unwrap_or("?"); + let reason = prop_str_field(envelope, "reason").unwrap_or("?"); + let condition = prop_str_field(envelope, "condition"); let detail = match condition { Some(value) => format!(" [{value}]"), None => String::new(), @@ -552,8 +629,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "sandbox.ready" => { - let provider = prop_str_field(&envelope, "provider").unwrap_or("?"); - let duration = format_duration_ms(prop_field(&envelope, "duration_ms")); + let provider = prop_str_field(envelope, "provider").unwrap_or("?"); + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); Some(format!( "{} Sandbox: {} {}", styles.dim.apply_to(&ts), @@ -562,8 +639,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "setup.completed" => { - let count = prop_field(&envelope, "command_count").and_then(serde_json::Value::as_u64); - let duration = format_duration_ms(prop_field(&envelope, "duration_ms")); + let count = prop_field(envelope, "command_count").and_then(serde_json::Value::as_u64); + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); Some(match count { Some(count) => format!( "{} Setup: {} commands {}", @@ -579,10 +656,10 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option }) } "agent.compaction.completed" => { - let original = prop_field(&envelope, "original_turn_count") + let original = prop_field(envelope, "original_turn_count") .and_then(serde_json::Value::as_u64) .unwrap_or(0); - let preserved = prop_field(&envelope, "preserved_turn_count") + let preserved = prop_field(envelope, "preserved_turn_count") .and_then(serde_json::Value::as_u64) .unwrap_or(0); Some(format!( @@ -594,7 +671,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "parallel.started" => { - let count = prop_field(&envelope, "branch_count") + let count = prop_field(envelope, "branch_count") .and_then(serde_json::Value::as_u64) .unwrap_or(0); Some(format!( @@ -605,7 +682,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "parallel.branch.started" => { - let label = str_field(&envelope, "node_label").unwrap_or("?"); + let label = str_field(envelope, "node_label").unwrap_or("?"); Some(format!( "{} {} {}", styles.dim.apply_to(&ts), @@ -614,7 +691,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "parallel.branch.completed" => { - let label = str_field(&envelope, "node_label").unwrap_or("?"); + let label = str_field(envelope, "node_label").unwrap_or("?"); Some(format!( "{} {} {}", styles.dim.apply_to(&ts), @@ -623,7 +700,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "parallel.completed" => { - let duration = format_duration_ms(prop_field(&envelope, "duration_ms")); + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); Some(format!( "{} {} Parallel {}", styles.dim.apply_to(&ts), @@ -632,8 +709,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "pull_request.created" => { - let url = prop_str_field(&envelope, "pr_url").unwrap_or("?"); - let draft = prop_field(&envelope, "draft") + let url = prop_str_field(envelope, "pr_url").unwrap_or("?"); + let draft = prop_field(envelope, "draft") .and_then(serde_json::Value::as_bool) .unwrap_or(false); let label = if draft { "Draft PR:" } else { "PR:" }; @@ -645,7 +722,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "pull_request.failed" => { - let error = prop_str_field(&envelope, "error").unwrap_or("unknown error"); + let error = prop_str_field(envelope, "error").unwrap_or("unknown error"); Some(format!( "{} {} {}", styles.dim.apply_to(&ts), @@ -654,7 +731,7 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "retro.completed" => { - let duration = format_duration_ms(prop_field(&envelope, "duration_ms")); + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); Some(format!( "{} {} Retro {}", styles.dim.apply_to(&ts), @@ -663,8 +740,8 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option )) } "retro.failed" => { - let error = prop_str_field(&envelope, "error").unwrap_or("unknown error"); - let duration = format_duration_ms(prop_field(&envelope, "duration_ms")); + let error = prop_str_field(envelope, "error").unwrap_or("unknown error"); + let duration = format_duration_ms(prop_field(envelope, "duration_ms")); Some(format!( "{} {} Retro {} {}", styles.dim.apply_to(&ts), @@ -682,6 +759,10 @@ pub(crate) fn format_event_pretty(line: &str, styles: &Styles) -> Option } } +fn is_metadata_snapshot_compat_notice(envelope: &serde_json::Value) -> bool { + prop_str_field(envelope, "code").is_some_and(is_metadata_snapshot_compat_notice_code) +} + fn str_field<'a>(value: &'a serde_json::Value, key: &str) -> Option<&'a str> { value.get(key)?.as_str() } @@ -1048,6 +1129,45 @@ mod tests { assert!(result.contains("[launch_failed]"), "got: {result}"); } + #[test] + fn pretty_metadata_snapshot_completed() { + let styles = no_color_styles(); + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"metadata.snapshot.completed","properties":{"phase":"checkpoint","branch":"fabro/meta","duration_ms":2800,"entry_count":2,"bytes":42,"commit_sha":"abc123"}}"#; + let result = format_event_pretty(line, &styles).unwrap(); + assert!(result.contains("Metadata checkpoint"), "got: {result}"); + assert!(result.contains("3s"), "got: {result}"); + } + + #[test] + fn pretty_metadata_snapshot_failed() { + let styles = no_color_styles(); + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"metadata.snapshot.failed","properties":{"phase":"finalize","branch":"fabro/meta","duration_ms":900,"failure_kind":"push","error":"push rejected","commit_sha":"abc123","entry_count":2,"bytes":42}}"#; + let result = format_event_pretty(line, &styles).unwrap(); + assert!(result.contains("Warning:"), "got: {result}"); + assert!( + result.contains("Metadata finalize failed: push rejected"), + "got: {result}" + ); + assert!(result.contains("[push]"), "got: {result}"); + } + + #[test] + fn pretty_stream_suppresses_metadata_compat_notice_only() { + let styles = no_color_styles(); + let failed = r#"{"ts":"2026-01-01T14:25:00Z","event":"metadata.snapshot.failed","properties":{"phase":"checkpoint","branch":"fabro/meta","duration_ms":900,"failure_kind":"write","error":"write failed"}}"#; + let compat_notice = r#"{"ts":"2026-01-01T14:25:01Z","event":"run.notice","properties":{"level":"warn","code":"checkpoint_metadata_write_failed","message":"legacy metadata warning"}}"#; + let degraded_notice = r#"{"ts":"2026-01-01T14:25:02Z","event":"run.notice","properties":{"level":"warn","code":"checkpoint_metadata_degraded","message":"metadata snapshots disabled"}}"#; + let mut state = PrettyEventState::default(); + + assert!(format_event_pretty_streamed(failed, &styles, &mut state).is_some()); + assert!(format_event_pretty_streamed(compat_notice, &styles, &mut state).is_none()); + let degraded = format_event_pretty_streamed(degraded_notice, &styles, &mut state).unwrap(); + assert!( + degraded.contains("metadata snapshots disabled"), + "got: {degraded}" + ); + } + #[test] fn pretty_workflow_run_failed() { let styles = no_color_styles(); diff --git a/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs b/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs index 911f48d6c..1b6344131 100644 --- a/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs +++ b/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs @@ -203,6 +203,15 @@ pub(super) enum ProgressEvent { RetroFailed { duration_ms: u64, }, + MetadataSnapshotCompleted { + phase: String, + duration_ms: u64, + }, + MetadataSnapshotFailed { + phase: String, + failure_kind: String, + error: String, + }, RunNotice { level: RunNoticeLevel, code: String, @@ -422,6 +431,17 @@ pub(super) fn from_run_event(stored: &RunEvent) -> Option { EventBody::RetroFailed(props) => Some(ProgressEvent::RetroFailed { duration_ms: props.duration_ms, }), + EventBody::MetadataSnapshotCompleted(props) => { + Some(ProgressEvent::MetadataSnapshotCompleted { + phase: props.phase.to_string(), + duration_ms: props.duration_ms, + }) + } + EventBody::MetadataSnapshotFailed(props) => Some(ProgressEvent::MetadataSnapshotFailed { + phase: props.phase.to_string(), + failure_kind: props.failure_kind.to_string(), + error: props.error.clone(), + }), EventBody::RunNotice(props) => Some(ProgressEvent::RunNotice { level: props.level, code: props.code.clone(), @@ -474,7 +494,7 @@ fn display_value(value: &Value) -> Option { #[cfg(test)] mod tests { use fabro_agent::AgentEvent; - use fabro_types::fixtures; + use fabro_types::{MetadataSnapshotFailureKind, MetadataSnapshotPhase, fixtures}; use fabro_workflow::event::{Event, to_run_event}; use super::*; @@ -701,4 +721,50 @@ mod tests { } if code == "sandbox_cleanup_failed" && message == "sandbox cleanup failed" )); } + + #[test] + fn round_trip_metadata_snapshot_completed() { + let event = Event::MetadataSnapshotCompleted { + phase: MetadataSnapshotPhase::Checkpoint, + branch: "fabro/meta".into(), + duration_ms: 2800, + entry_count: 2, + bytes: 42, + commit_sha: "abc123".into(), + }; + + let stored = to_run_event(&fixtures::RUN_1, &event); + let parsed = from_run_event(&stored).unwrap(); + assert!(matches!( + parsed, + ProgressEvent::MetadataSnapshotCompleted { phase, duration_ms } + if phase == "checkpoint" && duration_ms == 2800 + )); + } + + #[test] + fn round_trip_metadata_snapshot_failed() { + let event = Event::MetadataSnapshotFailed { + phase: MetadataSnapshotPhase::Finalize, + branch: "fabro/meta".into(), + duration_ms: 900, + failure_kind: MetadataSnapshotFailureKind::Push, + error: "push rejected".into(), + causes: vec!["remote rejected".into()], + commit_sha: Some("abc123".into()), + entry_count: Some(2), + bytes: Some(42), + }; + + let stored = to_run_event(&fixtures::RUN_1, &event); + let parsed = from_run_event(&stored).unwrap(); + assert!(matches!( + parsed, + ProgressEvent::MetadataSnapshotFailed { + phase, + failure_kind, + error, + } if phase == "finalize" && failure_kind == "push" && error == "push rejected" + )); + } } diff --git a/lib/crates/fabro-cli/src/commands/run/run_progress/info_display.rs b/lib/crates/fabro-cli/src/commands/run/run_progress/info_display.rs index 4d3c1bb35..4035d42ad 100644 --- a/lib/crates/fabro-cli/src/commands/run/run_progress/info_display.rs +++ b/lib/crates/fabro-cli/src/commands/run/run_progress/info_display.rs @@ -63,6 +63,38 @@ impl InfoDisplay { ); } + pub(super) fn on_metadata_snapshot_completed( + renderer: &ProgressRenderer, + phase: &str, + duration_ms: u64, + ) { + Self::insert_info_line( + renderer, + &format!("Metadata {phase} {}", format_duration_ms(duration_ms)), + ); + } + + pub(super) fn on_metadata_snapshot_failed( + renderer: &ProgressRenderer, + phase: &str, + failure_kind: &str, + error: &str, + ) { + let styles = renderer.styles(); + let kind_suffix = if failure_kind.is_empty() { + String::new() + } else { + format!(" {}", styles.dim.apply_to(format!("[{failure_kind}]"))) + }; + Self::insert_info_line( + renderer, + &format!( + "{} Metadata {phase} failed: {error}{kind_suffix}", + styles.yellow.apply_to("Warning:") + ), + ); + } + pub(super) fn on_edge_selected( &self, renderer: &ProgressRenderer, diff --git a/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs index 6d33477da..43e47bf21 100644 --- a/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs @@ -4,6 +4,7 @@ )] use fabro_types::RunEvent; +use fabro_types::run_event::is_metadata_snapshot_compat_notice_code; mod event; mod info_display; @@ -20,9 +21,10 @@ use stage_display::StageDisplay; pub(crate) struct ProgressUI { renderer: ProgressRenderer, - stage: StageDisplay, - setup: SetupDisplay, - info: InfoDisplay, + stage: StageDisplay, + setup: SetupDisplay, + info: InfoDisplay, + saw_metadata_snapshot_failure: bool, } impl ProgressUI { @@ -44,6 +46,7 @@ impl ProgressUI { stage: StageDisplay::new(verbose), setup: SetupDisplay::new(verbose), info: InfoDisplay::new(verbose), + saw_metadata_snapshot_failure: false, } } @@ -407,11 +410,27 @@ impl ProgressUI { ProgressEvent::RetroFailed { duration_ms } => { self.stage.on_retro_failed(renderer, duration_ms); } + ProgressEvent::MetadataSnapshotCompleted { phase, duration_ms } => { + InfoDisplay::on_metadata_snapshot_completed(renderer, &phase, duration_ms); + } + ProgressEvent::MetadataSnapshotFailed { + phase, + failure_kind, + error, + } => { + self.saw_metadata_snapshot_failure = true; + InfoDisplay::on_metadata_snapshot_failed(renderer, &phase, &failure_kind, &error); + } ProgressEvent::RunNotice { level, code, message, } => { + if self.saw_metadata_snapshot_failure + && is_metadata_snapshot_compat_notice_code(&code) + { + return; + } InfoDisplay::on_run_notice(renderer, level, &code, &message); } ProgressEvent::PullRequestCreated { pr_url, draft } => { @@ -443,7 +462,9 @@ mod tests { use fabro_agent::{AgentEvent, SandboxEvent}; use fabro_llm::types::TokenCounts; use fabro_model::Provider; - use fabro_types::{ParallelBranchId, StageId, fixtures}; + use fabro_types::{ + MetadataSnapshotFailureKind, MetadataSnapshotPhase, ParallelBranchId, StageId, fixtures, + }; use fabro_workflow::event::{Event, RunNoticeLevel, to_run_event, to_run_event_at}; use fabro_workflow::outcome::billed_model_usage_from_llm; @@ -1011,6 +1032,68 @@ mod tests { "); } + #[test] + fn plain_metadata_snapshot_snapshot() { + let (mut ui, buffer) = capture_ui(false); + + emit(&mut ui, Event::MetadataSnapshotCompleted { + phase: MetadataSnapshotPhase::Checkpoint, + branch: "fabro/meta".into(), + duration_ms: 2000, + entry_count: 2, + bytes: 42, + commit_sha: "abc123".into(), + }); + emit(&mut ui, Event::MetadataSnapshotFailed { + phase: MetadataSnapshotPhase::Finalize, + branch: "fabro/meta".into(), + duration_ms: 900, + failure_kind: MetadataSnapshotFailureKind::Push, + error: "push rejected".into(), + causes: Vec::new(), + commit_sha: Some("abc123".into()), + entry_count: Some(2), + bytes: Some(42), + }); + + insta::assert_snapshot!(rendered(&buffer), @r" + Metadata checkpoint 2s + Warning: Metadata finalize failed: push rejected [push] + "); + } + + #[test] + fn metadata_snapshot_failure_suppresses_compat_notice_only() { + let (mut ui, buffer) = capture_ui(false); + + emit(&mut ui, Event::MetadataSnapshotFailed { + phase: MetadataSnapshotPhase::Checkpoint, + branch: "fabro/meta".into(), + duration_ms: 900, + failure_kind: MetadataSnapshotFailureKind::Write, + error: "write failed".into(), + causes: Vec::new(), + commit_sha: None, + entry_count: None, + bytes: None, + }); + emit(&mut ui, Event::RunNotice { + level: RunNoticeLevel::Warn, + code: "checkpoint_metadata_write_failed".into(), + message: "legacy metadata warning".into(), + }); + emit(&mut ui, Event::RunNotice { + level: RunNoticeLevel::Warn, + code: "checkpoint_metadata_degraded".into(), + message: "metadata snapshots are disabled for this run".into(), + }); + + insta::assert_snapshot!(rendered(&buffer), @r" + Warning: Metadata checkpoint failed: write failed [write] + Warning: metadata snapshots are disabled for this run [checkpoint_metadata_degraded] + "); + } + #[test] fn tty_parallel_branch_completion_uses_recorded_duration() { let mut ui = ProgressUI::new(true, false); diff --git a/lib/crates/fabro-types/src/lib.rs b/lib/crates/fabro-types/src/lib.rs index d44fac1f6..603ae51c9 100644 --- a/lib/crates/fabro-types/src/lib.rs +++ b/lib/crates/fabro-types/src/lib.rs @@ -64,7 +64,10 @@ pub use run::{ RunProvenance, RunServerProvenance, RunSpec, RunSubjectProvenance, }; pub use run_blob_id::RunBlobId; -pub use run_event::{ActorKind, ActorRef, EventBody, InterviewOption, RunEvent, RunNoticeLevel}; +pub use run_event::{ + ActorKind, ActorRef, EventBody, InterviewOption, MetadataSnapshotFailureKind, + MetadataSnapshotPhase, RunEvent, RunNoticeLevel, +}; pub use run_id::{RunId, fixtures}; pub use run_projection::{NodeState, PendingInterviewRecord, RunProjection}; pub use run_summary::RunSummary; diff --git a/lib/crates/fabro-types/src/run_event/infra.rs b/lib/crates/fabro-types/src/run_event/infra.rs index 3d842b589..3dc726b77 100644 --- a/lib/crates/fabro-types/src/run_event/infra.rs +++ b/lib/crates/fabro-types/src/run_event/infra.rs @@ -1,5 +1,92 @@ use serde::{Deserialize, Serialize}; +/// Legacy `run.notice` codes paired with the new `metadata.snapshot.failed` +/// event for backward compatibility. Display layers suppress these so the +/// typed event renders without a duplicate raw warning. +pub const NOTICE_CODE_CHECKPOINT_METADATA_WRITE_FAILED: &str = "checkpoint_metadata_write_failed"; +pub const NOTICE_CODE_CHECKPOINT_METADATA_PUSH_FAILED: &str = "checkpoint_metadata_push_failed"; + +#[must_use] +pub fn is_metadata_snapshot_compat_notice_code(code: &str) -> bool { + matches!( + code, + NOTICE_CODE_CHECKPOINT_METADATA_WRITE_FAILED | NOTICE_CODE_CHECKPOINT_METADATA_PUSH_FAILED + ) +} + +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Serialize, + Deserialize, + strum::Display, + strum::EnumString, + strum::IntoStaticStr, +)] +#[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] +pub enum MetadataSnapshotPhase { + Init, + Checkpoint, + Finalize, +} + +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Serialize, + Deserialize, + strum::Display, + strum::EnumString, + strum::IntoStaticStr, +)] +#[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] +pub enum MetadataSnapshotFailureKind { + LoadState, + Write, + Push, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct MetadataSnapshotStartedProps { + pub phase: MetadataSnapshotPhase, + pub branch: String, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct MetadataSnapshotCompletedProps { + pub phase: MetadataSnapshotPhase, + pub branch: String, + pub duration_ms: u64, + pub entry_count: usize, + pub bytes: u64, + pub commit_sha: String, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct MetadataSnapshotFailedProps { + pub phase: MetadataSnapshotPhase, + pub branch: String, + pub duration_ms: u64, + pub failure_kind: MetadataSnapshotFailureKind, + pub error: String, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub causes: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub commit_sha: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub entry_count: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub bytes: Option, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct SandboxInitializingProps { pub provider: String, diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index d036d429b..cfbaeabd2 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -136,6 +136,12 @@ pub enum EventBody { RunFailed(RunFailedProps), #[serde(rename = "run.notice")] RunNotice(RunNoticeProps), + #[serde(rename = "metadata.snapshot.started")] + MetadataSnapshotStarted(MetadataSnapshotStartedProps), + #[serde(rename = "metadata.snapshot.completed")] + MetadataSnapshotCompleted(MetadataSnapshotCompletedProps), + #[serde(rename = "metadata.snapshot.failed")] + MetadataSnapshotFailed(MetadataSnapshotFailedProps), #[serde(rename = "stage.started")] StageStarted(StageStartedProps), #[serde(rename = "stage.completed")] @@ -396,6 +402,9 @@ impl EventBody { Self::RunCompleted(_) => "run.completed", Self::RunFailed(_) => "run.failed", Self::RunNotice(_) => "run.notice", + Self::MetadataSnapshotStarted(_) => "metadata.snapshot.started", + Self::MetadataSnapshotCompleted(_) => "metadata.snapshot.completed", + Self::MetadataSnapshotFailed(_) => "metadata.snapshot.failed", Self::StageStarted(_) => "stage.started", Self::StageCompleted(_) => "stage.completed", Self::StageFailed(_) => "stage.failed", @@ -527,6 +536,9 @@ fn is_known_event_name(event: &str) -> bool { | "run.completed" | "run.failed" | "run.notice" + | "metadata.snapshot.started" + | "metadata.snapshot.completed" + | "metadata.snapshot.failed" | "stage.started" | "stage.completed" | "stage.failed" @@ -1216,4 +1228,81 @@ mod tests { assert_eq!(parsed.to_value().unwrap()["event"], value["event"]); } } + + #[test] + fn metadata_snapshot_events_are_known_and_round_trip_json() { + let completed = RunEvent { + id: "evt_metadata_completed".to_string(), + ts: DateTime::parse_from_rfc3339("2026-04-29T12:00:00.000Z") + .unwrap() + .with_timezone(&Utc), + run_id: fixtures::RUN_1, + node_id: None, + node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, + session_id: None, + parent_session_id: None, + tool_call_id: None, + actor: None, + body: EventBody::MetadataSnapshotCompleted( + MetadataSnapshotCompletedProps { + phase: MetadataSnapshotPhase::Checkpoint, + branch: "fabro/metadata/run".to_string(), + duration_ms: 2800, + entry_count: 3, + bytes: 42, + commit_sha: "abc123".to_string(), + }, + ), + }; + + let serialized = completed.to_value().unwrap(); + assert_eq!(serialized["event"], "metadata.snapshot.completed"); + assert_eq!(serialized["properties"]["phase"], "checkpoint"); + assert_eq!(serialized["properties"]["branch"], "fabro/metadata/run"); + assert_eq!(serialized["properties"]["duration_ms"], 2800); + assert_eq!(serialized["properties"]["entry_count"], 3); + assert_eq!(serialized["properties"]["bytes"], 42); + assert_eq!(serialized["properties"]["commit_sha"], "abc123"); + + let parsed = RunEvent::from_value(serialized).unwrap(); + assert_eq!(parsed.event_name(), "metadata.snapshot.completed"); + assert!(matches!( + parsed.body, + EventBody::MetadataSnapshotCompleted(MetadataSnapshotCompletedProps { + phase: MetadataSnapshotPhase::Checkpoint, + .. + }) + )); + } + + #[test] + fn metadata_snapshot_failed_omits_empty_optional_fields() { + let body = EventBody::MetadataSnapshotFailed(MetadataSnapshotFailedProps { + phase: MetadataSnapshotPhase::Init, + branch: "fabro/metadata/run".to_string(), + duration_ms: 15, + failure_kind: MetadataSnapshotFailureKind::LoadState, + error: "state unavailable".to_string(), + causes: Vec::new(), + commit_sha: None, + entry_count: None, + bytes: None, + }); + + let value = serde_json::to_value(&body).unwrap(); + assert_eq!(value["event"], "metadata.snapshot.failed"); + assert_eq!( + value["properties"], + json!({ + "phase": "init", + "branch": "fabro/metadata/run", + "duration_ms": 15, + "failure_kind": "load_state", + "error": "state unavailable" + }) + ); + } } diff --git a/lib/crates/fabro-util/src/lib.rs b/lib/crates/fabro-util/src/lib.rs index 303cfa4a4..5993f2e39 100644 --- a/lib/crates/fabro-util/src/lib.rs +++ b/lib/crates/fabro-util/src/lib.rs @@ -13,6 +13,7 @@ pub mod run_log; pub mod session_secret; pub mod terminal; pub mod text; +pub mod time; pub mod version; pub mod warnings; diff --git a/lib/crates/fabro-util/src/time.rs b/lib/crates/fabro-util/src/time.rs new file mode 100644 index 000000000..fbfeaad39 --- /dev/null +++ b/lib/crates/fabro-util/src/time.rs @@ -0,0 +1,6 @@ +use std::time::Instant; + +#[must_use] +pub fn elapsed_ms(started: Instant) -> u64 { + u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX) +} diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index 8f8bc5297..b216f4711 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -146,6 +146,33 @@ pub enum Event { code: String, message: String, }, + MetadataSnapshotStarted { + phase: fabro_types::MetadataSnapshotPhase, + branch: String, + }, + MetadataSnapshotCompleted { + phase: fabro_types::MetadataSnapshotPhase, + branch: String, + duration_ms: u64, + entry_count: usize, + bytes: u64, + commit_sha: String, + }, + MetadataSnapshotFailed { + phase: fabro_types::MetadataSnapshotPhase, + branch: String, + duration_ms: u64, + failure_kind: fabro_types::MetadataSnapshotFailureKind, + error: String, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + causes: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + commit_sha: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + entry_count: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + bytes: Option, + }, StageStarted { node_id: String, name: String, @@ -679,6 +706,34 @@ impl Event { error!(code, message, "Run notice"); } }, + Self::MetadataSnapshotStarted { phase, branch } => { + debug!(%phase, branch, "Metadata snapshot started"); + } + Self::MetadataSnapshotCompleted { + phase, + branch, + duration_ms, + .. + } => { + debug!(%phase, branch, duration_ms, "Metadata snapshot completed"); + } + Self::MetadataSnapshotFailed { + phase, + branch, + duration_ms, + failure_kind, + error, + .. + } => { + warn!( + %phase, + branch, + duration_ms, + %failure_kind, + error, + "Metadata snapshot failed" + ); + } Self::StageStarted { node_id, name, @@ -1197,6 +1252,9 @@ pub fn event_name(event: &Event) -> &'static str { Event::WorkflowRunCompleted { .. } => "run.completed", Event::WorkflowRunFailed { .. } => "run.failed", Event::RunNotice { .. } => "run.notice", + Event::MetadataSnapshotStarted { .. } => "metadata.snapshot.started", + Event::MetadataSnapshotCompleted { .. } => "metadata.snapshot.completed", + Event::MetadataSnapshotFailed { .. } => "metadata.snapshot.failed", Event::StageStarted { .. } => "stage.started", Event::StageCompleted { .. } => "stage.completed", Event::StageFailed { .. } => "stage.failed", @@ -1660,6 +1718,48 @@ fn event_body_from_event(event: &Event) -> EventBody { code: code.clone(), message: message.clone(), }), + Event::MetadataSnapshotStarted { phase, branch } => { + EventBody::MetadataSnapshotStarted(fabro_types::MetadataSnapshotStartedProps { + phase: *phase, + branch: branch.clone(), + }) + } + Event::MetadataSnapshotCompleted { + phase, + branch, + duration_ms, + entry_count, + bytes, + commit_sha, + } => EventBody::MetadataSnapshotCompleted(fabro_types::MetadataSnapshotCompletedProps { + phase: *phase, + branch: branch.clone(), + duration_ms: *duration_ms, + entry_count: *entry_count, + bytes: *bytes, + commit_sha: commit_sha.clone(), + }), + Event::MetadataSnapshotFailed { + phase, + branch, + duration_ms, + failure_kind, + error, + causes, + commit_sha, + entry_count, + bytes, + } => EventBody::MetadataSnapshotFailed(fabro_types::MetadataSnapshotFailedProps { + phase: *phase, + branch: branch.clone(), + duration_ms: *duration_ms, + failure_kind: *failure_kind, + error: error.clone(), + causes: causes.clone(), + commit_sha: commit_sha.clone(), + entry_count: *entry_count, + bytes: *bytes, + }), Event::StageStarted { index, handler_type, @@ -3627,6 +3727,97 @@ mod tests { } } + #[test] + fn metadata_snapshot_events_map_to_typed_bodies() { + let started = to_run_event(&fixtures::RUN_1, &Event::MetadataSnapshotStarted { + phase: fabro_types::MetadataSnapshotPhase::Init, + branch: "fabro/metadata/run".to_string(), + }); + + assert_eq!(started.event_name(), "metadata.snapshot.started"); + assert!(started.node_id.is_none()); + assert!(started.stage_id.is_none()); + match started.body { + EventBody::MetadataSnapshotStarted(props) => { + assert_eq!(props.phase, fabro_types::MetadataSnapshotPhase::Init); + assert_eq!(props.branch, "fabro/metadata/run"); + } + other => panic!("expected MetadataSnapshotStarted body, got {other:?}"), + } + + let completed = to_run_event(&fixtures::RUN_1, &Event::MetadataSnapshotCompleted { + phase: fabro_types::MetadataSnapshotPhase::Finalize, + branch: "fabro/metadata/run".to_string(), + duration_ms: 2400, + entry_count: 4, + bytes: 512, + commit_sha: "abc123".to_string(), + }); + + assert_eq!(completed.event_name(), "metadata.snapshot.completed"); + match completed.body { + EventBody::MetadataSnapshotCompleted(props) => { + assert_eq!(props.phase, fabro_types::MetadataSnapshotPhase::Finalize); + assert_eq!(props.duration_ms, 2400); + assert_eq!(props.entry_count, 4); + assert_eq!(props.bytes, 512); + assert_eq!(props.commit_sha, "abc123"); + } + other => panic!("expected MetadataSnapshotCompleted body, got {other:?}"), + } + + let failed = to_run_event(&fixtures::RUN_1, &Event::MetadataSnapshotFailed { + phase: fabro_types::MetadataSnapshotPhase::Checkpoint, + branch: "fabro/metadata/run".to_string(), + duration_ms: 120, + failure_kind: fabro_types::MetadataSnapshotFailureKind::Push, + error: "push rejected".to_string(), + causes: vec!["permission denied".to_string()], + commit_sha: Some("def456".to_string()), + entry_count: Some(4), + bytes: Some(512), + }); + + assert_eq!(failed.event_name(), "metadata.snapshot.failed"); + match failed.body { + EventBody::MetadataSnapshotFailed(props) => { + assert_eq!( + props.failure_kind, + fabro_types::MetadataSnapshotFailureKind::Push + ); + assert_eq!(props.commit_sha.as_deref(), Some("def456")); + assert_eq!(props.entry_count, Some(4)); + assert_eq!(props.bytes, Some(512)); + } + other => panic!("expected MetadataSnapshotFailed body, got {other:?}"), + } + } + + #[test] + fn checkpoint_metadata_snapshot_events_can_be_stage_scoped() { + let scope = StageScope { + node_id: "build".to_string(), + visit: 2, + parallel_group_id: Some(StageId::new("fanout", 1)), + parallel_branch_id: Some(ParallelBranchId::new(StageId::new("fanout", 1), 0)), + }; + let stored = to_run_event_at( + &fixtures::RUN_1, + &Event::MetadataSnapshotStarted { + phase: fabro_types::MetadataSnapshotPhase::Checkpoint, + branch: "fabro/metadata/run".to_string(), + }, + Utc::now(), + Some(&scope), + ); + + assert_eq!(stored.node_id.as_deref(), Some("build")); + assert_eq!(stored.node_label.as_deref(), Some("build")); + assert_eq!(stored.stage_id, Some(StageId::new("build", 2))); + assert_eq!(stored.parallel_group_id, scope.parallel_group_id); + assert_eq!(stored.parallel_branch_id, scope.parallel_branch_id); + } + #[test] fn agent_assistant_message_populates_agent_actor() { let stored = to_run_event(&fixtures::RUN_1, &Event::Agent { diff --git a/lib/crates/fabro-workflow/src/lifecycle/git.rs b/lib/crates/fabro-workflow/src/lifecycle/git.rs index 3dc420676..d9aea0ea0 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/git.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/git.rs @@ -1,4 +1,5 @@ use std::sync::{Arc, Mutex}; +use std::time::Instant; use async_trait::async_trait; use fabro_core::error::{Error as CoreError, Result as CoreResult}; @@ -7,9 +8,12 @@ use fabro_core::lifecycle::RunLifecycle; use fabro_core::outcome::NodeResult; use fabro_core::state::ExecutionState; use fabro_types::RunId; +use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase}; +use fabro_util::error::collect_causes; +use fabro_util::time::elapsed_ms; use crate::artifact; -use crate::event::{Emitter, Event, RunNoticeLevel}; +use crate::event::{Emitter, Event, RunNoticeLevel, StageScope}; use crate::graph::{WorkflowGraph, WorkflowNode}; use crate::lifecycle::event::stage_scope_for; use crate::outcome::BilledModelUsage; @@ -17,7 +21,7 @@ use crate::run_dump::RunDump; use crate::run_options::RunOptions; use crate::runtime_store::RunStoreHandle; use crate::sandbox_git::{checked_git_checkpoint, git_diff}; -use crate::sandbox_metadata::{SandboxGitRuntime, SandboxMetadataWriter}; +use crate::sandbox_metadata::{MetadataSnapshot, SandboxGitRuntime, SandboxMetadataWriter}; type WfRunState = ExecutionState>; type WfNodeResult = NodeResult>; @@ -79,23 +83,42 @@ impl RunLifecycle for GitLifecycle { // Reset last_git_sha (diff base parity) *self.last_git_sha.lock().unwrap() = None; *self.checkpoint_git_result.lock().unwrap() = None; - if self - .run_options - .git - .as_ref() - .and_then(|g| g.meta_branch.as_ref()) - .is_some() - { + if let Some(meta_branch) = self.metadata_branch().map(str::to_string) { + if self.metadata_runtime.metadata_degraded() { + return Ok(()); + } + let phase = MetadataSnapshotPhase::Init; + let started = Instant::now(); + self.emit_metadata_snapshot_started(phase, &meta_branch, None); match self.run_store.state().await { Ok(state) => { let dump = RunDump::from_projection(&state); - let _ = self.write_metadata_snapshot(&dump, "init run").await; + let _ = self + .write_metadata_snapshot( + phase, + &meta_branch, + started, + &dump, + "init run", + None, + ) + .await; } Err(err) => { - self.emit_metadata_warning( - "checkpoint_metadata_write_failed", - format!("failed to load run state for metadata init: {err}"), + let message = format!("failed to load run state for metadata init: {err}"); + self.emit_metadata_snapshot_failed( + phase, + &meta_branch, + started, + MetadataSnapshotFailureKind::LoadState, + message.clone(), + collect_causes(err.as_ref()), + None, + None, + None, + None, ); + self.emit_metadata_warning("checkpoint_metadata_write_failed", message); } } } @@ -127,19 +150,50 @@ impl RunLifecycle for GitLifecycle { std::collections::HashMap::new(), None, ); - let shadow_sha = match self.run_store.state().await { - Ok(mut projection) => { - projection.checkpoint = Some(checkpoint); - let dump = RunDump::from_projection(&projection); - self.write_metadata_snapshot(&dump, "checkpoint").await - } - Err(err) => { - self.emit_metadata_warning( - "checkpoint_metadata_write_failed", - format!("failed to load run state for metadata checkpoint: {err}"), - ); + let shadow_sha = if let Some(meta_branch) = self.metadata_branch().map(str::to_string) { + if self.metadata_runtime.metadata_degraded() { None + } else { + let phase = MetadataSnapshotPhase::Checkpoint; + let started = Instant::now(); + let scope = stage_scope_for(state, node_id); + self.emit_metadata_snapshot_started(phase, &meta_branch, Some(&scope)); + match self.run_store.state().await { + Ok(mut projection) => { + projection.checkpoint = Some(checkpoint); + let dump = RunDump::from_projection(&projection); + self.write_metadata_snapshot( + phase, + &meta_branch, + started, + &dump, + "checkpoint", + Some(&scope), + ) + .await + } + Err(err) => { + let message = + format!("failed to load run state for metadata checkpoint: {err}"); + self.emit_metadata_snapshot_failed( + phase, + &meta_branch, + started, + MetadataSnapshotFailureKind::LoadState, + message.clone(), + collect_causes(err.as_ref()), + None, + None, + None, + Some(&scope), + ); + self.emit_metadata_warning("checkpoint_metadata_write_failed", message); + None + } + } } + } else { + None }; // Run branch commit via sandbox @@ -238,15 +292,25 @@ impl RunLifecycle for GitLifecycle { } impl GitLifecycle { - async fn write_metadata_snapshot(&self, dump: &RunDump, message: &str) -> Option { + fn metadata_branch(&self) -> Option<&str> { + self.run_options + .git + .as_ref() + .and_then(|git| git.meta_branch.as_deref()) + } + + async fn write_metadata_snapshot( + &self, + phase: MetadataSnapshotPhase, + meta_branch: &str, + started: Instant, + dump: &RunDump, + message: &str, + scope: Option<&StageScope>, + ) -> Option { if self.metadata_runtime.metadata_degraded() { return None; } - let meta_branch = self - .run_options - .git - .as_ref() - .and_then(|git| git.meta_branch.as_deref())?; let run_id = self.run_id.to_string(); let writer = SandboxMetadataWriter::new( @@ -259,23 +323,129 @@ impl GitLifecycle { match writer.write_snapshot(dump, message).await { Ok(snapshot) => { if let Some(detail) = snapshot.push_error.as_deref() { - self.emit_metadata_warning( - "checkpoint_metadata_push_failed", - format!("failed to push metadata ref refs/heads/{meta_branch}: {detail}"), + let message = + format!("failed to push metadata ref refs/heads/{meta_branch}: {detail}"); + self.emit_metadata_snapshot_failed( + phase, + meta_branch, + started, + MetadataSnapshotFailureKind::Push, + message.clone(), + Vec::new(), + Some(snapshot.commit_sha.clone()), + Some(snapshot.entry_count), + Some(snapshot.bytes), + scope, + ); + self.emit_metadata_warning("checkpoint_metadata_push_failed", message); + } else { + self.emit_metadata_snapshot_completed( + phase, + meta_branch, + started, + &snapshot, + scope, ); } Some(snapshot.commit_sha) } Err(err) => { - self.emit_metadata_warning( - "checkpoint_metadata_write_failed", - format!("failed to write checkpoint metadata: {err}"), + let message = format!("failed to write checkpoint metadata: {err}"); + self.emit_metadata_snapshot_failed( + phase, + meta_branch, + started, + MetadataSnapshotFailureKind::Write, + message.clone(), + collect_causes(&err), + None, + None, + None, + scope, ); + self.emit_metadata_warning("checkpoint_metadata_write_failed", message); None } } } + fn emit_metadata_snapshot_started( + &self, + phase: MetadataSnapshotPhase, + branch: &str, + scope: Option<&StageScope>, + ) { + self.emit_metadata_snapshot_event( + &Event::MetadataSnapshotStarted { + phase, + branch: branch.to_string(), + }, + scope, + ); + } + + fn emit_metadata_snapshot_completed( + &self, + phase: MetadataSnapshotPhase, + branch: &str, + started: Instant, + snapshot: &MetadataSnapshot, + scope: Option<&StageScope>, + ) { + self.emit_metadata_snapshot_event( + &Event::MetadataSnapshotCompleted { + phase, + branch: branch.to_string(), + duration_ms: elapsed_ms(started), + entry_count: snapshot.entry_count, + bytes: snapshot.bytes, + commit_sha: snapshot.commit_sha.clone(), + }, + scope, + ); + } + + #[allow( + clippy::too_many_arguments, + reason = "Metadata failure event carries the full event contract explicitly." + )] + fn emit_metadata_snapshot_failed( + &self, + phase: MetadataSnapshotPhase, + branch: &str, + started: Instant, + failure_kind: MetadataSnapshotFailureKind, + error: String, + causes: Vec, + commit_sha: Option, + entry_count: Option, + bytes: Option, + scope: Option<&StageScope>, + ) { + self.emit_metadata_snapshot_event( + &Event::MetadataSnapshotFailed { + phase, + branch: branch.to_string(), + duration_ms: elapsed_ms(started), + failure_kind, + error, + causes, + commit_sha, + entry_count, + bytes, + }, + scope, + ); + } + + fn emit_metadata_snapshot_event(&self, event: &Event, scope: Option<&StageScope>) { + if let Some(scope) = scope { + self.emitter.emit_scoped(event, scope); + } else { + self.emitter.emit(event); + } + } + fn emit_metadata_warning(&self, code: &str, message: String) { if self.metadata_runtime.mark_metadata_degraded() { self.emitter.emit(&Event::RunNotice { @@ -286,3 +456,514 @@ impl GitLifecycle { } } } + +#[cfg(test)] +mod tests { + use std::collections::{BTreeMap, HashMap}; + use std::path::Path; + use std::sync::Arc; + use std::time::Duration; + + use anyhow::Result; + use async_trait::async_trait; + use bytes::Bytes; + use fabro_core::graph::Graph as CoreGraph; + use fabro_core::lifecycle::RunLifecycle; + use fabro_core::state::ExecutionState; + use fabro_graphviz::graph::types::{AttrValue, Edge, Graph, Node}; + use fabro_store::{Database, EventEnvelope, RunDatabase, RunProjection}; + use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase}; + use fabro_types::{EventBody, RunBlobId, RunEvent, WorkflowSettings, fixtures}; + use object_store::memory::InMemory; + + use super::*; + use crate::event::append_event; + use crate::outcome::{Outcome, StageStatus}; + use crate::pipeline::write_finalize_commit; + use crate::records::Conclusion; + use crate::run_options::GitCheckpointOptions; + use crate::runtime_store::{RunStoreBackend, RunStoreHandle}; + use crate::services::RunServices; + + #[expect( + clippy::disallowed_methods, + reason = "metadata event tests use synchronous git commands to set up temporary repositories" + )] + fn init_git_repo(repo: &Path) { + let init = std::process::Command::new("git") + .args(["init", "-b", "main"]) + .current_dir(repo) + .output() + .unwrap(); + assert!(init.status.success()); + for (key, value) in [("user.name", "Test"), ("user.email", "test@test.com")] { + let config = std::process::Command::new("git") + .args(["config", key, value]) + .current_dir(repo) + .output() + .unwrap(); + assert!(config.status.success()); + } + let commit = std::process::Command::new("git") + .args(["commit", "--allow-empty", "-m", "initial"]) + .current_dir(repo) + .output() + .unwrap(); + assert!(commit.status.success()); + } + + fn workflow_graph() -> WorkflowGraph { + let mut graph = Graph::new("metadata"); + let mut start = Node::new("start"); + start.attrs.insert( + "shape".to_string(), + AttrValue::String("Mdiamond".to_string()), + ); + graph.nodes.insert("start".to_string(), start); + let mut build = Node::new("build"); + build + .attrs + .insert("shape".to_string(), AttrValue::String("box".to_string())); + graph.nodes.insert("build".to_string(), build); + let mut exit = Node::new("exit"); + exit.attrs.insert( + "shape".to_string(), + AttrValue::String("Msquare".to_string()), + ); + graph.nodes.insert("exit".to_string(), exit); + graph.edges.push(Edge::new("start", "build")); + graph.edges.push(Edge::new("build", "exit")); + WorkflowGraph(Arc::new(graph)) + } + + fn run_options(run_dir: &Path, meta_branch: &str) -> Arc { + Arc::new(RunOptions { + settings: WorkflowSettings::default(), + run_dir: run_dir.to_path_buf(), + cancel_token: None, + run_id: fixtures::RUN_1, + labels: HashMap::new(), + workflow_slug: Some("metadata".to_string()), + github_app: None, + pre_run_git: None, + fork_source_ref: None, + base_branch: None, + display_base_sha: None, + git: Some(GitCheckpointOptions { + base_sha: None, + run_branch: None, + meta_branch: Some(meta_branch.to_string()), + }), + }) + } + + async fn run_store(run_id: fabro_types::RunId) -> RunDatabase { + let store = Arc::new(Database::new( + Arc::new(InMemory::new()), + "", + Duration::from_millis(1), + None, + )); + let run_store = store.create_run(&run_id).await.unwrap(); + append_event(&run_store, &run_id, &Event::RunCreated { + run_id, + settings: serde_json::to_value(WorkflowSettings::default()).unwrap(), + graph: serde_json::to_value(fabro_types::Graph::new("metadata")).unwrap(), + workflow_source: None, + workflow_config: None, + labels: BTreeMap::new(), + run_dir: "/tmp/run".to_string(), + source_directory: Some("/tmp/project".to_string()), + workflow_slug: Some("metadata".to_string()), + db_prefix: None, + provenance: None, + manifest_blob: None, + git: None, + fork_source_ref: None, + in_place: false, + }) + .await + .unwrap(); + run_store + } + + fn record_events(emitter: &Arc) -> Arc>> { + let events = Arc::new(std::sync::Mutex::new(Vec::new())); + let captured = Arc::clone(&events); + emitter.on_event(move |event| { + captured.lock().unwrap().push(event.clone()); + }); + events + } + + fn git_lifecycle( + repo: &Path, + emitter: Arc, + run_store: RunStoreHandle, + run_options: Arc, + metadata_runtime: Arc, + ) -> GitLifecycle { + GitLifecycle { + sandbox: Arc::new(fabro_agent::LocalSandbox::new(repo.to_path_buf())), + emitter, + run_id: fixtures::RUN_1, + run_store, + run_options, + metadata_runtime, + start_node_id: Some("start".to_string()), + checkpoint_git_result: Arc::new(Mutex::new(None)), + last_git_sha: Arc::new(Mutex::new(None)), + } + } + + #[tokio::test] + async fn init_metadata_snapshot_success_emits_started_completed_unscoped() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let branch = "fabro/metadata/run"; + let run_store = run_store(fixtures::RUN_1).await; + let handle = RunStoreHandle::local(run_store.clone()); + let state = handle.state().await.unwrap(); + let expected_entries = RunDump::from_projection(&state).git_entries().unwrap(); + let expected_entry_count = expected_entries.len(); + let expected_bytes = expected_entries + .iter() + .map(|(_, bytes)| u64::try_from(bytes.len()).unwrap_or(u64::MAX)) + .sum::(); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + handle, + run_options(repo_dir.path(), branch), + Arc::new(SandboxGitRuntime::new()), + ); + let graph = workflow_graph(); + let state = ExecutionState::new(&graph).unwrap(); + + lifecycle.on_run_start(&graph, &state).await.unwrap(); + + let events = events.lock().unwrap(); + assert_eq!(events.len(), 2); + assert_eq!(events[0].event_name(), "metadata.snapshot.started"); + assert_eq!(events[1].event_name(), "metadata.snapshot.completed"); + assert!(events[0].node_id.is_none()); + assert!(events[1].node_id.is_none()); + match &events[1].body { + EventBody::MetadataSnapshotCompleted(props) => { + assert_eq!(props.phase, MetadataSnapshotPhase::Init); + assert_eq!(props.branch, branch); + assert_eq!(props.entry_count, expected_entry_count); + assert_eq!(props.bytes, expected_bytes); + assert!(!props.commit_sha.is_empty()); + } + other => panic!("expected metadata completed event, got {other:?}"), + } + } + + #[tokio::test] + async fn init_metadata_load_state_failure_emits_failed_before_notice() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let branch = "fabro/metadata/run"; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + RunStoreHandle::new(Arc::new(FailingStateStore)), + run_options(repo_dir.path(), branch), + Arc::new(SandboxGitRuntime::new()), + ); + let graph = workflow_graph(); + let state = ExecutionState::new(&graph).unwrap(); + + lifecycle.on_run_start(&graph, &state).await.unwrap(); + + let events = events.lock().unwrap(); + let names = events.iter().map(RunEvent::event_name).collect::>(); + assert_eq!(names, vec![ + "metadata.snapshot.started", + "metadata.snapshot.failed", + "run.notice", + ]); + match &events[1].body { + EventBody::MetadataSnapshotFailed(props) => { + assert_eq!(props.phase, MetadataSnapshotPhase::Init); + assert_eq!(props.failure_kind, MetadataSnapshotFailureKind::LoadState); + assert_eq!(props.commit_sha, None); + assert_eq!(props.entry_count, None); + assert_eq!(props.bytes, None); + } + other => panic!("expected metadata failed event, got {other:?}"), + } + } + + #[tokio::test] + #[expect( + clippy::disallowed_methods, + reason = "metadata push-failure test uses a synchronous git command to configure a temporary remote" + )] + async fn init_metadata_push_failure_emits_failed_with_snapshot_accounting() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let missing_origin = repo_dir.path().join("missing-origin.git"); + let remote = std::process::Command::new("git") + .args(["remote", "add", "origin", missing_origin.to_str().unwrap()]) + .current_dir(repo_dir.path()) + .output() + .unwrap(); + assert!(remote.status.success()); + let branch = "fabro/metadata/run"; + let run_store = run_store(fixtures::RUN_1).await; + let handle = RunStoreHandle::local(run_store.clone()); + let state = handle.state().await.unwrap(); + let expected_entries = RunDump::from_projection(&state).git_entries().unwrap(); + let expected_entry_count = expected_entries.len(); + let expected_bytes = expected_entries + .iter() + .map(|(_, bytes)| u64::try_from(bytes.len()).unwrap_or(u64::MAX)) + .sum::(); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let runtime = Arc::new(SandboxGitRuntime::new()); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + handle, + run_options(repo_dir.path(), branch), + Arc::clone(&runtime), + ); + let graph = workflow_graph(); + let state = ExecutionState::new(&graph).unwrap(); + + lifecycle.on_run_start(&graph, &state).await.unwrap(); + + assert!(runtime.metadata_degraded()); + let events = events.lock().unwrap(); + let names = events.iter().map(RunEvent::event_name).collect::>(); + assert_eq!(names, vec![ + "metadata.snapshot.started", + "metadata.snapshot.failed", + "run.notice", + ]); + match &events[1].body { + EventBody::MetadataSnapshotFailed(props) => { + assert_eq!(props.phase, MetadataSnapshotPhase::Init); + assert_eq!(props.failure_kind, MetadataSnapshotFailureKind::Push); + assert!(props.commit_sha.as_ref().is_some_and(|sha| !sha.is_empty())); + assert_eq!(props.entry_count, Some(expected_entry_count)); + assert_eq!(props.bytes, Some(expected_bytes)); + } + other => panic!("expected metadata failed event, got {other:?}"), + } + } + + #[tokio::test] + async fn checkpoint_metadata_load_state_failure_emits_scoped_failed_before_notice() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let branch = "fabro/metadata/run"; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + RunStoreHandle::new(Arc::new(FailingStateStore)), + run_options(repo_dir.path(), branch), + Arc::new(SandboxGitRuntime::new()), + ); + let graph = workflow_graph(); + let node = graph.get_node("build").unwrap(); + let mut state = ExecutionState::new(&graph).unwrap(); + state.increment_visits("build"); + let result = WfNodeResult::new(Outcome::success(), Duration::from_millis(10), 1, 1); + + lifecycle + .on_checkpoint(&node, &result, Some("exit"), &state) + .await + .unwrap(); + + let events = events.lock().unwrap(); + let names = events.iter().map(RunEvent::event_name).collect::>(); + assert_eq!(names, vec![ + "metadata.snapshot.started", + "metadata.snapshot.failed", + "run.notice", + ]); + assert_eq!(events[1].node_id.as_deref(), Some("build")); + match &events[1].body { + EventBody::MetadataSnapshotFailed(props) => { + assert_eq!(props.phase, MetadataSnapshotPhase::Checkpoint); + assert_eq!(props.failure_kind, MetadataSnapshotFailureKind::LoadState); + } + other => panic!("expected metadata failed event, got {other:?}"), + } + } + + #[tokio::test] + async fn checkpoint_metadata_snapshot_success_emits_scoped_events() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let branch = "fabro/metadata/run"; + let run_store = run_store(fixtures::RUN_1).await; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + RunStoreHandle::local(run_store), + run_options(repo_dir.path(), branch), + Arc::new(SandboxGitRuntime::new()), + ); + let graph = workflow_graph(); + let node = graph.get_node("build").unwrap(); + let mut state = ExecutionState::new(&graph).unwrap(); + state.increment_visits("build"); + let result = WfNodeResult::new(Outcome::success(), Duration::from_millis(10), 1, 1); + + lifecycle + .on_checkpoint(&node, &result, Some("exit"), &state) + .await + .unwrap(); + + let events = events.lock().unwrap(); + assert_eq!(events[0].event_name(), "metadata.snapshot.started"); + assert_eq!(events[1].event_name(), "metadata.snapshot.completed"); + assert_eq!(events[0].node_id.as_deref(), Some("build")); + assert_eq!( + events[0] + .stage_id + .as_ref() + .map(ToString::to_string) + .as_deref(), + Some("build@1") + ); + assert_eq!(events[1].node_id.as_deref(), Some("build")); + assert_eq!( + events[1] + .stage_id + .as_ref() + .map(ToString::to_string) + .as_deref(), + Some("build@1") + ); + } + + #[tokio::test] + async fn degraded_metadata_runtime_skips_snapshot_events() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let runtime = Arc::new(SandboxGitRuntime::new()); + runtime.mark_metadata_degraded(); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + RunStoreHandle::local(run_store(fixtures::RUN_1).await), + run_options(repo_dir.path(), "fabro/metadata/run"), + runtime, + ); + let graph = workflow_graph(); + let state = ExecutionState::new(&graph).unwrap(); + + lifecycle.on_run_start(&graph, &state).await.unwrap(); + + assert!(events.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn degraded_after_init_failure_skips_later_checkpoint_and_finalize_metadata_events() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let runtime = Arc::new(SandboxGitRuntime::new()); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + RunStoreHandle::new(Arc::new(FailingStateStore)), + run_options(repo_dir.path(), "fabro/metadata/run"), + runtime, + ); + let graph = workflow_graph(); + let state = ExecutionState::new(&graph).unwrap(); + + lifecycle.on_run_start(&graph, &state).await.unwrap(); + let after_init = events.lock().unwrap().len(); + let node = graph.get_node("build").unwrap(); + let mut checkpoint_state = ExecutionState::new(&graph).unwrap(); + checkpoint_state.increment_visits("build"); + let result = WfNodeResult::new(Outcome::success(), Duration::from_millis(10), 1, 1); + lifecycle + .on_checkpoint(&node, &result, Some("exit"), &checkpoint_state) + .await + .unwrap(); + let finalize_services = RunServices::new( + RunStoreHandle::new(Arc::new(FailingStateStore)), + Arc::clone(&lifecycle.emitter), + Arc::new(fabro_agent::LocalSandbox::new( + repo_dir.path().to_path_buf(), + )), + None, + None, + fabro_model::Provider::Anthropic, + Arc::new(fabro_auth::EnvCredentialSource::new()), + Arc::clone(&lifecycle.metadata_runtime), + ); + let conclusion = Conclusion { + timestamp: chrono::Utc::now(), + status: StageStatus::Success, + duration_ms: 10, + failure_reason: None, + final_git_commit_sha: None, + stages: Vec::new(), + billing: None, + total_retries: 0, + }; + write_finalize_commit( + lifecycle.run_options.as_ref(), + &finalize_services, + &conclusion, + ) + .await; + + let events = events.lock().unwrap(); + assert_eq!(events.len(), after_init); + assert_eq!( + events.iter().map(RunEvent::event_name).collect::>(), + vec![ + "metadata.snapshot.started", + "metadata.snapshot.failed", + "run.notice", + ] + ); + } + + struct FailingStateStore; + + #[async_trait] + impl RunStoreBackend for FailingStateStore { + async fn load_state(&self) -> Result { + Err(anyhow::anyhow!("state unavailable")) + } + + async fn list_events(&self) -> Result> { + Ok(Vec::new()) + } + + async fn append_run_event(&self, _event: &RunEvent) -> Result<()> { + Ok(()) + } + + async fn write_blob(&self, data: &[u8]) -> Result { + Ok(RunBlobId::new(data)) + } + + async fn read_blob(&self, _id: &RunBlobId) -> Result> { + Ok(None) + } + } +} diff --git a/lib/crates/fabro-workflow/src/pipeline/finalize.rs b/lib/crates/fabro-workflow/src/pipeline/finalize.rs index aa8c0b489..251ef28d0 100644 --- a/lib/crates/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/finalize.rs @@ -1,5 +1,10 @@ +use std::time::Instant; + use fabro_hooks::{HookContext, HookEvent}; +use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase}; use fabro_types::{BilledTokenCounts, EventBody}; +use fabro_util::error::collect_causes; +use fabro_util::time::elapsed_ms; use super::types::{Concluded, FinalizeOptions, Retroed}; use crate::error::Error; @@ -11,7 +16,7 @@ use crate::run_options::RunOptions; use crate::run_status::{FailureReason, RunStatus, SuccessReason}; use crate::runtime_store::RunStoreHandle; use crate::sandbox_git::git_diff_with_timeout; -use crate::sandbox_metadata::SandboxMetadataWriter; +use crate::sandbox_metadata::{MetadataSnapshot, SandboxMetadataWriter}; use crate::services::RunServices; pub fn classify_engine_result( @@ -153,47 +158,138 @@ pub async fn write_finalize_commit( return; }; + let phase = MetadataSnapshotPhase::Finalize; + let started = Instant::now(); + emit_metadata_snapshot_started(services, phase, meta_branch); + let mut projection = match services.run_store.state().await { Ok(state) => state, Err(err) => { - emit_metadata_warning( + let message = format!("failed to load run state for final metadata snapshot: {err}"); + emit_metadata_snapshot_failed( services, - "checkpoint_metadata_write_failed", - format!("failed to load run state for final metadata snapshot: {err}"), + phase, + meta_branch, + started, + MetadataSnapshotFailureKind::LoadState, + message.clone(), + collect_causes(err.as_ref()), + None, + None, + None, ); + emit_metadata_warning(services, "checkpoint_metadata_write_failed", message); return; } }; projection.conclusion = Some(conclusion.clone()); let dump = RunDump::from_projection(&projection); - let Some(spec_run_id) = projection.spec.as_ref().map(|spec| spec.run_id.to_string()) else { - return; - }; + let run_id = run_options.run_id.to_string(); let writer = SandboxMetadataWriter::new( &*services.sandbox, &services.metadata_runtime, - &spec_run_id, + &run_id, meta_branch, run_options.git_author(), ); match writer.write_snapshot(&dump, "finalize run").await { Ok(snapshot) => { if let Some(detail) = snapshot.push_error.as_deref() { - emit_metadata_warning( + let message = + format!("failed to push metadata ref refs/heads/{meta_branch}: {detail}"); + emit_metadata_snapshot_failed( services, - "checkpoint_metadata_push_failed", - format!("failed to push metadata ref refs/heads/{meta_branch}: {detail}"), + phase, + meta_branch, + started, + MetadataSnapshotFailureKind::Push, + message.clone(), + Vec::new(), + Some(snapshot.commit_sha.clone()), + Some(snapshot.entry_count), + Some(snapshot.bytes), ); + emit_metadata_warning(services, "checkpoint_metadata_push_failed", message); + } else { + emit_metadata_snapshot_completed(services, phase, meta_branch, started, &snapshot); } } - Err(err) => emit_metadata_warning( - services, - "checkpoint_metadata_write_failed", - format!("failed to write final checkpoint metadata: {err}"), - ), + Err(err) => { + let message = format!("failed to write final checkpoint metadata: {err}"); + emit_metadata_snapshot_failed( + services, + phase, + meta_branch, + started, + MetadataSnapshotFailureKind::Write, + message.clone(), + collect_causes(&err), + None, + None, + None, + ); + emit_metadata_warning(services, "checkpoint_metadata_write_failed", message); + } } } +fn emit_metadata_snapshot_started( + services: &RunServices, + phase: MetadataSnapshotPhase, + branch: &str, +) { + services.emitter.emit(&Event::MetadataSnapshotStarted { + phase, + branch: branch.to_string(), + }); +} + +fn emit_metadata_snapshot_completed( + services: &RunServices, + phase: MetadataSnapshotPhase, + branch: &str, + started: Instant, + snapshot: &MetadataSnapshot, +) { + services.emitter.emit(&Event::MetadataSnapshotCompleted { + phase, + branch: branch.to_string(), + duration_ms: elapsed_ms(started), + entry_count: snapshot.entry_count, + bytes: snapshot.bytes, + commit_sha: snapshot.commit_sha.clone(), + }); +} + +#[allow( + clippy::too_many_arguments, + reason = "Metadata failure event carries the full event contract explicitly." +)] +fn emit_metadata_snapshot_failed( + services: &RunServices, + phase: MetadataSnapshotPhase, + branch: &str, + started: Instant, + failure_kind: MetadataSnapshotFailureKind, + error: String, + causes: Vec, + commit_sha: Option, + entry_count: Option, + bytes: Option, +) { + services.emitter.emit(&Event::MetadataSnapshotFailed { + phase, + branch: branch.to_string(), + duration_ms: elapsed_ms(started), + failure_kind, + error, + causes, + commit_sha, + entry_count, + bytes, + }); +} + fn emit_metadata_warning(services: &RunServices, code: &str, message: String) { if services.metadata_runtime.mark_metadata_degraded() { services.emitter.notice(RunNoticeLevel::Warn, code, message); @@ -419,18 +515,25 @@ pub async fn finalize(retroed: Retroed, options: &FinalizeOptions) -> Result RunId { fixtures::RUN_1 @@ -453,6 +556,16 @@ mod tests { } } + fn test_git_run_options(run_dir: &std::path::Path, meta_branch: &str) -> RunOptions { + let mut options = test_run_options(run_dir); + options.git = Some(GitCheckpointOptions { + base_sha: None, + run_branch: None, + meta_branch: Some(meta_branch.to_string()), + }); + options + } + fn test_store() -> Arc { Arc::new(Database::new( Arc::new(InMemory::new()), @@ -462,6 +575,84 @@ mod tests { )) } + async fn seeded_run_store() -> RunDatabase { + let run_store = test_store().create_run(&test_run_id()).await.unwrap(); + append_event(&run_store, &test_run_id(), &Event::RunCreated { + run_id: test_run_id(), + settings: serde_json::to_value(WorkflowSettings::default()).unwrap(), + graph: serde_json::to_value(fabro_types::Graph::new("metadata")).unwrap(), + workflow_source: None, + workflow_config: None, + labels: std::collections::BTreeMap::new(), + run_dir: "/tmp/run".to_string(), + source_directory: Some("/tmp/project".to_string()), + workflow_slug: Some("metadata".to_string()), + db_prefix: None, + provenance: None, + manifest_blob: None, + git: None, + fork_source_ref: None, + in_place: false, + }) + .await + .unwrap(); + run_store + } + + #[expect( + clippy::disallowed_methods, + reason = "metadata event tests use synchronous git commands to set up temporary repositories" + )] + fn init_git_repo(repo: &Path) { + let init = std::process::Command::new("git") + .args(["init", "-b", "main"]) + .current_dir(repo) + .output() + .unwrap(); + assert!(init.status.success()); + for (key, value) in [("user.name", "Test"), ("user.email", "test@test.com")] { + let config = std::process::Command::new("git") + .args(["config", key, value]) + .current_dir(repo) + .output() + .unwrap(); + assert!(config.status.success()); + } + let commit = std::process::Command::new("git") + .args(["commit", "--allow-empty", "-m", "initial"]) + .current_dir(repo) + .output() + .unwrap(); + assert!(commit.status.success()); + } + + fn record_events(emitter: &Arc) -> Arc>> { + let events = Arc::new(std::sync::Mutex::new(Vec::new())); + let captured = Arc::clone(&events); + emitter.on_event(move |event| { + captured.lock().unwrap().push(event.clone()); + }); + events + } + + fn test_services( + run_store: RunStoreHandle, + emitter: Arc, + sandbox: Arc, + metadata_runtime: Arc, + ) -> Arc { + RunServices::new( + run_store, + emitter, + sandbox, + None, + None, + fabro_model::Provider::Anthropic, + Arc::new(fabro_auth::EnvCredentialSource::new()), + metadata_runtime, + ) + } + #[tokio::test] async fn finalize_persists_conclusion_in_projection() { let temp = tempfile::tempdir().unwrap(); @@ -506,4 +697,200 @@ mod tests { assert_eq!(concluded.conclusion.status, StageStatus::Success); } + + #[tokio::test] + async fn finalize_metadata_snapshot_success_emits_started_completed_unscoped() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let branch = "fabro/metadata/run"; + let run_store = seeded_run_store().await; + let handle = RunStoreHandle::local(run_store.clone()); + let conclusion = Conclusion { + timestamp: chrono::Utc::now(), + status: StageStatus::Success, + duration_ms: 10, + failure_reason: None, + final_git_commit_sha: None, + stages: Vec::new(), + billing: None, + total_retries: 0, + }; + let emitter = Arc::new(Emitter::new(test_run_id())); + let events = record_events(&emitter); + let services = test_services( + handle, + emitter, + Arc::new(fabro_agent::LocalSandbox::new( + repo_dir.path().to_path_buf(), + )), + Arc::new(SandboxGitRuntime::new()), + ); + let run_options = test_git_run_options(repo_dir.path(), branch); + + write_finalize_commit(&run_options, &services, &conclusion).await; + + let events = events.lock().unwrap(); + assert_eq!(events.len(), 2); + assert_eq!(events[0].event_name(), "metadata.snapshot.started"); + assert_eq!(events[1].event_name(), "metadata.snapshot.completed"); + assert!(events[0].node_id.is_none()); + match &events[1].body { + EventBody::MetadataSnapshotCompleted(props) => { + assert_eq!(props.phase, MetadataSnapshotPhase::Finalize); + assert_eq!(props.branch, branch); + assert!(!props.commit_sha.is_empty()); + } + other => panic!("expected metadata completed event, got {other:?}"), + } + } + + #[tokio::test] + async fn finalize_metadata_load_state_failure_emits_failed_before_notice() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let emitter = Arc::new(Emitter::new(test_run_id())); + let events = record_events(&emitter); + let services = test_services( + RunStoreHandle::new(Arc::new(FailingStateStore)), + emitter, + Arc::new(fabro_agent::LocalSandbox::new( + repo_dir.path().to_path_buf(), + )), + Arc::new(SandboxGitRuntime::new()), + ); + let run_options = test_git_run_options(repo_dir.path(), "fabro/metadata/run"); + let conclusion = Conclusion { + timestamp: chrono::Utc::now(), + status: StageStatus::Success, + duration_ms: 10, + failure_reason: None, + final_git_commit_sha: None, + stages: Vec::new(), + billing: None, + total_retries: 0, + }; + + write_finalize_commit(&run_options, &services, &conclusion).await; + + let events = events.lock().unwrap(); + let names = events.iter().map(RunEvent::event_name).collect::>(); + assert_eq!(names, vec![ + "metadata.snapshot.started", + "metadata.snapshot.failed", + "run.notice", + ]); + match &events[1].body { + EventBody::MetadataSnapshotFailed(props) => { + assert_eq!(props.phase, MetadataSnapshotPhase::Finalize); + assert_eq!(props.failure_kind, MetadataSnapshotFailureKind::LoadState); + } + other => panic!("expected metadata failed event, got {other:?}"), + } + } + + #[tokio::test] + async fn degraded_metadata_runtime_skips_finalize_metadata_events() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let run_store = seeded_run_store().await; + let emitter = Arc::new(Emitter::new(test_run_id())); + let events = record_events(&emitter); + let runtime = Arc::new(SandboxGitRuntime::new()); + runtime.mark_metadata_degraded(); + let services = test_services( + RunStoreHandle::local(run_store), + emitter, + Arc::new(fabro_agent::LocalSandbox::new( + repo_dir.path().to_path_buf(), + )), + runtime, + ); + let run_options = test_git_run_options(repo_dir.path(), "fabro/metadata/run"); + let conclusion = Conclusion { + timestamp: chrono::Utc::now(), + status: StageStatus::Success, + duration_ms: 10, + failure_reason: None, + final_git_commit_sha: None, + stages: Vec::new(), + billing: None, + total_retries: 0, + }; + + write_finalize_commit(&run_options, &services, &conclusion).await; + + assert!(events.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn finalize_emits_metadata_snapshot_before_run_completed() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let run_store = seeded_run_store().await; + let emitter = Arc::new(Emitter::new(test_run_id())); + let events = record_events(&emitter); + let services = test_services( + RunStoreHandle::local(run_store), + Arc::clone(&emitter), + Arc::new(fabro_agent::LocalSandbox::new( + repo_dir.path().to_path_buf(), + )), + Arc::new(SandboxGitRuntime::new()), + ); + let retroed = Retroed { + graph: Graph::new("test"), + outcome: Ok(Outcome::success()), + run_options: test_git_run_options(repo_dir.path(), "fabro/metadata/run"), + duration_ms: 5, + services, + retro: None, + }; + + finalize(retroed, &FinalizeOptions { + run_dir: repo_dir.path().to_path_buf(), + run_id: test_run_id(), + workflow_name: "test".to_string(), + preserve_sandbox: false, + last_git_sha: None, + }) + .await + .unwrap(); + + let names = events + .lock() + .unwrap() + .iter() + .map(|event| event.event_name().to_string()) + .collect::>(); + assert_eq!(names, vec![ + "metadata.snapshot.started", + "metadata.snapshot.completed", + "run.completed", + ]); + } + + struct FailingStateStore; + + #[async_trait] + impl RunStoreBackend for FailingStateStore { + async fn load_state(&self) -> Result { + Err(anyhow::anyhow!("state unavailable")) + } + + async fn list_events(&self) -> Result> { + Ok(Vec::new()) + } + + async fn append_run_event(&self, _event: &RunEvent) -> Result<()> { + Ok(()) + } + + async fn write_blob(&self, data: &[u8]) -> Result { + Ok(RunBlobId::new(data)) + } + + async fn read_blob(&self, _id: &RunBlobId) -> Result> { + Ok(None) + } + } } diff --git a/lib/crates/fabro-workflow/src/sandbox_git.rs b/lib/crates/fabro-workflow/src/sandbox_git.rs index 2df37bedc..09aa4a930 100644 --- a/lib/crates/fabro-workflow/src/sandbox_git.rs +++ b/lib/crates/fabro-workflow/src/sandbox_git.rs @@ -1350,8 +1350,16 @@ mod tests { crate::git::GitAuthor::default(), ); + let expected_entries = dump.git_entries().unwrap(); + let expected_entry_count = expected_entries.len(); + let expected_bytes = expected_entries + .iter() + .map(|(_, bytes)| u64::try_from(bytes.len()).unwrap_or(u64::MAX)) + .sum::(); let snapshot = writer.write_snapshot(&dump, "checkpoint").await.unwrap(); assert_eq!(snapshot.push_error, None); + assert_eq!(snapshot.entry_count, expected_entry_count); + assert_eq!(snapshot.bytes, expected_bytes); let commit_sha = snapshot.commit_sha; let current = std::process::Command::new("git") @@ -1399,7 +1407,15 @@ mod tests { assert!(String::from_utf8(status.stdout).unwrap().trim().is_empty()); dump.add_file_bytes("second.txt", b"second\n".to_vec()); + let second_expected_entries = dump.git_entries().unwrap(); + let second_expected_entry_count = second_expected_entries.len(); + let second_expected_bytes = second_expected_entries + .iter() + .map(|(_, bytes)| u64::try_from(bytes.len()).unwrap_or(u64::MAX)) + .sum::(); let second_snapshot = writer.write_snapshot(&dump, "checkpoint 2").await.unwrap(); + assert_eq!(second_snapshot.entry_count, second_expected_entry_count); + assert_eq!(second_snapshot.bytes, second_expected_bytes); let second_commit_sha = second_snapshot.commit_sha; let second_parent = std::process::Command::new("git") .args(["rev-list", "--parents", "-n", "1", &second_commit_sha]) @@ -1482,6 +1498,12 @@ mod tests { in_place: false, }); let dump = crate::run_dump::RunDump::from_projection(&projection); + let expected_entries = dump.git_entries().unwrap(); + let expected_entry_count = expected_entries.len(); + let expected_bytes = expected_entries + .iter() + .map(|(_, bytes)| u64::try_from(bytes.len()).unwrap_or(u64::MAX)) + .sum::(); let runtime = crate::sandbox_metadata::SandboxGitRuntime::new(); let writer = crate::sandbox_metadata::SandboxMetadataWriter::new( &sandbox, @@ -1492,6 +1514,8 @@ mod tests { ); let snapshot = writer.write_snapshot(&dump, "checkpoint").await.unwrap(); + assert_eq!(snapshot.entry_count, expected_entry_count); + assert_eq!(snapshot.bytes, expected_bytes); let push_error = snapshot.push_error.unwrap(); assert!(push_error.contains("git push origin")); diff --git a/lib/crates/fabro-workflow/src/sandbox_metadata.rs b/lib/crates/fabro-workflow/src/sandbox_metadata.rs index 9b83a9de6..f3b3c032b 100644 --- a/lib/crates/fabro-workflow/src/sandbox_metadata.rs +++ b/lib/crates/fabro-workflow/src/sandbox_metadata.rs @@ -76,8 +76,10 @@ pub(crate) struct SandboxMetadataWriter<'a> { } pub(crate) struct MetadataSnapshot { - pub commit_sha: String, - pub push_error: Option, + pub commit_sha: String, + pub push_error: Option, + pub entry_count: usize, + pub bytes: u64, } impl<'a> SandboxMetadataWriter<'a> { @@ -108,6 +110,8 @@ impl<'a> SandboxMetadataWriter<'a> { .map_err(SandboxMetadataError::GitUnavailable)?; let entries = dump.git_entries()?; + let entry_count = entries.len(); + let bytes = metadata_entries_bytes(&entries); let temp = sandbox_temp_dir(self.sandbox, self.run_id, "metadata"); exec_ok( self.sandbox, @@ -119,7 +123,9 @@ impl<'a> SandboxMetadataWriter<'a> { ) .await?; - let result = self.write_snapshot_in_temp(&entries, message, &temp).await; + let result = self + .write_snapshot_in_temp(&entries, message, &temp, entry_count, bytes) + .await; let _ = exec_ok( self.sandbox, &format!("rm -rf {}", shell_quote(&temp)), @@ -134,6 +140,8 @@ impl<'a> SandboxMetadataWriter<'a> { entries: &[(String, Vec)], message: &str, temp: &str, + entry_count: usize, + bytes: u64, ) -> Result { let full_ref = format!("refs/heads/{}", self.branch); let old_commit = exec_stdout( @@ -187,10 +195,18 @@ impl<'a> SandboxMetadataWriter<'a> { Ok(MetadataSnapshot { commit_sha: commit, push_error, + entry_count, + bytes, }) } } +fn metadata_entries_bytes(entries: &[(String, Vec)]) -> u64 { + entries.iter().fold(0, |total, (_, bytes)| { + total.saturating_add(u64::try_from(bytes.len()).unwrap_or(u64::MAX)) + }) +} + fn fast_import_stream( full_ref: &str, old_commit: Option<&str>,