mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Merge remote-tracking branch 'origin/main'
# Conflicts: # lib/crates/fabro-types/src/lib.rs
This commit is contained in:
commit
71ec9fdee6
16 changed files with 2183 additions and 117 deletions
|
|
@ -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
|
||||
|
|
|
|||
196
docs/superpowers/plans/2026-04-29-metadata-snapshot-events.md
Normal file
196
docs/superpowers/plans/2026-04-29-metadata-snapshot-events.md
Normal file
|
|
@ -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<String>, commit_sha: Option<String>, entry_count: Option<usize>, bytes: Option<u64> }`
|
||||
|
||||
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<String>` because push failures can have a local commit while load-state and write failures cannot.
|
||||
- Failed accounting fields are optional: `entry_count: Option<usize>` and `bytes: Option<u64>` 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<String>` 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.
|
||||
|
|
@ -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<u32> {
|
||||
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<String> {
|
||||
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<String> {
|
||||
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<String> {
|
||||
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<String>
|
|||
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<String>
|
|||
}
|
||||
}
|
||||
"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<String>
|
|||
_ => &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<String>
|
|||
)];
|
||||
|
||||
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<String>
|
|||
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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
})
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
))
|
||||
}
|
||||
"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<String>
|
|||
}
|
||||
}
|
||||
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<ProgressEvent> {
|
|||
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<String> {
|
|||
#[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"
|
||||
));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub commit_sha: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub entry_count: Option<usize>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub bytes: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct SandboxInitializingProps {
|
||||
pub provider: String,
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
})
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
|
|
|
|||
6
lib/crates/fabro-util/src/time.rs
Normal file
6
lib/crates/fabro-util/src/time.rs
Normal file
|
|
@ -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)
|
||||
}
|
||||
|
|
@ -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<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
commit_sha: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
entry_count: Option<usize>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
bytes: Option<u64>,
|
||||
},
|
||||
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 {
|
||||
|
|
|
|||
|
|
@ -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<Option<BilledModelUsage>>;
|
||||
type WfNodeResult = NodeResult<Option<BilledModelUsage>>;
|
||||
|
|
@ -79,23 +83,42 @@ impl RunLifecycle<WorkflowGraph> 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<WorkflowGraph> 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<WorkflowGraph> for GitLifecycle {
|
|||
}
|
||||
|
||||
impl GitLifecycle {
|
||||
async fn write_metadata_snapshot(&self, dump: &RunDump, message: &str) -> Option<String> {
|
||||
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<String> {
|
||||
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<String>,
|
||||
commit_sha: Option<String>,
|
||||
entry_count: Option<usize>,
|
||||
bytes: Option<u64>,
|
||||
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<RunOptions> {
|
||||
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<Emitter>) -> Arc<std::sync::Mutex<Vec<RunEvent>>> {
|
||||
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<Emitter>,
|
||||
run_store: RunStoreHandle,
|
||||
run_options: Arc<RunOptions>,
|
||||
metadata_runtime: Arc<SandboxGitRuntime>,
|
||||
) -> 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::<u64>();
|
||||
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::<Vec<_>>();
|
||||
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::<u64>();
|
||||
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::<Vec<_>>();
|
||||
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::<Vec<_>>();
|
||||
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<_>>(),
|
||||
vec![
|
||||
"metadata.snapshot.started",
|
||||
"metadata.snapshot.failed",
|
||||
"run.notice",
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
struct FailingStateStore;
|
||||
|
||||
#[async_trait]
|
||||
impl RunStoreBackend for FailingStateStore {
|
||||
async fn load_state(&self) -> Result<RunProjection> {
|
||||
Err(anyhow::anyhow!("state unavailable"))
|
||||
}
|
||||
|
||||
async fn list_events(&self) -> Result<Vec<EventEnvelope>> {
|
||||
Ok(Vec::new())
|
||||
}
|
||||
|
||||
async fn append_run_event(&self, _event: &RunEvent) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn write_blob(&self, data: &[u8]) -> Result<RunBlobId> {
|
||||
Ok(RunBlobId::new(data))
|
||||
}
|
||||
|
||||
async fn read_blob(&self, _id: &RunBlobId) -> Result<Option<Bytes>> {
|
||||
Ok(None)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
commit_sha: Option<String>,
|
||||
entry_count: Option<usize>,
|
||||
bytes: Option<u64>,
|
||||
) {
|
||||
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<Con
|
|||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::collections::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_graphviz::graph::Graph;
|
||||
use fabro_store::Database;
|
||||
use fabro_types::{RunId, WorkflowSettings, fixtures};
|
||||
use fabro_store::{Database, EventEnvelope, RunDatabase, RunProjection};
|
||||
use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase};
|
||||
use fabro_types::{EventBody, RunBlobId, RunEvent, RunId, WorkflowSettings, fixtures};
|
||||
use object_store::memory::InMemory;
|
||||
|
||||
use super::*;
|
||||
use crate::event::{Emitter, StoreProgressLogger};
|
||||
use crate::event::{Emitter, StoreProgressLogger, append_event};
|
||||
use crate::pipeline::types::Retroed;
|
||||
use crate::run_options::RunOptions;
|
||||
use crate::run_options::{GitCheckpointOptions, RunOptions};
|
||||
use crate::runtime_store::{RunStoreBackend, RunStoreHandle};
|
||||
use crate::sandbox_metadata::SandboxGitRuntime;
|
||||
|
||||
fn test_run_id() -> 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<Database> {
|
||||
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<Emitter>) -> Arc<std::sync::Mutex<Vec<RunEvent>>> {
|
||||
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<Emitter>,
|
||||
sandbox: Arc<dyn fabro_agent::Sandbox>,
|
||||
metadata_runtime: Arc<SandboxGitRuntime>,
|
||||
) -> Arc<RunServices> {
|
||||
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::<Vec<_>>();
|
||||
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::<Vec<_>>();
|
||||
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<RunProjection> {
|
||||
Err(anyhow::anyhow!("state unavailable"))
|
||||
}
|
||||
|
||||
async fn list_events(&self) -> Result<Vec<EventEnvelope>> {
|
||||
Ok(Vec::new())
|
||||
}
|
||||
|
||||
async fn append_run_event(&self, _event: &RunEvent) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn write_blob(&self, data: &[u8]) -> Result<RunBlobId> {
|
||||
Ok(RunBlobId::new(data))
|
||||
}
|
||||
|
||||
async fn read_blob(&self, _id: &RunBlobId) -> Result<Option<Bytes>> {
|
||||
Ok(None)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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::<u64>();
|
||||
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::<u64>();
|
||||
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::<u64>();
|
||||
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"));
|
||||
|
|
|
|||
|
|
@ -76,8 +76,10 @@ pub(crate) struct SandboxMetadataWriter<'a> {
|
|||
}
|
||||
|
||||
pub(crate) struct MetadataSnapshot {
|
||||
pub commit_sha: String,
|
||||
pub push_error: Option<String>,
|
||||
pub commit_sha: String,
|
||||
pub push_error: Option<String>,
|
||||
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<u8>)],
|
||||
message: &str,
|
||||
temp: &str,
|
||||
entry_count: usize,
|
||||
bytes: u64,
|
||||
) -> Result<MetadataSnapshot, SandboxMetadataError> {
|
||||
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<u8>)]) -> 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>,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue