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