Key the fold's stages by a typed FiringKey instead of a formatted string

stage_key formatted "<execution>:<firing>" and five places re-derived or
re-parsed that string: the engine, progress and platform folds, the
projector's ordering rule, and fork.rs's stage_labels, which split it
back apart. FiringKey is the fact itself, with of_event for the event
side and From<StagePosition> for the platform-record side. It still
serializes as "<execution>:<firing>", so the fold_json a stored view
holds keeps its shape and needs no migration; a unit test pins that.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-20 15:05:32 -04:00
parent fd0ae7f219
commit 1c30332636
No known key found for this signature in database
7 changed files with 136 additions and 63 deletions

View file

@ -509,16 +509,12 @@ pub async fn stage_labels(views: &DbPool, run_id: RunId) -> Result<StageLabels,
Ok(state
.stages
.iter()
.filter_map(|(key, stage)| {
let (execution, firing) = key.split_once(':')?;
Some((
(execution.parse().ok()?, firing.parse().ok()?),
StageLabel {
stage_id: stage.shown.then(|| stage.stage_id.to_string()),
node_name: stage.node_name.clone(),
visit: stage.visit,
},
))
.map(|(key, stage)| {
((key.execution, key.firing), StageLabel {
stage_id: stage.shown.then(|| stage.stage_id.to_string()),
node_name: stage.node_name.clone(),
visit: stage.visit,
})
})
.collect())
}

View file

@ -10,7 +10,7 @@ use fabro_types::{
use petri_execution::CoordinatorEvent;
use petri_execution::events::RunEvent;
use super::{InvocationRef, RunView, apply_status, settle_control, stage_key};
use super::{FiringKey, InvocationRef, RunView, apply_status, settle_control};
impl RunView {
pub(super) fn fold_coordinator(
@ -58,7 +58,7 @@ impl RunView {
let group = self
.state
.stages
.get(&stage_key(parent.execution.raw(), fork_firing))
.get(&FiringKey::new(parent.execution.raw(), fork_firing))
.map(|stage| stage.stage_id.clone());
if let Some(group) = group {
info.branch = Some((group, index));

View file

@ -19,7 +19,7 @@ use tracing::debug;
use super::model::{model_ref, usage_of};
use super::sandbox::{provider_kind, sandbox_instance, sandbox_plan_of};
use super::{RunView, StageRef, is_shown, node_meta_kind, stage_key, stage_label, visit_of};
use super::{FiringKey, RunView, StageRef, is_shown, node_meta_kind, stage_label, visit_of};
impl RunView {
pub(super) fn fold_engine(&mut self, engine: &Event, event: &RunEvent, at: DateTime<Utc>) {
@ -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;
}

View file

@ -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 `"<execution>:<firing>"`.
/// Stages by firing.
#[serde(default)]
pub stages: BTreeMap<String, StageRef>,
pub stages: BTreeMap<FiringKey, StageRef>,
/// 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<u64, u64>,
/// Open questions by id: the stage that asked.
/// Open questions by id: the firing that asked.
#[serde(default)]
pub questions: BTreeMap<String, String>,
pub questions: BTreeMap<String, FiringKey>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub root: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
@ -126,11 +129,10 @@ pub struct FoldState {
pub run_diff: Option<RunDiff>,
#[serde(default)]
pub health: RecordHealth,
/// Firings (`"<execution>:<firing>"`) 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<String>,
pub finished_firings: BTreeSet<FiringKey>,
/// 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 `<execution>:<firing>`, 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<Self> {
let execution = event.context.execution?;
let firing = event.subject.as_ref()?.firing?;
Some(Self::new(execution.raw(), firing.raw()))
}
}
impl From<StagePosition> 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 `<execution>:<firing>`.
#[derive(Debug, thiserror::Error)]
#[error("a firing key is `<execution>:<firing>`, not {0:?}")]
pub struct ParseFiringKeyError(String);
impl FromStr for FiringKey {
type Err = ParseFiringKeyError;
fn from_str(text: &str) -> Result<Self, Self::Err> {
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<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
serializer.collect_str(self)
}
}
impl<'de> Deserialize<'de> for FiringKey {
fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
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<FiringKey, u32> = 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<FiringKey, u32> = serde_json::from_str(&json).expect("the map decodes");
assert_eq!(back, stages);
assert!("3-7".parse::<FiringKey>().is_err());
assert_eq!(
FiringKey::from(StagePosition {
execution: 3,
firing: 7,
}),
FiringKey::new(3, 7)
);
}
}

View file

@ -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,

View file

@ -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<FiringKey>,
at: DateTime<Utc>,
) {
let closed: Vec<String> = match (question, firing_key) {
let closed: Vec<String> = 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(),

View file

@ -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<String>,
finished_before: &BTreeSet<FiringKey>,
run_finished: bool,
) -> (Vec<Item<'a>>, 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<String> = [projection::stage_key(0, 1)].into_iter().collect();
let finished_before: BTreeSet<FiringKey> = [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![