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.<node> 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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 22:49:15 -04:00
parent 7b466ad9f3
commit e0b546d465
No known key found for this signature in database
6 changed files with 112 additions and 56 deletions

View file

@ -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<Arc<AppS
);
let variables = Arc::new(VariableStore::new(db_pool.clone()));
let petri_runs = PetriRuns::new(db_pool.clone());
let petri_projector = Projector::new(db_pool.clone());
// Petri's records live on the shared pool; the view tables live where the
// run summary store keeps the `runs` row (the same database in the
// server, a fixture of its own in a test).
let petri_projector = Projector::new(db_pool.clone(), store.run_summary_store().pool());
{
let projector = Arc::clone(&petri_projector);
store.set_platform_record_hook(Arc::new(move |run_id| projector.signal(run_id)));

View file

@ -210,7 +210,7 @@ async fn settled_state(state: &AppState, app: &axum::Router, run_id: &str) -> 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!(

View file

@ -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.<node>` 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 {

View file

@ -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<Self> {
pub fn new(records: DbPool, views: DbPool) -> Arc<Self> {
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<StartupReport, ProjectError> {
let ids: Vec<String> = 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<String> = sqlx::query_scalar("SELECT run_id FROM petri_runs")
.fetch_all(&self.records)
.await
.map_err(ProjectError::Database)?;
let with_platform: Vec<String> =
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<T>(mutex: &Mutex<T>) -> 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<RunProjection>, 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<Option<(Positions, u64)>, 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<Option<RunProjection>, ProjectError> {
let json: Option<String> =
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<Vec<(u64, String, String)>, 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<Vec<StoredPlatformRecord>, ProjectError> {
PlatformRecordStore::new(pool.clone())
PlatformRecordStore::new(views.clone())
.read(&run_id)
.await
.map_err(ProjectError::Store)

View file

@ -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<Projector> {
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::<BTreeSet<_>>()
);
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

View file

@ -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,