From 7b466ad9f3afe374b054999afacebfecbdd559b0 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 22:32:23 -0400 Subject: [PATCH] Drive the Petri projector from the server The server holds one projector over its database and signals it after each committed worker append, after each committed platform record (through the run summary store's hook), at worker exit, and over every Petri run at startup after the restart reconcile. A run executing in the server process under the test override appends through the projector's observing store, so it is signalled the same way. The scenario tests read GET /runs/{id}/state after the view settles: the hello prompt stage with its response, the command stage with its output, and a two-branch parallel bundle whose branches are grouped under the fork with the fork's results. Co-Authored-By: Claude Fable 5.1 --- lib/apps/fabro-server/src/serve.rs | 5 + lib/apps/fabro-server/src/server.rs | 17 +++ .../fabro-server/src/server/handler/petri.rs | 6 +- .../fabro-server/src/server/petri_runs.rs | 4 +- .../fabro-server/tests/it/scenario/petri.rs | 134 +++++++++++++++++- 5 files changed, 162 insertions(+), 4 deletions(-) diff --git a/lib/apps/fabro-server/src/serve.rs b/lib/apps/fabro-server/src/serve.rs index 3296c877c..f9a581401 100644 --- a/lib/apps/fabro-server/src/serve.rs +++ b/lib/apps/fabro-server/src/serve.rs @@ -838,6 +838,11 @@ where "Reconciled stale in-flight runs on startup" ); } + state + .petri_projector + .startup_pass() + .await + .context("catching Petri projections up at startup")?; spawn_scheduler(Arc::clone(&state)); spawn_automation_scheduler(Arc::clone(&state)); let pull_request_creation_supervisor = diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 21a57cf24..78bdd8203 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -61,6 +61,7 @@ use fabro_llm::credentials::CredentialProvider; use fabro_llm::lithos_catalog::Catalog; use fabro_llm::{ClientOptions, FabroClient}; use fabro_mcp_store::McpServerStore; +use fabro_petri::projector::Projector; use fabro_redact::redact_jsonl_line; use fabro_sandbox::details::sandbox_details; use fabro_sandbox::driver::{DaytonaCredentials, ProviderAccess, ProviderConnectOptions}; @@ -1116,6 +1117,8 @@ pub struct AppState { pub(crate) worker_runtime: Arc, /// The Petri runs held open for workers over the API. pub(crate) petri_runs: PetriRuns, + /// The projector of Petri runs: signalled after each committed record. + pub(crate) petri_projector: Arc, scheduler_notify: Notify, automation_scheduler_notify: Notify, pull_request_scheduler_notify: Notify, @@ -1187,6 +1190,13 @@ impl AppState { self.petri_runs.store() } + /// The projector of Petri runs, so a test can wait for a run's view to + /// settle before it reads it. + #[cfg(any(test, feature = "test-support"))] + pub fn test_petri_projector(&self) -> &Arc { + &self.petri_projector + } + /// A worker token for `run_id` with the plain `run:worker` scope, as the /// server mints for the worker it launches. pub fn test_issue_worker_token(&self, run_id: &RunId) -> String { @@ -2490,6 +2500,11 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result anyhow::Result, run_id: RunId) { // The worker is gone: whatever Petri run handles it held open over the // API drop here, so its lease never outlives it. state.petri_runs.worker_exited(run_id); + state.petri_projector.signal(run_id); append_worker_exit_failure(&run_store, run_id, &worker_exit).await; let final_state = match run_store.state().await { diff --git a/lib/apps/fabro-server/src/server/handler/petri.rs b/lib/apps/fabro-server/src/server/handler/petri.rs index cc6824147..1e36fe378 100644 --- a/lib/apps/fabro-server/src/server/handler/petri.rs +++ b/lib/apps/fabro-server/src/server/handler/petri.rs @@ -145,7 +145,11 @@ async fn append_records( Err(err) => return store_error_response(id, &err), }; match writer.append(&log, &records).await { - Ok(()) => StatusCode::NO_CONTENT.into_response(), + Ok(()) => { + // The records are durable; the projection trails them from here. + state.petri_projector.signal(id); + StatusCode::NO_CONTENT.into_response() + } Err(err) => store_error_response(id, &err), } } diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 931a8a459..3d0e6ecf2 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -341,7 +341,9 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { run_id: run_id.to_string(), run_dir: run_dir.join("petri"), execution, - store: Arc::new(SqliteRunStore::new(state.db_pool.clone())), + store: state + .petri_projector + .observe_store(Arc::new(SqliteRunStore::new(state.db_pool.clone()))), runtime: runtime_spec(&state, &eligible, dry_run), provider: run_state.spec.settings.run.environment.provider.clone(), cancel, diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index ae2dae66c..369eeb53e 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -26,8 +26,8 @@ use std::sync::Arc; use axum::body::Body; use axum::http::{Request, StatusCode}; -use fabro_petri::SqliteRunStore; use fabro_petri::engine::{self, RunStatus}; +use fabro_petri::{SqliteRunStore, projector}; use fabro_server::server::AppState; use fabro_server::test_support::{ TestAppStateBuilder, llm_overlay_with_provider_base_url, test_app_db_pool, @@ -35,7 +35,7 @@ use fabro_server::test_support::{ }; use fabro_static::EnvVars; use fabro_test::{TwinScenario, TwinScenarios, twin_openai}; -use fabro_types::{WorkflowPath, WorkflowVersion}; +use fabro_types::{RunId, WorkflowPath, WorkflowVersion}; use tower::ServiceExt; use crate::helpers::{ @@ -87,6 +87,23 @@ const UNKNOWN_MODEL_DOT: &str = r#"digraph Bad { start -> work -> exit }"#; +/// Two command branches joined by a fan-in. +const PARALLEL_DOT: &str = r#"digraph Parallel { + graph [goal="Run two branches"] + start [shape=Mdiamond] + exit [shape=Msquare] + fork [shape=component] + a [shape=parallelogram, script="echo a"] + b [shape=parallelogram, script="echo b"] + merge [shape=tripleoctagon] + start -> fork + fork -> a + fork -> b + a -> merge + b -> merge + merge -> exit +}"#; + const PLAIN_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n"; const PETRI_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n"; @@ -168,6 +185,37 @@ async fn petri_outcome(state: &AppState, run_id: &str) -> engine::RunOutcome { .expect("the run's Petri record inspects") } +/// The run's projected state once its projector settled. +async fn settled_state(state: &AppState, app: &axum::Router, run_id: &str) -> serde_json::Value { + let id: RunId = run_id.parse().expect("the run id parses"); + state.test_petri_projector().settle(id).await; + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/state"))) + .body(Body::empty()) + .expect("state request should build"); + let response = app + .clone() + .oneshot(req) + .await + .expect("state request routes"); + response_json( + response, + StatusCode::OK, + format!("GET /api/v1/runs/{run_id}/state"), + ) + .await +} + +/// How many items the run's projected stream holds. +async fn petri_stream_len(state: &AppState, run_id: &str) -> usize { + let id: RunId = run_id.parse().expect("the run id parses"); + projector::stored_stream(&test_app_db_pool(state), id) + .await + .expect("the stream reads") + .len() +} + async fn run_engine(app: &axum::Router, run_id: &str) -> serde_json::Value { let req = Request::builder() .method("GET") @@ -260,6 +308,31 @@ async fn the_hello_bundle_runs_on_petri_when_the_version_names_the_engine() { let outcome = petri_outcome(&state, &run_id).await; assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); assert!(outcome.complete, "{:?}", outcome.incomplete); + let projection = settled_state(&state, &app, &run_id).await; + assert_eq!(projection["status"]["kind"], "succeeded", "{projection}"); + assert_eq!( + projection["conclusion"]["status"], "succeeded", + "{projection}" + ); + let stages = projection["stages"] + .as_object() + .expect("the state carries its stages"); + let prompt = stages + .values() + .find(|stage| stage["handler"] == "prompt") + .unwrap_or_else(|| panic!("the hello prompt stage is projected: {projection}")); + assert_eq!(prompt["state"], "succeeded", "{prompt}"); + assert!( + prompt["response"] + .as_str() + .is_some_and(|response| response.contains("A haiku, added.")), + "the prompt's response is projected: {prompt}" + ); + assert_eq!( + run["usage"]["tokens"]["input"].as_u64().is_some(), + true, + "{run}" + ); let logs = twin.request_logs(&namespace).await; let requests = logs["requests"] .as_array() @@ -302,6 +375,63 @@ async fn a_command_bundle_runs_on_petri_under_the_server_setting() { let outcome = petri_outcome(&state, &run_id).await; assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); assert!(outcome.complete, "{:?}", outcome.incomplete); + let projection = settled_state(&state, &app, &run_id).await; + let say = &projection["stages"]["say@1"]; + assert_eq!(say["state"], "succeeded", "{projection}"); + assert_eq!(say["handler"], "command", "{say}"); + assert!( + say["output"] + .as_str() + .is_some_and(|output| output.contains("hello from petri")), + "{say}" + ); + let stream = petri_stream_len(&state, &run_id).await; + assert!(stream > 0, "the run's stream holds its events"); +} + +/// A parallel bundle with two command branches runs on Petri through the +/// server: each branch is a child execution, projected as a stage grouped +/// under the fork, and the fork carries the branch results. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_parallel_bundle_projects_its_branches_through_the_server() { + if host_plugin().is_none() { + return; + } + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let settings = settings_from_toml( + "_version = 1\n\n[run.environment]\nid = \"local\"\n\n[server.execution]\nengine = \ + \"petri\"\n", + ); + let state = test_app_state_with_options(settings, 5); + let app = test_app_with_scheduler(Arc::clone(&state)); + + let version_id = register_version(&app, &[ + ("workflow.fabro", PARALLEL_DOT), + ("workflow.toml", PLAIN_SETTINGS), + ]) + .await; + let run_id = + create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await; + + let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await; + let run = run_json(&app, &run_id).await; + assert_eq!(status, "succeeded", "run: {run}"); + let projection = settled_state(&state, &app, &run_id).await; + for branch in ["a@1", "b@1"] { + let stage = &projection["stages"][branch]; + assert_eq!(stage["state"], "succeeded", "{branch}: {projection}"); + assert_eq!(stage["parallel_branch_id"]["group"], "fork@1", "{stage}"); + } + let fork = &projection["stages"]["fork@1"]; + assert_eq!( + fork["parallel_results"].as_array().map(Vec::len), + Some(2), + "{fork}" + ); + assert_eq!( + projection["conclusion"]["status"], "succeeded", + "{projection}" + ); } /// A version that names no engine on a server whose setting is the default