mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Fabro's hooks on a Petri run wrap the hooks the runtime installed for `[[run.hooks]]` and forward every point. In `prepare_result`, before the finish is recorded, they commit the stage's files on the run branch of its host workspace with Fabro's author identity and the run, execution, firing and attempt as trailers, and publish the commit to a snapshot repository beside the run's workspaces under a ref per checkpoint. A stage that failed on its own terms is committed like a successful one; a commit that fails is fatal: the outcome becomes a `checkpoint_failed` failure, the run is cancelled through the coordinator handle, and the transition refuses the firing's routes. In `transition` they write the platform checkpoint record, keyed on the Petri position and the checkpoint's operation identity, and a failed write is a recorded problem. On restart the server runs the recovery protocol before it relaunches a worker: a run with a failed checkpoint is reported failed; otherwise every live execution's last durable finish names the snapshot its workspace is verified against, reset to, or restored from, with a lost record reconciled from the snapshot repository, and a finish with no snapshot fails the run rather than resume it on stale files. The worker reaches the platform records over two new worker-scoped endpoints; the server reaches the table directly. A test gate directory lets the CLI scenarios hold a checkpoint at a named point. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
71 lines
2.2 KiB
Rust
71 lines
2.2 KiB
Rust
//! Petri's test kit, for Fabro crates that check a store implementation
|
|
//! against Petri's contract from their own tests, and an in-memory platform
|
|
//! record store for tests of the hooks and recovery. Compiled only with the
|
|
//! `test-support` feature, which a dev-dependency turns on.
|
|
|
|
use std::collections::HashMap;
|
|
use std::sync::{Mutex, MutexGuard, PoisonError};
|
|
|
|
use async_trait::async_trait;
|
|
use fabro_store::platform_records::now_ms;
|
|
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord};
|
|
use fabro_types::RunId;
|
|
pub use petri_testkit::run_store;
|
|
|
|
use crate::platform_records::{PlatformRecordError, PlatformRecords};
|
|
|
|
/// Platform records kept in memory, per run, in seq order.
|
|
#[derive(Debug, Default)]
|
|
pub struct MemoryPlatformRecords {
|
|
runs: Mutex<HashMap<RunId, Vec<StoredPlatformRecord>>>,
|
|
}
|
|
|
|
impl MemoryPlatformRecords {
|
|
#[must_use]
|
|
pub fn new() -> Self {
|
|
Self::default()
|
|
}
|
|
|
|
/// Every record of the run, in seq order.
|
|
#[must_use]
|
|
pub fn records(&self, run_id: &RunId) -> Vec<StoredPlatformRecord> {
|
|
lock(&self.runs).get(run_id).cloned().unwrap_or_default()
|
|
}
|
|
}
|
|
|
|
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
|
|
mutex.lock().unwrap_or_else(PoisonError::into_inner)
|
|
}
|
|
|
|
#[async_trait]
|
|
impl PlatformRecords for MemoryPlatformRecords {
|
|
async fn append(
|
|
&self,
|
|
run_id: &RunId,
|
|
record: &PlatformRecord,
|
|
position: Option<StagePosition>,
|
|
) -> Result<StoredPlatformRecord, PlatformRecordError> {
|
|
let mut runs = lock(&self.runs);
|
|
let records = runs.entry(*run_id).or_default();
|
|
let stored = StoredPlatformRecord {
|
|
seq: records.len() as u64 + 1,
|
|
recorded_at: now_ms(),
|
|
record: record.clone(),
|
|
position,
|
|
};
|
|
records.push(stored.clone());
|
|
Ok(stored)
|
|
}
|
|
|
|
async fn read_kind(
|
|
&self,
|
|
run_id: &RunId,
|
|
kind: PlatformRecordKind,
|
|
) -> Result<Vec<StoredPlatformRecord>, PlatformRecordError> {
|
|
Ok(self
|
|
.records(run_id)
|
|
.into_iter()
|
|
.filter(|record| record.record.kind() == kind)
|
|
.collect())
|
|
}
|
|
}
|