diff --git a/lib/components/fabro-petri/src/fork.rs b/lib/components/fabro-petri/src/fork.rs index 70f76dd95..764777029 100644 --- a/lib/components/fabro-petri/src/fork.rs +++ b/lib/components/fabro-petri/src/fork.rs @@ -509,16 +509,12 @@ pub async fn stage_labels(views: &DbPool, run_id: RunId) -> Result) { @@ -59,7 +59,7 @@ impl RunView { } => { self.state .finished_firings - .insert(stage_key(execution.raw(), firing.raw())); + .insert(FiringKey::new(execution.raw(), firing.raw())); let is_final = matches!( event.derived, Some(Derived::StepFinished { is_final: true, .. }) @@ -139,12 +139,11 @@ impl RunView { answer: Some(answer), }) = &event.derived { - let firing_key = event.subject.as_ref().and_then(|subject| { - subject - .firing - .map(|firing| stage_key(execution.raw(), firing.raw())) - }); - self.close_questions(answer.question.as_deref(), firing_key.as_deref(), at); + self.close_questions( + answer.question.as_deref(), + FiringKey::of_event(event), + at, + ); } } // ── Sandbox: the instance (VIEWS.md "Sandbox") ────────────────── @@ -271,7 +270,7 @@ impl RunView { results, .. } => { - let key = stage_key(occurrence.execution.raw(), occurrence.firing.raw()); + let key = FiringKey::new(occurrence.execution.raw(), occurrence.firing.raw()); let Some(stage_id) = self .state .stages @@ -317,7 +316,7 @@ impl RunView { let Some(firing) = subject.firing else { return; }; - let key = stage_key(execution.raw(), firing.raw()); + let key = FiringKey::new(execution.raw(), firing.raw()); if self.state.stages.contains_key(&key) { return; } diff --git a/lib/components/fabro-petri/src/projection/mod.rs b/lib/components/fabro-petri/src/projection/mod.rs index 6cc9f2cbf..e7a4b9cd8 100644 --- a/lib/components/fabro-petri/src/projection/mod.rs +++ b/lib/components/fabro-petri/src/projection/mod.rs @@ -33,15 +33,18 @@ mod progress; mod sandbox; use std::collections::{BTreeMap, BTreeSet}; +use std::fmt; +use std::str::FromStr; use chrono::{DateTime, TimeZone as _, Utc}; +use fabro_store::StagePosition; use fabro_store::platform_records::StoredPlatformRecord; use fabro_types::{ RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection, }; use petri_execution::ExecutionId; use petri_execution::events::{NodeRef, RunEvent, Subject}; -use serde::{Deserialize, Serialize}; +use serde::{Deserialize, Deserializer, Serialize, Serializer, de}; use serde_json::Value; use tracing::debug; @@ -91,9 +94,9 @@ pub struct RecordHealth { /// The fold's bookkeeping between items. #[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct FoldState { - /// Stages by `":"`. + /// Stages by firing. #[serde(default)] - pub stages: BTreeMap, + pub stages: BTreeMap, /// Labels taken, so a second firing with the same name and visit gets /// its own. #[serde(default)] @@ -103,9 +106,9 @@ pub struct FoldState { /// Which invocation each execution belongs to. #[serde(default)] pub executions: BTreeMap, - /// Open questions by id: the stage that asked. + /// Open questions by id: the firing that asked. #[serde(default)] - pub questions: BTreeMap, + pub questions: BTreeMap, #[serde(default, skip_serializing_if = "Option::is_none")] pub root: Option, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -126,11 +129,10 @@ pub struct FoldState { pub run_diff: Option, #[serde(default)] pub health: RecordHealth, - /// Firings (`":"`) whose attempt has recorded a - /// finish: what a position-keyed platform record may be streamed - /// behind. + /// Firings whose attempt has recorded a finish: what a position-keyed + /// platform record may be streamed behind. #[serde(default)] - pub finished_firings: BTreeSet, + pub finished_firings: BTreeSet, /// Whether the run's sandbox still exists after its release /// (`scope.released` `retained`): kept stopped, or deleted. Absent until /// the root invocation's lease was released. The view carries the same @@ -201,7 +203,7 @@ impl RunView { let stage = self .state .stages - .get(&stage_key(execution.raw(), firing.raw()))?; + .get(&FiringKey::new(execution.raw(), firing.raw()))?; if !stage.shown { return None; } @@ -240,10 +242,74 @@ fn settle_control(projection: &mut RunProjection, action: RunControlAction) { } } -/// The key of a stage: its execution and firing. -#[must_use] -pub fn stage_key(execution: u64, firing: u64) -> String { - format!("{execution}:{firing}") +/// The key of a stage: the execution and firing of the visit it shows. The +/// same fact a positioned platform record carries as its `StagePosition`. +/// It is written `:`, which is how the stored fold +/// state keys its maps. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct FiringKey { + pub execution: u64, + pub firing: u64, +} + +impl FiringKey { + #[must_use] + pub fn new(execution: u64, firing: u64) -> Self { + Self { execution, firing } + } + + /// The firing an event belongs to: its context's execution and its + /// subject's firing, when it has both. + #[must_use] + pub fn of_event(event: &RunEvent) -> Option { + let execution = event.context.execution?; + let firing = event.subject.as_ref()?.firing?; + Some(Self::new(execution.raw(), firing.raw())) + } +} + +impl From for FiringKey { + fn from(position: StagePosition) -> Self { + Self::new(position.execution, position.firing) + } +} + +impl fmt::Display for FiringKey { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}:{}", self.execution, self.firing) + } +} + +/// A firing key that is not `:`. +#[derive(Debug, thiserror::Error)] +#[error("a firing key is `:`, not {0:?}")] +pub struct ParseFiringKeyError(String); + +impl FromStr for FiringKey { + type Err = ParseFiringKeyError; + + fn from_str(text: &str) -> Result { + let invalid = || ParseFiringKeyError(text.to_string()); + let (execution, firing) = text.split_once(':').ok_or_else(invalid)?; + Ok(Self::new( + execution.parse().map_err(|_| invalid())?, + firing.parse().map_err(|_| invalid())?, + )) + } +} + +impl Serialize for FiringKey { + fn serialize(&self, serializer: S) -> Result { + serializer.collect_str(self) + } +} + +impl<'de> Deserialize<'de> for FiringKey { + fn deserialize>(deserializer: D) -> Result { + String::deserialize(deserializer)? + .parse() + .map_err(de::Error::custom) + } } /// Which firing of its node a subject is, 1-based. @@ -351,4 +417,22 @@ mod tests { assert_eq!(labels, vec!["build@1", "build/e2@1"]); assert_eq!(view.state.stages.len(), 2); } + + #[test] + fn a_firing_key_is_stored_as_execution_colon_firing() { + let mut stages: BTreeMap = BTreeMap::new(); + stages.insert(FiringKey::new(3, 7), 1); + let json = serde_json::to_string(&stages).expect("the map encodes"); + assert_eq!(json, r#"{"3:7":1}"#); + let back: BTreeMap = serde_json::from_str(&json).expect("the map decodes"); + assert_eq!(back, stages); + assert!("3-7".parse::().is_err()); + assert_eq!( + FiringKey::from(StagePosition { + execution: 3, + firing: 7, + }), + FiringKey::new(3, 7) + ); + } } diff --git a/lib/components/fabro-petri/src/projection/platform.rs b/lib/components/fabro-petri/src/projection/platform.rs index ebff0cd24..7109e25fd 100644 --- a/lib/components/fabro-petri/src/projection/platform.rs +++ b/lib/components/fabro-petri/src/projection/platform.rs @@ -17,7 +17,7 @@ use fabro_types::{ use tracing::debug; use super::sandbox::sandbox_plan; -use super::{RunView, apply_status, millis, settle_control, stage_key, touch}; +use super::{FiringKey, RunView, apply_status, millis, settle_control, touch}; impl RunView { pub(super) fn fold_platform(&mut self, stored: &StoredPlatformRecord, stream_seq: u64) { @@ -76,7 +76,7 @@ impl RunView { let stage = self .state .stages - .get(&stage_key(record.execution, record.firing)); + .get(&FiringKey::new(record.execution, record.firing)); let current_node = stage.map_or_else(String::new, |stage| stage.node_name.clone()); let stage_id = stage .filter(|stage| stage.shown) @@ -107,7 +107,7 @@ impl RunView { let stage = self .state .stages - .get(&stage_key(record.execution, record.firing)); + .get(&FiringKey::new(record.execution, record.firing)); let Some(stage_id) = stage.map(|stage| stage.stage_id.clone()) else { debug!( seq = stored.seq, diff --git a/lib/components/fabro-petri/src/projection/progress.rs b/lib/components/fabro-petri/src/projection/progress.rs index 09121fac1..f5a0eac25 100644 --- a/lib/components/fabro-petri/src/projection/progress.rs +++ b/lib/components/fabro-petri/src/projection/progress.rs @@ -20,7 +20,7 @@ use serde_json::Value; use tracing::debug; use super::model::{model_ref, split_model}; -use super::{RunView, apply_status, stage_key}; +use super::{FiringKey, RunView, apply_status}; use crate::interview::question_type; impl RunView { @@ -115,7 +115,7 @@ impl RunView { let group = self .state .stages - .get(&stage_key(execution.raw(), occurrence.firing)) + .get(&FiringKey::new(execution.raw(), occurrence.firing)) .map(|stage| stage.stage_id.clone()); if let Some(group) = group { self.state.invocations.entry(invocation).or_default().branch = @@ -143,7 +143,7 @@ impl RunView { let Some(firing) = subject.firing else { return; }; - let key = stage_key(execution.raw(), firing.raw()); + let key = FiringKey::new(execution.raw(), firing.raw()); let label = self.state.stages.get(&key).map_or_else( || subject.node.name.to_string(), |stage| stage.stage_id.to_string(), @@ -201,16 +201,16 @@ impl RunView { pub(super) fn close_questions( &mut self, question: Option<&str>, - firing_key: Option<&str>, + firing: Option, at: DateTime, ) { - let closed: Vec = match (question, firing_key) { + let closed: Vec = match (question, firing) { (Some(question), _) => vec![question.to_string()], - (None, Some(key)) => self + (None, Some(firing)) => self .state .questions .iter() - .filter(|(_, asked_by)| asked_by.as_str() == key) + .filter(|(_, asked_by)| **asked_by == firing) .map(|(id, _)| id.clone()) .collect(), (None, None) => Vec::new(), diff --git a/lib/components/fabro-petri/src/projector.rs b/lib/components/fabro-petri/src/projector.rs index b9b3e274d..533a590c5 100644 --- a/lib/components/fabro-petri/src/projector.rs +++ b/lib/components/fabro-petri/src/projector.rs @@ -77,7 +77,7 @@ use tracing::{debug, info, warn}; use self::cache::{Caches, IDLE, RunCache}; use crate::SqliteRunStore; -use crate::projection::{self, FoldState, Item, RecordHealth, RunView}; +use crate::projection::{self, FiringKey, FoldState, Item, RecordHealth, RunView}; /// The positions a view committed: the last event consumed per Petri log, /// and the last platform record consumed. @@ -891,23 +891,16 @@ impl petri_execution::RunLogs for SignallingLogs { fn order_items<'a>( events: &'a [RunEvent], platform_records: &'a [StoredPlatformRecord], - finished_before: &BTreeSet, + finished_before: &BTreeSet, run_finished: bool, ) -> (Vec>, usize) { - let firing_of = |event: &RunEvent| -> Option<(u64, u64)> { - let execution = event.context.execution?; - let firing = event.subject.as_ref()?.firing?; - Some((execution.raw(), firing.raw())) - }; - let finished_in_pass = |at: (u64, u64)| { + let finished_in_pass = |at: FiringKey| { events.iter().any(|event| { - firing_of(event) == Some(at) + FiringKey::of_event(event) == Some(at) && matches!(event.engine(), Some(Event::StepFinished { .. })) }) }; - let finished = |at: (u64, u64)| { - finished_in_pass(at) || finished_before.contains(&projection::stage_key(at.0, at.1)) - }; + let finished = |at: FiringKey| finished_in_pass(at) || finished_before.contains(&at); // Platform records are consumed in seq order: the first one whose firing // has not finished holds itself and everything after it. let consumed = if run_finished { @@ -918,7 +911,7 @@ fn order_items<'a>( .position(|record| { record .position - .is_some_and(|position| !finished((position.execution, position.firing))) + .is_some_and(|position| !finished(FiringKey::from(position))) }) .unwrap_or(platform_records.len()) }; @@ -940,7 +933,7 @@ fn order_items<'a>( items.sort_by_key(|(recorded_at, rank, _)| (*recorded_at, *rank)); let item_firing = |item: &Item<'a>| match item { - Item::Petri(event) => firing_of(event), + Item::Petri(event) => FiringKey::of_event(event), Item::Platform(_) => None, }; let is_routing = |item: &Item<'a>| { @@ -959,7 +952,7 @@ fn order_items<'a>( let Some(position) = record.position else { continue; }; - let at = (position.execution, position.firing); + let at = FiringKey::from(position); let first_routing = items .iter() .position(|(_, _, other)| item_firing(other) == Some(at) && is_routing(other)); @@ -967,7 +960,8 @@ fn order_items<'a>( .iter() .rposition(|(_, _, other)| item_firing(other) == Some(at)); let first_later = items.iter().position(|(_, _, other)| { - item_firing(other).is_some_and(|(execution, firing)| execution == at.0 && firing > at.1) + item_firing(other) + .is_some_and(|key| key.execution == at.execution && key.firing > at.firing) }); keys[index] = if let Some(before) = first_routing { (before, 0) @@ -1323,7 +1317,7 @@ mod tests { fn a_positioned_record_precedes_later_firings_and_an_unpositioned_one_keeps_its_clock() { let events = vec![started(13, 2, 103), finished(14, 2, 104)]; let records = vec![checkpoint(1, 1, 250)]; - let finished_before: BTreeSet = [projection::stage_key(0, 1)].into_iter().collect(); + let finished_before: BTreeSet = [FiringKey::new(0, 1)].into_iter().collect(); let (items, held) = order_items(&events, &records, &finished_before, false); assert_eq!(held, 0); assert_eq!(names(&items), vec![