mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
feat(events): add run.archived and run.unarchived event variants
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<RunStatus>` — `RunArchived` captures the current status before switching to Archived; `RunUnarchived` applies the event's `restored_status` payload (authoritative) and clears `prior_status`.
This commit is contained in:
parent
ed5e3f1792
commit
beda6b00d8
4 changed files with 331 additions and 4 deletions
|
|
@ -24,6 +24,8 @@ pub struct RunProjection {
|
|||
pub graph_source: Option<String>,
|
||||
pub start: Option<StartRecord>,
|
||||
pub status: Option<RunStatusRecord>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub prior_status: Option<RunStatus>,
|
||||
pub pending_control: Option<RunControlAction>,
|
||||
pub checkpoint: Option<Checkpoint>,
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 [
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct RunArchivedProps {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub actor: Option<ActorRef>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct RunUnarchivedProps {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub actor: Option<ActorRef>,
|
||||
pub restored_status: RunStatus,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct RunCompletedProps {
|
||||
pub duration_ms: u64,
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
},
|
||||
RunArchived {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
actor: Option<ActorRef>,
|
||||
},
|
||||
RunUnarchived {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
actor: Option<ActorRef>,
|
||||
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 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue