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,