From d986fc97512364aaf655306abd9f91521ce1dffa Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 11 Sep 2026 13:32:54 -0600 Subject: [PATCH] Keep the sandbox driver's events whole as run events A bridge translated the driver's events into thirteen lifecycle variants of fabro's own (start, stop, and delete phases, image pulls, snapshot builds) and dropped everything else the driver reported, pairing an image pull's first progress report with the create's completion to invent a duration. The driver's event is now stored as the run event itself, under a name derived from it: subject, action, and phase (sandbox.stop.completed, sandbox.create.progress for an image pull, snapshot.create.started), or .state and .notice. Every operation the driver performs on the run's sandbox lands on the run, including creates and state observations the bridge skipped. The CLI reads image pulls and snapshot builds from the driver's event for its setup progress and pretty output, the thirteen variants and their props go, and a run stored under the old names still reads as an unknown body. Checkpoint file numbers in a dump shift because the run records more events before each checkpoint. Co-Authored-By: Claude Fable 5.1 --- Cargo.lock | 1 + docs/internal/events-strategy.md | 13 +- docs/internal/events.md | 96 +++------ lib/apps/fabro-cli/Cargo.toml | 1 + lib/apps/fabro-cli/src/commands/run/events.rs | 65 ++++-- .../src/commands/run/run_progress/event.rs | 175 ++++++++++++---- .../src/commands/run/run_progress/mod.rs | 101 +++++---- lib/apps/fabro-cli/tests/it/cmd/dump.rs | 6 +- lib/components/fabro-workflow/src/event.rs | 4 +- .../fabro-workflow/src/event/convert.rs | 136 ++++--------- .../fabro-workflow/src/event/driver_events.rs | 36 ++++ .../fabro-workflow/src/event/events.rs | 155 +++----------- .../fabro-workflow/src/event/names.rs | 25 +-- .../src/event/sandbox_bridge.rs | 192 ------------------ .../fabro-workflow/src/pipeline/initialize.rs | 14 +- lib/foundation/fabro-types/src/lib.rs | 1 + .../fabro-types/src/run_event/infra.rs | 76 ------- .../fabro-types/src/run_event/mod.rs | 165 +++++++++------ 18 files changed, 529 insertions(+), 733 deletions(-) create mode 100644 lib/components/fabro-workflow/src/event/driver_events.rs delete mode 100644 lib/components/fabro-workflow/src/event/sandbox_bridge.rs diff --git a/Cargo.lock b/Cargo.lock index 40252a362..b31e186b8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2480,6 +2480,7 @@ dependencies = [ "reqwest 0.13.4", "ring", "rustls", + "sandbox-driver", "scopeguard", "semver", "serde", diff --git a/docs/internal/events-strategy.md b/docs/internal/events-strategy.md index b3aea81c0..ad7a63823 100644 --- a/docs/internal/events-strategy.md +++ b/docs/internal/events-strategy.md @@ -126,10 +126,15 @@ Never build the same `RunEvent` twice if multiple sinks receive it. ### 1. Add the typed event Add a variant to `Event`, `AgentEvent`, or `SandboxLifecycle` as appropriate. Sandbox -lifecycle facts come from two places: the pipeline emits `Initializing`, `Ready`, and -`InitializeFailed` around bringing the sandbox up, and `SandboxEventBridge` (in the -`fabro-workflow::event` module) translates the sandbox driver's own events — start, -stop, delete, image pulls, snapshot builds — into the rest. Fabro-sandbox emits no +facts come from two places: the pipeline emits `Initializing`, `Ready`, and +`InitializeFailed` around bringing the sandbox up, and the sandbox driver's own events +(operations and their outcome, progress inside a create such as an image pull, snapshot +builds, state observations, notices) are stored whole as `Event::SandboxDriver` by the +`DriverEventRecorder` in the `fabro-workflow::event` module. Their names derive from the +event (`fabro_types::sandbox_driver_event_name`): `..` such as +`sandbox.stop.completed` or `snapshot.create.started`, `.state`, and +`.notice`; their `properties` are the driver's event as the driver serializes +it, so the driver's `Event` is part of fabro's stored format. Fabro-sandbox emits no events of its own. ### 2. Add tracing diff --git a/docs/internal/events.md b/docs/internal/events.md index b2f9d7ad0..c58f041d2 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -1807,82 +1807,52 @@ Emitted after the engine completes sandbox initialization (distinct from `sandbo | `provider` | string | Sandbox provider name | | `error` | string | Error message | -### `sandbox.snapshot.pulling` +### Sandbox driver events -Emitted only when the Docker image cache misses and Fabro starts pulling the image. +Everything the sandbox driver reports about a run's sandbox is stored whole. The +event name derives from the driver's event: `..` for an +operation (`sandbox.start.started`, `sandbox.stop.completed`, `sandbox.delete.failed`, +`sandbox.create.progress` for an image pull inside the create, `snapshot.create.started` +and `snapshot.create.completed` for a snapshot build), `.state` for a state +observation, and `.notice` for a notice. `properties` is the driver's event as +the driver serializes it. ```json { "id": "...", "ts": "...", "run_id": "...", - "event": "sandbox.snapshot.pulling", + "event": "sandbox.stop.completed", "properties": { - "name": "my-image:latest" + "id": {"source_id": "9b2f…", "sequence": 4}, + "occurred_at": "2026-08-31T20:00:00Z", + "provider": "docker", + "subject": {"type": "sandbox", "id": "container-abc123"}, + "operation_id": "58a1…", + "correlation_id": "01JQ…", + "type": "operation_completed", + "action": "stop", + "duration": {"secs": 1, "nanos": 250000000} } } ``` | Property | Type | Description | |----------|------|-------------| -| `name` | string | Image/snapshot name | +| `id` | object | The driver's event id: `source_id` and `sequence` within that source | +| `occurred_at` | string | When the driver observed the event (RFC 3339) | +| `provider` | string | The driver's provider kind (`host`, `docker`, `daytona`, a plugin's kind) | +| `subject` | object | `type` (`sandbox`, `snapshot`, `volume`, `provider`) with the resource's `id` and `name` when known | +| `operation_id` | string | Groups the started, progress, and completed or failed events of one operation | +| `correlation_id` | string | The run id fabro attached | +| `type` | string | `operation_started`, `operation_progress`, `operation_completed`, `operation_failed`, `state_observed`, or `notice` | +| `action` | string | The operation (`create`, `start`, `stop`, `delete`, `snapshot`, …) on operation events | +| `progress` | object | `code` (`image.pull`, `snapshot.build`, …), `message`, and optional `completed`, `total`, `unit` on progress events | +| `duration` | object | `secs` and `nanos` on completed and failed events | +| `error` | object | `kind`, `message`, `retryable`, `causes` on failed events | -### `sandbox.snapshot.creating` - -Emitted only when a Daytona snapshot cache miss or inactive snapshot requires Fabro to create or wait for the snapshot. - -```json -{ - "id": "...", "ts": "...", "run_id": "...", - "event": "sandbox.snapshot.creating", - "properties": { - "name": "my-snapshot" - } -} -``` - -| Property | Type | Description | -|----------|------|-------------| -| `name` | string | Snapshot name | - -### `sandbox.snapshot.ready` - -Emitted when an image or snapshot ensure step succeeds. Cache hits still emit this event with a near-zero `duration_ms`; explicit no-op paths such as Docker `auto_pull = false` and the Daytona default snapshot path do not. - -```json -{ - "id": "...", "ts": "...", "run_id": "...", - "event": "sandbox.snapshot.ready", - "properties": { - "name": "my-snapshot", - "duration_ms": 30000 - } -} -``` - -| Property | Type | Description | -|----------|------|-------------| -| `name` | string | Snapshot name | -| `duration_ms` | number | Ensure duration | - -### `sandbox.snapshot.failed` - -Emitted when an image or snapshot ensure step fails. - -```json -{ - "id": "...", "ts": "...", "run_id": "...", - "event": "sandbox.snapshot.failed", - "properties": { - "name": "my-snapshot", - "error": "disk quota exceeded" - } -} -``` - -| Property | Type | Description | -|----------|------|-------------| -| `name` | string | Snapshot name | -| `error` | string | Error message | -| `causes` | string[] | Optional error cause chain | +Events stored under `sandbox.start.*`, `sandbox.stop.*`, `sandbox.delete.*`, and +`sandbox.snapshot.*` before the driver's events were kept whole carry fabro's earlier +`provider`, `name`, `duration_ms`, and `error` properties instead; readers treat them as +unknown bodies. ### `sandbox.git.started` diff --git a/lib/apps/fabro-cli/Cargo.toml b/lib/apps/fabro-cli/Cargo.toml index b81f7e019..3b997722f 100644 --- a/lib/apps/fabro-cli/Cargo.toml +++ b/lib/apps/fabro-cli/Cargo.toml @@ -25,6 +25,7 @@ fabro-llm = { path = "../../components/fabro-llm" } fabro-oauth = { path = "../../foundation/fabro-oauth" } fabro-github = { path = "../../components/fabro-github" } fabro-agent = { path = "../../components/fabro-agent" } +sandbox-driver.workspace = true fabro-dump = { path = "../../components/fabro-dump" } fabro-hooks = { path = "../../components/fabro-hooks" } fabro-install = { path = "../../components/fabro-install" } diff --git a/lib/apps/fabro-cli/src/commands/run/events.rs b/lib/apps/fabro-cli/src/commands/run/events.rs index 247a22c6f..2149bfdcf 100644 --- a/lib/apps/fabro-cli/src/commands/run/events.rs +++ b/lib/apps/fabro-cli/src/commands/run/events.rs @@ -654,25 +654,35 @@ fn format_event_pretty_value(envelope: &serde_json::Value, styles: &Styles) -> O styles.dim.apply_to(&duration), )) } - "sandbox.snapshot.pulling" => { - let name = prop_str_field(envelope, "name").unwrap_or("?"); + "sandbox.create.progress" => { + let code = envelope + .pointer("/properties/progress/code") + .and_then(serde_json::Value::as_str)?; + if code != "image.pull" { + return None; + } + let message = envelope + .pointer("/properties/progress/message") + .and_then(serde_json::Value::as_str) + .unwrap_or("image"); + let name = message.strip_prefix("pulling image ").unwrap_or(message); Some(format!( "{} Sandbox: pulling {}", styles.dim.apply_to(&ts), name, )) } - "sandbox.snapshot.creating" => { - let name = prop_str_field(envelope, "name").unwrap_or("?"); + "snapshot.create.started" => { + let name = driver_subject_name(envelope); Some(format!( "{} Sandbox: building {}", styles.dim.apply_to(&ts), name, )) } - "sandbox.snapshot.ready" => { - let name = prop_str_field(envelope, "name").unwrap_or("?"); - let duration = format_duration_ms(prop_field(envelope, "duration_ms")); + "snapshot.create.completed" => { + let name = driver_subject_name(envelope); + let duration = format_duration_ms(driver_duration_ms(envelope).as_ref()); Some(format!( "{} Sandbox snapshot: {} {}", styles.dim.apply_to(&ts), @@ -680,9 +690,12 @@ fn format_event_pretty_value(envelope: &serde_json::Value, styles: &Styles) -> O styles.dim.apply_to(&duration), )) } - "sandbox.snapshot.failed" => { - let name = prop_str_field(envelope, "name").unwrap_or("?"); - let error = prop_str_field(envelope, "error").unwrap_or("unknown error"); + "snapshot.create.failed" => { + let name = driver_subject_name(envelope); + let error = envelope + .pointer("/properties/error/message") + .and_then(serde_json::Value::as_str) + .unwrap_or("unknown error"); Some(format!( "{} {} Sandbox snapshot {} failed: {}", styles.dim.apply_to(&ts), @@ -809,6 +822,30 @@ fn str_field<'a>(value: &'a serde_json::Value, key: &str) -> Option<&'a str> { value.get(key)?.as_str() } +/// The name of the resource a sandbox driver event is about, falling back +/// to its id. +fn driver_subject_name(envelope: &serde_json::Value) -> &str { + envelope + .pointer("/properties/subject/name") + .or_else(|| envelope.pointer("/properties/subject/id")) + .and_then(serde_json::Value::as_str) + .unwrap_or("?") +} + +/// A sandbox driver operation's duration, in milliseconds, as the number +/// [`format_duration_ms`] reads. +fn driver_duration_ms(envelope: &serde_json::Value) -> Option { + let duration = envelope.pointer("/properties/duration")?; + let secs = duration.get("secs").and_then(serde_json::Value::as_u64)?; + let nanos = duration + .get("nanos") + .and_then(serde_json::Value::as_u64) + .unwrap_or(0); + Some(serde_json::Value::from( + secs.saturating_mul(1000).saturating_add(nanos / 1_000_000), + )) +} + fn prop_field<'a>(value: &'a serde_json::Value, key: &str) -> Option<&'a serde_json::Value> { value.get("properties")?.get(key) } @@ -1261,7 +1298,7 @@ mod tests { #[test] fn pretty_sandbox_snapshot_pulling() { let styles = no_color_styles(); - let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"sandbox.snapshot.pulling","properties":{"name":"buildpack-deps:noble"}}"#; + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"sandbox.create.progress","properties":{"id":{"source_id":"t","sequence":1},"occurred_at":"2026-01-01T14:25:00Z","provider":"docker","subject":{"type":"sandbox"},"type":"operation_progress","action":"create","progress":{"code":"image.pull","message":"pulling image buildpack-deps:noble"}}}"#; let result = format_event_pretty(line, &styles).unwrap(); assert!(result.contains("Sandbox: pulling"), "got: {result}"); assert!(result.contains("buildpack-deps:noble"), "got: {result}"); @@ -1270,7 +1307,7 @@ mod tests { #[test] fn pretty_sandbox_snapshot_creating() { let styles = no_color_styles(); - let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"sandbox.snapshot.creating","properties":{"name":"fabro-v9-test"}}"#; + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"snapshot.create.started","properties":{"id":{"source_id":"t","sequence":1},"occurred_at":"2026-01-01T14:25:00Z","provider":"daytona","subject":{"type":"snapshot","name":"fabro-v9-test"},"type":"operation_started","action":"create"}}"#; let result = format_event_pretty(line, &styles).unwrap(); assert!(result.contains("Sandbox: building"), "got: {result}"); assert!(result.contains("fabro-v9-test"), "got: {result}"); @@ -1279,7 +1316,7 @@ mod tests { #[test] fn pretty_sandbox_snapshot_ready() { let styles = no_color_styles(); - let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"sandbox.snapshot.ready","properties":{"name":"buildpack-deps:noble","duration_ms":8200}}"#; + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"snapshot.create.completed","properties":{"id":{"source_id":"t","sequence":1},"occurred_at":"2026-01-01T14:25:00Z","provider":"daytona","subject":{"type":"snapshot","name":"buildpack-deps:noble"},"type":"operation_completed","action":"create","duration":{"secs":8,"nanos":200000000}}}"#; let result = format_event_pretty(line, &styles).unwrap(); assert!(result.contains("Sandbox snapshot:"), "got: {result}"); assert!(result.contains("buildpack-deps:noble"), "got: {result}"); @@ -1289,7 +1326,7 @@ mod tests { #[test] fn pretty_sandbox_snapshot_failed() { let styles = no_color_styles(); - let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"sandbox.snapshot.failed","properties":{"name":"buildpack-deps:noble","error":"pull failed"}}"#; + let line = r#"{"ts":"2026-01-01T14:25:00Z","event":"snapshot.create.failed","properties":{"id":{"source_id":"t","sequence":1},"occurred_at":"2026-01-01T14:25:00Z","provider":"docker","subject":{"type":"snapshot","name":"buildpack-deps:noble"},"type":"operation_failed","action":"create","duration":{"secs":1,"nanos":0},"error":{"kind":"provider","message":"pull failed","retryable":false,"causes":[]}}}"#; let result = format_event_pretty(line, &styles).unwrap(); assert!( result.contains("Sandbox snapshot buildpack-deps:noble failed: pull failed"), diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs index 57167154b..92f4e5459 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs @@ -257,20 +257,7 @@ pub(super) fn from_run_event(stored: &RunEvent) -> Option { provider: props.provider.clone(), error: props.error.clone(), }), - EventBody::SnapshotPulling(props) => Some(ProgressEvent::SnapshotPulling { - name: props.name.clone(), - }), - EventBody::SnapshotCreating(props) => Some(ProgressEvent::SnapshotCreating { - name: props.name.clone(), - }), - EventBody::SnapshotReady(props) => Some(ProgressEvent::SnapshotReady { - name: props.name.clone(), - duration_ms: props.duration_ms, - }), - EventBody::SnapshotFailed(props) => Some(ProgressEvent::SnapshotFailed { - name: props.name.clone(), - error: props.error.clone(), - }), + EventBody::SandboxDriver { event, .. } => driver_progress_event(event), EventBody::SshAccessReady(props) => Some(ProgressEvent::SshAccessReady { ssh_command: props.ssh_command.clone(), }), @@ -525,6 +512,65 @@ fn display_value(value: &Value) -> Option { } } +/// The setup progress a sandbox driver event stands for: the image pull +/// inside the sandbox's create, or a snapshot build. Every other driver +/// event is stored on the run but renders nothing here. +fn driver_progress_event(event: &sandbox_driver::Event) -> Option { + use sandbox_driver::{Action, EventBody as Body, EventSubject, ProgressCode}; + + match (&event.subject, &event.body) { + ( + EventSubject::Sandbox { .. }, + Body::OperationProgress { + action: Action::Create, + progress, + }, + ) if progress.code.as_str() == ProgressCode::IMAGE_PULL => { + Some(ProgressEvent::SnapshotPulling { + name: pulled_image_name(progress.message.as_deref()), + }) + } + (EventSubject::Snapshot { id, name }, body) => { + let name = name + .clone() + .or_else(|| id.as_ref().map(ToString::to_string)) + .unwrap_or_default(); + match body { + Body::OperationStarted { + action: Action::Create, + } => Some(ProgressEvent::SnapshotCreating { name }), + Body::OperationCompleted { + action: Action::Create, + duration, + } => Some(ProgressEvent::SnapshotReady { + name, + duration_ms: u64::try_from(duration.as_millis()).unwrap_or(u64::MAX), + }), + Body::OperationFailed { + action: Action::Create, + error, + .. + } => Some(ProgressEvent::SnapshotFailed { + name, + error: error.message.clone(), + }), + _ => None, + } + } + _ => None, + } +} + +/// The image an image pull progress report names. The Docker provider +/// says `pulling image `; the reference alone reads better. +fn pulled_image_name(message: Option<&str>) -> String { + let message = message.unwrap_or("image"); + message + .strip_prefix("pulling image ") + .unwrap_or(message) + .to_owned() +} + #[cfg(test)] mod tests { use fabro_agent::AgentEvent; @@ -780,32 +826,76 @@ mod tests { )); } - #[test] - fn round_trip_snapshot_lifecycle_events() { - let pulling = to_run_event(&fixtures::RUN_1, &Event::Sandbox { - event: SandboxLifecycle::SnapshotPulling { - name: "buildpack-deps:noble".into(), - }, - }); - let creating = to_run_event(&fixtures::RUN_1, &Event::Sandbox { - event: SandboxLifecycle::SnapshotCreating { - name: "fabro-v9".into(), - }, - }); - let ready = to_run_event(&fixtures::RUN_1, &Event::Sandbox { - event: SandboxLifecycle::SnapshotReady { - name: "buildpack-deps:noble".into(), - duration_ms: 1200, - }, - }); - let failed = to_run_event(&fixtures::RUN_1, &Event::Sandbox { - event: SandboxLifecycle::SnapshotFailed { - name: "fabro-v9".into(), - error: "build failed".into(), - causes: Vec::new(), - }, - }); + fn driver_event(value: serde_json::Value) -> Event { + Event::SandboxDriver { + event: serde_json::from_value(value).expect("a driver event"), + } + } + #[test] + fn round_trip_driver_events_that_render_setup_progress() { + let pulling = to_run_event( + &fixtures::RUN_1, + &driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 1}, + "occurred_at": "2026-01-01T00:00:00Z", + "provider": "docker", + "subject": {"type": "sandbox"}, + "type": "operation_progress", + "action": "create", + "progress": {"code": "image.pull", "message": "pulling image buildpack-deps:noble"} + })), + ); + let creating = to_run_event( + &fixtures::RUN_1, + &driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 2}, + "occurred_at": "2026-01-01T00:00:00Z", + "provider": "daytona", + "subject": {"type": "snapshot", "name": "fabro-v9"}, + "type": "operation_started", + "action": "create" + })), + ); + let ready = to_run_event( + &fixtures::RUN_1, + &driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 3}, + "occurred_at": "2026-01-01T00:00:01Z", + "provider": "daytona", + "subject": {"type": "snapshot", "name": "fabro-v9"}, + "type": "operation_completed", + "action": "create", + "duration": {"secs": 1, "nanos": 200_000_000} + })), + ); + let failed = to_run_event( + &fixtures::RUN_1, + &driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 4}, + "occurred_at": "2026-01-01T00:00:02Z", + "provider": "daytona", + "subject": {"type": "snapshot", "name": "fabro-v9"}, + "type": "operation_failed", + "action": "create", + "duration": {"secs": 2, "nanos": 0}, + "error": {"kind": "provider", "message": "build failed", "retryable": false, "causes": []} + })), + ); + let stopped = to_run_event( + &fixtures::RUN_1, + &driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 5}, + "occurred_at": "2026-01-01T00:00:03Z", + "provider": "docker", + "subject": {"type": "sandbox", "id": "c1"}, + "type": "operation_completed", + "action": "stop", + "duration": {"secs": 0, "nanos": 0} + })), + ); + + assert_eq!(pulling.event_name(), "sandbox.create.progress"); assert!(matches!( from_run_event(&pulling).unwrap(), ProgressEvent::SnapshotPulling { name } if name == "buildpack-deps:noble" @@ -817,13 +907,18 @@ mod tests { assert!(matches!( from_run_event(&ready).unwrap(), ProgressEvent::SnapshotReady { name, duration_ms } - if name == "buildpack-deps:noble" && duration_ms == 1200 + if name == "fabro-v9" && duration_ms == 1200 )); assert!(matches!( from_run_event(&failed).unwrap(), ProgressEvent::SnapshotFailed { name, error } if name == "fabro-v9" && error == "build failed" )); + assert_eq!(stopped.event_name(), "sandbox.stop.completed"); + assert!( + from_run_event(&stopped).is_none(), + "a stop is stored on the run but renders no setup progress" + ); } #[test] diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs index 330ca4cbd..d8d87e03a 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs @@ -504,6 +504,55 @@ mod tests { .expect("valid utf-8") } + fn driver_event(value: serde_json::Value) -> Event { + Event::SandboxDriver { + event: serde_json::from_value(value).expect("a driver event"), + } + } + + /// A snapshot build reported by the driver: started, or completed after + /// `secs`. + fn snapshot_build_event(name: &str, kind: &str, secs: Option) -> Event { + let mut value = serde_json::json!({ + "id": {"source_id": "test", "sequence": 1}, + "occurred_at": "2026-01-01T00:00:00Z", + "provider": "daytona", + "subject": {"type": "snapshot", "name": name}, + "type": kind, + "action": "create" + }); + if let Some(secs) = secs { + value["duration"] = serde_json::json!({"secs": secs, "nanos": 0}); + } + driver_event(value) + } + + fn snapshot_build_failed_event(name: &str, error: &str) -> Event { + driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 1}, + "occurred_at": "2026-01-01T00:00:00Z", + "provider": "docker", + "subject": {"type": "snapshot", "name": name}, + "type": "operation_failed", + "action": "create", + "duration": {"secs": 1, "nanos": 0}, + "error": {"kind": "provider", "message": error, "retryable": false, "causes": []} + })) + } + + /// The Docker provider pulling the sandbox's image inside its create. + fn image_pull_event(image: &str) -> Event { + driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 1}, + "occurred_at": "2026-01-01T00:00:00Z", + "provider": "docker", + "subject": {"type": "sandbox"}, + "type": "operation_progress", + "action": "create", + "progress": {"code": "image.pull", "message": format!("pulling image {image}")} + })) + } + fn emit(ui: &mut ProgressUI, event: Event) { let stored = to_run_event(&fixtures::RUN_1, &event); ui.handle_event(&stored); @@ -1078,17 +1127,14 @@ mod tests { provider: "daytona".into(), }, }); - emit(&mut ui, Event::Sandbox { - event: SandboxLifecycle::SnapshotCreating { - name: "fabro-v9-test".into(), - }, - }); - emit(&mut ui, Event::Sandbox { - event: SandboxLifecycle::SnapshotReady { - name: "fabro-v9-test".into(), - duration_ms: 210_000, - }, - }); + emit( + &mut ui, + snapshot_build_event("fabro-v9-test", "operation_started", None), + ); + emit( + &mut ui, + snapshot_build_event("fabro-v9-test", "operation_completed", Some(210)), + ); emit(&mut ui, Event::Sandbox { event: SandboxLifecycle::Ready { provider: "daytona".into(), @@ -1114,17 +1160,7 @@ mod tests { provider: "docker".into(), }, }); - emit(&mut ui, Event::Sandbox { - event: SandboxLifecycle::SnapshotPulling { - name: "buildpack-deps:noble".into(), - }, - }); - emit(&mut ui, Event::Sandbox { - event: SandboxLifecycle::SnapshotReady { - name: "buildpack-deps:noble".into(), - duration_ms: 8_200, - }, - }); + emit(&mut ui, image_pull_event("buildpack-deps:noble")); emit(&mut ui, Event::Sandbox { event: SandboxLifecycle::Ready { provider: "docker".into(), @@ -1170,13 +1206,10 @@ mod tests { provider: "docker".into(), }, }); - emit(&mut ui, Event::Sandbox { - event: SandboxLifecycle::SnapshotFailed { - name: "buildpack-deps:noble".into(), - error: "pull failed".into(), - causes: Vec::new(), - }, - }); + emit( + &mut ui, + snapshot_build_failed_event("buildpack-deps:noble", "pull failed"), + ); emit(&mut ui, Event::Sandbox { event: SandboxLifecycle::InitializeFailed { provider: "docker".into(), @@ -1203,12 +1236,10 @@ mod tests { }); assert!(ui.setup.sandbox_bar.is_some()); - emit(&mut ui, Event::Sandbox { - event: SandboxLifecycle::SnapshotReady { - name: "buildpack-deps:noble".into(), - duration_ms: 10, - }, - }); + emit( + &mut ui, + snapshot_build_event("buildpack-deps:noble", "operation_completed", Some(0)), + ); assert!(ui.setup.sandbox_bar.is_some()); emit(&mut ui, Event::Sandbox { diff --git a/lib/apps/fabro-cli/tests/it/cmd/dump.rs b/lib/apps/fabro-cli/tests/it/cmd/dump.rs index 532257f18..f56f590b2 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/dump.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/dump.rs @@ -261,9 +261,9 @@ fn dump_exports_completed_run_snapshot() { "); assert_snapshot!(dump_file_summary(&output_dir), @" - checkpoints/0014.json - checkpoints/0018.json - checkpoints/0022.json + checkpoints/0017.json + checkpoints/0021.json + checkpoints/0025.json events.jsonl graph.fabro run.json diff --git a/lib/components/fabro-workflow/src/event.rs b/lib/components/fabro-workflow/src/event.rs index e20385e9b..31d8255af 100644 --- a/lib/components/fabro-workflow/src/event.rs +++ b/lib/components/fabro-workflow/src/event.rs @@ -1,9 +1,9 @@ mod convert; +mod driver_events; mod emitter; mod events; mod names; mod redaction; -mod sandbox_bridge; mod sink; mod stored_fields; #[cfg(test)] @@ -12,13 +12,13 @@ mod test_support; pub use fabro_types::{EventBody, RunNoticeCode, RunNoticeLevel}; pub use self::convert::{to_run_event, to_run_event_at}; +pub use self::driver_events::DriverEventRecorder; pub use self::emitter::Emitter; pub use self::events::{Event, SandboxLifecycle}; pub use self::names::event_name; pub use self::redaction::{ build_redacted_event_payload, event_payload_from_redacted_json, redacted_event_json, }; -pub use self::sandbox_bridge::SandboxEventBridge; pub use self::sink::{ RunEventLogger, RunEventPersistenceError, RunEventSink, StoreProgressLogger, append_event, append_event_if, append_event_to_sink, create_run, diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index 14246d6d7..83df942cf 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -997,91 +997,8 @@ fn event_body_from_event(event: &Event) -> EventBody { causes: causes.clone(), duration_ms: *duration_ms, }), - SandboxLifecycle::StartStarted { provider } => { - EventBody::SandboxStartStarted(fabro_types::SandboxStartStartedProps { - provider: provider.clone(), - }) - } - SandboxLifecycle::StartCompleted { - provider, - duration_ms, - } => EventBody::SandboxStartCompleted(fabro_types::SandboxStartCompletedProps { - provider: provider.clone(), - duration_ms: *duration_ms, - }), - SandboxLifecycle::StartFailed { - provider, - error, - causes, - } => EventBody::SandboxStartFailed(fabro_types::SandboxStartFailedProps { - provider: provider.clone(), - error: error.clone(), - causes: causes.clone(), - }), - SandboxLifecycle::StopStarted { provider } => { - EventBody::SandboxStopStarted(fabro_types::SandboxStopStartedProps { - provider: provider.clone(), - }) - } - SandboxLifecycle::StopCompleted { - provider, - duration_ms, - } => EventBody::SandboxStopCompleted(fabro_types::SandboxStopCompletedProps { - provider: provider.clone(), - duration_ms: *duration_ms, - }), - SandboxLifecycle::StopFailed { - provider, - error, - causes, - } => EventBody::SandboxStopFailed(fabro_types::SandboxStopFailedProps { - provider: provider.clone(), - error: error.clone(), - causes: causes.clone(), - }), - SandboxLifecycle::DeleteStarted { provider } => { - EventBody::SandboxDeleteStarted(fabro_types::SandboxDeleteStartedProps { - provider: provider.clone(), - }) - } - SandboxLifecycle::DeleteCompleted { - provider, - duration_ms, - } => EventBody::SandboxDeleteCompleted(fabro_types::SandboxDeleteCompletedProps { - provider: provider.clone(), - duration_ms: *duration_ms, - }), - SandboxLifecycle::DeleteFailed { - provider, - error, - causes, - } => EventBody::SandboxDeleteFailed(fabro_types::SandboxDeleteFailedProps { - provider: provider.clone(), - error: error.clone(), - causes: causes.clone(), - }), - SandboxLifecycle::SnapshotPulling { name } => { - EventBody::SnapshotPulling(fabro_types::SnapshotNameProps { name: name.clone() }) - } - SandboxLifecycle::SnapshotCreating { name } => { - EventBody::SnapshotCreating(fabro_types::SnapshotNameProps { name: name.clone() }) - } - SandboxLifecycle::SnapshotReady { name, duration_ms } => { - EventBody::SnapshotReady(fabro_types::SnapshotCompletedProps { - name: name.clone(), - duration_ms: *duration_ms, - }) - } - SandboxLifecycle::SnapshotFailed { - name, - error, - causes, - } => EventBody::SnapshotFailed(fabro_types::SnapshotFailedProps { - name: name.clone(), - error: error.clone(), - causes: causes.clone(), - }), }, + Event::SandboxDriver { event } => EventBody::sandbox_driver(event.clone()), Event::SandboxInitialized { working_directory, provider, @@ -1674,22 +1591,49 @@ mod tests { } #[test] - fn run_event_sandbox_stop_and_delete_use_distinct_event_names() { - let stopped = to_run_event(&fixtures::RUN_5, &Event::Sandbox { - event: SandboxLifecycle::StopCompleted { - provider: "docker".to_string(), - duration_ms: 10, - }, + fn run_event_driver_events_are_named_from_the_subject_action_and_phase() { + let stopped = to_run_event(&fixtures::RUN_5, &Event::SandboxDriver { + event: driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 1}, + "occurred_at": "2026-05-09T12:00:00Z", + "provider": "docker", + "subject": {"type": "sandbox", "id": "container-1"}, + "type": "operation_completed", + "action": "stop", + "duration": {"secs": 0, "nanos": 10_000_000} + })), }); - let deleted = to_run_event(&fixtures::RUN_5, &Event::Sandbox { - event: SandboxLifecycle::DeleteCompleted { - provider: "docker".to_string(), - duration_ms: 20, - }, + let building = to_run_event(&fixtures::RUN_5, &Event::SandboxDriver { + event: driver_event(serde_json::json!({ + "id": {"source_id": "test", "sequence": 2}, + "occurred_at": "2026-05-09T12:00:01Z", + "provider": "daytona", + "subject": {"type": "snapshot", "name": "sandbox-driver-abc"}, + "type": "operation_started", + "action": "create" + })), }); assert_eq!(stopped.event_name(), "sandbox.stop.completed"); - assert_eq!(deleted.event_name(), "sandbox.delete.completed"); + assert_eq!(building.event_name(), "snapshot.create.started"); + let properties = stopped.properties().unwrap(); + assert_eq!(properties["action"], "stop"); + assert_eq!(properties["subject"]["id"], "container-1"); + assert_eq!(properties["duration"]["nanos"], 10_000_000); + + // The stored form reads back as the driver's event. + let round_trip: RunEvent = serde_json::from_value(serde_json::to_value(&stopped).unwrap()) + .expect("a stored driver event decodes"); + assert!(matches!( + &round_trip.body, + EventBody::SandboxDriver { name, event } + if name == "sandbox.stop.completed" + && matches!(event.body, sandbox_driver::EventBody::OperationCompleted { .. }) + )); + } + + fn driver_event(value: serde_json::Value) -> sandbox_driver::Event { + serde_json::from_value(value).expect("a driver event") } #[test] diff --git a/lib/components/fabro-workflow/src/event/driver_events.rs b/lib/components/fabro-workflow/src/event/driver_events.rs new file mode 100644 index 000000000..a14637685 --- /dev/null +++ b/lib/components/fabro-workflow/src/event/driver_events.rs @@ -0,0 +1,36 @@ +//! The sandbox driver's events for a run's sandbox, kept as run events. +//! +//! A run's sandbox is created or attached with a driver [`EventContext`] +//! whose observer is a [`DriverEventRecorder`]. Everything the driver +//! reports about the sandbox — the operations it performs and their +//! outcome, progress inside a create such as an image pull, snapshot +//! builds, state observations, notices — is stored whole as an +//! [`Event::SandboxDriver`], named from the event (see +//! `fabro_types::sandbox_driver_event_name`). +//! +//! [`EventContext`]: sandbox_driver::EventContext + +use std::sync::Arc; + +use async_trait::async_trait; +use sandbox_driver::{Event as DriverEvent, EventObserver}; + +use super::{Emitter, Event}; + +/// Records every event the driver reports as a run event. +pub struct DriverEventRecorder { + emitter: Arc, +} + +impl DriverEventRecorder { + pub fn new(emitter: Arc) -> Self { + Self { emitter } + } +} + +#[async_trait] +impl EventObserver for DriverEventRecorder { + async fn observe(&self, event: DriverEvent) { + self.emitter.emit(&Event::SandboxDriver { event }); + } +} diff --git a/lib/components/fabro-workflow/src/event/events.rs b/lib/components/fabro-workflow/src/event/events.rs index f8ffebf86..18af4fd24 100644 --- a/lib/components/fabro-workflow/src/event/events.rs +++ b/lib/components/fabro-workflow/src/event/events.rs @@ -517,11 +517,17 @@ pub enum Event { status: String, duration_ms: u64, }, - /// A fact about the run's sandbox: the pipeline bringing it up, or a - /// driver operation on it. + /// A fact about the run's sandbox from the pipeline bringing it up. Sandbox { event: SandboxLifecycle, }, + /// An event the sandbox driver reported about the run's sandbox (an + /// operation and its outcome, progress inside a create, a state + /// observation, a notice), kept whole. Named from the event; see + /// `fabro_types::sandbox_driver_event_name`. + SandboxDriver { + event: sandbox_driver::Event, + }, /// Emitted after the sandbox has been initialized (by engine lifecycle). SandboxInitialized { working_directory: String, @@ -768,7 +774,7 @@ pub enum Event { /// Initializing, ready, and failed are the pipeline's view of bringing the /// sandbox up — create, activate, and prepare the workspace as one step. /// The rest are the sandbox driver's own operations and snapshot work, -/// translated from its events by [`super::SandboxEventBridge`]. +/// the driver's own events are kept whole as [`Event::SandboxDriver`]. #[derive(Debug, Clone, Serialize, Deserialize)] pub enum SandboxLifecycle { Initializing { @@ -789,68 +795,11 @@ pub enum SandboxLifecycle { causes: Vec, duration_ms: u64, }, - StartStarted { - provider: String, - }, - StartCompleted { - provider: String, - duration_ms: u64, - }, - StartFailed { - provider: String, - error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - causes: Vec, - }, - StopStarted { - provider: String, - }, - StopCompleted { - provider: String, - duration_ms: u64, - }, - StopFailed { - provider: String, - error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - causes: Vec, - }, - DeleteStarted { - provider: String, - }, - DeleteCompleted { - provider: String, - duration_ms: u64, - }, - DeleteFailed { - provider: String, - error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - causes: Vec, - }, - /// The provider is pulling the image the sandbox is created from. - SnapshotPulling { - name: String, - }, - /// The provider is building or activating the snapshot. - SnapshotCreating { - name: String, - }, - SnapshotReady { - name: String, - duration_ms: u64, - }, - SnapshotFailed { - name: String, - error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - causes: Vec, - }, } impl SandboxLifecycle { pub fn trace(&self) { - use tracing::{debug, error, info, warn}; + use tracing::{debug, error, info}; match self { Self::Initializing { provider } => { debug!(provider, "Sandbox initializing"); @@ -870,70 +819,6 @@ impl SandboxLifecycle { } => { error!(provider, error, causes = ?causes, duration_ms, "Sandbox init failed"); } - Self::StartStarted { provider } => { - info!(provider, "Sandbox start started"); - } - Self::StartCompleted { - provider, - duration_ms, - } => { - info!(provider, duration_ms, "Sandbox start completed"); - } - Self::StartFailed { - provider, - error, - causes, - } => { - warn!(provider, error, causes = ?causes, "Sandbox start failed"); - } - Self::StopStarted { provider } => { - info!(provider, "Sandbox stop started"); - } - Self::StopCompleted { - provider, - duration_ms, - } => { - info!(provider, duration_ms, "Sandbox stop completed"); - } - Self::StopFailed { - provider, - error, - causes, - } => { - warn!(provider, error, causes = ?causes, "Sandbox stop failed"); - } - Self::DeleteStarted { provider } => { - info!(provider, "Sandbox delete started"); - } - Self::DeleteCompleted { - provider, - duration_ms, - } => { - info!(provider, duration_ms, "Sandbox delete completed"); - } - Self::DeleteFailed { - provider, - error, - causes, - } => { - warn!(provider, error, causes = ?causes, "Sandbox delete failed"); - } - Self::SnapshotPulling { name } => { - debug!(name, "Snapshot pulling"); - } - Self::SnapshotCreating { name } => { - debug!(name, "Snapshot creating"); - } - Self::SnapshotReady { name, duration_ms } => { - info!(name, duration_ms, "Snapshot ready"); - } - Self::SnapshotFailed { - name, - error, - causes, - } => { - error!(name, error, causes = ?causes, "Snapshot failed"); - } } } } @@ -1489,6 +1374,7 @@ impl Event { } Self::Agent { .. } => {} Self::Sandbox { event } => event.trace(), + Self::SandboxDriver { event } => trace_driver_event(event), Self::SandboxInitialized { working_directory, provider, @@ -1760,3 +1646,22 @@ impl Event { } } } + +/// Traces a sandbox driver event under its run event name. +fn trace_driver_event(event: &sandbox_driver::Event) { + use sandbox_driver::EventBody as Body; + use tracing::{debug, info, warn}; + + let name = fabro_types::sandbox_driver_event_name(event); + match &event.body { + Body::OperationFailed { error, .. } => { + warn!(event = %name, error = %error.message, "Sandbox driver operation failed"); + } + Body::OperationCompleted { duration, .. } => { + let duration_ms = u64::try_from(duration.as_millis()).unwrap_or(u64::MAX); + info!(event = %name, duration_ms, "Sandbox driver operation completed"); + } + Body::OperationStarted { .. } => info!(event = %name, "Sandbox driver operation started"), + _ => debug!(event = %name, "Sandbox driver event"), + } +} diff --git a/lib/components/fabro-workflow/src/event/names.rs b/lib/components/fabro-workflow/src/event/names.rs index a63d4f647..5a9cdcf5e 100644 --- a/lib/components/fabro-workflow/src/event/names.rs +++ b/lib/components/fabro-workflow/src/event/names.rs @@ -1,10 +1,15 @@ +use std::borrow::Cow; + use fabro_agent::AgentEvent; use super::{Event, SandboxLifecycle}; #[must_use] -pub fn event_name(event: &Event) -> &'static str { - match event { +pub fn event_name(event: &Event) -> Cow<'static, str> { + let name: &'static str = match event { + Event::SandboxDriver { event } => { + return Cow::Owned(fabro_types::sandbox_driver_event_name(event)); + } Event::RunCreated { .. } => "run.created", Event::WorkflowRunStarted { .. } => "run.started", Event::RunSubmitted { .. } => "run.submitted", @@ -105,19 +110,6 @@ pub fn event_name(event: &Event) -> &'static str { SandboxLifecycle::Initializing { .. } => "sandbox.initializing", SandboxLifecycle::Ready { .. } => "sandbox.ready", SandboxLifecycle::InitializeFailed { .. } => "sandbox.failed", - SandboxLifecycle::StartStarted { .. } => "sandbox.start.started", - SandboxLifecycle::StartCompleted { .. } => "sandbox.start.completed", - SandboxLifecycle::StartFailed { .. } => "sandbox.start.failed", - SandboxLifecycle::StopStarted { .. } => "sandbox.stop.started", - SandboxLifecycle::StopCompleted { .. } => "sandbox.stop.completed", - SandboxLifecycle::StopFailed { .. } => "sandbox.stop.failed", - SandboxLifecycle::DeleteStarted { .. } => "sandbox.delete.started", - SandboxLifecycle::DeleteCompleted { .. } => "sandbox.delete.completed", - SandboxLifecycle::DeleteFailed { .. } => "sandbox.delete.failed", - SandboxLifecycle::SnapshotPulling { .. } => "sandbox.snapshot.pulling", - SandboxLifecycle::SnapshotCreating { .. } => "sandbox.snapshot.creating", - SandboxLifecycle::SnapshotReady { .. } => "sandbox.snapshot.ready", - SandboxLifecycle::SnapshotFailed { .. } => "sandbox.snapshot.failed", }, Event::SandboxInitialized { .. } => "sandbox.initialized", Event::SetupStarted { .. } => "setup.started", @@ -150,7 +142,8 @@ pub fn event_name(event: &Event) -> &'static str { Event::PullRequestLinked { .. } => "pull_request.linked", Event::PullRequestUnlinked { .. } => "pull_request.unlinked", Event::PullRequestFailed { .. } => "pull_request.failed", - } + }; + Cow::Borrowed(name) } #[cfg(test)] diff --git a/lib/components/fabro-workflow/src/event/sandbox_bridge.rs b/lib/components/fabro-workflow/src/event/sandbox_bridge.rs deleted file mode 100644 index d00446a57..000000000 --- a/lib/components/fabro-workflow/src/event/sandbox_bridge.rs +++ /dev/null @@ -1,192 +0,0 @@ -//! The sandbox driver's events for a run's sandbox, as workflow events. -//! -//! A run's sandbox is created or attached with a driver [`EventContext`] -//! whose observer is a [`SandboxEventBridge`]. The driver reports every -//! operation it performs — start, stop, delete, the image pull inside a -//! create, snapshot builds — and the bridge turns the ones fabro records -//! on a run into [`SandboxLifecycle`] events. Everything else the driver -//! reports (state observations, notices, other operations) is not a run -//! event and is dropped here. - -use std::collections::HashMap; -use std::sync::{Arc, Mutex, PoisonError}; -use std::time::{Duration, Instant}; - -use async_trait::async_trait; -use sandbox_driver::{ - Action, ErrorReport, Event as DriverEvent, EventBody as DriverEventBody, EventObserver, - EventSubject, OperationId, ProgressCode, -}; - -use super::{Emitter, Event, SandboxLifecycle}; - -/// Emits the workflow's sandbox lifecycle events from the driver's. -pub struct SandboxEventBridge { - emitter: Arc, - /// Fabro's name for the provider, which is what the run records; the - /// driver's own kind name can differ (`host` for a `local` run). - provider: String, - /// The image the sandbox is created from, named on pull events. - image: Option, - /// Creates that pulled an image, by operation, with when the pull began. - pulls: Mutex>, -} - -impl SandboxEventBridge { - pub fn new(emitter: Arc, provider: impl Into, image: Option) -> Self { - Self { - emitter, - provider: provider.into(), - image, - pulls: Mutex::new(HashMap::new()), - } - } - - /// The lifecycle event a driver event stands for, if fabro records one. - fn translate(&self, event: &DriverEvent) -> Option { - match &event.subject { - EventSubject::Sandbox { .. } => self.translate_sandbox(event), - EventSubject::Snapshot { id, name } => { - let name = name - .clone() - .or_else(|| id.as_ref().map(ToString::to_string)) - .unwrap_or_default(); - match &event.body { - DriverEventBody::OperationStarted { .. } => { - Some(SandboxLifecycle::SnapshotCreating { name }) - } - DriverEventBody::OperationCompleted { duration, .. } => { - Some(SandboxLifecycle::SnapshotReady { - name, - duration_ms: duration_ms(*duration), - }) - } - DriverEventBody::OperationFailed { error, .. } => { - Some(SandboxLifecycle::SnapshotFailed { - name, - error: error.message.clone(), - causes: error.causes.clone(), - }) - } - _ => None, - } - } - _ => None, - } - } - - fn translate_sandbox(&self, event: &DriverEvent) -> Option { - let provider = self.provider.clone(); - match &event.body { - DriverEventBody::OperationStarted { action } => match action { - Action::Start => Some(SandboxLifecycle::StartStarted { provider }), - Action::Stop => Some(SandboxLifecycle::StopStarted { provider }), - Action::Delete => Some(SandboxLifecycle::DeleteStarted { provider }), - _ => None, - }, - DriverEventBody::OperationProgress { action, progress } => { - if *action != Action::Create || progress.code.as_str() != ProgressCode::IMAGE_PULL { - return None; - } - // The first pull report of a create opens the pull; later - // ones are the same pull's progress. - let operation_id = event.operation_id.clone()?; - let mut pulls = self.pulls.lock().unwrap_or_else(PoisonError::into_inner); - if pulls.contains_key(&operation_id) { - return None; - } - pulls.insert(operation_id, Instant::now()); - Some(SandboxLifecycle::SnapshotPulling { - name: self - .image - .clone() - .or_else(|| progress.message.clone()) - .unwrap_or_default(), - }) - } - DriverEventBody::OperationCompleted { action, duration } => match action { - Action::Create => { - let pulled = self.take_pull(event.operation_id.as_ref())?; - Some(SandboxLifecycle::SnapshotReady { - name: self.image.clone().unwrap_or_default(), - duration_ms: duration_ms(pulled.elapsed()), - }) - } - Action::Start => Some(SandboxLifecycle::StartCompleted { - provider, - duration_ms: duration_ms(*duration), - }), - Action::Stop => Some(SandboxLifecycle::StopCompleted { - provider, - duration_ms: duration_ms(*duration), - }), - Action::Delete => Some(SandboxLifecycle::DeleteCompleted { - provider, - duration_ms: duration_ms(*duration), - }), - _ => None, - }, - DriverEventBody::OperationFailed { action, error, .. } => match action { - Action::Create => { - self.take_pull(event.operation_id.as_ref())?; - Some(SandboxLifecycle::SnapshotFailed { - name: self.image.clone().unwrap_or_default(), - error: error.message.clone(), - causes: error.causes.clone(), - }) - } - Action::Start => Some(failed(error, |error, causes| { - SandboxLifecycle::StartFailed { - provider, - error, - causes, - } - })), - Action::Stop => Some(failed(error, |error, causes| { - SandboxLifecycle::StopFailed { - provider, - error, - causes, - } - })), - Action::Delete => Some(failed(error, |error, causes| { - SandboxLifecycle::DeleteFailed { - provider, - error, - causes, - } - })), - _ => None, - }, - _ => None, - } - } - - /// When the create `operation_id` began pulling its image, if it did. - fn take_pull(&self, operation_id: Option<&OperationId>) -> Option { - self.pulls - .lock() - .unwrap_or_else(PoisonError::into_inner) - .remove(operation_id?) - } -} - -fn failed( - error: &ErrorReport, - build: impl FnOnce(String, Vec) -> SandboxLifecycle, -) -> SandboxLifecycle { - build(error.message.clone(), error.causes.clone()) -} - -fn duration_ms(duration: Duration) -> u64 { - u64::try_from(duration.as_millis()).unwrap_or(u64::MAX) -} - -#[async_trait] -impl EventObserver for SandboxEventBridge { - async fn observe(&self, event: DriverEvent) { - if let Some(lifecycle) = self.translate(&event) { - self.emitter.emit(&Event::Sandbox { event: lifecycle }); - } - } -} diff --git a/lib/components/fabro-workflow/src/pipeline/initialize.rs b/lib/components/fabro-workflow/src/pipeline/initialize.rs index 4688f4721..1b3d849ab 100644 --- a/lib/components/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/components/fabro-workflow/src/pipeline/initialize.rs @@ -24,7 +24,7 @@ use tokio::sync::RwLock as AsyncRwLock; use super::types::{InitOptions, Initialized, LlmSpec, Persisted, SandboxEnvSpec}; use crate::error::Error; -use crate::event::{Event, RunNoticeCode, RunNoticeLevel, SandboxEventBridge, SandboxLifecycle}; +use crate::event::{DriverEventRecorder, Event, RunNoticeCode, RunNoticeLevel, SandboxLifecycle}; use crate::git::GitAuthor; use crate::git_bridge; use crate::handler::llm::{AgentAcpBackend, AgentApiBackend, BackendRouter, routing}; @@ -385,14 +385,12 @@ pub async fn initialize( ); } - // The driver reports what it does to the run's sandbox; the bridge - // records the operations fabro keeps as run events. + // The driver reports what it does to the run's sandbox; every event is + // kept as a run event. let provider_name = options.sandbox.provider_name(); - let sandbox_events = EventContext::new(Arc::new(SandboxEventBridge::new( - Arc::clone(&options.emitter), - provider_name.clone(), - options.sandbox.image(), - ))) + let sandbox_events = EventContext::new(Arc::new(DriverEventRecorder::new(Arc::clone( + &options.emitter, + )))) .correlation_id(CorrelationId::new(options.run_options.run_id.to_string())); let attach_instance = if is_resume { let record = options diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index d6fb7b9b6..4b3c173bb 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -130,6 +130,7 @@ pub use run_event::{ LlmRetryPhase, MetadataSnapshotFailureKind, MetadataSnapshotPhase, RunEvent, RunNoticeCode, RunNoticeLevel, RunPairEndedReason, RunPairFailedReason, RunRunnableSource, SessionCapability, TodoCreatedProps, TodoDeletedProps, TodoUpdatedProps, initial_subagent_generation, + sandbox_driver_event_name, }; pub use run_failure::RunFailure; pub use run_id::{RunId, fixtures}; diff --git a/lib/foundation/fabro-types/src/run_event/infra.rs b/lib/foundation/fabro-types/src/run_event/infra.rs index 203375b12..9e93b833b 100644 --- a/lib/foundation/fabro-types/src/run_event/infra.rs +++ b/lib/foundation/fabro-types/src/run_event/infra.rs @@ -200,82 +200,6 @@ pub struct SandboxReadyProps { pub type SandboxFailedProps = RunSandboxFailure; -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxStartStartedProps { - pub provider: String, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxStartCompletedProps { - pub provider: String, - pub duration_ms: u64, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxStartFailedProps { - pub provider: String, - pub error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub causes: Vec, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxStopStartedProps { - pub provider: String, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxStopCompletedProps { - pub provider: String, - pub duration_ms: u64, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxStopFailedProps { - pub provider: String, - pub error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub causes: Vec, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxDeleteStartedProps { - pub provider: String, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxDeleteCompletedProps { - pub provider: String, - pub duration_ms: u64, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SandboxDeleteFailedProps { - pub provider: String, - pub error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub causes: Vec, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SnapshotNameProps { - pub name: String, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SnapshotCompletedProps { - pub name: String, - pub duration_ms: u64, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct SnapshotFailedProps { - pub name: String, - pub error: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub causes: Vec, -} - #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct SandboxInitializedProps { pub working_directory: String, diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index a118c223f..0eeafc05f 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -283,32 +283,17 @@ pub enum EventBody { SandboxReady(SandboxReadyProps), #[serde(rename = "sandbox.failed")] SandboxFailed(SandboxFailedProps), - #[serde(rename = "sandbox.start.started")] - SandboxStartStarted(SandboxStartStartedProps), - #[serde(rename = "sandbox.start.completed")] - SandboxStartCompleted(SandboxStartCompletedProps), - #[serde(rename = "sandbox.start.failed")] - SandboxStartFailed(SandboxStartFailedProps), - #[serde(rename = "sandbox.stop.started")] - SandboxStopStarted(SandboxStopStartedProps), - #[serde(rename = "sandbox.stop.completed")] - SandboxStopCompleted(SandboxStopCompletedProps), - #[serde(rename = "sandbox.stop.failed")] - SandboxStopFailed(SandboxStopFailedProps), - #[serde(rename = "sandbox.delete.started")] - SandboxDeleteStarted(SandboxDeleteStartedProps), - #[serde(rename = "sandbox.delete.completed")] - SandboxDeleteCompleted(SandboxDeleteCompletedProps), - #[serde(rename = "sandbox.delete.failed")] - SandboxDeleteFailed(SandboxDeleteFailedProps), - #[serde(rename = "sandbox.snapshot.pulling")] - SnapshotPulling(SnapshotNameProps), - #[serde(rename = "sandbox.snapshot.creating")] - SnapshotCreating(SnapshotNameProps), - #[serde(rename = "sandbox.snapshot.ready")] - SnapshotReady(SnapshotCompletedProps), - #[serde(rename = "sandbox.snapshot.failed")] - SnapshotFailed(SnapshotFailedProps), + /// An event the sandbox driver reported about the run's sandbox, a + /// snapshot, a volume, or the provider, stored as the driver's own event + /// under a name derived from it (`sandbox.stop.completed`, + /// `snapshot.create.started`, `sandbox.state`); see + /// [`sandbox_driver_event_name`]. The derive never sees this variant: + /// the run event writes the name and the driver's event itself. + #[serde(skip)] + SandboxDriver { + name: String, + event: sandbox_driver::Event, + }, #[serde(rename = "sandbox.initialized")] SandboxInitialized(SandboxInitializedProps), #[serde(rename = "setup.started")] @@ -412,6 +397,88 @@ struct RunEventParts<'a> { properties: &'a Value, } +impl EventBody { + /// The sandbox driver's event as a run event body, named by + /// [`sandbox_driver_event_name`]. + #[must_use] + pub fn sandbox_driver(event: sandbox_driver::Event) -> Self { + Self::SandboxDriver { + name: sandbox_driver_event_name(&event), + event, + } + } + + /// A stored driver event: `name` has the shape the driver's events are + /// stored under and `properties` decode to a driver event that yields + /// that name. Anything else, including an event stored under one of + /// these names before the driver's events were kept whole, is left to + /// the other variants. + fn sandbox_driver_from_stored(name: &str, properties: &Value) -> Option { + if !is_sandbox_driver_event_name(name) { + return None; + } + let event: sandbox_driver::Event = serde_json::from_value(properties.clone()).ok()?; + (sandbox_driver_event_name(&event) == name).then(|| Self::SandboxDriver { + name: name.to_owned(), + event, + }) + } +} + +/// Whether `name` has the shape the sandbox driver's events are stored +/// under: `..`, `.state`, `.notice`, +/// or `.event`, for the subjects the driver reports on. +fn is_sandbox_driver_event_name(name: &str) -> bool { + let Some((subject, rest)) = name.split_once('.') else { + return false; + }; + matches!(subject, "sandbox" | "snapshot" | "volume" | "provider") + && (matches!(rest, "state" | "notice" | "event") + || rest.split_once('.').is_some_and(|(_, phase)| { + matches!(phase, "started" | "progress" | "completed" | "failed") + })) +} + +/// The run event name for a sandbox driver event: the subject kind, the +/// action, and the phase, so a stop on the sandbox is `sandbox.stop.started`, +/// `sandbox.stop.completed`, or `sandbox.stop.failed`, an image pull inside +/// a create is `sandbox.create.progress`, and a snapshot build is +/// `snapshot.create.*`. A state observation is `.state`, a notice +/// `.notice`, and an event kind this build does not know +/// `.event`. +#[must_use] +pub fn sandbox_driver_event_name(event: &sandbox_driver::Event) -> String { + use sandbox_driver::{EventBody as Body, EventSubject}; + + let subject = match &event.subject { + EventSubject::Snapshot { .. } => "snapshot", + EventSubject::Volume { .. } => "volume", + EventSubject::Provider => "provider", + _ => "sandbox", + }; + let (action, phase) = match &event.body { + Body::OperationStarted { action } => (Some(*action), "started"), + Body::OperationProgress { action, .. } => (Some(*action), "progress"), + Body::OperationCompleted { action, .. } => (Some(*action), "completed"), + Body::OperationFailed { action, .. } => (Some(*action), "failed"), + Body::StateObserved { .. } => (None, "state"), + Body::Notice { .. } => (None, "notice"), + _ => (None, "event"), + }; + match action { + Some(action) => format!("{subject}.{}.{phase}", driver_action_name(action)), + None => format!("{subject}.{phase}"), + } +} + +/// The driver action's wire name (`stop`, `refresh_activity`). +fn driver_action_name(action: sandbox_driver::Action) -> String { + match serde_json::to_value(action) { + Ok(Value::String(name)) => name, + _ => "unknown".to_owned(), + } +} + impl EventBody { pub fn event_name(&self) -> &str { match self { @@ -526,19 +593,6 @@ impl EventBody { Self::SandboxInitializing(_) => "sandbox.initializing", Self::SandboxReady(_) => "sandbox.ready", Self::SandboxFailed(_) => "sandbox.failed", - Self::SandboxStartStarted(_) => "sandbox.start.started", - Self::SandboxStartCompleted(_) => "sandbox.start.completed", - Self::SandboxStartFailed(_) => "sandbox.start.failed", - Self::SandboxStopStarted(_) => "sandbox.stop.started", - Self::SandboxStopCompleted(_) => "sandbox.stop.completed", - Self::SandboxStopFailed(_) => "sandbox.stop.failed", - Self::SandboxDeleteStarted(_) => "sandbox.delete.started", - Self::SandboxDeleteCompleted(_) => "sandbox.delete.completed", - Self::SandboxDeleteFailed(_) => "sandbox.delete.failed", - Self::SnapshotPulling(_) => "sandbox.snapshot.pulling", - Self::SnapshotCreating(_) => "sandbox.snapshot.creating", - Self::SnapshotReady(_) => "sandbox.snapshot.ready", - Self::SnapshotFailed(_) => "sandbox.snapshot.failed", Self::SandboxInitialized(_) => "sandbox.initialized", Self::SetupStarted(_) => "setup.started", Self::SetupCommandStarted(_) => "setup.command.started", @@ -563,7 +617,7 @@ impl EventBody { Self::PullRequestLinked(_) => "pull_request.linked", Self::PullRequestUnlinked(_) => "pull_request.unlinked", Self::PullRequestFailed(_) => "pull_request.failed", - Self::Unknown { name, .. } => name.as_str(), + Self::SandboxDriver { name, .. } | Self::Unknown { name, .. } => name.as_str(), } } @@ -575,6 +629,9 @@ impl EventBody { if let Self::Unknown { properties, .. } = self { return Ok(properties.clone()); } + if let Self::SandboxDriver { event, .. } = self { + return serde_json::to_value(event); + } match serde_json::to_value(self)? { Value::Object(mut map) => { @@ -697,19 +754,6 @@ fn is_known_event_name(event: &str) -> bool { | "sandbox.cleanup.started" | "sandbox.cleanup.completed" | "sandbox.cleanup.failed" - | "sandbox.start.started" - | "sandbox.start.completed" - | "sandbox.start.failed" - | "sandbox.stop.started" - | "sandbox.stop.completed" - | "sandbox.stop.failed" - | "sandbox.delete.started" - | "sandbox.delete.completed" - | "sandbox.delete.failed" - | "sandbox.snapshot.pulling" - | "sandbox.snapshot.creating" - | "sandbox.snapshot.ready" - | "sandbox.snapshot.failed" | "sandbox.git.started" | "sandbox.git.completed" | "sandbox.git.failed" @@ -818,12 +862,15 @@ impl RunEvent { "event": parts.event, "properties": parts.properties, }); - let body: EventBody = match serde_json::from_value(body_payload) { - Ok(body) => body, - Err(err) if is_known_event_name(parts.event) => return Err(err), - Err(_) => EventBody::Unknown { - name: parts.event.to_string(), - properties: parts.properties.clone(), + let body = match EventBody::sandbox_driver_from_stored(parts.event, parts.properties) { + Some(body) => body, + None => match serde_json::from_value(body_payload) { + Ok(body) => body, + Err(err) if is_known_event_name(parts.event) => return Err(err), + Err(_) => EventBody::Unknown { + name: parts.event.to_string(), + properties: parts.properties.clone(), + }, }, }; Ok(Self {