diff --git a/Cargo.lock b/Cargo.lock index 670f22e8b..4c5223449 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2891,6 +2891,7 @@ dependencies = [ "anyhow", "async-trait", "bytes", + "chrono", "fabro-api", "fabro-auth", "fabro-client", @@ -2899,6 +2900,7 @@ dependencies = [ "fabro-llm", "fabro-store", "fabro-types", + "fabro-util", "lithos-llm", "petri-attractor-steps", "petri-execution", diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 1933209e3..60b07f507 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -25,6 +25,7 @@ fabro-db = { path = "../../foundation/fabro-db" } fabro-http.workspace = true fabro-store = { path = "../fabro-store" } fabro-types = { path = "../../foundation/fabro-types" } +fabro-util = { path = "../../foundation/fabro-util" } petri_runtime.workspace = true petri_execution.workspace = true petri_store.workspace = true @@ -36,6 +37,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 @@ -49,5 +51,6 @@ tracing.workspace = true 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-types = { path = "../../foundation/fabro-types", features = ["test-support"] } petri_testkit.workspace = true tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 9b8134b11..5c71575c9 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -46,8 +46,39 @@ 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. +- `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, interviews over - Fabro's API, secrets, output storage, run tools, the event projection. + Fabro's API, secrets, output storage, 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 @@ -79,6 +110,16 @@ Integration tests live under `tests/`: operator release, lease exclusivity, a crash between appends, and blob interoperation with Fabro's `BlobStore`. +- `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 @@ -91,10 +132,12 @@ 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, and Petri's diagnostics refuse a run at create. The +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, 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. diff --git a/lib/components/fabro-petri/VIEWS.md b/lib/components/fabro-petri/VIEWS.md index 0954969c7..2df4e3c4f 100644 --- a/lib/components/fabro-petri/VIEWS.md +++ b/lib/components/fabro-petri/VIEWS.md @@ -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 `.` name from diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index 443032ac8..0d13b9d5f 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -23,8 +23,11 @@ //! - [`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; +//! - [`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, interviews over Fabro's API, -//! secrets, output storage, the run tools, the event projection. +//! secrets, output storage, the run tools. //! //! The Petri packages are pinned by revision in the workspace `Cargo.toml` //! under `petri_*` keys. @@ -35,6 +38,8 @@ pub mod engine; pub mod http_store; pub mod interviewer; pub mod petri; +pub mod projection; +pub mod projector; pub mod run_store; pub mod runtime; #[cfg(feature = "test-support")] diff --git a/lib/components/fabro-petri/src/projection.rs b/lib/components/fabro-petri/src/projection.rs new file mode 100644 index 000000000..e98156021 --- /dev/null +++ b/lib/components/fabro-petri/src/projection.rs @@ -0,0 +1,1375 @@ +//! The projection of a Petri run: Petri's public events and Fabro's platform +//! records folded into the view Fabro's read side serves. +//! +//! The fold is pure. [`RunView`] holds the [`RunProjection`] the API serves +//! (`GET /runs/{id}/state`, the run list through its summary) and the +//! bookkeeping the fold needs between items ([`FoldState`]): which Petri +//! firing each stage is, which invocation each execution belongs to and +//! whether it is a parallel branch, which stage asked each open question. +//! Both halves are stored by the projector and reloaded for the next pass, +//! so a pass folds only the items past the committed positions. +//! +//! The mapping follows `VIEWS.md`, row by row. The stage key is `(execution, +//! firing)`; Fabro's `StageId` (`node@visit`) is the display label the +//! `RunProjection` keys stages by, and a label two firings would share (two +//! child invocations with the same node name and visit) is made unique by +//! naming the execution. What the matrix leaves default is left default +//! here and named in the crate's README. +//! +//! Every item the fold sees carries the delivery sequence the projector +//! assigned it (`stream_seq`), which a checkpoint keeps as its `seq`. A +//! stage's `first_event_seq`, the key the stage list sorts by, is not the +//! delivery sequence: two logs' records can be committed in an order that +//! differs from their recording times by a few positions, and the view +//! built live must equal the view rebuilt from the records alone. It is the +//! milliseconds from the run's creation to the stage's `visit.started`, +//! plus one, which is the same however the records were delivered. + +use std::collections::{BTreeMap, BTreeSet, HashMap}; + +use chrono::{DateTime, TimeZone as _, Utc}; +use fabro_store::platform_records::{ + PlatformRecord, RunLifecycleKind, RunLifecycleRecord, StoredPlatformRecord, +}; +use fabro_types::settings::run::RunEnvironmentSettings; +use fabro_types::{ + BlockedReason, CheckpointRecord as ViewCheckpoint, CodingAgentEvent, CodingEvent, Conclusion, + FailureCategory, FailureDetail, FailureReason, InterviewOption, InterviewQuestionRecord, + ModelRef, ModelUsage, ParallelBranchId, ParallelBranchResult, PendingInterviewRecord, + PullRequestLink, QuestionType, RunApproval, RunApprovalState, RunControlAction, RunDiff, + RunFailure, RunId, RunProjection, RunSandbox, RunSandboxPlan, RunStatus, RunTiming, + SandboxProviderKind, StageCompletion, StageHandler, StageId, StageInferenceProjection, + StageModelUsage, StageOutcome, StageProjection, StageState, StageTiming, StartRecord, + SuccessReason, first_event_seq, timing, usage_rollup, +}; +use lithos_llm::catalog::{ModelId, ProviderId}; +use lithos_llm::types::Usage; +use petri_execution::events::{Derived, Parsed, RunEvent, Subject, ViewEvent, WaitState}; +use petri_execution::{CoordinatorEvent, ExecutionId}; +use petri_runtime::engine::{Admission, Event}; +use petri_runtime::ir::{Metrics, Status, StepEvent}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use tracing::debug; + +/// One item the projector hands the fold, with its delivery sequence. +pub enum Item<'a> { + Petri(&'a RunEvent), + Platform(&'a StoredPlatformRecord), +} + +/// A stage as the fold knows it: its label in the projection, and what it +/// learned about it. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct StageRef { + pub stage_id: StageId, + /// Whether the stage is a logical one the projection shows, or a + /// lowering node it keeps off the list. + pub shown: bool, + /// The node's instance name and visit, for the collision rule. + pub node_name: String, + pub visit: u32, +} + +/// What the fold knows about one invocation. +#[derive(Clone, Debug, Default, Serialize, Deserialize)] +pub struct InvocationRef { + /// The calling execution and firing, for a nested invocation. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub parent: Option<(u64, u64)>, + /// The parallel group and branch index, for a branch child. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub branch: Option<(StageId, u32)>, + /// The result the invocation recorded, for the root. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub failure: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output: Option, +} + +/// Whether the run's durable record is whole, as the projector last read it. +#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct RecordHealth { + pub complete: bool, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub incomplete: Vec, +} + +/// The fold's bookkeeping between items. +#[derive(Clone, Debug, Default, Serialize, Deserialize)] +pub struct FoldState { + /// Stages by `":"`. + #[serde(default)] + pub stages: BTreeMap, + /// Labels taken, so a second firing with the same name and visit gets + /// its own. + #[serde(default)] + pub labels: BTreeSet, + #[serde(default)] + pub invocations: BTreeMap, + /// Which invocation each execution belongs to. + #[serde(default)] + pub executions: BTreeMap, + /// Open questions by id: the stage that asked. + #[serde(default)] + pub questions: BTreeMap, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub root: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub started_at: Option, + /// The run's recorded finish, when Petri recorded one. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub finished: Option, + /// The run branch and base sha, when they arrive before `run.started`. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub run_branch: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub base_sha: Option, + #[serde(default)] + pub checkpoints: u32, + #[serde(default)] + pub health: RecordHealth, +} + +/// The view of one run: what the API serves and what the fold keeps. +#[derive(Clone, Debug)] +pub struct RunView { + pub projection: Option, + pub state: FoldState, +} + +impl RunView { + #[must_use] + pub fn new() -> Self { + Self { + projection: None, + state: FoldState::default(), + } + } + + /// Fold one item at its delivery sequence. + pub fn fold(&mut self, item: &Item<'_>, stream_seq: u64) { + match item { + Item::Platform(record) => self.fold_platform(record, stream_seq), + Item::Petri(event) => self.fold_petri(event), + } + } + + /// The run's projection, once its `run.created` record was folded. + #[must_use] + pub fn projection(&self) -> Option<&RunProjection> { + self.projection.as_ref() + } + + // ── Platform records ──────────────────────────────────────────────── + + fn fold_platform(&mut self, stored: &StoredPlatformRecord, stream_seq: u64) { + let at = millis(stored.recorded_at); + if let PlatformRecord::RunCreated(created) = &stored.record { + let title = created + .title + .clone() + .unwrap_or_else(|| fabro_types::infer_run_title(created.spec.graph.goal())); + let mut projection = RunProjection::new(title, created.spec.clone(), at); + projection.parent_id = created.parent_id; + projection.retried_from = created.retried_from; + projection.web_url.clone_from(&created.web_url); + projection.sandbox = Some(RunSandbox::planned(sandbox_plan( + &projection.spec.settings.run.environment, + ))); + self.projection = Some(projection); + return; + } + let Some(projection) = self.projection.as_mut() else { + debug!( + seq = stored.seq, + kind = %stored.record.kind(), + "platform record before run.created; not folded" + ); + return; + }; + touch(projection, at); + match &stored.record { + PlatformRecord::RunLifecycle(record) => fold_lifecycle(projection, record, at), + PlatformRecord::RunTitle(record) => projection.title.clone_from(&record.title), + PlatformRecord::RunParent(record) => projection.parent_id = record.parent_id, + PlatformRecord::RunArchived => projection.archived_at = Some(at), + PlatformRecord::RunUnarchived => projection.archived_at = None, + PlatformRecord::RunSuperseded(record) => { + projection.superseded_by = Some(record.new_run_id); + } + PlatformRecord::RunCreated(_) + | PlatformRecord::RunNotice(_) + | PlatformRecord::InterviewAnswered(_) + | PlatformRecord::NotificationSent(_) + | PlatformRecord::RunPaired(_) => {} + PlatformRecord::RunBranch(record) => { + self.state.run_branch.clone_from(&record.run_branch); + self.state.base_sha.clone_from(&record.base_sha); + if let Some(start) = projection.start.as_mut() { + start.run_branch.clone_from(&record.run_branch); + start.base_sha.clone_from(&record.base_sha); + } + } + PlatformRecord::GitIdentity(record) => { + projection.git_identity = Some(record.identity.clone()); + } + PlatformRecord::Checkpoint(record) => { + self.state.checkpoints = self.state.checkpoints.saturating_add(1); + let stage = self + .state + .stages + .get(&stage_key(record.execution, record.firing)); + let current_node = stage.map_or_else(String::new, |stage| stage.node_name.clone()); + let checkpoint = fabro_types::Checkpoint { + timestamp: at, + current_node: current_node.clone(), + completed_nodes: Vec::new(), + node_retries: HashMap::default(), + context_values: HashMap::default(), + node_outcomes: HashMap::default(), + next_node_id: None, + git_commit_sha: record.git_commit_sha.clone(), + loop_failure_signatures: HashMap::default(), + restart_failure_signatures: HashMap::default(), + node_visits: HashMap::default(), + }; + projection.checkpoints.push(ViewCheckpoint { + seq: u32::try_from(stream_seq).unwrap_or(u32::MAX), + checkpoint, + diff: RunDiff { + patch: None, + summary: record.diff_summary, + }, + }); + } + PlatformRecord::PullRequestCreated(record) => { + projection.pull_request = Some(PullRequestLink { + owner: record.owner.clone(), + repo: record.repo.clone(), + number: record.number, + }); + } + } + } + + // ── Petri events ──────────────────────────────────────────────────── + + fn fold_petri(&mut self, event: &RunEvent) { + let at = millis(event.recorded_at); + if let Some(record) = event.coordinator() { + self.fold_coordinator(record, event, at); + } else if let Some(engine) = event.engine() { + self.fold_engine(engine, event, at); + } else if let Some(view) = event.view() { + self.fold_view(view, event, at); + } + if let Some(projection) = self.projection.as_mut() { + touch(projection, at); + } + } + + fn fold_coordinator(&mut self, record: &CoordinatorEvent, event: &RunEvent, at: DateTime) { + match record { + CoordinatorEvent::RunStarted { root, .. } => { + self.state.root = Some(root.raw()); + self.state.started_at = Some(event.recorded_at); + if let Some(projection) = self.projection.as_mut() { + apply_status(projection, RunStatus::Running, at); + projection.start = Some(StartRecord { + start_time: at, + run_branch: self.state.run_branch.clone(), + base_sha: self.state.base_sha.clone(), + }); + } + } + CoordinatorEvent::InvocationDeclared { invocation, .. } => { + let mut info = InvocationRef::default(); + if let Some(parent) = &event.context.parent { + info.parent = Some((parent.execution.raw(), parent.firing.raw())); + if let Some((fork_firing, index)) = branch_slot(&parent.slot) { + let group = self + .state + .stages + .get(&stage_key(parent.execution.raw(), fork_firing)) + .map(|stage| stage.stage_id.clone()); + if let Some(group) = group { + info.branch = Some((group, index)); + } + } + } + self.state.invocations.insert(invocation.raw(), info); + } + CoordinatorEvent::ExecutionDeclared { + execution, + invocation, + .. + } => { + self.state + .executions + .insert(execution.raw(), invocation.raw()); + } + CoordinatorEvent::InvocationFinished { invocation, result } => { + let info = self.state.invocations.entry(invocation.raw()).or_default(); + info.failure = result + .failure + .as_ref() + .map(|failure| failure.message.clone()); + info.output = Some(result.output.clone()); + } + CoordinatorEvent::RunPaused => { + if let Some(projection) = self.projection.as_mut() { + let prior_block = match projection.status { + RunStatus::Blocked { blocked_reason } => Some(blocked_reason), + _ => None, + }; + apply_status(projection, RunStatus::Paused { prior_block }, at); + if projection.pending_control == Some(RunControlAction::Pause) { + projection.pending_control = None; + } + } + } + CoordinatorEvent::RunUnpaused => { + if let Some(projection) = self.projection.as_mut() { + let next = match projection.status { + RunStatus::Paused { + prior_block: Some(blocked_reason), + } => RunStatus::Blocked { blocked_reason }, + _ => RunStatus::Running, + }; + apply_status(projection, next, at); + if projection.pending_control == Some(RunControlAction::Unpause) { + projection.pending_control = None; + } + } + } + CoordinatorEvent::RunFinished { status } => { + self.state.finished = Some(status.to_string()); + self.conclude(status.to_string().as_str(), at); + } + CoordinatorEvent::GraphRegistered { .. } + | CoordinatorEvent::ExecutionFinished { .. } + | CoordinatorEvent::InvocationCancelRequested { .. } + | CoordinatorEvent::RunNoteRecorded { .. } => {} + } + } + + /// The run's conclusion, from its recorded finish and what the stages + /// summed to. + fn conclude(&mut self, status: &str, at: DateTime) { + let Some(projection) = self.projection.as_mut() else { + return; + }; + let root = self + .state + .root + .and_then(|root| self.state.invocations.get(&root)); + let failure_message = root.and_then(|root| root.failure.clone()); + let (run_status, outcome, failure) = match status { + "success" => ( + RunStatus::Succeeded { + reason: SuccessReason::Completed, + }, + StageOutcome::Succeeded, + None, + ), + "cancelled" => ( + RunStatus::Failed { + reason: FailureReason::Cancelled, + }, + StageOutcome::Failed { + retry_requested: false, + }, + Some(RunFailure { + reason: FailureReason::Cancelled, + detail: FailureDetail::new( + failure_message + .clone() + .unwrap_or_else(|| "the run was cancelled".to_string()), + FailureCategory::Canceled, + ), + }), + ), + _ => ( + RunStatus::Failed { + reason: FailureReason::WorkflowError, + }, + StageOutcome::Failed { + retry_requested: false, + }, + Some(RunFailure { + reason: FailureReason::WorkflowError, + detail: FailureDetail::new( + failure_message + .clone() + .unwrap_or_else(|| "the run failed".to_string()), + FailureCategory::Deterministic, + ), + }), + ), + }; + apply_status(projection, run_status, at); + projection.pending_control = None; + projection.pending_interviews.clear(); + let rollup = usage_rollup::usage_rollup_from_projection(projection); + let (stages, total_retries) = rollup.conclusion_stages(projection); + let wall_time_ms = self.state.started_at.map_or(0, |started| { + u64::try_from(at.timestamp_millis()) + .unwrap_or(0) + .saturating_sub(started) + }); + let timing = RunTiming::new( + wall_time_ms, + rollup.timing.inference_time_ms, + rollup.timing.tool_time_ms, + ); + let last_checkpoint = projection.checkpoints.last(); + projection.conclusion = Some(Conclusion { + timestamp: at, + status: outcome, + timing, + failure, + final_git_commit_sha: last_checkpoint + .and_then(|checkpoint| checkpoint.checkpoint.git_commit_sha.clone()), + stages, + usage: rollup.usage_if_present(), + total_retries, + diff: last_checkpoint + .map(|checkpoint| checkpoint.diff.clone()) + .unwrap_or_default(), + }); + } + + fn fold_engine(&mut self, engine: &Event, event: &RunEvent, at: DateTime) { + let Some(execution) = event.context.execution else { + return; + }; + match engine { + Event::AdmissionDecided { decision, .. } => { + if let Admission::Skip { outcome } = decision { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + stage.state = StageState::Skipped; + stage.completion = Some(StageCompletion { + outcome: StageOutcome::Skipped, + notes: None, + failure_reason: failure_message(&outcome.status), + timestamp: at, + }); + } + } + } + Event::StepStarted { attempt, .. } => { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + if attempt.raw() > 1 { + stage.clear_live_timing(); + stage.output = None; + stage.output_bytes = None; + } + stage.state = StageState::Running; + stage.live_streaming = Some(true); + } + } + Event::StepProgressRecorded { ev, .. } => { + self.fold_progress(execution, event, ev, at); + } + Event::StepFinished { + attempt, outcome, .. + } => { + let is_final = matches!( + event.derived, + Some(Derived::StepFinished { is_final: true, .. }) + ); + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + if let Some(output) = outcome.output.as_str() { + stage.output = Some(output.to_string()); + stage.output_bytes = Some(output.len() as u64); + } + stage.live_streaming = Some(false); + apply_metrics(stage, &outcome.metrics); + if is_final { + stage.completion = Some(StageCompletion { + outcome: stage_outcome(&outcome.status), + notes: None, + failure_reason: failure_message(&outcome.status), + timestamp: at, + }); + stage.termination = Some(match outcome.status { + Status::TimedOut => fabro_types::CommandTermination::TimedOut, + Status::Cancelled => fabro_types::CommandTermination::Cancelled, + Status::Success + | Status::PartialSuccess { .. } + | Status::Failure(_) + | Status::Skipped => fabro_types::CommandTermination::Exited, + }); + } else { + stage.state = StageState::Retrying; + debug!(attempt = attempt.raw(), "attempt returned; a retry follows"); + } + } + } + Event::ControlRequested { .. } => { + if let Some(Derived::ControlRequested { + deliverable: true, + 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); + } + } + Event::ExecutionStarted { .. } + | Event::TokenEmitted { .. } + | Event::RoutingResolved { .. } + | Event::RouteApplied { .. } + | Event::RetryElapsed { .. } + | Event::NodeExpanded { .. } + | Event::CancelRequested { .. } + | Event::KillRequested { .. } => {} + } + } + + fn fold_progress( + &mut self, + execution: ExecutionId, + event: &RunEvent, + ev: &StepEvent, + at: DateTime, + ) { + match ev { + StepEvent::Log { line, .. } => { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + let output = stage.output.get_or_insert_default(); + output.push_str(line); + output.push('\n'); + stage.output_bytes = Some(output.len() as u64); + stage.live_streaming = Some(true); + } + } + StepEvent::Artifact { .. } => {} + StepEvent::Custom(payload) => { + if let Some(parsed) = event.parsed() { + self.fold_parsed(execution, event, parsed, at); + return; + } + let kind = payload.get("kind").and_then(Value::as_str).unwrap_or(""); + match kind { + "pebble" => self.fold_pebble(execution, event, payload, at), + "attractor.prompt" => { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + stage.prompt = payload + .get("prompt") + .and_then(Value::as_str) + .map(str::to_string); + let model = payload.get("model").and_then(Value::as_str); + if let Some(model) = model { + let (provider, model_id) = split_model(model); + stage.provider_used = Some(StageModelUsage { + mode: StageModelUsage::MODE_PROMPT.to_string(), + provider: provider.map(str::to_string), + model: Some(model_id.to_string()), + reasoning_effort: None, + speed: None, + }); + stage.model = model_ref(provider, model_id); + } + } + } + "attractor.prompt.completed" => { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + stage.response = payload + .get("response") + .and_then(Value::as_str) + .map(str::to_string); + if let Some(usage) = usage_of(payload.get("usage")) { + stage.usage = usage; + } + } + } + "attractor.fallback.plan" => { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + let route = payload + .get("routes") + .and_then(Value::as_array) + .and_then(|routes| routes.first()); + if let Some(route) = route { + let provider = route.get("provider").and_then(Value::as_str); + let model = route.get("model").and_then(Value::as_str); + stage.provider_used = Some(StageModelUsage { + mode: StageModelUsage::MODE_AGENT.to_string(), + provider: provider.map(str::to_string), + model: model.map(str::to_string), + reasoning_effort: None, + speed: None, + }); + if let Some(model) = model { + stage.model = model_ref(provider, model); + } + } + } + } + "attractor.parallel.branch.started" => { + let invocation = payload.get("invocation").and_then(Value::as_u64); + let index = payload + .get("index") + .and_then(Value::as_u64) + .and_then(|index| u32::try_from(index).ok()); + let fork_firing = payload + .get("occurrence") + .and_then(|occurrence| occurrence.get("firing")) + .and_then(Value::as_u64); + if let (Some(invocation), Some(index), Some(fork_firing)) = + (invocation, index, fork_firing) + { + let group = self + .state + .stages + .get(&stage_key(execution.raw(), fork_firing)) + .map(|stage| stage.stage_id.clone()); + if let Some(group) = group { + self.state.invocations.entry(invocation).or_default().branch = + Some((group, index)); + } + } + } + _ => {} + } + } + } + } + + fn fold_parsed( + &mut self, + execution: ExecutionId, + event: &RunEvent, + parsed: &Parsed, + at: DateTime, + ) { + match parsed { + Parsed::Question { question } => { + let Some(subject) = event.subject.as_ref() else { + return; + }; + let Some(firing) = subject.firing else { + return; + }; + let key = stage_key(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(), + ); + self.state.questions.insert(question.id.clone(), key); + let Some(projection) = self.projection.as_mut() else { + return; + }; + projection + .pending_interviews + .insert(question.id.clone(), PendingInterviewRecord { + question: InterviewQuestionRecord { + id: question.id.clone(), + text: question.text.clone(), + stage: label, + question_type: question + .kind + .as_deref() + .and_then(|kind| kind.parse::().ok()) + .unwrap_or_default(), + options: question + .options + .iter() + .map(|option| InterviewOption { + key: option.key.clone(), + label: option.label.clone(), + description: None, + preview: None, + }) + .collect(), + allow_freeform: question.freeform, + timeout_seconds: question + .timeout_ms + .map(|timeout| timeout as f64 / 1000.0), + context_display: None, + review_target: None, + }, + started_at: at, + }); + apply_status( + projection, + RunStatus::Blocked { + blocked_reason: BlockedReason::HumanInputRequired, + }, + at, + ); + } + Parsed::QuestionExpired { expired } => { + self.close_questions(Some(expired.question.as_str()), None, at); + } + Parsed::Note { .. } => {} + } + } + + /// Close one question by id, or every question of a firing, and unblock + /// the run when none is left. + fn close_questions( + &mut self, + question: Option<&str>, + firing_key: Option<&str>, + at: DateTime, + ) { + let closed: Vec = match (question, firing_key) { + (Some(question), _) => vec![question.to_string()], + (None, Some(key)) => self + .state + .questions + .iter() + .filter(|(_, asked_by)| asked_by.as_str() == key) + .map(|(id, _)| id.clone()) + .collect(), + (None, None) => Vec::new(), + }; + for id in &closed { + self.state.questions.remove(id); + } + let Some(projection) = self.projection.as_mut() else { + return; + }; + for id in &closed { + projection.pending_interviews.remove(id); + } + if projection.pending_interviews.is_empty() + && matches!(projection.status, RunStatus::Blocked { .. }) + { + apply_status(projection, RunStatus::Running, at); + } + } + + fn fold_pebble( + &mut self, + execution: ExecutionId, + event: &RunEvent, + payload: &Value, + at: DateTime, + ) { + let Some(envelope) = payload.get("event") else { + return; + }; + let envelope: CodingAgentEvent = match serde_json::from_value(envelope.clone()) { + Ok(envelope) => envelope, + Err(error) => { + debug!(error = %error, "a pebble envelope did not decode; skipped"); + return; + } + }; + let Some(stage) = self.stage_of(execution, event.subject.as_ref()) else { + return; + }; + let agent = stage.agent.get_or_insert_default(); + agent.apply(&envelope); + if stage.completion.is_none() { + stage.usage = agent.usage.saturating_add(agent.descendant_usage()); + } + let is_root = envelope.parent_session_id.is_none(); + #[expect( + clippy::wildcard_enum_match_arm, + reason = "pebble's event vocabulary is non-exhaustive and only some events project" + )] + match &envelope.event { + CodingEvent::SessionStarted { + provider, model, .. + } if is_root => { + stage.provider_used = Some(StageModelUsage { + mode: StageModelUsage::MODE_AGENT.to_string(), + provider: provider.clone(), + model: model.clone(), + reasoning_effort: None, + speed: None, + }); + if let Some(model) = model.as_deref() { + stage.model = model_ref(provider.as_deref(), model); + } + } + CodingEvent::LlmRequestStarted { requested_model } if is_root => { + stage.inference = Some(StageInferenceProjection { + session_id: envelope.session_id.clone(), + started_at: at, + requested_model: requested_model.clone(), + first_output_at: None, + first_output_kind: None, + retries: 0, + }); + } + CodingEvent::LlmFirstOutput { kind } => { + if let Some(inference) = stage.inference.as_mut() { + if inference.session_id == envelope.session_id { + inference.first_output_at = Some(at); + inference.first_output_kind = Some(*kind); + } + } + } + CodingEvent::LlmRetry { .. } => { + if let Some(inference) = stage.inference.as_mut() { + if inference.session_id == envelope.session_id { + inference.retries = inference.retries.saturating_add(1); + inference.first_output_at = None; + inference.first_output_kind = None; + } + } + } + CodingEvent::AssistantMessage { model, .. } => { + if is_root { + if let Some(provider) = stage + .provider_used + .as_ref() + .and_then(|used| used.provider.as_deref()) + { + stage.model = model_ref(Some(provider), model); + } + } + close_inference(stage, &envelope.session_id, at); + } + CodingEvent::Error { .. } | CodingEvent::RoundInterrupted { .. } => { + close_inference(stage, &envelope.session_id, at); + } + CodingEvent::SessionEnded => { + close_inference(stage, &envelope.session_id, at); + stage.close_tool_batch_for_session(&envelope.session_id, at); + } + CodingEvent::ToolCallStarted { tool_call_id, .. } if is_root => { + stage.open_tool_call(envelope.session_id.clone(), tool_call_id.clone(), at); + } + CodingEvent::ToolCallCompleted { tool_call_id, .. } if is_root => { + stage.close_tool_call(&envelope.session_id, tool_call_id, at); + } + _ => {} + } + } + + fn fold_view(&mut self, view: &ViewEvent, event: &RunEvent, at: DateTime) { + let Some(execution) = event.context.execution else { + return; + }; + match view { + ViewEvent::VisitStarted { .. } => { + let Some(subject) = event.subject.as_ref() else { + return; + }; + self.start_visit(execution, subject, at); + } + ViewEvent::WaitStateChanged { state } => { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + match state { + WaitState::AwaitingAdmission => { + if stage.state == StageState::Running { + stage.state = StageState::Pending; + } + } + WaitState::Running | WaitState::AwaitingAnswer | WaitState::Cancelling => { + stage.state = StageState::Running; + } + WaitState::AwaitingRetry => stage.state = StageState::Retrying, + } + } + } + ViewEvent::RetryScheduled { .. } => { + if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) { + stage.state = StageState::Retrying; + } + } + ViewEvent::VisitCompleted { + outcome, + executed, + attempts, + } => { + let Some(stage) = self.stage_of(execution, event.subject.as_ref()) else { + return; + }; + stage.state = match outcome.status { + Status::Success => StageState::Succeeded, + Status::PartialSuccess { .. } => StageState::PartiallySucceeded, + Status::Failure(_) | Status::TimedOut => StageState::Failed, + Status::Skipped => StageState::Skipped, + Status::Cancelled => StageState::Cancelled, + }; + if stage.completion.is_none() || !*executed { + stage.completion = Some(StageCompletion { + outcome: stage_outcome(&outcome.status), + notes: None, + failure_reason: failure_message(&outcome.status), + timestamp: at, + }); + } + if stage.timing.is_none() { + let wall = stage + .started_at + .map_or(0, |started| timing::elapsed_ms(started, at)); + stage.set_authoritative_timing(StageTiming::new(wall, 0, 0)); + } + debug!(attempts, "visit completed"); + } + ViewEvent::ForkCompleted { + occurrence, + results, + .. + } => { + let key = stage_key(occurrence.execution.raw(), occurrence.firing.raw()); + let Some(stage_id) = self + .state + .stages + .get(&key) + .map(|stage| stage.stage_id.clone()) + else { + return; + }; + let Some(projection) = self.projection.as_mut() else { + return; + }; + if let Some(stage) = projection.stage_mut(&stage_id) { + stage.parallel_results = Some( + results + .iter() + .map(|result| ParallelBranchResult { + id: result.node.name.to_string(), + index: Some(result.branch.index as usize), + item_label: None, + status: stage_outcome(&result.status), + context_updates: BTreeMap::new(), + }) + .collect(), + ); + } + } + ViewEvent::ForkStarted { .. } + | ViewEvent::BranchCompleted { .. } + | ViewEvent::RunStalled { .. } => {} + } + } + + /// A firing exists: register its stage and, when it is a logical stage, + /// show it. + fn start_visit(&mut self, execution: ExecutionId, subject: &Subject, at: DateTime) { + let Some(firing) = subject.firing else { + return; + }; + let key = stage_key(execution.raw(), firing.raw()); + if self.state.stages.contains_key(&key) { + return; + } + let node_name = subject.node.name.to_string(); + let visit = subject.visit.unwrap_or(1).max(1); + let meta_kind = subject + .node + .meta + .get("kind") + .and_then(Value::as_str) + .unwrap_or(""); + let synthetic = subject + .node + .meta + .get("synthetic") + .and_then(Value::as_bool) + .unwrap_or(false); + let shown = !synthetic && meta_kind != "parallel.branch"; + // Only a shown stage takes a label: a lowering node (a branch's + // parent-side delegate shares its target's name) never competes with + // the stage it stands for. + let mut stage_id = StageId::new(node_name.clone(), visit); + if shown { + if self.state.labels.contains(&stage_id.to_string()) { + stage_id = StageId::new(format!("{node_name}/e{}", execution.raw()), visit); + } + self.state.labels.insert(stage_id.to_string()); + } + self.state.stages.insert(key, StageRef { + stage_id: stage_id.clone(), + shown, + node_name, + visit, + }); + if !shown { + return; + } + let branch = self + .state + .executions + .get(&execution.raw()) + .and_then(|invocation| self.state.invocations.get(invocation)) + .and_then(|invocation| invocation.branch.clone()); + let Some(projection) = self.projection.as_mut() else { + return; + }; + let since_created = at + .signed_duration_since(projection.spec.run_id.created_at()) + .num_milliseconds() + .max(0); + let ordinal = u32::try_from(since_created) + .unwrap_or(u32::MAX - 1) + .saturating_add(1); + let stage = projection.stage_entry(stage_id.node_id(), visit, first_event_seq(ordinal)); + stage.handler = Some(StageHandler::from_handler_type(Some(meta_kind))); + stage.started_at = Some(at); + stage.graph_visit = Some(visit); + stage.state = StageState::Pending; + stage.parallel_branch_id = branch.map(|(group, index)| ParallelBranchId::new(group, index)); + } + + /// The shown stage an event's subject firing belongs to. + fn stage_of( + &mut self, + execution: ExecutionId, + subject: Option<&Subject>, + ) -> Option<&mut StageProjection> { + let firing = subject?.firing?; + let stage = self + .state + .stages + .get(&stage_key(execution.raw(), firing.raw()))?; + if !stage.shown { + return None; + } + let stage_id = stage.stage_id.clone(); + self.projection.as_mut()?.stage_mut(&stage_id) + } +} + +impl Default for RunView { + fn default() -> Self { + Self::new() + } +} + +// ── Lifecycle ─────────────────────────────────────────────────────────── + +fn fold_lifecycle(projection: &mut RunProjection, record: &RunLifecycleRecord, at: DateTime) { + use RunLifecycleKind as Kind; + match record.transition { + Kind::Submitted => apply_status(projection, RunStatus::Submitted, at), + Kind::StartRequested | Kind::Unpaused => {} + Kind::Pending => { + if let Some(status) = record.status { + apply_status(projection, status, at); + } + projection.approval = Some(RunApproval { + state: RunApprovalState::Pending, + requested_at: at, + decided_at: None, + denial_reason: None, + }); + } + Kind::Approved => { + if let Some(approval) = projection.approval.as_mut() { + approval.state = RunApprovalState::Approved; + approval.decided_at = Some(at); + } + } + Kind::Denied => { + if let Some(approval) = projection.approval.as_mut() { + approval.state = RunApprovalState::Denied; + approval.decided_at = Some(at); + approval.denial_reason.clone_from(&record.reason); + } + apply_status( + projection, + RunStatus::Failed { + reason: FailureReason::ApprovalDenied, + }, + at, + ); + } + Kind::Runnable + | Kind::Starting + | Kind::Running + | Kind::Blocked + | Kind::Unblocked + | Kind::Removing + | Kind::Dead => { + if let Some(status) = record.status { + apply_status(projection, status, at); + } + } + Kind::Paused => { + let prior_block = match projection.status { + RunStatus::Blocked { blocked_reason } => Some(blocked_reason), + _ => None, + }; + apply_status(projection, RunStatus::Paused { prior_block }, at); + } + Kind::Succeeded | Kind::Failed => { + if let Some(status) = record.status { + apply_status(projection, status, at); + } + projection.pending_control = None; + if projection.conclusion.is_none() { + let (outcome, failure) = match record.status { + Some(RunStatus::Failed { reason }) => ( + StageOutcome::Failed { + retry_requested: false, + }, + Some(RunFailure { + reason, + detail: FailureDetail::new( + record + .reason + .clone() + .unwrap_or_else(|| "the run failed".to_string()), + FailureCategory::Deterministic, + ), + }), + ), + _ => (StageOutcome::Succeeded, None), + }; + projection.conclusion = Some(Conclusion { + timestamp: at, + status: outcome, + timing: RunTiming::default(), + failure, + final_git_commit_sha: None, + stages: Vec::new(), + usage: None, + total_retries: 0, + diff: RunDiff::default(), + }); + } + } + Kind::CancelRequested => projection.pending_control = Some(RunControlAction::Cancel), + Kind::PauseRequested => projection.pending_control = Some(RunControlAction::Pause), + Kind::UnpauseRequested => projection.pending_control = Some(RunControlAction::Unpause), + } +} + +/// Apply a status transition; one the lifecycle refuses is logged and +/// skipped, since the view never fails the run. +fn apply_status(projection: &mut RunProjection, status: RunStatus, at: DateTime) { + if let Err(error) = projection.try_apply_status(status, at) { + debug!(error = %error, "status transition not applied to the Petri projection"); + } +} + +fn touch(projection: &mut RunProjection, at: DateTime) { + if at > projection.last_event_at { + projection.last_event_at = at; + } +} + +// ── Helpers ───────────────────────────────────────────────────────────── + +/// The key of a stage: its execution and firing. +#[must_use] +pub fn stage_key(execution: u64, firing: u64) -> String { + format!("{execution}:{firing}") +} + +/// The fork firing and branch index a branch child's call slot names: +/// `branch:@::`. +fn branch_slot(slot: &str) -> Option<(u64, u32)> { + let rest = slot.strip_prefix("branch:")?; + let mut parts = rest.splitn(3, ':'); + let fork = parts.next()?; + let index = parts.next()?.parse::().ok()?; + let firing = fork.rsplit_once('@')?.1.parse::().ok()?; + Some((firing, index)) +} + +fn millis(recorded_at: u64) -> DateTime { + Utc.timestamp_millis_opt(i64::try_from(recorded_at).unwrap_or(i64::MAX)) + .single() + .unwrap_or_default() +} + +fn sandbox_plan(settings: &RunEnvironmentSettings) -> RunSandboxPlan { + RunSandboxPlan { + provider: settings.provider.clone(), + image: (settings.provider == SandboxProviderKind::DOCKER) + .then(|| settings.image.docker.clone()) + .flatten() + .filter(|image| !image.is_empty()), + snapshot: None, + } +} + +fn stage_outcome(status: &Status) -> StageOutcome { + match status { + Status::Success => StageOutcome::Succeeded, + Status::PartialSuccess { .. } => StageOutcome::PartiallySucceeded, + Status::Failure(info) => StageOutcome::Failed { + retry_requested: info.class.as_str() == "retry_requested", + }, + Status::Skipped => StageOutcome::Skipped, + Status::Cancelled | Status::TimedOut => StageOutcome::Failed { + retry_requested: false, + }, + } +} + +fn failure_message(status: &Status) -> Option { + match status { + Status::Failure(info) + | Status::PartialSuccess { + underlying: Some(info), + } => Some(info.message.clone()), + Status::TimedOut => Some("the step timed out".to_string()), + Status::Cancelled => Some("the step was cancelled".to_string()), + Status::Success | Status::PartialSuccess { underlying: None } | Status::Skipped => None, + } +} + +/// The finished attempt's metrics onto its stage: the timing and the usage +/// the backend reported. +fn apply_metrics(stage: &mut StageProjection, metrics: &Metrics) { + let custom = &metrics.custom; + let inference = custom + .get("pebble.inference_ms") + .and_then(Value::as_u64) + .unwrap_or(0); + let tool = custom + .get("pebble.tool_ms") + .and_then(Value::as_u64) + .unwrap_or(0); + let wall = metrics.duration_ms.unwrap_or(0); + let (inference, tool) = match stage.handler { + Some(StageHandler::Prompt) => (wall, 0), + Some(StageHandler::Command) => (0, wall), + _ => (inference, tool), + }; + stage.set_authoritative_timing(StageTiming::new(wall, inference, tool).clamped_to_wall()); + if let Some(usage) = + usage_of(custom.get("pebble.usage")).or_else(|| usage_of(custom.get("prompt.usage"))) + { + stage.usage = usage; + } + if let Some(sessions) = custom + .get("pebble.subagents") + .and_then(|subagents| subagents.get("sessions")) + .and_then(Value::as_array) + { + let mut by_model: Vec = Vec::new(); + for session in sessions { + let provider = session.get("provider").and_then(Value::as_str); + let model = session.get("model").and_then(Value::as_str); + let Some(usage) = usage_of(session.get("usage")) else { + continue; + }; + let Some(model) = model.and_then(|model| model_ref(provider, model)) else { + continue; + }; + if let Some(entry) = by_model.iter_mut().find(|entry| entry.model == model) { + entry.usage = entry.usage.saturating_add(usage); + } else { + by_model.push(ModelUsage::new(model, usage)); + } + } + if !by_model.is_empty() { + stage.usage_by_model = by_model; + } + } +} + +fn usage_of(value: Option<&Value>) -> Option { + serde_json::from_value(value?.clone()).ok() +} + +/// `provider/model` into its parts, or the model alone. +fn split_model(model: &str) -> (Option<&str>, &str) { + match model.split_once('/') { + Some((provider, model)) if !provider.is_empty() && !model.is_empty() => { + (Some(provider), model) + } + _ => (None, model), + } +} + +fn model_ref(provider: Option<&str>, model: &str) -> Option { + let provider = provider.filter(|provider| !provider.is_empty())?; + Some(ModelRef::new( + ProviderId::new(provider), + ModelId::new(model), + )) +} + +fn close_inference(stage: &mut StageProjection, session_id: &str, at: DateTime) { + let open = stage + .inference + .as_ref() + .is_some_and(|inference| inference.session_id == session_id); + if !open { + return; + } + if let Some(inference) = stage.inference.take() { + stage.accumulate_inference_ms(timing::elapsed_ms(inference.started_at, at)); + } +} + +/// The run id a Petri run key names. +#[must_use] +pub fn run_id_of(key: &str) -> Option { + key.parse().ok() +} + +#[cfg(test)] +mod tests { + use fabro_store::platform_records::RunCreatedRecord; + use fabro_types::test_support as types_support; + use petri_execution::events::NodeRef; + use petri_runtime::driver::BranchRole; + use petri_runtime::ir::{FiringId, NodeId}; + + use super::*; + + #[test] + fn a_branch_slot_names_the_fork_firing_and_the_index() { + assert_eq!(branch_slot("branch:fan@7:2:review"), Some((7, 2))); + assert_eq!(branch_slot("branch:fan@7:x:review"), None); + assert_eq!(branch_slot("child:0"), None); + } + + #[test] + fn a_model_selector_splits_into_provider_and_model() { + assert_eq!(split_model("openai/gpt-5.4"), (Some("openai"), "gpt-5.4")); + assert_eq!(split_model("gpt-5.4"), (None, "gpt-5.4")); + assert!(model_ref(None, "gpt-5.4").is_none()); + assert!(model_ref(Some("openai"), "gpt-5.4").is_some()); + } + + #[test] + fn a_taken_label_is_made_unique_by_the_execution() { + let mut view = RunView::new(); + let created = StoredPlatformRecord { + seq: 1, + recorded_at: 1_000, + record: PlatformRecord::RunCreated(RunCreatedRecord { + spec: types_support::test_run_spec(), + title: Some("A run".to_string()), + parent_id: None, + retried_from: None, + web_url: None, + }), + position: None, + }; + view.fold(&Item::Platform(&created), 1); + let subject = |name: &str| Subject { + node: NodeRef { + id: NodeId::new(1), + name: name.into(), + kind: "attractor/command".into(), + meta: serde_json::json!({ "kind": "command" }), + }, + firing: Some(FiringId::new(4)), + visit: Some(1), + attempt: None, + generation: None, + branch: BranchRole::None, + }; + view.start_visit(ExecutionId::new(1), &subject("build"), millis(2_000)); + view.start_visit(ExecutionId::new(2), &subject("build"), millis(3_000)); + let labels: Vec = view + .projection() + .expect("the run was created") + .iter_stages() + .map(|(id, _)| id.to_string()) + .collect(); + assert_eq!(labels, vec!["build@1", "build/e2@1"]); + assert_eq!(view.state.stages.len(), 2); + } +} diff --git a/lib/components/fabro-petri/src/projector.rs b/lib/components/fabro-petri/src/projector.rs new file mode 100644 index 000000000..39bc98c3c --- /dev/null +++ b/lib/components/fabro-petri/src/projector.rs @@ -0,0 +1,761 @@ +//! 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::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, + #[serde(default)] + pub platform_seq: u64, +} + +impl Positions { + fn held(&self) -> BTreeMap { + 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, +} + +/// 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. +pub struct Projector { + pool: DbPool, + store: SqliteRunStore, + platform: PlatformRecordStore, + slots: Mutex>, + /// 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 a pool whose migrations have run. + #[must_use] + pub fn new(pool: DbPool) -> Arc { + Arc::new(Self { + store: SqliteRunStore::new(pool.clone()), + platform: PlatformRecordStore::new(pool.clone()), + pool, + slots: 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, 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 { + let ids: Vec = sqlx::query_scalar( + "SELECT run_id FROM petri_runs UNION SELECT run_id FROM platform_records ORDER BY 1", + ) + .fetch_all(&self.pool) + .await + .map_err(ProjectError::Database)?; + 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; + let pass = self.project_run(run_id).await?; + if !pass.skipped { + report.projected += 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. + pub async fn project_run(&self, run_id: RunId) -> Result { + 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 = 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 { + 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 = + 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, 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.pool) + .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, + ) -> Result { + 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, + inner: Arc, + ) -> Arc { + Arc::new(SignallingStore { + inner, + projector: Arc::clone(self), + }) + } +} + +/// A run store that signals a projector after each append. +struct SignallingStore { + inner: Arc, + projector: Arc, +} + +#[async_trait::async_trait] +impl petri_execution::RunStore for SignallingStore { + async fn open( + &self, + key: &RunKey, + access: Access, + ) -> Result, 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, + run_id: Option, + projector: Arc, +} + +#[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, StoreError> { + self.inner.read(log).await + } + + async fn put_blob(&self, bytes: &[u8]) -> Result { + self.inner.put_blob(bytes).await + } + + async fn get_blob(&self, digest: petri_store::Digest) -> Result>, 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: `//`. +#[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(mutex: &Mutex) -> 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. +pub async fn rebuild( + pool: &DbPool, + run_id: RunId, +) -> Result<(Option, Positions, u64), ProjectError> { + let store = SqliteRunStore::new(pool.clone()); + let platform = PlatformRecordStore::new(pool.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. +pub async fn stored_positions( + pool: &DbPool, + run_id: RunId, +) -> Result, 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(pool) + .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( + pool: &DbPool, + run_id: RunId, +) -> Result, ProjectError> { + let json: Option = + sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?") + .bind(run_id.to_string()) + .fetch_optional(pool) + .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( + pool: &DbPool, + run_id: RunId, +) -> Result, 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(pool) + .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( + pool: &DbPool, + run_id: RunId, +) -> Result, ProjectError> { + PlatformRecordStore::new(pool.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() +} diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs new file mode 100644 index 000000000..58d2ee545 --- /dev/null +++ b/lib/components/fabro-petri/tests/projection.rs @@ -0,0 +1,853 @@ +//! 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 { + 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, + 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 { + let projector = Projector::new(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 { + 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(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, 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 = stream.iter().map(|(seq, _, _)| *seq).collect(); + assert_eq!( + seqs, + (1..=stream_seq).collect::>(), + "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 = 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()); + 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()); + 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"); +} + +/// 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)> = 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()); + 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()); + 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()); + 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::>() + .len() + >= 3, + "the parallel run has child executions: {:?}", + rows.iter().map(|row| &row.0).collect::>() + ); + let staged = copy_run_without_records(&scenario.pool, scenario.run_id).await; + let first = Projector::new(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()); + 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()); + 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()); + 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" + ); +}