From e0b546d465c43afad01ba606c0932383a5f631c4 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 22:49:15 -0400 Subject: [PATCH] Read the view tables where the run summary store keeps them The projector takes two pools: the one Petri's records live in and the one the view tables live in. In the server both are the one database; a test fixture keeps the runs row, the platform records and the projection tables in the run summary store's own pool, which the projector was not reading, so a run projected in a test server folded its Petri events before its run.created record. The startup run-history verification checks only a Petri run's identity and legacy guard, since its row is the projector's. An agent stage's response is the response. its outcome wrote into the run context, as the prompt step writes it. The scenario tests assert each branch's own index. Co-Authored-By: Claude Fable 5.1 --- lib/apps/fabro-server/src/server.rs | 11 ++- .../fabro-server/tests/it/scenario/petri.rs | 31 ++++----- lib/components/fabro-petri/src/projection.rs | 20 ++++++ lib/components/fabro-petri/src/projector.rs | 69 ++++++++++++------- .../fabro-petri/tests/projection.rs | 22 +++--- .../fabro-store/src/run_summary_store.rs | 15 ++++ 6 files changed, 112 insertions(+), 56 deletions(-) diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 78bdd8203..83959d509 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -1197,6 +1197,12 @@ impl AppState { &self.petri_projector } + /// The pool the Petri view tables live in, so a test can read them. + #[cfg(any(test, feature = "test-support"))] + pub fn test_petri_view_pool(&self) -> DbPool { + self.stores.runs.run_summary_store().pool() + } + /// A worker token for `run_id` with the plain `run:worker` scope, as the /// server mints for the worker it launches. pub fn test_issue_worker_token(&self, run_id: &RunId) -> String { @@ -2500,7 +2506,10 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result se /// 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) + projector::stored_stream(&state.test_petri_view_pool(), id) .await .expect("the stream reads") .len() @@ -314,25 +314,16 @@ async fn the_hello_bundle_runs_on_petri_when_the_version_names_the_engine() { 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}"); + let greet = &projection["stages"]["greet@1"]; + assert_eq!(greet["state"], "succeeded", "{projection}"); + assert_eq!(greet["handler"], "agent", "{greet}"); assert!( - prompt["response"] + greet["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}" + "the agent's answer is projected as the stage's response: {greet}" ); + assert!(run["usage"]["tokens"]["input"].as_u64().is_some(), "{run}"); let logs = twin.request_logs(&namespace).await; let requests = logs["requests"] .as_array() @@ -417,10 +408,14 @@ async fn a_parallel_bundle_projects_its_branches_through_the_server() { 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"] { + for (branch, index) in [("a@1", 0), ("b@1", 1)] { let stage = &projection["stages"][branch]; assert_eq!(stage["state"], "succeeded", "{branch}: {projection}"); - assert_eq!(stage["parallel_branch_id"]["group"], "fork@1", "{stage}"); + assert_eq!( + stage["parallel_branch_id"], + format!("fork@1:{index}"), + "{stage}" + ); } let fork = &projection["stages"]["fork@1"]; assert_eq!( diff --git a/lib/components/fabro-petri/src/projection.rs b/lib/components/fabro-petri/src/projection.rs index e98156021..dc0f03262 100644 --- a/lib/components/fabro-petri/src/projection.rs +++ b/lib/components/fabro-petri/src/projection.rs @@ -479,11 +479,31 @@ impl RunView { event.derived, Some(Derived::StepFinished { is_final: true, .. }) ); + let node_name = event + .subject + .as_ref() + .map(|subject| subject.node.name.to_string()); 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); } + // An agent's answer: the `response.` the step wrote + // into the run context, as the prompt step writes it. + if stage.handler == Some(StageHandler::Agent) { + let response = node_name + .as_deref() + .and_then(|name| { + outcome + .context_updates + .get(format!("response.{name}").as_str()) + }) + .and_then(Value::as_str) + .or_else(|| outcome.output.as_str()); + if let Some(response) = response { + stage.response = Some(response.to_string()); + } + } stage.live_streaming = Some(false); apply_metrics(stage, &outcome.metrics); if is_final { diff --git a/lib/components/fabro-petri/src/projector.rs b/lib/components/fabro-petri/src/projector.rs index 39bc98c3c..6421180a3 100644 --- a/lib/components/fabro-petri/src/projector.rs +++ b/lib/components/fabro-petri/src/projector.rs @@ -138,8 +138,13 @@ struct Slot { pending: bool, } -/// The projector over one database. +/// The projector over one database: the pool Petri's records are read +/// from, and the pool the view tables (`platform_records`, +/// `petri_projection`, `petri_stream`, `runs`) are read and written on. In +/// the server both are the one database; a test may hand it the run +/// summary store's own pool for the views. pub struct Projector { + records: DbPool, pool: DbPool, store: SqliteRunStore, platform: PlatformRecordStore, @@ -156,13 +161,16 @@ impl std::fmt::Debug for Projector { } impl Projector { - /// A projector over a pool whose migrations have run. + /// A projector over `records`, the pool Petri's records live in, and + /// `views`, the pool the view tables live in; both migrated. The server + /// passes its one pool twice. #[must_use] - pub fn new(pool: DbPool) -> Arc { + pub fn new(records: DbPool, views: DbPool) -> Arc { Arc::new(Self { - store: SqliteRunStore::new(pool.clone()), - platform: PlatformRecordStore::new(pool.clone()), - pool, + store: SqliteRunStore::new(records.clone()), + platform: PlatformRecordStore::new(views.clone()), + records, + pool: views, slots: Mutex::default(), fault: AtomicBool::new(false), }) @@ -233,12 +241,18 @@ impl Projector { /// 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 ids: Vec = sqlx::query_scalar("SELECT run_id FROM petri_runs") + .fetch_all(&self.records) + .await + .map_err(ProjectError::Database)?; + let with_platform: Vec = + sqlx::query_scalar("SELECT DISTINCT run_id FROM platform_records") + .fetch_all(&self.pool) + .await + .map_err(ProjectError::Database)?; + ids.extend(with_platform); + ids.sort(); + ids.dedup(); let mut report = StartupReport::default(); for id in ids { let Some(run_id) = projection::run_id_of(&id) else { @@ -485,7 +499,7 @@ impl Projector { OR log LIKE 'execution %') GROUP BY log", ) .bind(run_id.to_string()) - .fetch_all(&self.pool) + .fetch_all(&self.records) .await .map_err(ProjectError::Database)?; Ok(rows @@ -645,13 +659,15 @@ fn lock(mutex: &Mutex) -> MutexGuard<'_, T> { /// 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. +/// compares it with the live view. `records` and `views` are the two pools +/// [`Projector::new`] takes. pub async fn rebuild( - pool: &DbPool, + records: &DbPool, + views: &DbPool, run_id: RunId, ) -> Result<(Option, Positions, u64), ProjectError> { - let store = SqliteRunStore::new(pool.clone()); - let platform = PlatformRecordStore::new(pool.clone()); + let store = SqliteRunStore::new(records.clone()); + let platform = PlatformRecordStore::new(views.clone()); let key = RunKey::new(run_id.to_string()); let platform_records = platform.read(&run_id).await.map_err(ProjectError::Store)?; let events = match store.open(&key, Access::Read).await { @@ -690,15 +706,16 @@ pub async fn rebuild( Ok((view.projection, positions, stream_seq)) } -/// The stored view's positions and stream sequence, for a test. +/// The stored view's positions and stream sequence, for a test; `views` is +/// the pool the view tables live in. pub async fn stored_positions( - pool: &DbPool, + views: &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) + .fetch_optional(views) .await .map_err(ProjectError::Database)?; row.map(|(positions, stream_seq)| { @@ -712,13 +729,13 @@ pub async fn stored_positions( /// The stored view's projection, for a test or a reader outside the store. pub async fn stored_projection( - pool: &DbPool, + views: &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) + .fetch_optional(views) .await .map_err(ProjectError::Database)?; json.map(|json| serde_json::from_str(&json).map_err(ProjectError::Encode)) @@ -727,14 +744,14 @@ pub async fn stored_projection( /// The stream rows of a run: `(stream_seq, item_kind, item_id)`, in order. pub async fn stored_stream( - pool: &DbPool, + views: &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) + .fetch_all(views) .await .map_err(ProjectError::Database)?; Ok(rows @@ -745,10 +762,10 @@ pub async fn stored_stream( /// Every stored platform record of a run, for a reader outside the store. pub async fn stored_platform_records( - pool: &DbPool, + views: &DbPool, run_id: RunId, ) -> Result, ProjectError> { - PlatformRecordStore::new(pool.clone()) + PlatformRecordStore::new(views.clone()) .read(&run_id) .await .map_err(ProjectError::Store) diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs index 58d2ee545..e231515f7 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -274,7 +274,7 @@ async fn parallel_scenario() -> Scenario { /// 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()); + let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone()); projector.signal(scenario.run_id); let store = projector.observe_store(Arc::new(SqliteRunStore::new(scenario.pool.clone()))); run_workflow( @@ -344,7 +344,7 @@ async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) { .await .expect("the positions read") .expect("the run has positions"); - let (rebuilt, positions, stream_seq) = projector::rebuild(pool, run_id) + let (rebuilt, positions, stream_seq) = projector::rebuild(pool, pool, run_id) .await .expect("the run rebuilds"); let rebuilt = rebuilt.expect("the rebuild has a projection"); @@ -491,7 +491,7 @@ async fn dropped_wake_ups_are_caught_up_by_the_next_signal() { .is_none(), "nothing woke the view" ); - let projector = Projector::new(scenario.pool.clone()); + let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone()); projector.signal(scenario.run_id); projector.settle(scenario.run_id).await; assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; @@ -511,7 +511,7 @@ async fn the_startup_pass_catches_up_a_view_nobody_signalled() { } let scenario = command_scenario().await; run_unobserved(&scenario).await; - let projector = Projector::new(scenario.pool.clone()); + let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone()); let report = projector .startup_pass() .await @@ -632,7 +632,7 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix( for row in &first { insert_petri_row(&replayed, scenario.run_id, row).await; } - let before = Projector::new(replayed.clone()); + let before = Projector::new(replayed.clone(), replayed.clone()); let pass = before .project_run(scenario.run_id) .await @@ -669,7 +669,7 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix( ); // A new projector, as a restarted server builds one. - let after = Projector::new(replayed.clone()); + let after = Projector::new(replayed.clone(), replayed.clone()); let report = after.startup_pass().await.expect("the restart catches up"); assert_eq!((report.runs, report.projected), (1, 1)); let (positions_after, stream_after) = projector::stored_positions(&replayed, scenario.run_id) @@ -702,7 +702,7 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix( 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()); + let source = Projector::new(scenario.pool.clone(), scenario.pool.clone()); source.startup_pass().await.expect("the source projects"); let whole = projector::stored_projection(&scenario.pool, scenario.run_id) .await @@ -737,7 +737,7 @@ async fn a_restarted_projector_agrees_over_nested_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()); + let first = Projector::new(staged.clone(), staged.clone()); // The parent execution and the coordinator log up to the first child's // declaration go in first; a restart then sees the children. let (early, late): (Vec<_>, Vec<_>) = rows @@ -754,14 +754,14 @@ async fn a_restarted_projector_agrees_over_nested_child_executions() { insert_petri_row(&staged, scenario.run_id, row).await; } drop(first); - let second = Projector::new(staged.clone()); + let second = Projector::new(staged.clone(), staged.clone()); second .startup_pass() .await .expect("the second projector passes"); assert_view_equals_rebuild(&staged, scenario.run_id).await; - let whole = Projector::new(scenario.pool.clone()); + let whole = Projector::new(scenario.pool.clone(), scenario.pool.clone()); whole.startup_pass().await.expect("the source projects"); let one_go = projector::stored_projection(&scenario.pool, scenario.run_id) .await @@ -790,7 +790,7 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() { } let scenario = command_scenario().await; run_unobserved(&scenario).await; - let projector = Projector::new(scenario.pool.clone()); + let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone()); let clean = projector .project_run(scenario.run_id) .await diff --git a/lib/components/fabro-store/src/run_summary_store.rs b/lib/components/fabro-store/src/run_summary_store.rs index 802415cef..3b0254c1c 100644 --- a/lib/components/fabro-store/src/run_summary_store.rs +++ b/lib/components/fabro-store/src/run_summary_store.rs @@ -238,6 +238,14 @@ impl RunSummaryStore { } } + /// The pool this store's tables live in: the `runs` row, the run events, + /// the platform records and the Petri projection tables. The server's + /// one database; a test fixture's own. + #[must_use] + pub fn pool(&self) -> SqlitePool { + self.pool.clone() + } + /// The platform records over the same pool. #[must_use] pub fn platform_records(&self) -> PlatformRecordStore { @@ -902,6 +910,13 @@ WHERE id = ? let diff = run.diff.unwrap_or_default(); verify_run_field(&row, run, "id", &run.id.to_string())?; verify_run_field(&row, run, "source_last_seq", &i64::from(record.last_seq))?; + if entry.projection.spec.engine.is_petri() { + // A Petri run's row is written by its projector from Petri's + // records and the platform records; the legacy fold knows the + // lifecycle alone, so only the identity and the legacy guard + // are checked here. + return Ok(()); + } verify_run_field( &row, run,