Merge branch 'petri-integration-read' into petri-integration

# Conflicts:
#	Cargo.lock
#	lib/components/fabro-petri/Cargo.toml
#	lib/components/fabro-petri/README.md
#	lib/components/fabro-petri/src/lib.rs
This commit is contained in:
Bryan Helmkamp 2026-09-17 23:48:36 -04:00
commit 5dae891a98
No known key found for this signature in database
22 changed files with 4693 additions and 30 deletions

2
Cargo.lock generated
View file

@ -2891,6 +2891,7 @@ dependencies = [
"anyhow",
"async-trait",
"bytes",
"chrono",
"fabro-api",
"fabro-auth",
"fabro-client",
@ -2901,6 +2902,7 @@ dependencies = [
"fabro-store",
"fabro-test",
"fabro-types",
"fabro-util",
"fabro-vault",
"fabro-workflow",
"lithos-llm",

View file

@ -838,6 +838,11 @@ where
"Reconciled stale in-flight runs on startup"
);
}
state
.petri_projector
.startup_pass()
.await
.context("catching Petri projections up at startup")?;
spawn_scheduler(Arc::clone(&state));
spawn_automation_scheduler(Arc::clone(&state));
let pull_request_creation_supervisor =

View file

@ -61,6 +61,7 @@ use fabro_llm::credentials::CredentialProvider;
use fabro_llm::lithos_catalog::Catalog;
use fabro_llm::{ClientOptions, FabroClient};
use fabro_mcp_store::McpServerStore;
use fabro_petri::projector::Projector;
use fabro_redact::redact_jsonl_line;
use fabro_sandbox::details::sandbox_details;
use fabro_sandbox::driver::{DaytonaCredentials, ProviderAccess, ProviderConnectOptions};
@ -1116,6 +1117,8 @@ pub struct AppState {
pub(crate) worker_runtime: Arc<dyn WorkerRuntime>,
/// The Petri runs held open for workers over the API.
pub(crate) petri_runs: PetriRuns,
/// The projector of Petri runs: signalled after each committed record.
pub(crate) petri_projector: Arc<Projector>,
scheduler_notify: Notify,
automation_scheduler_notify: Notify,
pull_request_scheduler_notify: Notify,
@ -1187,6 +1190,19 @@ impl AppState {
self.petri_runs.store()
}
/// The projector of Petri runs, so a test can wait for a run's view to
/// settle before it reads it.
#[cfg(any(test, feature = "test-support"))]
pub fn test_petri_projector(&self) -> &Arc<Projector> {
&self.petri_projector
}
/// The pool the Petri view tables live in, so a test can read them.
#[cfg(any(test, feature = "test-support"))]
pub fn test_petri_view_pool(&self) -> DbPool {
self.stores.runs.run_summary_store().pool()
}
/// A worker token for `run_id` with the plain `run:worker` scope, as the
/// server mints for the worker it launches.
pub fn test_issue_worker_token(&self, run_id: &RunId) -> String {
@ -2490,6 +2506,14 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
);
let variables = Arc::new(VariableStore::new(db_pool.clone()));
let petri_runs = PetriRuns::new(db_pool.clone());
// Petri's records live on the shared pool; the view tables live where the
// run summary store keeps the `runs` row (the same database in the
// server, a fixture of its own in a test).
let petri_projector = Projector::new(db_pool.clone(), store.run_summary_store().pool());
{
let projector = Arc::clone(&petri_projector);
store.set_platform_record_hook(Arc::new(move |run_id| projector.signal(run_id)));
}
let session_records = Arc::new(RunSessionRecordStore::new(db_pool.clone()));
let secret_store = Arc::new(SecretStore::new(db_pool));
let vault = preloaded_vault;
@ -2607,6 +2631,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
worker_control_bus,
worker_runtime,
petri_runs,
petri_projector,
scheduler_notify: Notify::new(),
automation_scheduler_notify: Notify::new(),
pull_request_scheduler_notify: Notify::new(),
@ -4539,6 +4564,7 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
// The worker is gone: whatever Petri run handles it held open over the
// API drop here, so its lease never outlives it.
state.petri_runs.worker_exited(run_id);
state.petri_projector.signal(run_id);
append_worker_exit_failure(&run_store, run_id, &worker_exit).await;
let final_state = match run_store.state().await {

View file

@ -145,7 +145,11 @@ async fn append_records(
Err(err) => return store_error_response(id, &err),
};
match writer.append(&log, &records).await {
Ok(()) => StatusCode::NO_CONTENT.into_response(),
Ok(()) => {
// The records are durable; the projection trails them from here.
state.petri_projector.signal(id);
StatusCode::NO_CONTENT.into_response()
}
Err(err) => store_error_response(id, &err),
}
}

View file

@ -382,7 +382,9 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
run_id: run_id.to_string(),
run_dir: run_dir.join("petri"),
execution,
store: Arc::new(SqliteRunStore::new(state.db_pool.clone())),
store: state
.petri_projector
.observe_store(Arc::new(SqliteRunStore::new(state.db_pool.clone()))),
runtime: runtime_spec(&state, &eligible, dry_run),
provider: run_state.spec.settings.run.environment.provider.clone(),
cancel,

View file

@ -26,8 +26,8 @@ use std::sync::Arc;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use fabro_petri::SqliteRunStore;
use fabro_petri::engine::{self, RunStatus};
use fabro_petri::{SqliteRunStore, projector};
use fabro_server::server::AppState;
use fabro_server::test_support::{
TestAppStateBuilder, llm_overlay_with_provider_base_url, test_app_db_pool,
@ -35,7 +35,7 @@ use fabro_server::test_support::{
};
use fabro_static::EnvVars;
use fabro_test::{TwinScenario, TwinScenarios, twin_openai};
use fabro_types::{WorkflowPath, WorkflowVersion};
use fabro_types::{RunId, WorkflowPath, WorkflowVersion};
use tower::ServiceExt;
use crate::helpers::{
@ -87,6 +87,23 @@ const UNKNOWN_MODEL_DOT: &str = r#"digraph Bad {
start -> work -> exit
}"#;
/// Two command branches joined by a fan-in.
const PARALLEL_DOT: &str = r#"digraph Parallel {
graph [goal="Run two branches"]
start [shape=Mdiamond]
exit [shape=Msquare]
fork [shape=component]
a [shape=parallelogram, script="echo a"]
b [shape=parallelogram, script="echo b"]
merge [shape=tripleoctagon]
start -> fork
fork -> a
fork -> b
a -> merge
b -> merge
merge -> exit
}"#;
const PLAIN_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
const PETRI_SETTINGS: &str =
"_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n";
@ -168,6 +185,37 @@ async fn petri_outcome(state: &AppState, run_id: &str) -> engine::RunOutcome {
.expect("the run's Petri record inspects")
}
/// The run's projected state once its projector settled.
async fn settled_state(state: &AppState, app: &axum::Router, run_id: &str) -> serde_json::Value {
let id: RunId = run_id.parse().expect("the run id parses");
state.test_petri_projector().settle(id).await;
let req = Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}/state")))
.body(Body::empty())
.expect("state request should build");
let response = app
.clone()
.oneshot(req)
.await
.expect("state request routes");
response_json(
response,
StatusCode::OK,
format!("GET /api/v1/runs/{run_id}/state"),
)
.await
}
/// How many items the run's projected stream holds.
async fn petri_stream_len(state: &AppState, run_id: &str) -> usize {
let id: RunId = run_id.parse().expect("the run id parses");
projector::stored_stream(&state.test_petri_view_pool(), id)
.await
.expect("the stream reads")
.len()
}
async fn run_engine(app: &axum::Router, run_id: &str) -> serde_json::Value {
let req = Request::builder()
.method("GET")
@ -260,6 +308,22 @@ async fn the_hello_bundle_runs_on_petri_when_the_version_names_the_engine() {
let outcome = petri_outcome(&state, &run_id).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(outcome.complete, "{:?}", outcome.incomplete);
let projection = settled_state(&state, &app, &run_id).await;
assert_eq!(projection["status"]["kind"], "succeeded", "{projection}");
assert_eq!(
projection["conclusion"]["status"], "succeeded",
"{projection}"
);
let greet = &projection["stages"]["greet@1"];
assert_eq!(greet["state"], "succeeded", "{projection}");
assert_eq!(greet["handler"], "agent", "{greet}");
assert!(
greet["response"]
.as_str()
.is_some_and(|response| response.contains("A haiku, added.")),
"the agent's answer is projected as the stage's response: {greet}"
);
assert!(run["usage"]["tokens"]["input"].as_u64().is_some(), "{run}");
let logs = twin.request_logs(&namespace).await;
let requests = logs["requests"]
.as_array()
@ -302,6 +366,67 @@ async fn a_command_bundle_runs_on_petri_under_the_server_setting() {
let outcome = petri_outcome(&state, &run_id).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(outcome.complete, "{:?}", outcome.incomplete);
let projection = settled_state(&state, &app, &run_id).await;
let say = &projection["stages"]["say@1"];
assert_eq!(say["state"], "succeeded", "{projection}");
assert_eq!(say["handler"], "command", "{say}");
assert!(
say["output"]
.as_str()
.is_some_and(|output| output.contains("hello from petri")),
"{say}"
);
let stream = petri_stream_len(&state, &run_id).await;
assert!(stream > 0, "the run's stream holds its events");
}
/// A parallel bundle with two command branches runs on Petri through the
/// server: each branch is a child execution, projected as a stage grouped
/// under the fork, and the fork carries the branch results.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_parallel_bundle_projects_its_branches_through_the_server() {
if host_plugin().is_none() {
return;
}
let workspace = tempfile::tempdir().expect("workspace tempdir");
let settings = settings_from_toml(
"_version = 1\n\n[run.environment]\nid = \"local\"\n\n[server.execution]\nengine = \
\"petri\"\n",
);
let state = test_app_state_with_options(settings, 5);
let app = test_app_with_scheduler(Arc::clone(&state));
let version_id = register_version(&app, &[
("workflow.fabro", PARALLEL_DOT),
("workflow.toml", PLAIN_SETTINGS),
])
.await;
let run_id =
create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await;
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
let run = run_json(&app, &run_id).await;
assert_eq!(status, "succeeded", "run: {run}");
let projection = settled_state(&state, &app, &run_id).await;
for (branch, index) in [("a@1", 0), ("b@1", 1)] {
let stage = &projection["stages"][branch];
assert_eq!(stage["state"], "succeeded", "{branch}: {projection}");
assert_eq!(
stage["parallel_branch_id"],
format!("fork@1:{index}"),
"{stage}"
);
}
let fork = &projection["stages"]["fork@1"];
assert_eq!(
fork["parallel_results"].as_array().map(Vec::len),
Some(2),
"{fork}"
);
assert_eq!(
projection["conclusion"]["status"], "succeeded",
"{projection}"
);
}
/// A version that names no engine on a server whose setting is the default

View file

@ -28,6 +28,7 @@ fabro-store = { path = "../fabro-store" }
fabro-types = { path = "../../foundation/fabro-types" }
fabro-vault = { path = "../../foundation/fabro-vault" }
fabro-workflow = { path = "../fabro-workflow" }
fabro-util = { path = "../../foundation/fabro-util" }
petri_runtime.workspace = true
petri_execution.workspace = true
petri_store.workspace = true
@ -39,6 +40,7 @@ petri_testkit = { workspace = true, optional = true }
anyhow.workspace = true
bytes.workspace = true
async-trait.workspace = true
chrono = { workspace = true, features = ["serde"] }
serde.workspace = true
serde_json.workspace = true
sqlx.workspace = true
@ -52,6 +54,7 @@ fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"]
fabro-llm = { path = "../fabro-llm", features = ["test-support"] }
fabro-store = { path = "../fabro-store", features = ["test-support"] }
fabro-test.workspace = true
fabro-types = { path = "../../foundation/fabro-types", features = ["test-support"] }
petri_testkit.workspace = true
tempfile = "3"
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }

