From beda6b00d8a3f949675c7518521be991cfa0b838 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 19 Apr 2026 15:42:39 -0400 Subject: [PATCH] feat(events): add run.archived and run.unarchived event variants MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds `RunArchived` and `RunUnarchived` events end-to-end through the engine. Internal `Event` carries `actor` (and `restored_status` on unarchive); wire `EventBody` serializes as `run.archived`/`run.unarchived` with typed props. Projection gains `prior_status: Option` — `RunArchived` captures the current status before switching to Archived; `RunUnarchived` applies the event's `restored_status` payload (authoritative) and clears `prior_status`. --- lib/crates/fabro-store/src/run_state.rs | 149 ++++++++++++++++++++ lib/crates/fabro-types/src/run_event/mod.rs | 78 ++++++++++ lib/crates/fabro-types/src/run_event/run.rs | 17 ++- lib/crates/fabro-workflow/src/event.rs | 91 +++++++++++- 4 files changed, 331 insertions(+), 4 deletions(-) diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index c31454b45..da05d3ac3 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -24,6 +24,8 @@ pub struct RunProjection { pub graph_source: Option, pub start: Option, pub status: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub prior_status: Option, pub pending_control: Option, pub checkpoint: Option, pub checkpoints: Vec<(u32, Checkpoint)>, @@ -212,6 +214,14 @@ impl RunProjection { EventBody::RunRewound(_) => { self.reset_for_rewind(); } + EventBody::RunArchived(_props) => { + self.prior_status = self.status.as_ref().map(|record| record.status); + self.status = Some(run_status_record(RunStatus::Archived, None, ts)); + } + EventBody::RunUnarchived(props) => { + self.status = Some(run_status_record(props.restored_status, None, ts)); + self.prior_status = None; + } EventBody::CheckpointCompleted(props) => { let checkpoint = checkpoint_from_props(props, ts); if let Some(node_id) = stored.node_id.as_deref() { @@ -1097,4 +1107,143 @@ mod tests { events[1].payload.as_value()["properties"]["definition_blob"] ); } + + #[test] + fn run_archived_captures_prior_status_and_sets_archived() { + use fabro_types::RunStatus; + use fabro_types::run_event::{RunArchivedProps, RunCompletedProps}; + + let mut state = RunProjection::default(); + state + .apply_event(&test_event( + 1, + EventBody::RunCompleted(RunCompletedProps { + duration_ms: 10, + artifact_count: 0, + status: "success".to_string(), + reason: None, + total_usd_micros: None, + final_git_commit_sha: None, + final_patch: None, + billing: None, + }), + None, + )) + .unwrap(); + assert_eq!( + state.status.as_ref().map(|record| record.status), + Some(RunStatus::Succeeded) + ); + assert_eq!(state.prior_status, None); + + state + .apply_event(&test_event( + 2, + EventBody::RunArchived(RunArchivedProps { actor: None }), + None, + )) + .unwrap(); + + assert_eq!( + state.status.as_ref().map(|record| record.status), + Some(RunStatus::Archived) + ); + assert_eq!(state.prior_status, Some(RunStatus::Succeeded)); + } + + #[test] + fn run_unarchived_restores_status_and_clears_prior_status() { + use fabro_types::RunStatus; + use fabro_types::run_event::{RunArchivedProps, RunCompletedProps, RunUnarchivedProps}; + + let mut state = RunProjection::default(); + state + .apply_event(&test_event( + 1, + EventBody::RunCompleted(RunCompletedProps { + duration_ms: 10, + artifact_count: 0, + status: "success".to_string(), + reason: None, + total_usd_micros: None, + final_git_commit_sha: None, + final_patch: None, + billing: None, + }), + None, + )) + .unwrap(); + state + .apply_event(&test_event( + 2, + EventBody::RunArchived(RunArchivedProps { actor: None }), + None, + )) + .unwrap(); + state + .apply_event(&test_event( + 3, + EventBody::RunUnarchived(RunUnarchivedProps { + actor: None, + restored_status: RunStatus::Succeeded, + }), + None, + )) + .unwrap(); + + assert_eq!( + state.status.as_ref().map(|record| record.status), + Some(RunStatus::Succeeded) + ); + assert_eq!(state.prior_status, None); + } + + #[test] + fn run_unarchived_uses_event_payload_even_when_prior_status_differs() { + // The event payload is authoritative: the unarchive apply arm sets status + // from `restored_status`, ignoring whatever `prior_status` was captured. + use fabro_types::RunStatus; + use fabro_types::run_event::{RunArchivedProps, RunCompletedProps, RunUnarchivedProps}; + + let mut state = RunProjection::default(); + state + .apply_event(&test_event( + 1, + EventBody::RunCompleted(RunCompletedProps { + duration_ms: 10, + artifact_count: 0, + status: "success".to_string(), + reason: None, + total_usd_micros: None, + final_git_commit_sha: None, + final_patch: None, + billing: None, + }), + None, + )) + .unwrap(); + state + .apply_event(&test_event( + 2, + EventBody::RunArchived(RunArchivedProps { actor: None }), + None, + )) + .unwrap(); + state + .apply_event(&test_event( + 3, + EventBody::RunUnarchived(RunUnarchivedProps { + actor: None, + restored_status: RunStatus::Failed, + }), + None, + )) + .unwrap(); + + assert_eq!( + state.status.as_ref().map(|record| record.status), + Some(RunStatus::Failed) + ); + assert_eq!(state.prior_status, None); + } } diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index ecc6a37a3..b973e3457 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -114,6 +114,10 @@ pub enum EventBody { RunUnpaused(RunControlEffectProps), #[serde(rename = "run.rewound")] RunRewound(RunRewoundProps), + #[serde(rename = "run.archived")] + RunArchived(RunArchivedProps), + #[serde(rename = "run.unarchived")] + RunUnarchived(RunUnarchivedProps), #[serde(rename = "run.completed")] RunCompleted(RunCompletedProps), #[serde(rename = "run.failed")] @@ -375,6 +379,8 @@ impl EventBody { Self::RunPaused(_) => "run.paused", Self::RunUnpaused(_) => "run.unpaused", Self::RunRewound(_) => "run.rewound", + Self::RunArchived(_) => "run.archived", + Self::RunUnarchived(_) => "run.unarchived", Self::RunCompleted(_) => "run.completed", Self::RunFailed(_) => "run.failed", Self::RunNotice(_) => "run.notice", @@ -504,6 +510,8 @@ fn is_known_event_name(event: &str) -> bool { | "run.unblocked" | "run.removing" | "run.rewound" + | "run.archived" + | "run.unarchived" | "run.completed" | "run.failed" | "run.notice" @@ -1106,6 +1114,76 @@ mod tests { ); } + #[test] + fn run_archived_serializes_with_dotted_event_name_and_actor_property() { + let body = EventBody::RunArchived(RunArchivedProps { + actor: Some(ActorRef::user("alice".to_string())), + }); + let value = serde_json::to_value(&body).unwrap(); + assert_eq!(value["event"], "run.archived"); + assert_eq!(value["properties"]["actor"]["kind"], "user"); + assert_eq!(value["properties"]["actor"]["id"], "alice"); + } + + #[test] + fn run_unarchived_serializes_with_restored_status() { + let body = EventBody::RunUnarchived(RunUnarchivedProps { + actor: None, + restored_status: crate::RunStatus::Failed, + }); + let value = serde_json::to_value(&body).unwrap(); + assert_eq!(value["event"], "run.unarchived"); + assert_eq!(value["properties"]["restored_status"], "failed"); + } + + #[test] + fn run_archived_round_trips_through_from_value() { + let value = json!({ + "id": "evt_archived", + "ts": "2026-04-19T12:00:00.000Z", + "run_id": fixtures::RUN_1, + "event": "run.archived", + "properties": { + "actor": { + "kind": "user", + "id": "alice", + "display": "alice" + } + } + }); + + let parsed = RunEvent::from_value(value.clone()).unwrap(); + assert!(matches!(parsed.body, EventBody::RunArchived(_))); + let serialized = parsed.to_value().unwrap(); + assert_eq!(serialized["event"], "run.archived"); + assert_eq!( + serialized["properties"]["actor"], + value["properties"]["actor"] + ); + } + + #[test] + fn run_unarchived_round_trips_through_from_value() { + let value = json!({ + "id": "evt_unarchived", + "ts": "2026-04-19T12:00:00.000Z", + "run_id": fixtures::RUN_1, + "event": "run.unarchived", + "properties": { + "restored_status": "succeeded" + } + }); + + let parsed = RunEvent::from_value(value.clone()).unwrap(); + match &parsed.body { + EventBody::RunUnarchived(props) => { + assert_eq!(props.restored_status, crate::RunStatus::Succeeded); + assert!(props.actor.is_none()); + } + other => panic!("expected RunUnarchived body, got {other:?}"), + } + } + #[test] fn run_queued_and_unblocked_round_trip_as_typed_events() { for value in [ diff --git a/lib/crates/fabro-types/src/run_event/run.rs b/lib/crates/fabro-types/src/run_event/run.rs index e7948284f..2bb51302b 100644 --- a/lib/crates/fabro-types/src/run_event/run.rs +++ b/lib/crates/fabro-types/src/run_event/run.rs @@ -2,9 +2,9 @@ use std::collections::BTreeMap; use serde::{Deserialize, Serialize}; -use super::{BilledTokenCounts, RunNoticeLevel}; +use super::{ActorRef, BilledTokenCounts, RunNoticeLevel}; use crate::settings::SettingsLayer; -use crate::status::BlockedReason; +use crate::status::{BlockedReason, RunStatus}; use crate::{Graph, RunBlobId, RunControlAction, RunProvenance, StatusReason}; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -93,6 +93,19 @@ pub struct RunRewoundProps { pub run_commit_sha: Option, } +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct RunArchivedProps { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub actor: Option, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct RunUnarchivedProps { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub actor: Option, + pub restored_status: RunStatus, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct RunCompletedProps { pub duration_ms: u64, diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index 7db5a75e5..ebdac9976 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -6,7 +6,8 @@ use std::sync::atomic::{AtomicI64, Ordering}; use ::fabro_types::{ ActorRef, BilledTokenCounts, BlockedReason, ParallelBranchId, RunBlobId, RunControlAction, - RunEvent, RunId, RunProvenance, StageId, StageStatus, StatusReason, run_event as fabro_types, + RunEvent, RunId, RunProvenance, RunStatus, StageId, StageStatus, StatusReason, + run_event as fabro_types, }; use anyhow::{Context, Result}; use chrono::Utc; @@ -118,6 +119,15 @@ pub enum Event { #[serde(default, skip_serializing_if = "Option::is_none")] run_commit_sha: Option, }, + RunArchived { + #[serde(default, skip_serializing_if = "Option::is_none")] + actor: Option, + }, + RunUnarchived { + #[serde(default, skip_serializing_if = "Option::is_none")] + actor: Option, + restored_status: RunStatus, + }, WorkflowRunCompleted { duration_ms: u64, artifact_count: usize, @@ -628,6 +638,15 @@ impl Event { "Run rewound" ); } + Self::RunArchived { actor } => { + info!(?actor, "Run archived"); + } + Self::RunUnarchived { + actor, + restored_status, + } => { + info!(?actor, ?restored_status, "Run unarchived"); + } Self::WorkflowRunCompleted { duration_ms, artifact_count, @@ -1172,6 +1191,8 @@ pub fn event_name(event: &Event) -> &'static str { Event::RunPaused => "run.paused", Event::RunUnpaused => "run.unpaused", Event::RunRewound { .. } => "run.rewound", + Event::RunArchived { .. } => "run.archived", + Event::RunUnarchived { .. } => "run.unarchived", Event::WorkflowRunCompleted { .. } => "run.completed", Event::WorkflowRunFailed { .. } => "run.failed", Event::RunNotice { .. } => "run.notice", @@ -1357,7 +1378,9 @@ fn stored_event_fields_for_variant(event: &Event) -> StoredEventFields { }, Event::RunCancelRequested { actor } | Event::RunPauseRequested { actor } - | Event::RunUnpauseRequested { actor } => StoredEventFields { + | Event::RunUnpauseRequested { actor } + | Event::RunArchived { actor } + | Event::RunUnarchived { actor, .. } => StoredEventFields { actor: actor.clone(), ..StoredEventFields::default() }, @@ -1584,6 +1607,16 @@ fn event_body_from_event(event: &Event) -> EventBody { previous_status: previous_status.clone(), run_commit_sha: run_commit_sha.clone(), }), + Event::RunArchived { actor } => EventBody::RunArchived(fabro_types::RunArchivedProps { + actor: actor.clone(), + }), + Event::RunUnarchived { + actor, + restored_status, + } => EventBody::RunUnarchived(fabro_types::RunUnarchivedProps { + actor: actor.clone(), + restored_status: *restored_status, + }), Event::WorkflowRunCompleted { duration_ms, artifact_count, @@ -3395,6 +3428,60 @@ mod tests { assert!(unpause.actor.is_none()); } + #[test] + fn run_archived_event_name_matches_dot_notation() { + assert_eq!( + event_name(&Event::RunArchived { actor: None }), + "run.archived" + ); + assert_eq!( + event_name(&Event::RunUnarchived { + actor: None, + restored_status: RunStatus::Succeeded, + }), + "run.unarchived" + ); + } + + #[test] + fn run_archived_round_trips_actor_in_envelope() { + let actor = ActorRef { + kind: ActorKind::User, + id: Some("alice".to_string()), + display: Some("alice".to_string()), + }; + + let archived = to_run_event(&fixtures::RUN_1, &Event::RunArchived { + actor: Some(actor.clone()), + }); + assert_eq!(archived.event_name(), "run.archived"); + assert_eq!(archived.actor.as_ref().expect("actor set"), &actor); + assert!(matches!(archived.body, EventBody::RunArchived(_))); + } + + #[test] + fn run_unarchived_round_trips_actor_and_restored_status() { + let actor = ActorRef { + kind: ActorKind::User, + id: Some("bob".to_string()), + display: Some("bob".to_string()), + }; + + let unarchived = to_run_event(&fixtures::RUN_1, &Event::RunUnarchived { + actor: Some(actor.clone()), + restored_status: RunStatus::Failed, + }); + assert_eq!(unarchived.event_name(), "run.unarchived"); + assert_eq!(unarchived.actor.as_ref().expect("actor set"), &actor); + match &unarchived.body { + EventBody::RunUnarchived(props) => { + assert_eq!(props.restored_status, RunStatus::Failed); + assert_eq!(props.actor.as_ref().expect("actor set"), &actor); + } + other => panic!("expected RunUnarchived body, got {other:?}"), + } + } + #[test] fn agent_assistant_message_populates_agent_actor() { let stored = to_run_event(&fixtures::RUN_1, &Event::Agent {