View file

@ -67,8 +67,38 @@ Every adapter the integration plan describes lands here.
- `petri`: the Petri store vocabulary re-exported for the server, which
answers the worker endpoints from a `SqliteRunStore` without naming a Petri
package in its own manifest.
- The platform adapters the plan adds after it: hooks, run tools, the
event projection.
- `projection`: the fold of a Petri run's public events (`replay_since` over
its records) and Fabro's platform records (`fabro-store`'s
`platform_records`) into the `RunProjection` the API serves, row by row as
`VIEWS.md` maps them. The stage key is `(execution, firing)`; the
`StageId` label is `node@visit`, made unique with the execution when two
child invocations would share one.
- `projector`: the view pass and its wake-up. Records first: Petri's append
and a platform record's insert return before any view work; a pass reads
what is committed, folds the items past the committed positions, and
writes the projection document (`petri_projection`), the ordered stream
(`petri_stream`, one `stream_seq` per Petri event or platform record) and
the narrowed `runs` row in one later transaction. The server signals the
projector after each committed worker append, after each committed
platform record (the run summary store's hook), at worker exit and, over
every Petri run, at startup. A run that executes in the server process
goes through `Projector::observe_store`, which signals after each append.
A torn tail (a record Petri cannot read) holds the view where it stands
and reports the run incomplete with the reason.
- The platform adapters the plan adds after it: hooks and the run tools.
### What the projection leaves default
`VIEWS.md` rows with no source yet, or whose source this crate does not read
yet, keep their default value in the projection: `StageProjection.diff` and
`Conclusion.diff.patch` (the checkpoint's `patch_blob` is not resolved),
`Checkpoint`'s engine-derived maps (`completed_nodes`, `node_retries`,
`context_values`, `node_outcomes`, `next_node_id`), `agent_tools`,
`permission_level`, `script_invocation` and `script_timing`, a stage's
`notes`, `StageCompletion` details for a `parsed.note`, the sandbox instance
(the matrix's two gaps), `Run.ask_fabro`, an interview option's
`description` and `preview`, the pull request `creation` state, and the
run's notices, notifications and pairings (recorded, not shown).
A run goes to Petri when its workflow version's `workflow.toml` names
`engine = "petri"` in `[workflow]`, or when the server's
@ -118,6 +148,16 @@ Integration tests live under `tests/`:
Those four need the host plugin like `runs.rs` does, and `model.rs` also
starts the twin.
- `projection.rs` builds the view live (every append signals the
projector) for the `hello` bundle on the stub registry, a command-only
workflow and a two-branch parallel workflow, and checks it equals the view
rebuilt from the records alone (`projector::rebuild`); catches a view up
after every wake-up was dropped, by a signal and by the startup pass;
recovers a crash between the record commit and the view transaction by
applying only the missing suffix, with the positions and `stream_seq`
continuing; runs two projectors over one store with child executions; and
holds the view at a torn tail. All skip without the host plugin.
The conformance suite over `HttpRunStore` needs a server to talk to, so it
lives with the server's integration tests
(`lib/apps/fabro-server/tests/it/api/petri_store.rs`), which reach the suite
@ -130,13 +170,14 @@ ulimit -n 4096 && cargo nextest run -p fabro-petri
```
The server's end-to-end coverage is `lib/apps/fabro-server/tests/it/scenario/petri.rs`:
the `hello` bundle on the OpenAI twin and a command-only bundle run to
completion through the create handler and the scheduler, in the server
process under its test override, under the version flag and under the
server setting, a human gate is answered through the questions API, and
Petri's diagnostics refuse a run at create. The server's `petri_runs` unit
tests cover the lease ending at worker exit and the restart reconcile that
relaunches a worker in resume mode.
the `hello` bundle on the OpenAI twin, a command-only bundle and a
two-branch parallel bundle run to completion through the create handler and
the scheduler, in the server process under its test override, under the
version flag and under the server setting, with `GET /runs/{id}/state`
serving the projection over Petri's records; a human gate is answered
through the questions API; and Petri's diagnostics refuse a run at create.
The server's `petri_runs` unit tests cover the lease ending at worker exit
and the restart reconcile that relaunches a worker in resume mode.
The worker path is covered with the real binary in
`lib/apps/fabro-cli/tests/it/scenario/petri.rs`: a command-only Petri run

View file

@ -6,6 +6,11 @@ comes from once Petri's records are the store. It is written before any view
changes. F2.2 (the projection), F2.3 (platform records) and F2.4 (API, CLI,
web) build from it.
F2.2 and F2.3 implement this matrix: `src/projection.rs` is the fold,
`src/projector.rs` the view pass and its wake-up, and `fabro-store`'s
`platform_records` module the platform record kinds and their table. The
crate README names the rows the fold still leaves default.
Sources are named three ways:
- A Petri event, by its `<subject>.<verb>` name from

View file

@ -30,8 +30,10 @@
//! - [`HttpRunStore`]: the same store as a run's worker process reaches it,
//! over the server's API with the worker's token and its launch id as the
//! lease owner;
//! - the platform adapters still to come: hooks, the run tools, the event
//! projection.
//! - [`projection`] and [`projector`]: the view of a Petri run, folded from its
//! records and Fabro's platform records, and the pass that writes it after
//! each committed record;
//! - the platform adapters still to come: hooks and the run tools.
//!
//! The Petri packages are pinned by revision in the workspace `Cargo.toml`
//! under `petri_*` keys.
@ -43,6 +45,8 @@ pub mod engine;
pub mod http_store;
pub mod interview;
pub mod petri;
pub mod projection;
pub mod projector;
pub mod run_store;
pub mod runtime;
pub mod secrets;

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,800 @@
//! The projector: the view pass that folds a Petri run's committed records
//! into its stored projection, and the wake-up that drives it.
//!
//! # Commit rule
//!
//! Records first. Petri's append (the worker's append endpoint, then
//! `SqliteRunStore::append`) and a platform record's insert are the
//! durability boundaries, and both return before any view work. A view pass
//! then reads what is committed, folds the items past the positions the
//! view last committed, and writes the derived rows in one later
//! transaction together with the new positions: the last event consumed per
//! Petri log, the last platform record consumed, and the delivery sequence
//! (`stream_seq`) it assigned to each item. The view therefore trails a
//! committed record and never leads one. No projection state of Petri's is
//! checkpointed: each pass replays the run through `replay_since`, which
//! rebuilds the engine and invocation state the derivation needs and
//! delivers only the events past the held positions.
//!
//! A pass that finds new platform records committed between its read and
//! its write leaves the view alone and runs again, so the `runs` row never
//! moves backwards behind a concurrent lifecycle write.
//!
//! # Where it runs
//!
//! In the server. [`Projector::signal`] schedules a pass for a run: the
//! server calls it after each committed worker append and, through the run
//! summary store's hook, after each committed platform record; signals
//! that arrive while a pass runs coalesce into one more pass. A signal is a
//! wake-up only, never a source of facts: a signal that is lost costs
//! nothing but latency, because the next signal or the startup pass
//! ([`Projector::startup_pass`]) folds everything the view still trails.
//!
//! # A torn tail
//!
//! A record the store holds that Petri cannot read (a gap in a log, a line
//! that does not decode) fails the replay. The pass then advances no Petri
//! position, folds only the platform records, and reports the run's record
//! as incomplete with the replay's error; `inspect_run` decides
//! completeness once the run has recorded its finish.
use std::collections::{BTreeMap, HashMap};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::Duration;
use fabro_db::DbPool;
use fabro_store::platform_records::{PlatformRecordStore, StoredPlatformRecord, now_ms};
use fabro_store::{RunProjection, RunSummaryStore};
use fabro_types::RunId;
use fabro_util::error::collect_chain;
use petri_execution::events::{self, EventId, EventSource, RunEvent};
use petri_execution::{Access, RunKey, RunStore as _, inspect};
use petri_store::StoreError;
use serde::{Deserialize, Serialize};
use tokio::sync::Mutex as AsyncMutex;
use tokio::time;
use tracing::{debug, info, warn};
use crate::SqliteRunStore;
use crate::projection::{self, FoldState, Item, RecordHealth, RunView};
/// The positions a view committed: the last event consumed per Petri log,
/// and the last platform record consumed.
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Positions {
#[serde(default)]
pub petri: Vec<EventId>,
#[serde(default)]
pub platform_seq: u64,
}
impl Positions {
fn held(&self) -> BTreeMap<EventSource, EventId> {
self.petri.iter().map(|id| (id.source, *id)).collect()
}
fn advance(&mut self, id: EventId) {
match self.petri.iter_mut().find(|held| held.source == id.source) {
Some(held) => {
if id > *held {
*held = id;
}
}
None => self.petri.push(id),
}
}
}
/// What one pass did.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct PassReport {
pub run_id: RunId,
/// The pass found nothing past the committed positions and wrote nothing.
pub skipped: bool,
/// The view was left alone because a platform record landed during the
/// pass; the projector runs the pass again.
pub contended: bool,
pub petri_events: usize,
pub platform_records: usize,
/// The last delivery sequence the view holds.
pub stream_seq: u64,
pub positions: Positions,
pub health: RecordHealth,
}
/// What the startup pass did.
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct StartupReport {
pub runs: usize,
pub projected: usize,
/// Runs whose pass failed and was left for the next signal.
pub failed: usize,
}
/// Why a pass could not run or commit.
#[derive(Debug, thiserror::Error)]
pub enum ProjectError {
#[error("the run's Petri record could not be opened")]
Open(#[source] StoreError),
#[error("the projection tables could not be read or written")]
Database(#[source] sqlx::Error),
#[error("the platform records could not be read or written")]
Store(#[source] fabro_store::Error),
#[error("the view could not be encoded")]
Encode(#[source] serde_json::Error),
#[error("the pass was stopped before its view transaction (injected)")]
Injected,
}
/// The stored view of a run, as the projection tables hold it.
struct StoredView {
view: RunView,
positions: Positions,
stream_seq: u64,
}
/// A run's pass state under the projector's lock.
#[derive(Default)]
struct Slot {
running: bool,
pending: bool,
}
/// The projector over one database: the pool Petri's records are read
/// from, and the pool the view tables (`platform_records`,
/// `petri_projection`, `petri_stream`, `runs`) are read and written on. In
/// the server both are the one database; a test may hand it the run
/// summary store's own pool for the views.
pub struct Projector {
records: DbPool,
pool: DbPool,
store: SqliteRunStore,
platform: PlatformRecordStore,
slots: Mutex<HashMap<RunId, Slot>>,
/// One pass at a time per run: a signalled pass and the startup pass
/// over the same run never interleave their reads and writes.
passes: Mutex<HashMap<RunId, Arc<AsyncMutex<()>>>>,
/// Test-only: stop the next pass after its reads, before its view
/// transaction, as a crash there would.
fault: AtomicBool,
}
impl std::fmt::Debug for Projector {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Projector").finish_non_exhaustive()
}
}
impl Projector {
/// A projector over `records`, the pool Petri's records live in, and
/// `views`, the pool the view tables live in; both migrated. The server
/// passes its one pool twice.
#[must_use]
pub fn new(records: DbPool, views: DbPool) -> Arc<Self> {
Arc::new(Self {
store: SqliteRunStore::new(records.clone()),
platform: PlatformRecordStore::new(views.clone()),
records,
pool: views,
slots: Mutex::default(),
passes: Mutex::default(),
fault: AtomicBool::new(false),
})
}
/// Schedule a pass for the run. A pass already running for it runs once
/// more when it ends; any number of signals in between coalesce.
pub fn signal(self: &Arc<Self>, run_id: RunId) {
{
let mut slots = lock(&self.slots);
let slot = slots.entry(run_id).or_default();
if slot.running {
slot.pending = true;
return;
}
slot.running = true;
}
let projector = Arc::clone(self);
tokio::spawn(async move {
loop {
let again = match projector.project_run(run_id).await {
Ok(report) => report.contended,
Err(error) => {
warn!(
run_id = %run_id,
error = %collect_chain(&error).join(": "),
"Petri projection pass failed; the next signal retries it"
);
false
}
};
let mut slots = lock(&projector.slots);
let slot = slots.entry(run_id).or_default();
if again || slot.pending {
slot.pending = false;
continue;
}
slot.running = false;
return;
}
});
}
/// Wait until no pass is running or pending for the run: a test's way
/// to observe the view after its signals.
pub async fn settle(&self, run_id: RunId) {
loop {
let idle = {
let slots = lock(&self.slots);
slots
.get(&run_id)
.is_none_or(|slot| !slot.running && !slot.pending)
};
if idle {
return;
}
time::sleep(Duration::from_millis(5)).await;
}
}
/// Stop the next pass after its reads and before its view transaction,
/// as a crash there would, once.
pub fn fail_before_view(&self) {
self.fault.store(true, Ordering::SeqCst);
}
/// One pass over every Petri run the database holds: the runs with a
/// Petri record, and the runs with platform records. Runs whose view
/// already covers every committed record are skipped cheaply.
pub async fn startup_pass(&self) -> Result<StartupReport, ProjectError> {
let mut ids: Vec<String> = sqlx::query_scalar("SELECT run_id FROM petri_runs")
.fetch_all(&self.records)
.await
.map_err(ProjectError::Database)?;
let with_platform: Vec<String> =
sqlx::query_scalar("SELECT DISTINCT run_id FROM platform_records")
.fetch_all(&self.pool)
.await
.map_err(ProjectError::Database)?;
ids.extend(with_platform);
ids.sort();
ids.dedup();
let mut report = StartupReport::default();
for id in ids {
let Some(run_id) = projection::run_id_of(&id) else {
debug!(run_key = %id, "Petri run key is not a Fabro run id; not projected");
continue;
};
report.runs += 1;
match self.project_run(run_id).await {
Ok(pass) => {
if !pass.skipped {
report.projected += 1;
}
}
// One run's view trailing never stops the server: the next
// signal for the run retries its pass.
Err(error) => {
warn!(
run_id = %run_id,
error = %collect_chain(&error).join(": "),
"Petri projection pass failed at startup; the next signal retries it"
);
report.failed += 1;
}
}
}
if report.projected > 0 {
info!(
runs = report.runs,
projected = report.projected,
"Petri projections caught up at startup"
);
}
Ok(report)
}
/// One view pass for the run. Passes over one run run one at a time.
pub async fn project_run(&self, run_id: RunId) -> Result<PassReport, ProjectError> {
let pass = Arc::clone(lock(&self.passes).entry(run_id).or_default());
let _one_at_a_time = pass.lock().await;
let stored = self.load_view(&run_id).await?;
let key = RunKey::new(run_id.to_string());
let platform_head = self
.platform
.head(&run_id)
.await
.map_err(ProjectError::Store)?
.unwrap_or(0);
let petri_heads = self.petri_heads(&run_id).await?;
let at_head = platform_head == stored.positions.platform_seq
&& petri_heads.iter().all(|(log, head)| {
stored
.positions
.petri
.iter()
.any(|held| log_text(&held.source) == *log && held.seq == *head)
});
if at_head && stored.view.projection.is_some() {
return Ok(PassReport {
run_id,
skipped: true,
contended: false,
petri_events: 0,
platform_records: 0,
stream_seq: stored.stream_seq,
positions: stored.positions,
health: stored.view.state.health,
});
}
let StoredView {
mut view,
mut positions,
mut stream_seq,
} = stored;
let platform_records = self
.platform
.read_after(&run_id, positions.platform_seq)
.await
.map_err(ProjectError::Store)?;
let (events, replay_failure) = match self.store.open(&key, Access::Read).await {
Ok(logs) => match events::replay_since(&*logs, &positions.held()).await {
Ok(events) => (events, None),
Err(error) => {
let chain = collect_chain(&error).join(": ");
warn!(run_id = %run_id, error = %chain, "Petri run does not replay; the view holds");
(Vec::new(), Some(chain))
}
},
Err(StoreError::NotFound { .. }) => (Vec::new(), None),
Err(error) => return Err(ProjectError::Open(error)),
};
let mut items: Vec<(u64, u8, Item<'_>)> =
Vec::with_capacity(events.len() + platform_records.len());
for event in &events {
let rank = match event.id.source {
EventSource::Coordinator => 0,
EventSource::Execution { .. } => 1,
};
items.push((event.recorded_at, rank, Item::Petri(event)));
}
for record in &platform_records {
items.push((record.recorded_at, 2, Item::Platform(record)));
}
items.sort_by_key(|(recorded_at, rank, _)| (*recorded_at, *rank));
let mut rows: Vec<StreamRow> = Vec::with_capacity(items.len());
for (_, _, item) in &items {
stream_seq += 1;
view.fold(item, stream_seq);
let row = match item {
Item::Petri(event) => {
positions.advance(event.id);
StreamRow {
stream_seq,
item_kind: "petri",
item_id: event_id_text(&event.id),
event_json: serde_json::to_string(event).map_err(ProjectError::Encode)?,
}
}
Item::Platform(record) => {
positions.platform_seq = record.seq;
StreamRow {
stream_seq,
item_kind: "platform",
item_id: record.seq.to_string(),
event_json: serde_json::to_string(record).map_err(ProjectError::Encode)?,
}
}
};
rows.push(row);
}
view.state.health = self.health(&key, &view.state, replay_failure).await?;
if self.fault.swap(false, Ordering::SeqCst) {
return Err(ProjectError::Injected);
}
let mut tx = self
.pool
.begin_with("BEGIN IMMEDIATE")
.await
.map_err(ProjectError::Database)?;
let head_now: i64 = sqlx::query_scalar(
"SELECT COALESCE(MAX(seq), 0) FROM platform_records WHERE run_id = ?",
)
.bind(run_id.to_string())
.fetch_one(&mut *tx)
.await
.map_err(ProjectError::Database)?;
if u64::try_from(head_now).unwrap_or(0) != positions.platform_seq {
debug!(run_id = %run_id, "platform records landed during the pass; running it again");
drop(tx);
return Ok(PassReport {
run_id,
skipped: false,
contended: true,
petri_events: 0,
platform_records: 0,
stream_seq: 0,
positions: Positions::default(),
health: RecordHealth::default(),
});
}
let projection_json =
serde_json::to_string(&view.projection).map_err(ProjectError::Encode)?;
let fold_json = serde_json::to_string(&view.state).map_err(ProjectError::Encode)?;
let positions_json = serde_json::to_string(&positions).map_err(ProjectError::Encode)?;
sqlx::query(
"INSERT INTO petri_projection (run_id, projection_json, fold_json, positions_json, \
stream_seq, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(run_id) DO UPDATE \
SET projection_json = excluded.projection_json, fold_json = excluded.fold_json, \
positions_json = excluded.positions_json, stream_seq = excluded.stream_seq, \
updated_at_ms = excluded.updated_at_ms",
)
.bind(run_id.to_string())
.bind(projection_json)
.bind(fold_json)
.bind(positions_json)
.bind(column(stream_seq))
.bind(column(now_ms()))
.execute(&mut *tx)
.await
.map_err(ProjectError::Database)?;
for row in &rows {
sqlx::query(
"INSERT INTO petri_stream (run_id, stream_seq, item_kind, item_id, event_json) \
VALUES (?, ?, ?, ?, ?)",
)
.bind(run_id.to_string())
.bind(column(row.stream_seq))
.bind(row.item_kind)
.bind(&row.item_id)
.bind(&row.event_json)
.execute(&mut *tx)
.await
.map_err(ProjectError::Database)?;
}
if let Some(projection) = view.projection.as_ref() {
RunSummaryStore::write_petri_run_row_on_connection(&mut tx, &run_id, projection)
.await
.map_err(ProjectError::Store)?;
}
tx.commit().await.map_err(ProjectError::Database)?;
debug!(
run_id = %run_id,
petri_events = events.len(),
platform_records = platform_records.len(),
stream_seq,
"Petri projection pass committed"
);
Ok(PassReport {
run_id,
skipped: false,
contended: false,
petri_events: events.len(),
platform_records: platform_records.len(),
stream_seq,
positions,
health: view.state.health.clone(),
})
}
/// The stored view of the run, or an empty one.
async fn load_view(&self, run_id: &RunId) -> Result<StoredView, ProjectError> {
let row: Option<(String, String, String, i64)> = sqlx::query_as(
"SELECT projection_json, fold_json, positions_json, stream_seq FROM petri_projection \
WHERE run_id = ?",
)
.bind(run_id.to_string())
.fetch_optional(&self.pool)
.await
.map_err(ProjectError::Database)?;
let Some((projection_json, fold_json, positions_json, stream_seq)) = row else {
return Ok(StoredView {
view: RunView::new(),
positions: Positions::default(),
stream_seq: 0,
});
};
let projection: Option<RunProjection> =
serde_json::from_str(&projection_json).map_err(ProjectError::Encode)?;
let state: FoldState = serde_json::from_str(&fold_json).map_err(ProjectError::Encode)?;
let positions: Positions =
serde_json::from_str(&positions_json).map_err(ProjectError::Encode)?;
Ok(StoredView {
view: RunView { projection, state },
positions,
stream_seq: u64::try_from(stream_seq).unwrap_or(0),
})
}
/// The last seq of every Petri log of the run, by the log column's text.
async fn petri_heads(&self, run_id: &RunId) -> Result<Vec<(String, u64)>, ProjectError> {
// The coordinator log and the execution logs are what the projection
// reads; the resources log is the sandbox ledger and has no events.
let rows: Vec<(String, i64)> = sqlx::query_as(
"SELECT log, MAX(seq) FROM petri_records WHERE run_id = ? AND (log = 'coordinator' \
OR log LIKE 'execution %') GROUP BY log",
)
.bind(run_id.to_string())
.fetch_all(&self.records)
.await
.map_err(ProjectError::Database)?;
Ok(rows
.into_iter()
.map(|(log, seq)| (log, u64::try_from(seq).unwrap_or(0)))
.collect())
}
/// Whether the run's record is whole: a replay failure says no with its
/// reason; a run that has not recorded its finish is not yet; a finished
/// run is what `inspect_run` says, checked until it says complete.
async fn health(
&self,
key: &RunKey,
state: &FoldState,
replay_failure: Option<String>,
) -> Result<RecordHealth, ProjectError> {
if let Some(failure) = replay_failure {
return Ok(RecordHealth {
complete: false,
incomplete: vec![failure],
});
}
if state.finished.is_none() {
return Ok(RecordHealth {
complete: false,
incomplete: vec!["the run has not recorded its finish".to_string()],
});
}
if state.health.complete {
return Ok(state.health.clone());
}
let logs = match self.store.open(key, Access::Read).await {
Ok(logs) => logs,
Err(StoreError::NotFound { .. }) => return Ok(state.health.clone()),
Err(error) => return Err(ProjectError::Open(error)),
};
match inspect::inspect_run(&*logs).await {
Ok(inspection) => Ok(RecordHealth {
complete: inspection.complete,
incomplete: inspection.incomplete,
}),
Err(error) => Ok(RecordHealth {
complete: false,
incomplete: vec![collect_chain(&error).join(": ")],
}),
}
}
}
impl Projector {
/// A run store whose appends signal this projector: for a run that
/// executes in the same process as the projector, over the SQLite store
/// directly, where no append endpoint is there to signal. The signal is
/// sent after the store's append returned, so the records it covers are
/// durable before the view sees them.
pub fn observe_store(
self: &Arc<Self>,
inner: Arc<dyn petri_execution::RunStore>,
) -> Arc<dyn petri_execution::RunStore> {
Arc::new(SignallingStore {
inner,
projector: Arc::clone(self),
})
}
}
/// A run store that signals a projector after each append.
struct SignallingStore {
inner: Arc<dyn petri_execution::RunStore>,
projector: Arc<Projector>,
}
#[async_trait::async_trait]
impl petri_execution::RunStore for SignallingStore {
async fn open(
&self,
key: &RunKey,
access: Access,
) -> Result<Arc<dyn petri_execution::RunLogs>, StoreError> {
let logs = self.inner.open(key, access).await?;
Ok(Arc::new(SignallingLogs {
inner: logs,
run_id: projection::run_id_of(key.as_str()),
projector: Arc::clone(&self.projector),
}))
}
}
struct SignallingLogs {
inner: Arc<dyn petri_execution::RunLogs>,
run_id: Option<RunId>,
projector: Arc<Projector>,
}
#[async_trait::async_trait]
impl petri_execution::RunLogs for SignallingLogs {
fn locator(&self) -> String {
self.inner.locator()
}
async fn append(
&self,
log: &petri_execution::LogId,
records: &[petri_execution::Record],
) -> Result<(), StoreError> {
self.inner.append(log, records).await?;
if let Some(run_id) = self.run_id {
self.projector.signal(run_id);
}
Ok(())
}
async fn read(
&self,
log: &petri_execution::LogId,
) -> Result<Vec<petri_execution::Record>, StoreError> {
self.inner.read(log).await
}
async fn put_blob(&self, bytes: &[u8]) -> Result<petri_store::Digest, StoreError> {
self.inner.put_blob(bytes).await
}
async fn get_blob(&self, digest: petri_store::Digest) -> Result<Option<Vec<u8>>, StoreError> {
self.inner.get_blob(digest).await
}
}
struct StreamRow {
stream_seq: u64,
item_kind: &'static str,
item_id: String,
event_json: String,
}
/// A Petri event id as the stream names it: `<log>/<seq>/<index>`.
#[must_use]
pub fn event_id_text(id: &EventId) -> String {
format!("{}/{}/{}", log_text(&id.source), id.seq, id.index)
}
fn log_text(source: &EventSource) -> String {
match source {
EventSource::Coordinator => "coordinator".to_string(),
EventSource::Execution { execution } => format!("execution {execution}"),
}
}
fn column(value: u64) -> i64 {
i64::try_from(value).unwrap_or(i64::MAX)
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
/// The run's projection rebuilt from its records alone, with nothing
/// stored: what a fresh projector would commit over the same records. A test
/// compares it with the live view. `records` and `views` are the two pools
/// [`Projector::new`] takes.
pub async fn rebuild(
records: &DbPool,
views: &DbPool,
run_id: RunId,
) -> Result<(Option<RunProjection>, Positions, u64), ProjectError> {
let store = SqliteRunStore::new(records.clone());
let platform = PlatformRecordStore::new(views.clone());
let key = RunKey::new(run_id.to_string());
let platform_records = platform.read(&run_id).await.map_err(ProjectError::Store)?;
let events = match store.open(&key, Access::Read).await {
Ok(logs) => events::replay_run(&*logs)
.await
.inspect_err(|error| {
warn!(error = %collect_chain(error).join(": "), "rebuild: the run does not replay");
})
.unwrap_or_default(),
Err(StoreError::NotFound { .. }) => Vec::new(),
Err(error) => return Err(ProjectError::Open(error)),
};
let mut items: Vec<(u64, u8, Item<'_>)> = Vec::new();
for event in &events {
let rank = match event.id.source {
EventSource::Coordinator => 0,
EventSource::Execution { .. } => 1,
};
items.push((event.recorded_at, rank, Item::Petri(event)));
}
for record in &platform_records {
items.push((record.recorded_at, 2, Item::Platform(record)));
}
items.sort_by_key(|(recorded_at, rank, _)| (*recorded_at, *rank));
let mut view = RunView::new();
let mut positions = Positions::default();
let mut stream_seq = 0;
for (_, _, item) in &items {
stream_seq += 1;
view.fold(item, stream_seq);
match item {
Item::Petri(event) => positions.advance(event.id),
Item::Platform(record) => positions.platform_seq = record.seq,
}
}
Ok((view.projection, positions, stream_seq))
}
/// The stored view's positions and stream sequence, for a test; `views` is
/// the pool the view tables live in.
pub async fn stored_positions(
views: &DbPool,
run_id: RunId,
) -> Result<Option<(Positions, u64)>, ProjectError> {
let row: Option<(String, i64)> =
sqlx::query_as("SELECT positions_json, stream_seq FROM petri_projection WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_optional(views)
.await
.map_err(ProjectError::Database)?;
row.map(|(positions, stream_seq)| {
Ok((
serde_json::from_str(&positions).map_err(ProjectError::Encode)?,
u64::try_from(stream_seq).unwrap_or(0),
))
})
.transpose()
}
/// The stored view's projection, for a test or a reader outside the store.
pub async fn stored_projection(
views: &DbPool,
run_id: RunId,
) -> Result<Option<RunProjection>, ProjectError> {
let json: Option<String> =
sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_optional(views)
.await
.map_err(ProjectError::Database)?;
json.map(|json| serde_json::from_str(&json).map_err(ProjectError::Encode))
.transpose()
}
/// The stream rows of a run: `(stream_seq, item_kind, item_id)`, in order.
pub async fn stored_stream(
views: &DbPool,
run_id: RunId,
) -> Result<Vec<(u64, String, String)>, ProjectError> {
let rows: Vec<(i64, String, String)> = sqlx::query_as(
"SELECT stream_seq, item_kind, item_id FROM petri_stream WHERE run_id = ? ORDER BY stream_seq",
)
.bind(run_id.to_string())
.fetch_all(views)
.await
.map_err(ProjectError::Database)?;
Ok(rows
.into_iter()
.map(|(seq, kind, id)| (u64::try_from(seq).unwrap_or(0), kind, id))
.collect())
}
/// Every stored platform record of a run, for a reader outside the store.
pub async fn stored_platform_records(
views: &DbPool,
run_id: RunId,
) -> Result<Vec<StoredPlatformRecord>, ProjectError> {
PlatformRecordStore::new(views.clone())
.read(&run_id)
.await
.map_err(ProjectError::Store)
}
/// A recorded event's projection is what `RunEvent` serializes to.
#[must_use]
pub fn event_json(event: &RunEvent) -> serde_json::Value {
serde_json::to_value(event).unwrap_or_default()
}

View file

@ -0,0 +1,886 @@
//! The projection of a Petri run: the view built live, as the run appends
//! its records, equals the view rebuilt from the records alone; a view that
//! missed its wake-ups catches up on the next signal; a crash between the
//! record commit and the view transaction is recovered by applying only the
//! missing suffix; two projectors over one store agree over nested child
//! executions; and a torn tail holds the view where it stands.
//!
//! Every run here takes its scope's environment through the sandbox-driver
//! host plugin, so the tests skip, and say why, when the executable is not
//! found, unless `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
#![expect(
clippy::disallowed_methods,
reason = "the tests locate the plugin executable through the process environment"
)]
#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")]
use std::collections::BTreeSet;
use std::env;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use fabro_db::DbPool;
use fabro_petri::SqliteRunStore;
use fabro_petri::projector::{self, Projector};
use fabro_store::platform_records::{
PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord,
};
use fabro_store::test_support;
use fabro_types::{
BlobHash, PetriAdmission, PetriGraphRef, RunEngine, RunId, RunStatus, StageHandler, StageId,
StageState, test_support as types_support,
};
use petri_execution::host::{self, HostRun};
use petri_frontend_fabro::Fabro;
use petri_runtime::executor::Retention;
use petri_runtime::frontend::CompileInputs;
use petri_runtime::ir::RunStatus as PetriRunStatus;
use petri_runtime::{RunOptions, Runtime};
use petri_store::{RunKey, RunStore};
use tokio::fs;
const HOST_PLUGIN: &str = "sandbox-driver-host";
const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN";
const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS";
const COMMAND_WORKFLOW: &str = r#"digraph Command {
graph [goal="Run one command"]
start [shape=Mdiamond]
exit [shape=Msquare]
say [shape=parallelogram, script="echo hello from petri"]
start -> say -> exit
}"#;
/// Two branches, each a command, joined by a fan-in.
const PARALLEL_WORKFLOW: &str = r#"digraph Parallel {
graph [goal="Run two branches"]
start [shape=Mdiamond]
exit [shape=Msquare]
fork [shape=component]
a [shape=parallelogram, script="echo a"]
b [shape=parallelogram, script="echo b"]
merge [shape=tripleoctagon]
report [shape=parallelogram, script="echo done"]
start -> fork
fork -> a
fork -> b
a -> merge
b -> merge
merge -> report -> exit
}"#;
const SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
fn host_plugin() -> Option<PathBuf> {
let found = env::var_os(HOST_PLUGIN_OVERRIDE)
.map(PathBuf::from)
.or_else(|| {
env::split_paths(&env::var_os("PATH")?)
.map(|dir| dir.join(HOST_PLUGIN))
.find(|candidate| candidate.is_file())
});
if found.is_none() {
assert!(
env::var_os(REQUIRE_ENV).is_none(),
"{REQUIRE_ENV} is set, but {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset"
);
eprintln!("skipping: {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset");
}
found
}
/// A fresh in-memory database with every table the projection touches.
fn pool() -> DbPool {
test_support::in_memory_pool_with(&[
fabro_db::BLOBS_MIGRATION_SQL,
fabro_db::RUNS_MIGRATION_SQL,
fabro_db::PETRI_RECORDS_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
])
}
fn hello_bundle() -> PathBuf {
Path::new(env!("CARGO_MANIFEST_DIR")).join("../../../.fabro/workflows/hello")
}
async fn install_bundle(root: &Path, name: &str, files: &[(&str, &str)]) -> PathBuf {
let bundle = root.join(".fabro").join("workflows").join(name);
fs::create_dir_all(&bundle)
.await
.expect("the bundle directory is creatable");
for (file, text) in files {
fs::write(bundle.join(file), text)
.await
.expect("the bundle file is writable");
}
bundle.join("workflow.fabro")
}
fn run_options(run_dir: &Path, run_id: RunId) -> RunOptions {
let mut options = RunOptions::new(run_dir);
options.grace = Duration::from_secs(2);
options.retention = Retention::Never;
options.echo = false;
options.run_key = Some(RunKey::new(run_id.to_string()));
options
}
/// The run's `run.created` platform record, as the create handler writes it,
/// and the `running` lifecycle record the execute path writes.
async fn create_run(pool: &DbPool, run_id: RunId, goal: &str) {
let store = PlatformRecordStore::new(pool.clone());
let mut spec = types_support::test_run_spec();
spec.run_id = run_id;
spec.engine = RunEngine::Petri(PetriAdmission {
graph: PetriGraphRef {
blob: BlobHash::new(b"graph"),
digest: "digest".to_string(),
},
children: Vec::new(),
});
store
.append(
&run_id,
&PlatformRecord::RunCreated(RunCreatedRecord {
spec,
title: Some(goal.to_string()),
parent_id: None,
retried_from: None,
web_url: None,
}),
None,
)
.await
.expect("the created record stores");
for (transition, status) in [
(RunLifecycleKind::Runnable, RunStatus::Runnable),
(RunLifecycleKind::Starting, RunStatus::Starting),
(RunLifecycleKind::Running, RunStatus::Running),
] {
store
.append(
&run_id,
&PlatformRecord::RunLifecycle(
RunLifecycleRecord::new(transition).with_status(status),
),
None,
)
.await
.expect("the lifecycle record stores");
}
}
/// Run `workflow` to completion on the real registry over `store`.
async fn run_workflow(
store: Arc<dyn RunStore>,
run_dir: &Path,
run_id: RunId,
workflow: &Path,
stubs: bool,
) {
let runtime = Runtime::standard().frontend(Fabro::new());
let runtime = if stubs {
petri_attractor_steps::register_stubs(runtime)
} else {
petri_attractor_steps::register(runtime)
};
let rt = runtime.store(store).options(run_options(run_dir, run_id));
let lowered = rt
.check(workflow, None, None, &CompileInputs::new())
.expect("the workflow file loads");
let graph = lowered
.graph
.unwrap_or_else(|| panic!("the workflow lowers: {:?}", lowered.diagnostics));
let host_run = HostRun::new(graph).with_children(lowered.children);
let report = host::run_configured(&rt, host_run, |_, _| {})
.await
.expect("the run completes");
assert_eq!(
report.status,
PetriRunStatus::Success,
"errors: {:?}",
report.state.errors()
);
}
/// A scenario: its bundle installed, its run created in the database.
struct Scenario {
pool: DbPool,
run_id: RunId,
workflow: PathBuf,
run_dir: PathBuf,
stubs: bool,
_root: tempfile::TempDir,
}
async fn scenario(name: &str, files: &[(&str, &str)], stubs: bool) -> Scenario {
let root = tempfile::tempdir().expect("a temp dir");
let workflow = install_bundle(root.path(), name, files).await;
let pool = pool();
let run_id = RunId::new();
create_run(&pool, run_id, name).await;
Scenario {
pool,
run_id,
workflow,
run_dir: root.path().join("run"),
stubs,
_root: root,
}
}
async fn hello_scenario() -> Scenario {
let bundle = hello_bundle();
let workflow = fs::read_to_string(bundle.join("workflow.fabro"))
.await
.expect("the hello workflow is checked in");
let settings = fs::read_to_string(bundle.join("workflow.toml"))
.await
.expect("the hello settings are checked in");
scenario(
"hello",
&[("workflow.fabro", &workflow), ("workflow.toml", &settings)],
true,
)
.await
}
async fn command_scenario() -> Scenario {
scenario(
"command",
&[
("workflow.fabro", COMMAND_WORKFLOW),
("workflow.toml", SETTINGS),
],
false,
)
.await
}
async fn parallel_scenario() -> Scenario {
scenario(
"parallel",
&[
("workflow.fabro", PARALLEL_WORKFLOW),
("workflow.toml", SETTINGS),
],
false,
)
.await
}
/// Run the scenario live: every append signals the projector, and the view
/// settles before the run is compared with its rebuild.
async fn run_live(scenario: &Scenario) -> Arc<Projector> {
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
projector.signal(scenario.run_id);
let store = projector.observe_store(Arc::new(SqliteRunStore::new(scenario.pool.clone())));
run_workflow(
store,
&scenario.run_dir,
scenario.run_id,
&scenario.workflow,
scenario.stubs,
)
.await;
projector.settle(scenario.run_id).await;
projector
}
/// Run the scenario with no projector attached: the records land and
/// nothing wakes the view.
async fn run_unobserved(scenario: &Scenario) {
run_workflow(
Arc::new(SqliteRunStore::new(scenario.pool.clone())),
&scenario.run_dir,
scenario.run_id,
&scenario.workflow,
scenario.stubs,
)
.await;
}
/// Every path where two JSON values differ, with both sides.
fn diff_json(path: &str, left: &serde_json::Value, right: &serde_json::Value) -> Vec<String> {
use serde_json::Value;
match (left, right) {
(Value::Object(left), Value::Object(right)) => {
let keys: BTreeSet<&String> = left.keys().chain(right.keys()).collect();
keys.into_iter()
.flat_map(|key| {
diff_json(
&format!("{path}/{key}"),
left.get(key).unwrap_or(&Value::Null),
right.get(key).unwrap_or(&Value::Null),
)
})
.collect()
}
(Value::Array(left), Value::Array(right)) if left.len() == right.len() => left
.iter()
.zip(right)
.enumerate()
.flat_map(|(index, (left, right))| diff_json(&format!("{path}[{index}]"), left, right))
.collect(),
_ if left == right => Vec::new(),
_ => vec![format!("{path}: live {left} != rebuilt {right}")],
}
}
fn json<T: serde::Serialize>(value: &T) -> serde_json::Value {
serde_json::to_value(value).expect("the value serializes")
}
/// The stored view equals the view rebuilt from the records alone: the
/// projection, the positions and the delivery sequence.
async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) {
let stored = projector::stored_projection(pool, run_id)
.await
.expect("the stored projection reads")
.expect("the run has a stored projection");
let (stored_positions, stored_stream_seq) = projector::stored_positions(pool, run_id)
.await
.expect("the positions read")
.expect("the run has positions");
let (rebuilt, positions, stream_seq) = projector::rebuild(pool, pool, run_id)
.await
.expect("the run rebuilds");
let rebuilt = rebuilt.expect("the rebuild has a projection");
let differences = diff_json("", &json(&stored), &json(&rebuilt));
assert!(
differences.is_empty(),
"live view differs from the rebuild at:\n{}",
differences.join("\n")
);
let mut stored_positions = stored_positions;
let mut positions = positions;
stored_positions.petri.sort();
positions.petri.sort();
assert_eq!(stored_positions, positions);
assert_eq!(stored_stream_seq, stream_seq);
let stream = projector::stored_stream(pool, run_id)
.await
.expect("the stream reads");
let seqs: Vec<u64> = stream.iter().map(|(seq, _, _)| *seq).collect();
assert_eq!(
seqs,
(1..=stream_seq).collect::<Vec<_>>(),
"contiguous stream"
);
}
async fn stage_states(pool: &DbPool, run_id: RunId) -> Vec<(String, StageState)> {
let stored = projector::stored_projection(pool, run_id)
.await
.expect("the stored projection reads")
.expect("the run has a stored projection");
stored
.iter_stages()
.map(|(id, stage)| (id.to_string(), stage.state))
.collect()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_hello_bundle_projects_live_as_it_rebuilds() {
if host_plugin().is_none() {
return;
}
let scenario = hello_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
assert!(
matches!(stored.status, RunStatus::Succeeded { .. }),
"{:?}",
stored.status
);
assert!(stored.conclusion.is_some(), "the run concluded");
let states = stage_states(&scenario.pool, scenario.run_id).await;
assert!(
states
.iter()
.any(|(label, state)| label.starts_with("start@") && *state == StageState::Succeeded),
"{states:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_command_workflow_projects_live_as_it_rebuilds() {
if host_plugin().is_none() {
return;
}
let scenario = command_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let say = stored
.stage(&StageId::new("say", 1))
.expect("the command stage is shown");
assert_eq!(say.state, StageState::Succeeded);
assert_eq!(say.handler, Some(StageHandler::Command));
assert!(
say.output
.as_deref()
.is_some_and(|output| output.contains("hello from petri")),
"{:?}",
say.output
);
assert!(say.timing.is_some());
let states = stage_states(&scenario.pool, scenario.run_id).await;
assert_eq!(states.len(), 3, "start, say, exit: {states:?}");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_parallel_workflow_projects_its_branches_as_child_executions() {
if host_plugin().is_none() {
return;
}
let scenario = parallel_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let fork = StageId::new("fork", 1);
for branch in ["a", "b"] {
let stage = stored
.stage(&StageId::new(branch, 1))
.unwrap_or_else(|| panic!("branch {branch} is a stage"));
assert_eq!(stage.state, StageState::Succeeded);
let branch_id = stage
.parallel_branch_id
.as_ref()
.unwrap_or_else(|| panic!("branch {branch} is grouped under the fork"));
assert_eq!(branch_id.group(), &fork);
}
let fork_stage = stored.stage(&fork).expect("the fork is a stage");
let results = fork_stage
.parallel_results
.as_ref()
.expect("the fork carries its branch results");
assert_eq!(results.len(), 2, "{results:?}");
let labels: Vec<String> = stored.iter_stages().map(|(id, _)| id.to_string()).collect();
assert!(
!labels.iter().any(|label| label.contains("fan_in")),
"synthetic nodes stay off the list: {labels:?}"
);
}
/// The projector is not signalled for any append; one signal at the end
/// folds everything.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn dropped_wake_ups_are_caught_up_by_the_next_signal() {
if host_plugin().is_none() {
return;
}
let scenario = command_scenario().await;
run_unobserved(&scenario).await;
assert!(
projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.is_none(),
"nothing woke the view"
);
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
projector.signal(scenario.run_id);
projector.settle(scenario.run_id).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let report = projector
.project_run(scenario.run_id)
.await
.expect("a pass over a caught-up view");
assert!(report.skipped, "nothing is left to fold: {report:?}");
assert!(report.health.complete, "{:?}", report.health.incomplete);
}
/// The same, through the startup pass.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_startup_pass_catches_up_a_view_nobody_signalled() {
if host_plugin().is_none() {
return;
}
let scenario = command_scenario().await;
run_unobserved(&scenario).await;
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
let report = projector
.startup_pass()
.await
.expect("the startup pass runs");
assert_eq!((report.runs, report.projected), (1, 1));
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let again = projector
.startup_pass()
.await
.expect("a second startup pass");
assert_eq!((again.runs, again.projected), (1, 0), "nothing left to do");
}
/// Passes over one run never interleave: the startup pass and a signalled
/// pass racing over the same run commit one stream, contiguous and without
/// a duplicate.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn concurrent_passes_over_one_run_commit_one_contiguous_stream() {
if host_plugin().is_none() {
return;
}
let scenario = parallel_scenario().await;
run_unobserved(&scenario).await;
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
let mut passes = Vec::new();
for _ in 0..4 {
let projector = Arc::clone(&projector);
let run_id = scenario.run_id;
passes.push(tokio::spawn(
async move { projector.project_run(run_id).await },
));
}
projector.signal(scenario.run_id);
projector
.startup_pass()
.await
.expect("the startup pass runs");
for pass in passes {
pass.await
.expect("the pass task joins")
.expect("a concurrent pass commits or skips");
}
projector.settle(scenario.run_id).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
}
/// Every Petri record of the run, as `(log, seq, recorded_at, record_json)`.
async fn petri_rows(pool: &DbPool, run_id: RunId) -> Vec<(String, i64, i64, String)> {
sqlx::query_as(
"SELECT log, seq, recorded_at, record_json FROM petri_records WHERE run_id = ? ORDER BY \
log, seq",
)
.bind(run_id.to_string())
.fetch_all(pool)
.await
.expect("the records read")
}
async fn insert_petri_row(pool: &DbPool, run_id: RunId, row: &(String, i64, i64, String)) {
sqlx::query(
"INSERT INTO petri_records (run_id, log, seq, recorded_at, record_json) VALUES (?, ?, ?, \
?, ?)",
)
.bind(run_id.to_string())
.bind(&row.0)
.bind(row.1)
.bind(row.2)
.bind(&row.3)
.execute(pool)
.await
.expect("the record inserts");
}
/// A copy of the run in a fresh database: its blobs, its Petri run row and
/// its platform records, but none of its Petri records yet.
async fn copy_run_without_records(source: &DbPool, run_id: RunId) -> DbPool {
let target = pool();
let blobs: Vec<(String, Vec<u8>)> = sqlx::query_as("SELECT hash, data FROM blobs")
.fetch_all(source)
.await
.expect("the blobs read");
for (hash, data) in blobs {
sqlx::query("INSERT INTO blobs (hash, data) VALUES (?, ?)")
.bind(hash)
.bind(data)
.execute(&target)
.await
.expect("the blob inserts");
}
sqlx::query("INSERT INTO petri_runs (run_id, created_at_ms, owner_id, acquired_at_ms) VALUES (?, 0, NULL, NULL)")
.bind(run_id.to_string())
.execute(&target)
.await
.expect("the run row inserts");
let platform: Vec<(i64, i64, String, String)> = sqlx::query_as(
"SELECT seq, recorded_at, kind, record_json FROM platform_records WHERE run_id = ? ORDER \
BY seq",
)
.bind(run_id.to_string())
.fetch_all(source)
.await
.expect("the platform records read");
for (seq, recorded_at, kind, record_json) in platform {
sqlx::query(
"INSERT INTO platform_records (run_id, seq, recorded_at, kind, record_json) VALUES \
(?, ?, ?, ?, ?)",
)
.bind(run_id.to_string())
.bind(seq)
.bind(recorded_at)
.bind(kind)
.bind(record_json)
.execute(&target)
.await
.expect("the platform record inserts");
}
target
}
/// The records commit in two halves and the view runs between them, then
/// the process dies before the view catches the second half: the rebuilt
/// view applies only the suffix, with the positions and the delivery
/// sequence continuing from where the committed view stood.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix() {
if host_plugin().is_none() {
return;
}
let scenario = parallel_scenario().await;
run_unobserved(&scenario).await;
let rows = petri_rows(&scenario.pool, scenario.run_id).await;
let replayed = copy_run_without_records(&scenario.pool, scenario.run_id).await;
// The first half of every log: a prefix per log, the coordinator log
// short of its finish.
let mut first: Vec<&(String, i64, i64, String)> = Vec::new();
let mut second: Vec<&(String, i64, i64, String)> = Vec::new();
for row in &rows {
let head = rows
.iter()
.filter(|other| other.0 == row.0)
.map(|other| other.1)
.max()
.expect("the log has a head");
if row.1 <= head / 2 {
first.push(row);
} else {
second.push(row);
}
}
for row in &first {
insert_petri_row(&replayed, scenario.run_id, row).await;
}
let before = Projector::new(replayed.clone(), replayed.clone());
let pass = before
.project_run(scenario.run_id)
.await
.expect("the first pass commits");
assert!(!pass.skipped);
assert!(!pass.health.complete, "the run has not finished");
let (positions_before, stream_before) = projector::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions");
assert_eq!(pass.stream_seq, stream_before);
let stream_rows_before = projector::stored_stream(&replayed, scenario.run_id)
.await
.expect("reads")
.len();
// The rest of the records commit; the view transaction never runs.
for row in &second {
insert_petri_row(&replayed, scenario.run_id, row).await;
}
before.fail_before_view();
let crashed = before.project_run(scenario.run_id).await;
assert!(
matches!(crashed, Err(projector::ProjectError::Injected)),
"{crashed:?}"
);
assert_eq!(
projector::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions"),
(positions_before.clone(), stream_before),
"the crash left the committed view alone"
);
// A new projector, as a restarted server builds one.
let after = Projector::new(replayed.clone(), replayed.clone());
let report = after.startup_pass().await.expect("the restart catches up");
assert_eq!((report.runs, report.projected), (1, 1));
let (positions_after, stream_after) = projector::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions");
let stream_rows_after = projector::stored_stream(&replayed, scenario.run_id)
.await
.expect("reads");
// Only the suffix was applied: the stream grew by the suffix's events,
// numbered on from the committed sequence, and every earlier row stayed.
assert_eq!(
stream_rows_after.len(),
stream_rows_before + usize::try_from(stream_after - stream_before).expect("a small count")
);
assert!(stream_after > stream_before);
assert_eq!(
stream_rows_after[stream_rows_before].0,
stream_before + 1,
"the suffix starts right after the committed sequence"
);
for held in &positions_before.petri {
let now = positions_after
.petri
.iter()
.find(|after| after.source == held.source)
.expect("a held log is still held");
assert!(now >= held, "{now:?} >= {held:?}");
}
assert_eq!(positions_after.platform_seq, positions_before.platform_seq);
assert_view_equals_rebuild(&replayed, scenario.run_id).await;
// And the copy agrees with the run projected in one go over the source.
let source = Projector::new(scenario.pool.clone(), scenario.pool.clone());
source.startup_pass().await.expect("the source projects");
let whole = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let pieced = projector::stored_projection(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("stored");
assert_eq!(json(&whole), json(&pieced));
}
/// Two projectors over one store, one after the other, over a run with
/// child executions: the second continues where the first stopped and both
/// agree with a projector that saw the run whole.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_restarted_projector_agrees_over_nested_child_executions() {
if host_plugin().is_none() {
return;
}
let scenario = parallel_scenario().await;
run_unobserved(&scenario).await;
let rows = petri_rows(&scenario.pool, scenario.run_id).await;
assert!(
rows.iter()
.filter(|row| row.0.starts_with("execution "))
.map(|row| &row.0)
.collect::<BTreeSet<_>>()
.len()
>= 3,
"the parallel run has child executions: {:?}",
rows.iter().map(|row| &row.0).collect::<BTreeSet<_>>()
);
let staged = copy_run_without_records(&scenario.pool, scenario.run_id).await;
let first = Projector::new(staged.clone(), staged.clone());
// The parent execution and the coordinator log up to the first child's
// declaration go in first; a restart then sees the children.
let (early, late): (Vec<_>, Vec<_>) = rows
.iter()
.partition(|row| row.0 == "execution 0" || (row.0 == "coordinator" && row.1 < 6));
for row in &early {
insert_petri_row(&staged, scenario.run_id, row).await;
}
first
.startup_pass()
.await
.expect("the first projector passes");
for row in &late {
insert_petri_row(&staged, scenario.run_id, row).await;
}
drop(first);
let second = Projector::new(staged.clone(), staged.clone());
second
.startup_pass()
.await
.expect("the second projector passes");
assert_view_equals_rebuild(&staged, scenario.run_id).await;
let whole = Projector::new(scenario.pool.clone(), scenario.pool.clone());
whole.startup_pass().await.expect("the source projects");
let one_go = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let restarted = projector::stored_projection(&staged, scenario.run_id)
.await
.expect("reads")
.expect("stored");
assert_eq!(json(&one_go), json(&restarted));
let states = stage_states(&staged, scenario.run_id).await;
assert!(
states.iter().any(|(label, _)| label == "a@1")
&& states.iter().any(|(label, _)| label == "b@1"),
"{states:?}"
);
}
/// A record at seq n+2 of an execution log, past a gap: Petri cannot read
/// the log, the view does not advance past what it held, and the run is
/// reported incomplete with the reason.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() {
if host_plugin().is_none() {
return;
}
let scenario = command_scenario().await;
run_unobserved(&scenario).await;
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
let clean = projector
.project_run(scenario.run_id)
.await
.expect("the clean pass commits");
assert!(clean.health.complete, "{:?}", clean.health.incomplete);
let (positions, stream_seq) = projector::stored_positions(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("positions");
let before = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
// A record two past the head of the execution log.
let rows = petri_rows(&scenario.pool, scenario.run_id).await;
let last = rows
.iter()
.filter(|row| row.0 == "execution 0")
.max_by_key(|row| row.1)
.expect("the execution log has records");
let mut torn: serde_json::Value = serde_json::from_str(&last.3).expect("the record is JSON");
torn["seq"] = serde_json::json!(last.1 + 2);
insert_petri_row(
&scenario.pool,
scenario.run_id,
&(last.0.clone(), last.1 + 2, last.2, torn.to_string()),
)
.await;
let held = projector
.project_run(scenario.run_id)
.await
.expect("the pass over the torn log still commits its health");
assert!(!held.health.complete, "the torn log is incomplete");
assert!(
!held.health.incomplete.is_empty(),
"the reason is reported: {:?}",
held.health
);
let (positions_after, stream_after) =
projector::stored_positions(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("positions");
assert_eq!(
positions_after, positions,
"the view did not advance past the tear"
);
assert_eq!(stream_after, stream_seq);
let after = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
assert_eq!(
json(&before),
json(&after),
"the projection stands where it was"
);
}

View file

@ -7,6 +7,7 @@ mod keyed_mutex;
mod keys;
mod legacy_blob_import;
mod legacy_run_history_import;
pub mod platform_records;
#[cfg(test)]
mod record;
mod run_session_record_store;
@ -43,9 +44,13 @@ pub use legacy_run_history_import::{
LegacyRunHistorySourceIdentity, LegacyRunHistorySourceIdentityError,
LegacyRunHistoryVerificationError, LegacyRunHistoryVerificationReport,
};
pub use platform_records::{
PlatformRecord, PlatformRecordHook, PlatformRecordKind, PlatformRecordStore, StagePosition,
StoredPlatformRecord,
};
pub use run_session_record_store::{RunSessionRecordStore, StoredSessionRecord};
pub use run_sessions::{ProjectedRunSession, project_run_session, project_run_sessions};
pub use run_state::RunProjectionReducer;
pub use run_state::{RunProjectionReducer, build_summary, projected_usage};
pub use run_summary_store::{
RunSummaryIdentity, RunSummaryListQuery, RunSummaryPage, RunSummarySort,
RunSummarySortDirection, RunSummaryStore, RunSummaryVisibility,

File diff suppressed because it is too large Load diff

View file

@ -1145,7 +1145,10 @@ fn stage_at_completed_visit<'a>(
Some(state.stage_entry(node_id, visit, first_event_seq(seq)))
}
pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> Run {
/// The run summary (`Run`) a projection stands for: what the run list, the
/// board and the scheduler read.
#[must_use]
pub fn build_summary(state: &RunProjection, run_id: &RunId) -> Run {
let goal = state.spec.graph.goal().to_string();
let diff_summary = state
.conclusion
@ -1241,7 +1244,8 @@ pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> Run {
/// The run's usage: the conclusion's total once the run ended, else the sum
/// of every non-boundary stage's usage so far.
pub(crate) fn projected_usage(state: &RunProjection) -> Usage {
#[must_use]
pub fn projected_usage(state: &RunProjection) -> Usage {
if let Some(usage) = state
.conclusion
.as_ref()

View file

@ -1,5 +1,5 @@
use std::fmt::Write as _;
use std::sync::LazyLock;
use std::sync::{Arc, LazyLock, PoisonError, RwLock};
use chrono::{DateTime, Utc};
use fabro_types::{
@ -12,8 +12,9 @@ use sqlx::sqlite::{SqliteArguments, SqliteConnection, SqliteRow};
use sqlx::{Connection as _, QueryBuilder, Row as _, Sqlite, SqlitePool, Transaction};
use strum::VariantArray as _;
use crate::platform_records::{self, PlatformRecordHook, PlatformRecordStore};
use crate::run_state::{ProjectedRun, build_summary, projected_usage};
use crate::{Error, EventPayload, Result, keys};
use crate::{Error, EventPayload, Result, RunProjection, keys};
const INSERT_RUN_SQL: &str = r"
INSERT INTO runs (
@ -65,6 +66,38 @@ ON CONFLICT(id) DO UPDATE SET
WHERE excluded.source_last_seq > runs.source_last_seq
";
/// The `runs` row of a Petri run, written by its projector: every column the
/// list views and the scheduler read, and never `source_last_seq`, which the
/// legacy event path owns while it still writes the row.
const UPSERT_PETRI_RUN_SQL: &str = r"
INSERT INTO runs (
id, source_last_seq, created_at_ms, started_at_ms, last_event_at_ms, completed_at_ms,
status, archived_at_ms, parent_id, title, workflow_slug, workflow_name,
repository_name, automation_id, diff_files_changed, diff_additions, diff_deletions,
input_tokens, output_tokens, reasoning_tokens, cache_read_tokens, cache_write_tokens,
total_usd_micros, summary_json
) VALUES (
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
)
ON CONFLICT(id) DO UPDATE SET
created_at_ms = excluded.created_at_ms,
started_at_ms = excluded.started_at_ms,
last_event_at_ms = excluded.last_event_at_ms,
completed_at_ms = excluded.completed_at_ms,
status = excluded.status,
archived_at_ms = excluded.archived_at_ms,
parent_id = excluded.parent_id,
title = excluded.title,
workflow_slug = excluded.workflow_slug,
workflow_name = excluded.workflow_name,
repository_name = excluded.repository_name,
automation_id = excluded.automation_id,
diff_additions = excluded.diff_additions,
diff_deletions = excluded.diff_deletions,
total_usd_micros = excluded.total_usd_micros,
summary_json = excluded.summary_json
";
const UPDATE_RUN_SQL: &str = r"
UPDATE runs SET
source_last_seq = ?,
@ -184,7 +217,10 @@ pub struct RunSummaryPage {
#[derive(Clone)]
pub struct RunSummaryStore {
pool: SqlitePool,
pool: SqlitePool,
/// Called after a platform record for a Petri run is committed beside
/// its legacy event: the projector's wake-up.
platform_hook: Arc<RwLock<Option<PlatformRecordHook>>>,
}
impl std::fmt::Debug for RunSummaryStore {
@ -196,7 +232,81 @@ impl std::fmt::Debug for RunSummaryStore {
impl RunSummaryStore {
#[must_use]
pub fn new(pool: SqlitePool) -> Self {
Self { pool }
Self {
pool,
platform_hook: Arc::new(RwLock::new(None)),
}
}
/// The pool this store's tables live in: the `runs` row, the run events,
/// the platform records and the Petri projection tables. The server's
/// one database; a test fixture's own.
#[must_use]
pub fn pool(&self) -> SqlitePool {
self.pool.clone()
}
/// The platform records over the same pool.
#[must_use]
pub fn platform_records(&self) -> PlatformRecordStore {
PlatformRecordStore::new(self.pool.clone())
}
/// Install the wake-up called after a platform record of a Petri run is
/// committed beside its legacy event.
pub fn set_platform_record_hook(&self, hook: PlatformRecordHook) {
*self
.platform_hook
.write()
.unwrap_or_else(PoisonError::into_inner) = Some(hook);
}
pub(crate) fn notify_platform_record(&self, run_id: RunId) {
let hook = self
.platform_hook
.read()
.unwrap_or_else(PoisonError::into_inner)
.clone();
if let Some(hook) = hook {
hook(run_id);
}
}
/// The stored projection of a Petri run, as the run's projector last
/// committed it, or `None` when no view pass has run for it yet.
pub async fn load_petri_projection(
&self,
run_id: &RunId,
) -> Result<Option<Arc<RunProjection>>> {
let json: Option<String> =
sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_optional(&self.pool)
.await?;
json.map(|json| Ok(Arc::new(serde_json::from_str(&json)?)))
.transpose()
}
/// Write the `runs` row of a Petri run from its projection, on a
/// connection the caller holds a transaction on: the columns the list
/// views and the scheduler read, and the summary JSON. The legacy
/// concurrency guard `source_last_seq` is left as the legacy path set it
/// (or `1` when this write creates the row), so both writers keep
/// working until the legacy events go.
pub async fn write_petri_run_row_on_connection(
connection: &mut SqliteConnection,
run_id: &RunId,
projection: &RunProjection,
) -> Result<()> {
let entry = ProjectedRun::new(*run_id, Arc::new(projection.clone()), 1);
let record = PreparedRunSummary::from_entry(&entry);
bind_run_columns(
sqlx::query(UPSERT_PETRI_RUN_SQL).bind(run_id.to_string()),
&record,
)?
.execute(connection)
.await?;
Ok(())
}
#[cfg(test)]
@ -657,6 +767,7 @@ impl RunSummaryStore {
insert_run_on_connection(connection, &record).await?;
insert_event_on_connection(connection, &record, payload, &envelope).await?;
insert_platform_record_on_connection(connection, entry, &envelope).await?;
Ok(envelope)
}
@ -674,6 +785,7 @@ impl RunSummaryStore {
update_run_on_connection(connection, &record, expected_last_seq).await?;
insert_event_on_connection(connection, &record, payload, &envelope).await?;
insert_platform_record_on_connection(connection, entry, &envelope).await?;
Ok(envelope)
}
@ -798,6 +910,13 @@ WHERE id = ?
let diff = run.diff.unwrap_or_default();
verify_run_field(&row, run, "id", &run.id.to_string())?;
verify_run_field(&row, run, "source_last_seq", &i64::from(record.last_seq))?;
if entry.projection.spec.engine.is_petri() {
// A Petri run's row is written by its projector from Petri's
// records and the platform records; the legacy fold knows the
// lifecycle alone, so only the identity and the legacy guard
// are checked here.
return Ok(());
}
verify_run_field(
&row,
run,
@ -1113,6 +1232,43 @@ async fn insert_event_json_on_connection(
Ok(())
}
/// For a Petri run, the platform record the legacy event stands for, stored
/// in the event's transaction so the projection over Petri's records reads
/// the lifecycle from platform records alone. Whether one was written is
/// what [`platform_record_written`] answers after the commit.
async fn insert_platform_record_on_connection(
connection: &mut SqliteConnection,
entry: &ProjectedRun,
envelope: &EventEnvelope,
) -> Result<()> {
let Some(record) = platform_record_written(entry, envelope) else {
return Ok(());
};
let recorded_at = u64::try_from(envelope.event.ts.timestamp_millis()).unwrap_or(0);
PlatformRecordStore::append_on_connection(
connection,
&entry.run_id,
recorded_at,
&record,
None,
)
.await?;
Ok(())
}
/// The platform record a committed legacy event of a Petri run produced,
/// if any: the same derivation the insert makes, for the caller that
/// notifies after the commit.
pub(crate) fn platform_record_written(
entry: &ProjectedRun,
envelope: &EventEnvelope,
) -> Option<platform_records::PlatformRecord> {
if !entry.projection.spec.engine.is_petri() {
return None;
}
platform_records::platform_record_for(&envelope.event)
}
fn sql_limit(limit: usize) -> i64 {
i64::try_from(limit.saturating_add(1)).unwrap_or(i64::MAX)
}

View file

@ -13,7 +13,9 @@ use run_store::RunDatabaseInner;
use slatedb::config::{CompressionCodec, Settings};
use tokio::sync::{Mutex, MutexGuard, OnceCell};
use crate::{BlobStore, Error, EventPayload, Result, RunProjection, RunSummaryStore, keys};
use crate::{
BlobStore, Error, EventPayload, Result, RunProjection, RunSummaryStore, keys, run_summary_store,
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UnreadableRun {
@ -128,9 +130,13 @@ impl Database {
) -> Result<RunDatabase> {
let (mut active_runs, run_store) = self.reserve_new_run(run_id).await?;
let (envelope, projected) = run_store.commit_first_event(payload).await?;
let platform_record = run_summary_store::platform_record_written(&projected, &envelope);
run_store.install_in_memory_state(projected);
Self::cache_active_run(&mut active_runs, &run_store);
run_store.publish(&envelope);
if platform_record.is_some() {
self.run_summary_store.notify_platform_record(*run_id);
}
Ok(run_store)
}
@ -232,15 +238,32 @@ impl Database {
Ok(())
}
/// The run's projection: for a legacy run the reducer's fold of its
/// events; for a Petri run the projection its projector last committed
/// over Petri's records and the platform records, falling back to the
/// legacy fold (the lifecycle alone) until the first view pass commits.
pub async fn load_run_projection(&self, run_id: &RunId) -> Result<Option<Arc<RunProjection>>> {
if let Some(active) = self.get_active_run(run_id).await {
return active.projection_snapshot().await.map(Some);
}
match self.run_summary_store.load_projection(run_id).await {
Ok(projected) => Ok(Some(projected.projection)),
Err(Error::RunNotFound(_)) => Ok(None),
Err(error) => Err(error),
let legacy = if let Some(active) = self.get_active_run(run_id).await {
active.projection_snapshot().await?
} else {
match self.run_summary_store.load_projection(run_id).await {
Ok(projected) => projected.projection,
Err(Error::RunNotFound(_)) => return Ok(None),
Err(error) => return Err(error),
}
};
if legacy.spec.engine.is_petri() {
if let Some(petri) = self.run_summary_store.load_petri_projection(run_id).await? {
return Ok(Some(petri));
}
}
Ok(Some(legacy))
}
/// Install the wake-up called after a platform record of a Petri run is
/// committed beside its legacy event.
pub fn set_platform_record_hook(&self, hook: crate::PlatformRecordHook) {
self.run_summary_store.set_platform_record_hook(hook);
}
/// Resolves the run that owns `session_id` from the canonical typed

View file

@ -220,8 +220,14 @@ impl RunDatabase {
let (envelope, projected) = self.commit_event_locked(payload, event).await?;
// Keep post-commit propagation await-free: cancellation after SQLite
// commits must not leave in-memory state stale or omit the broadcast.
let platform_record = run_summary_store::platform_record_written(&projected, &envelope);
self.install_in_memory_state(projected);
self.publish(&envelope);
if platform_record.is_some() {
self.inner
.run_summary_store
.notify_platform_record(self.inner.run_id);
}
Ok(envelope)
}
@ -504,6 +510,86 @@ mod tests {
);
}
/// A Petri run's legacy lifecycle events leave platform records beside
/// them, in the same commit, and the hook fires after each; a legacy
/// run's events leave none.
#[tokio::test]
async fn a_petri_runs_lifecycle_events_become_platform_records_and_wake_the_hook() {
use std::sync::atomic::{AtomicUsize, Ordering};
use crate::platform_records::{PlatformRecord, PlatformRecordKind, RunLifecycleKind};
let store = store();
let woken = Arc::new(AtomicUsize::new(0));
let counter = Arc::clone(&woken);
store.set_platform_record_hook(Arc::new(move |_| {
counter.fetch_add(1, Ordering::SeqCst);
}));
let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65E".parse().unwrap();
let mut created = run_created_payload(&run_id);
let mut petri = serde_json::to_value(&created).unwrap();
petri["properties"]["engine"] = json!({
"kind": "petri",
"graph": { "blob": fabro_types::BlobHash::new(b"graph").to_string(), "digest": "d" },
});
created = EventPayload::new(petri, &run_id).unwrap();
let run = store
.create_run_with_first_event(&run_id, &created)
.await
.unwrap();
run.append_event(
&EventPayload::new(
json!({
"id": "evt-starting",
"ts": "2026-04-09T12:00:00Z",
"run_id": run_id.to_string(),
"event": "run.start_requested",
"properties": { "resume": false },
}),
&run_id,
)
.unwrap(),
)
.await
.unwrap();
run.append_event(&stage_payload(&run_id, 3)).await.unwrap();
let records = store
.run_summary_store()
.platform_records()
.read(&run_id)
.await
.unwrap();
let kinds: Vec<PlatformRecordKind> = records.iter().map(|r| r.record.kind()).collect();
assert_eq!(kinds, vec![
PlatformRecordKind::RunCreated,
PlatformRecordKind::RunLifecycle
]);
let PlatformRecord::RunLifecycle(lifecycle) = &records[1].record else {
panic!("the second record is the lifecycle");
};
assert_eq!(lifecycle.transition, RunLifecycleKind::StartRequested);
assert_eq!(lifecycle.source.as_deref(), Some("start"));
assert_eq!(woken.load(Ordering::SeqCst), 2, "one wake-up per record");
let legacy_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65F".parse().unwrap();
store
.create_run_with_first_event(&legacy_id, &run_created_payload(&legacy_id))
.await
.unwrap();
assert!(
store
.run_summary_store()
.platform_records()
.read(&legacy_id)
.await
.unwrap()
.is_empty(),
"a legacy run leaves no platform records"
);
assert_eq!(woken.load(Ordering::SeqCst), 2);
}
#[tokio::test]
async fn watcher_catches_up_from_sql_without_duplicates() {
let store = store();

View file

@ -34,6 +34,7 @@ pub fn test_run_summary_store() -> Arc<RunSummaryStore> {
fabro_db::RUN_EVENTS_MIGRATION_SQL,
fabro_db::RUN_HISTORY_ACTIVATION_MIGRATION_SQL,
fabro_db::RUN_EVENT_SESSION_OWNER_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
])))
}
@ -115,6 +116,7 @@ pub fn test_run_summary_store_at(store_dir: &Path) -> Arc<RunSummaryStore> {
fabro_db::RUN_EVENTS_MIGRATION_SQL,
fabro_db::RUN_HISTORY_ACTIVATION_MIGRATION_SQL,
fabro_db::RUN_EVENT_SESSION_OWNER_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
],
)))
}

View file

@ -0,0 +1,60 @@
-- Fabro's own facts about a Petri run, beside Petri's records.
--
-- `platform_records` holds every fact Fabro records about a run that Petri
-- does not: the lifecycle before and after the engine, a checkpoint commit,
-- a pull request, a notification, a pairing. `seq` is per run and assigned
-- by the store; `kind` is the record's kind and `record_json` the typed
-- record with its kind tag; `execution` and `firing` name the Petri stage a
-- record belongs to, when it belongs to one.
CREATE TABLE platform_records (
run_id TEXT NOT NULL,
seq INTEGER NOT NULL,
recorded_at INTEGER NOT NULL,
kind TEXT NOT NULL,
record_json TEXT NOT NULL,
execution INTEGER NULL,
firing INTEGER NULL,
PRIMARY KEY (run_id, seq),
CHECK (seq >= 1),
CHECK (json_valid(record_json))
);
CREATE INDEX platform_records_by_kind
ON platform_records(run_id, kind, seq);
-- The projection of a Petri run: the view document Fabro's read side serves,
-- derived from the run's Petri records and platform records, rewritten in
-- one transaction per view pass together with the positions it covers.
-- `projection_json` is the `RunProjection`; `fold_json` is the projector's
-- own bookkeeping; `positions_json` is the last event consumed per Petri
-- log and the last platform record consumed; `stream_seq` is the last
-- delivery sequence assigned to `petri_stream`.
CREATE TABLE petri_projection (
run_id TEXT PRIMARY KEY NOT NULL,
projection_json TEXT NOT NULL,
fold_json TEXT NOT NULL,
positions_json TEXT NOT NULL,
stream_seq INTEGER NOT NULL,
updated_at_ms INTEGER NOT NULL,
CHECK (stream_seq >= 0),
CHECK (json_valid(projection_json)),
CHECK (json_valid(fold_json)),
CHECK (json_valid(positions_json))
);
-- One ordered stream per run of everything the projection consumed: each
-- Petri event and each platform record, in the order the view committed
-- them. `stream_seq` is the cursor a client resumes from; `item_kind` and
-- `item_id` are the item's own identity (a Petri event id as
-- `<log>/<seq>/<index>`, or a platform record's `seq`), for deduplication.
CREATE TABLE petri_stream (
run_id TEXT NOT NULL,
stream_seq INTEGER NOT NULL,
item_kind TEXT NOT NULL,
item_id TEXT NOT NULL,
event_json TEXT NOT NULL,
PRIMARY KEY (run_id, stream_seq),
CHECK (stream_seq >= 1),
CHECK (item_kind IN ('petri', 'platform')),
CHECK (json_valid(event_json))
);

View file

@ -47,6 +47,12 @@ pub const RUN_SESSION_RECORDS_MIGRATION_SQL: &str =
pub const PETRI_RECORDS_MIGRATION_SQL: &str =
include_str!("../migrations/2026091701_petri_records.sql");
/// The Petri projection migration (`platform_records`, `petri_projection`,
/// `petri_stream`), exposed so fixtures in other crates can install the
/// production schema without a filesystem path into this crate.
pub const PETRI_PROJECTION_MIGRATION_SQL: &str =
include_str!("../migrations/2026091801_petri_projection.sql");
/// The temporary run-history activation migration, exposed so fixtures in
/// other crates can install the production compatibility schema.
pub const RUN_HISTORY_ACTIVATION_MIGRATION_SQL: &str =