From 22844300af25787cf08a598db70c146061aefd70 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 11 Jul 2026 14:54:04 -0400 Subject: [PATCH 1/2] Add SQLite runs read model --- Cargo.lock | 2 + .../2026-07-11-sqlite-runs-read-model-plan.md | 68 ++ .../fabro-db/migrations/2026071104_runs.sql | 58 ++ lib/crates/fabro-db/tests/sqlite.rs | 59 ++ lib/crates/fabro-server/src/serve.rs | 10 +- lib/crates/fabro-server/src/server.rs | 18 +- .../src/server/handler/automations.rs | 43 +- .../fabro-server/src/server/handler/runs.rs | 226 ++--- lib/crates/fabro-store/Cargo.toml | 2 + lib/crates/fabro-store/src/error.rs | 7 + lib/crates/fabro-store/src/lib.rs | 5 + lib/crates/fabro-store/src/run_state.rs | 47 +- .../fabro-store/src/run_summary_store.rs | 789 ++++++++++++++++++ lib/crates/fabro-store/src/slate/mod.rs | 104 ++- .../fabro-store/src/slate/projection_cache.rs | 31 +- lib/crates/fabro-store/src/slate/run_store.rs | 113 ++- 16 files changed, 1326 insertions(+), 256 deletions(-) create mode 100644 docs/plans/2026-07-11-sqlite-runs-read-model-plan.md create mode 100644 lib/crates/fabro-db/migrations/2026071104_runs.sql create mode 100644 lib/crates/fabro-store/src/run_summary_store.rs diff --git a/Cargo.lock b/Cargo.lock index 2ac91d62f..2198d929b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3114,6 +3114,7 @@ dependencies = [ "bytes", "chrono", "dashmap", + "fabro-db", "fabro-types", "fabro-util", "futures", @@ -3124,6 +3125,7 @@ dependencies = [ "serde", "serde_json", "slatedb", + "sqlx", "tempfile", "thiserror 2.0.18", "tokio", diff --git a/docs/plans/2026-07-11-sqlite-runs-read-model-plan.md b/docs/plans/2026-07-11-sqlite-runs-read-model-plan.md new file mode 100644 index 000000000..53a15c30e --- /dev/null +++ b/docs/plans/2026-07-11-sqlite-runs-read-model-plan.md @@ -0,0 +1,68 @@ +# SQLite Runs Read Model Plan + +## Goal + +Add SQLite `runs` projection. Power run summary reads, board, table. SlateDB events remain authoritative. + +## Decisions + +- `runs` rebuildable from events; never directly mutated. +- Store query columns + canonical base `Run` JSON. +- Store five token buckets + `total_usd_micros`; derive `total_tokens`, `RunSize`. +- Dynamic fields at read: children count, live wall time, Ask Fabro readiness, queue position. +- Preserve API wire contract. No event schema changes. + +## 1. Schema + +- Add `fabro-db/migrations/2026071101_runs.sql`. +- Columns: `id`, `source_last_seq`, timestamps, status, archive, parent, title, workflow, repo, automation, diff totals, token buckets, cost, `summary_json`. +- Index: created, updated, status/archive, parent, workflow, repo, automation. +- Constraints: JSON valid; booleans/counts valid; known status strings. +- Migration tests: table, indexes, constraints. + +## 2. Concrete store + +- Add `fabro-store::RunSummaryStore` over `SqlitePool`; no trait/backend enum. +- Methods: upsert projected row, get, list/count, children count, delete. +- `source_last_seq` monotonic; duplicate/older projection no-op. +- One transaction writes columns + `summary_json`. +- Row decode validates `Run`; parity-check indexed columns against JSON in tests. + +## 3. Projection + reconciliation + +- Convert existing `CachedRunProjection` to SQL row via existing `build_summary`. +- After durable SlateDB append and in-memory projection update, synchronously upsert SQLite. +- Backfill/reconcile at startup from warmed SlateDB projections. +- Idempotent restart; newer SQLite watermark never overwritten. +- Projection failure observable with run id/seq; no payload logging. +- Delete SQLite row after authoritative SlateDB run deletion succeeds. + +## 4. Shadow verification + +- Keep current reads. +- Compare SQLite vs current cache for list/get in integration tests. +- Cover create, lifecycle, title, parent, archive, retry, billing, diff, delete, restart/backfill. +- Cover failed/interrupted reconciliation and resume. + +## 5. Read cutover + +- `GET /runs`: SQLite filtering, sorting, count, pagination. +- `GET /runs/{id}` and resolve: SQLite summary. +- Automation/parent-child summary lists: SQLite. +- Apply dynamic decorations after row decode. +- Keep `/state`, stages, detailed billing, settings, questions, events on full projection/event store. +- Remove list-path dependency on global projection-cache scan; retain detailed projection cache. + +## 6. Verification + +- `cargo nextest run -p fabro-db` +- `cargo nextest run -p fabro-store` +- `cargo nextest run -p fabro-server` +- `cargo build --workspace` +- Pinned fmt + Clippy. +- Representative multi-run append/list benchmark; record p50/p95, DB size. + +## Unresolved questions + +- Keep `RunSize` cost-derived (recommended), or redefine from total tokens? +- Cut over reads in same release after parity tests (recommended), or shadow for one release? diff --git a/lib/crates/fabro-db/migrations/2026071104_runs.sql b/lib/crates/fabro-db/migrations/2026071104_runs.sql new file mode 100644 index 000000000..42d2b7a7d --- /dev/null +++ b/lib/crates/fabro-db/migrations/2026071104_runs.sql @@ -0,0 +1,58 @@ +CREATE TABLE runs ( + id TEXT PRIMARY KEY NOT NULL, + source_last_seq INTEGER NOT NULL, + created_at_ms INTEGER NOT NULL, + started_at_ms INTEGER, + last_event_at_ms INTEGER NOT NULL, + completed_at_ms INTEGER, + status TEXT NOT NULL, + archived_at_ms INTEGER, + parent_id TEXT, + title TEXT NOT NULL, + workflow_slug TEXT, + workflow_name TEXT, + repository_name TEXT, + automation_id TEXT, + diff_files_changed INTEGER NOT NULL DEFAULT 0, + diff_additions INTEGER NOT NULL DEFAULT 0, + diff_deletions INTEGER NOT NULL DEFAULT 0, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + reasoning_tokens INTEGER NOT NULL DEFAULT 0, + cache_read_tokens INTEGER NOT NULL DEFAULT 0, + cache_write_tokens INTEGER NOT NULL DEFAULT 0, + total_usd_micros INTEGER, + summary_json TEXT NOT NULL, + CHECK (source_last_seq >= 1), + CHECK (status IN ( + 'submitted', + 'pending', + 'runnable', + 'starting', + 'running', + 'blocked', + 'paused', + 'removing', + 'succeeded', + 'failed', + 'dead' + )), + CHECK (diff_files_changed >= 0), + CHECK (diff_additions >= 0), + CHECK (diff_deletions >= 0), + CHECK (input_tokens >= 0), + CHECK (output_tokens >= 0), + CHECK (reasoning_tokens >= 0), + CHECK (cache_read_tokens >= 0), + CHECK (cache_write_tokens >= 0), + CHECK (total_usd_micros IS NULL OR total_usd_micros >= 0), + CHECK (json_valid(summary_json)) +); + +CREATE INDEX runs_by_created_at ON runs(created_at_ms DESC, id DESC); +CREATE INDEX runs_by_updated_at ON runs(last_event_at_ms DESC, id DESC); +CREATE INDEX runs_by_status ON runs(archived_at_ms, status, last_event_at_ms DESC, id DESC); +CREATE INDEX runs_by_parent ON runs(parent_id, created_at_ms DESC, id DESC); +CREATE INDEX runs_by_workflow ON runs(workflow_slug, created_at_ms DESC, id DESC); +CREATE INDEX runs_by_repository ON runs(repository_name, created_at_ms DESC, id DESC); +CREATE INDEX runs_by_automation ON runs(automation_id, created_at_ms DESC, id DESC); diff --git a/lib/crates/fabro-db/tests/sqlite.rs b/lib/crates/fabro-db/tests/sqlite.rs index 384dbcec2..e875209b8 100644 --- a/lib/crates/fabro-db/tests/sqlite.rs +++ b/lib/crates/fabro-db/tests/sqlite.rs @@ -66,6 +66,13 @@ async fn connect_creates_parent_directory_and_migrate_is_idempotent() -> anyhow: assert_eq!(count, 1, "{table} table should exist"); } + let runs_table_count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'runs'", + ) + .fetch_one(database.pool()) + .await?; + assert_eq!(runs_table_count, 1); + let legacy_import_table_count: i64 = sqlx::query_scalar( "SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'legacy_imports'", ) @@ -331,6 +338,58 @@ async fn insert_minimal_automation( Ok(()) } +#[tokio::test] +async fn runs_schema_creates_indexes_and_rejects_invalid_rows() -> anyhow::Result<()> { + let dir = tempfile::tempdir()?; + let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?; + database.migrate().await?; + + let index_count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' AND name LIKE 'runs_by_%'", + ) + .fetch_one(database.pool()) + .await?; + assert_eq!(index_count, 7); + + insert_minimal_run(database.pool(), "submitted", 0, r#"{"id":"run"}"#).await?; + for (status, input_tokens, summary_json) in [ + ("unknown", 0, r#"{"id":"run-2"}"#), + ("submitted", -1, r#"{"id":"run-3"}"#), + ("submitted", 0, "not-json"), + ] { + assert!( + insert_minimal_run(database.pool(), status, input_tokens, summary_json) + .await + .is_err() + ); + } + + Ok(()) +} + +async fn insert_minimal_run( + pool: &fabro_db::DbPool, + status: &str, + input_tokens: i64, + summary_json: &str, +) -> Result<(), sqlx::Error> { + sqlx::query( + r" +INSERT INTO runs ( + id, source_last_seq, created_at_ms, last_event_at_ms, status, title, + input_tokens, summary_json +) VALUES (?, 1, 0, 0, ?, 'title', ?, ?) +", + ) + .bind(format!("run-{status}-{input_tokens}")) + .bind(status) + .bind(input_tokens) + .bind(summary_json) + .execute(pool) + .await?; + Ok(()) +} + #[tokio::test] async fn environments_schema_rejects_invalid_rows() -> anyhow::Result<()> { let dir = tempfile::tempdir()?; diff --git a/lib/crates/fabro-server/src/serve.rs b/lib/crates/fabro-server/src/serve.rs index 6325e04a3..f46bb7e0a 100644 --- a/lib/crates/fabro-server/src/serve.rs +++ b/lib/crates/fabro-server/src/serve.rs @@ -779,10 +779,6 @@ where flush_interval, cache_path, )); - store - .warm_projection_cache() - .await - .context("warming run projection cache")?; let auth_code_store = store.auth_codes().await?; let auth_token_store = store.refresh_tokens().await?; let (artifact_object_store, artifact_prefix) = build_artifact_object_store_with_server_secrets( @@ -815,6 +811,12 @@ where #[cfg(any(test, feature = "test-support"))] automation_materializer_override: None, })?; + state + .stores + .runs + .warm_projection_cache() + .await + .context("warming run projection cache and reconciling run summaries")?; let reconciled = reconcile_incomplete_runs_on_startup(&state).await?; if reconciled > 0 { info!( diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index d31506ab6..1b37763d6 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -86,7 +86,7 @@ use fabro_slack::{blocks as slack_blocks, connection as slack_connection}; use fabro_static::EnvVars; use fabro_store::{ ArtifactKey, ArtifactStore, Database, EventEnvelope, EventPayload, NodeArtifact, - PendingInterviewRecord, StageArtifactEntry, StageId, + PendingInterviewRecord, RunSummaryStore, StageArtifactEntry, StageId, }; #[cfg(test)] use fabro_types::BlockedReason; @@ -1107,12 +1107,13 @@ pub struct AppState { } pub(crate) struct AppStores { - pub(crate) runs: Arc, - pub(crate) automations: Arc, - pub(crate) environments: Arc, - pub(crate) mcp_servers: Arc, - pub(crate) vault: Arc, - pub(crate) variables: Arc, + pub(crate) runs: Arc, + pub(crate) run_summaries: Arc, + pub(crate) automations: Arc, + pub(crate) environments: Arc, + pub(crate) mcp_servers: Arc, + pub(crate) vault: Arc, + pub(crate) variables: Arc, } type PullRequestCreateLocks = Arc>>>>; @@ -2386,6 +2387,8 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result anyhow::Result return ApiError::from(err).into_response(), } - let entries = match state - .stores - .runs - .list_cached_runs(&fabro_store::ListRunsQuery::default(), Utc::now()) - .await - { - Ok(entries) => entries, + let query = RunSummaryListQuery { + automation_id: Some(id.to_string()), + visibility: RunSummaryVisibility::All, + limit: pagination.limit.clamp(1, 100), + offset: pagination.offset.min(MAX_PAGE_OFFSET), + ..RunSummaryListQuery::default() + }; + let page = match state.stores.run_summaries.list(&query, Utc::now()).await { + Ok(page) => page, Err(err) => { return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) .into_response(); } }; - let mut runs: Vec = entries - .into_iter() - .map(|entry| entry.summary) - .filter(|run| { - run.automation - .as_ref() - .is_some_and(|automation| automation.id == id.as_str()) - }) - .collect(); - runs.sort_by(|a, b| { - b.timestamps - .created_at - .cmp(&a.timestamps.created_at) - .then_with(|| b.id.cmp(&a.id)) - }); - - let total = runs.len() as u64; - let (page, has_more) = paginate_items(runs, &pagination); - let data = state.decorate_run_summaries(page).await; + let data = state.decorate_run_summaries(page.data).await; ( StatusCode::OK, Json(serde_json::json!({ "data": data, - "meta": { "has_more": has_more, "total": total } + "meta": { "has_more": page.has_more, "total": page.total } })), ) .into_response() diff --git a/lib/crates/fabro-server/src/server/handler/runs.rs b/lib/crates/fabro-server/src/server/handler/runs.rs index 89bec59a8..cfdfbc2b0 100644 --- a/lib/crates/fabro-server/src/server/handler/runs.rs +++ b/lib/crates/fabro-server/src/server/handler/runs.rs @@ -11,13 +11,16 @@ use axum_extra::extract::Query as ExtraQuery; use base64::Engine as _; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use bytes::Bytes; -use chrono::{DateTime, Utc}; +use chrono::Utc; use fabro_api::types::{ BoardColumn, RunManifest, SubmitAnswerRequest, UpdateRunParentRequest, UpdateRunRequest, }; use fabro_config::Storage; use fabro_interview::AnswerSubmission; use fabro_llm::client::Client as LlmClient; +use fabro_store::{ + RunSummaryListQuery, RunSummarySort, RunSummarySortDirection, RunSummaryVisibility, +}; use fabro_types::settings::ResolveError; use fabro_types::{ AutomationRef, Principal, RunClientProvenance, RunId, RunProvenance, RunServerProvenance, @@ -33,9 +36,9 @@ use tokio::fs; use tracing::info; use super::super::{ - AppState, DeleteRunOutcome, ListResponse, PaginationParams, RunExecutionMode, VariableError, + AppState, DeleteRunOutcome, ListResponse, MAX_PAGE_OFFSET, RunExecutionMode, VariableError, answer_from_request, api_question_from_pending_interview, default_page_limit, - delete_run_internal, load_pending_interview, managed_run, paginate_items, parse_run_id_path, + delete_run_internal, load_pending_interview, managed_run, parse_run_id_path, parse_stage_id_path, reject_if_archived, submit_pending_interview_answer, workflow_event, }; use crate::error::ApiError; @@ -127,107 +130,83 @@ struct ListRunsParams { } impl ListRunsParams { - fn pagination(&self) -> PaginationParams { - PaginationParams { - limit: self.limit, - offset: self.offset, - } - } - - fn status_filter(&self) -> Option> { - if self.status.is_empty() { - None - } else { - Some(self.status.iter().copied().collect()) + fn summary_query(&self) -> RunSummaryListQuery { + RunSummaryListQuery { + parent_id: self.parent_id, + visibility: summary_visibility(&self.status, self.include_archived), + sort: summary_sort(self.sort), + direction: summary_sort_direction(self.direction), + limit: self.limit.clamp(1, 100), + offset: self.offset.min(MAX_PAGE_OFFSET), + ..RunSummaryListQuery::default() } } } -pub(crate) fn board_column(status: RunStatus, archived: bool) -> BoardColumn { - if archived { - return BoardColumn::Archived; +fn summary_visibility(selected: &[BoardColumn], include_archived: bool) -> RunSummaryVisibility { + if selected.is_empty() { + return RunSummaryVisibility::Default { include_archived }; } - match status { - RunStatus::Submitted | RunStatus::Pending { .. } => BoardColumn::Pending, - RunStatus::Runnable => BoardColumn::Runnable, - RunStatus::Starting => BoardColumn::Initializing, - RunStatus::Running | RunStatus::Paused { .. } => BoardColumn::Running, - RunStatus::Blocked { .. } => BoardColumn::Blocked, - RunStatus::Succeeded { .. } => BoardColumn::Succeeded, - RunStatus::Failed { .. } | RunStatus::Dead => BoardColumn::Failed, - RunStatus::Removing => BoardColumn::Removing, + + let mut statuses = HashSet::new(); + let mut archived = false; + for column in selected { + match column { + BoardColumn::Pending => { + statuses.insert(fabro_types::RunStatusKind::Submitted); + statuses.insert(fabro_types::RunStatusKind::Pending); + } + BoardColumn::Runnable => { + statuses.insert(fabro_types::RunStatusKind::Runnable); + } + BoardColumn::Initializing => { + statuses.insert(fabro_types::RunStatusKind::Starting); + } + BoardColumn::Running => { + statuses.insert(fabro_types::RunStatusKind::Running); + statuses.insert(fabro_types::RunStatusKind::Paused); + } + BoardColumn::Blocked => { + statuses.insert(fabro_types::RunStatusKind::Blocked); + } + BoardColumn::Succeeded => { + statuses.insert(fabro_types::RunStatusKind::Succeeded); + } + BoardColumn::Failed => { + statuses.insert(fabro_types::RunStatusKind::Failed); + statuses.insert(fabro_types::RunStatusKind::Dead); + } + BoardColumn::Archived => archived = true, + BoardColumn::Removing => { + statuses.insert(fabro_types::RunStatusKind::Removing); + } + } + } + RunSummaryVisibility::Selected { + statuses: statuses.into_iter().collect(), + archived, } } -fn run_elapsed_ms(run: &fabro_types::Run, now: DateTime) -> i64 { - let start = run - .timestamps - .started_at - .unwrap_or(run.timestamps.created_at); - let end = run.timestamps.completed_at.unwrap_or(now); - (end - start).num_milliseconds().max(0) +fn summary_sort(sort: RunsSortKey) -> RunSummarySort { + match sort { + RunsSortKey::CreatedAt => RunSummarySort::CreatedAt, + RunsSortKey::UpdatedAt => RunSummarySort::UpdatedAt, + RunsSortKey::Status => RunSummarySort::Status, + RunsSortKey::Elapsed => RunSummarySort::Elapsed, + RunsSortKey::Repo => RunSummarySort::Repository, + RunsSortKey::Title => RunSummarySort::Title, + RunsSortKey::Workflow => RunSummarySort::Workflow, + RunsSortKey::Changes => RunSummarySort::Changes, + RunsSortKey::Size => RunSummarySort::Size, + } } -fn sort_runs(runs: &mut [fabro_types::Run], key: RunsSortKey, direction: RunsSortDirection) { - let now = Utc::now(); - let asc = matches!(direction, RunsSortDirection::Asc); - runs.sort_by(|a, b| { - let primary = match key { - RunsSortKey::CreatedAt => a.timestamps.created_at.cmp(&b.timestamps.created_at), - RunsSortKey::UpdatedAt => { - let av = a - .timestamps - .last_event_at - .unwrap_or(a.timestamps.created_at); - let bv = b - .timestamps - .last_event_at - .unwrap_or(b.timestamps.created_at); - av.cmp(&bv) - } - RunsSortKey::Status => { - let ac = board_column(a.lifecycle.status, a.lifecycle.archived); - let bc = board_column(b.lifecycle.status, b.lifecycle.archived); - ac.cmp(&bc) - } - RunsSortKey::Elapsed => run_elapsed_ms(a, now).cmp(&run_elapsed_ms(b, now)), - RunsSortKey::Repo => run_repo_key(a).cmp(&run_repo_key(b)), - RunsSortKey::Title => run_title_key(a).cmp(&run_title_key(b)), - RunsSortKey::Workflow => run_workflow_key(a).cmp(&run_workflow_key(b)), - RunsSortKey::Changes => run_changes_total(a).cmp(&run_changes_total(b)), - RunsSortKey::Size => a.size.cmp(&b.size), - }; - let primary = if asc { primary } else { primary.reverse() }; - // Stable tiebreak: newer ULIDs (and thus newer runs) first. - primary.then_with(|| b.id.cmp(&a.id)) - }); -} - -fn run_repo_key(run: &fabro_types::Run) -> String { - run.repository - .as_ref() - .map(|repo| repo.name.to_lowercase()) - .unwrap_or_default() -} - -fn run_title_key(run: &fabro_types::Run) -> String { - run.title.trim().to_lowercase() -} - -fn run_workflow_key(run: &fabro_types::Run) -> String { - let wf = &run.workflow; - wf.name - .as_deref() - .or(wf.graph_name.as_deref()) - .or(wf.slug.as_deref()) - .map(str::to_lowercase) - .unwrap_or_default() -} - -fn run_changes_total(run: &fabro_types::Run) -> i64 { - run.diff - .as_ref() - .map_or(0, |diff| diff.additions + diff.deletions) +fn summary_sort_direction(direction: RunsSortDirection) -> RunSummarySortDirection { + match direction { + RunsSortDirection::Asc => RunSummarySortDirection::Asc, + RunsSortDirection::Desc => RunSummarySortDirection::Desc, + } } async fn link_run_parent( @@ -364,12 +343,7 @@ async fn validate_parent_link( } async fn updated_run_response(state: &AppState, run_id: &RunId) -> Response { - match state - .stores - .runs - .get_cached_summary(run_id, Utc::now()) - .await - { + match state.stores.run_summaries.get(run_id, Utc::now()).await { Ok(Some(summary)) => ( StatusCode::OK, Json(state.decorate_run_summary(summary).await), @@ -387,53 +361,26 @@ async fn list_runs( State(state): State>, ExtraQuery(params): ExtraQuery, ) -> Response { - let entries = match state + let page = match state .stores - .runs - .list_cached_runs( - &fabro_store::ListRunsQuery { - parent_id: params.parent_id, - ..fabro_store::ListRunsQuery::default() - }, - Utc::now(), - ) + .run_summaries + .list(¶ms.summary_query(), Utc::now()) .await { - Ok(entries) => entries, + Ok(page) => page, Err(err) => { return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) .into_response(); } }; - let status_filter = params.status_filter(); - let include_archived = params.include_archived; - - let filtered: Vec = entries - .into_iter() - .map(|entry| entry.summary) - .filter(|run| { - let column = board_column(run.lifecycle.status, run.lifecycle.archived); - match &status_filter { - Some(set) => set.contains(&column), - None => { - column != BoardColumn::Removing - && (include_archived || column != BoardColumn::Archived) - } - } - }) - .collect(); - - let mut decorated = state.decorate_run_summaries(filtered).await; - sort_runs(&mut decorated, params.sort, params.direction); - let total = decorated.len() as u64; - let (data, has_more) = paginate_items(decorated, ¶ms.pagination()); + let data = state.decorate_run_summaries(page.data).await; ( StatusCode::OK, Json(serde_json::json!({ "data": data, - "meta": { "has_more": has_more, "total": total } + "meta": { "has_more": page.has_more, "total": page.total } })), ) .into_response() @@ -478,13 +425,18 @@ async fn resolve_run( State(state): State>, Query(query): Query, ) -> Response { + let summary_query = RunSummaryListQuery { + visibility: RunSummaryVisibility::All, + limit: u32::MAX, + ..RunSummaryListQuery::default() + }; let runs = match state .stores - .runs - .list_runs(&fabro_store::ListRunsQuery::default(), Utc::now()) + .run_summaries + .list(&summary_query, Utc::now()) .await { - Ok(runs) => runs, + Ok(page) => page.data, Err(err) => { return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) .into_response(); @@ -1031,7 +983,7 @@ async fn get_run_status( RequireRunManagementTarget(id, _actor): RequireRunManagementTarget, State(state): State>, ) -> Response { - match state.stores.runs.get_cached_summary(&id, Utc::now()).await { + match state.stores.run_summaries.get(&id, Utc::now()).await { Ok(Some(run)) => { (StatusCode::OK, Json(state.decorate_run_summary(run).await)).into_response() } diff --git a/lib/crates/fabro-store/Cargo.toml b/lib/crates/fabro-store/Cargo.toml index c81bca853..1fd4c10dd 100644 --- a/lib/crates/fabro-store/Cargo.toml +++ b/lib/crates/fabro-store/Cargo.toml @@ -24,6 +24,7 @@ tokio-stream.workspace = true dashmap.workspace = true serde.workspace = true serde_json.workspace = true +sqlx.workspace = true chrono = { workspace = true, features = ["serde"] } bytes.workspace = true thiserror.workspace = true @@ -37,3 +38,4 @@ tokio = { workspace = true, features = ["test-util", "macros"] } tempfile = "3" ulid.workspace = true insta = { workspace = true } +fabro-db = { path = "../fabro-db" } diff --git a/lib/crates/fabro-store/src/error.rs b/lib/crates/fabro-store/src/error.rs index 068a8d035..33126d474 100644 --- a/lib/crates/fabro-store/src/error.rs +++ b/lib/crates/fabro-store/src/error.rs @@ -8,6 +8,8 @@ pub enum Error { ObjectStore(#[from] object_store::Error), #[error("Serialization error: {0}")] Serde(#[from] serde_json::Error), + #[error("SQLite error: {0}")] + Sqlite(#[from] sqlx::Error), #[error("I/O error: {0}")] Io(#[from] std::io::Error), #[error("Invalid event payload: {0}")] @@ -26,6 +28,11 @@ pub enum Error { InvalidKeySegment { segment: String }, #[error("failed to parse key: {0}")] KeyParse(String), + #[error("stored run summary {run_id} has inconsistent field {field}")] + RunSummaryMismatch { + run_id: String, + field: &'static str, + }, #[error("invalid status transition: {0}")] InvalidTransition(#[from] fabro_types::InvalidTransition), #[error("{0}")] diff --git a/lib/crates/fabro-store/src/lib.rs b/lib/crates/fabro-store/src/lib.rs index cfc8e1c0c..d33d8aa91 100644 --- a/lib/crates/fabro-store/src/lib.rs +++ b/lib/crates/fabro-store/src/lib.rs @@ -7,6 +7,7 @@ mod keys; mod record; mod run_sessions; mod run_state; +mod run_summary_store; mod serializable_projection; mod slate; mod types; @@ -25,6 +26,10 @@ pub use run_sessions::{ project_run_sessions, }; pub use run_state::RunProjectionReducer; +pub use run_summary_store::{ + RunSummaryListQuery, RunSummaryPage, RunSummarySort, RunSummarySortDirection, RunSummaryStore, + RunSummaryVisibility, +}; pub use serializable_projection::SerializableProjection; pub use slate::{ AuthCode, AuthCodeStore, Blob, BlobStore, CachedRunProjection, ConsumeOutcome, Database, diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index 8335d277c..e0275a795 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -8,16 +8,16 @@ use fabro_types::run_event::{ }; use fabro_types::settings::run::{EnvironmentProvider, RunEnvironmentSettings}; use fabro_types::{ - ActivatedSkill, AskFabro, BilledModelUsage, Checkpoint, CheckpointRecord, CommandTermination, - Conclusion, EventBody, FailureCategory, FailureSignature, InterviewQuestionRecord, - McpServerProjection, McpServerStatus, Outcome, PendingInterviewRecord, PendingReason, - PullRequestLink, RepositoryRef, Run, RunApproval, RunApprovalState, RunBillingSummary, - RunControlAction, RunDiff, RunEvent, RunId, RunLifecycle, RunLinks, RunModel, RunOrigin, - RunProjection, RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxPlan, - RunSandboxRuntime, RunSize, RunSpec, RunStatus, RunTimestamps, SandboxProviderKind, - StageCompletion, StageHandler, StageId, StageModelUsage, StageOutcome, StageProjection, - StageState, StartRecord, SubAgentProjection, SubAgentStatus, TodoListKind, TodoListProjection, - TodoProjection, WorkflowRef, first_event_seq, + ActivatedSkill, AskFabro, BilledModelUsage, BilledTokenCounts, Checkpoint, CheckpointRecord, + CommandTermination, Conclusion, EventBody, FailureCategory, FailureSignature, + InterviewQuestionRecord, McpServerProjection, McpServerStatus, Outcome, PendingInterviewRecord, + PendingReason, PullRequestLink, RepositoryRef, Run, RunApproval, RunApprovalState, + RunBillingSummary, RunControlAction, RunDiff, RunEvent, RunId, RunLifecycle, RunLinks, + RunModel, RunOrigin, RunProjection, RunSandbox, RunSandboxFailure, RunSandboxInstance, + RunSandboxPlan, RunSandboxRuntime, RunSize, RunSpec, RunStatus, RunTimestamps, + SandboxProviderKind, StageCompletion, StageHandler, StageId, StageModelUsage, StageOutcome, + StageProjection, StageState, StartRecord, SubAgentProjection, SubAgentStatus, TodoListKind, + TodoListProjection, TodoProjection, WorkflowRef, first_event_seq, }; use fabro_util::error::render_compact_with_causes; @@ -1004,20 +1004,25 @@ fn terminal_total_usd_micros(state: &RunProjection) -> Option { } fn projected_total_usd_micros(state: &RunProjection) -> Option { - let mut total_usd_micros = 0_i64; - let mut has_total = false; + projected_billing(state).total_usd_micros +} - for (stage_id, stage) in state.iter_stages() { - if is_boundary_stage(state, stage_id.node_id()) { - continue; - } - if let Some(value) = stage.usage.total_usd_micros { - total_usd_micros = total_usd_micros.saturating_add(value); - has_total = true; - } +pub(crate) fn projected_billing(state: &RunProjection) -> BilledTokenCounts { + if let Some(billing) = state + .conclusion + .as_ref() + .and_then(|conclusion| conclusion.billing.as_ref()) + { + return billing.clone(); } - has_total.then_some(total_usd_micros) + let mut billing = BilledTokenCounts::default(); + for (stage_id, stage) in state.iter_stages() { + if !is_boundary_stage(state, stage_id.node_id()) { + billing.add_counts(&stage.usage); + } + } + billing } fn is_boundary_stage(projection: &RunProjection, node_id: &str) -> bool { diff --git a/lib/crates/fabro-store/src/run_summary_store.rs b/lib/crates/fabro-store/src/run_summary_store.rs new file mode 100644 index 000000000..289530d1c --- /dev/null +++ b/lib/crates/fabro-store/src/run_summary_store.rs @@ -0,0 +1,789 @@ +use std::collections::HashSet; + +use chrono::{DateTime, Utc}; +use fabro_types::{Run, RunId, RunStatusKind, RunTiming}; +use sqlx::sqlite::{SqliteConnection, SqliteRow}; +use sqlx::{QueryBuilder, Row as _, Sqlite, SqlitePool}; + +use crate::run_state::projected_billing; +use crate::slate::CachedRunProjection; +use crate::{Error, Result}; + +const UPSERT_RUN_SQL: &str = r" +INSERT INTO runs ( + id, source_last_seq, created_at_ms, started_at_ms, last_event_at_ms, completed_at_ms, + status, archived_at_ms, parent_id, title, workflow_slug, workflow_name, + repository_name, automation_id, diff_files_changed, diff_additions, diff_deletions, + input_tokens, output_tokens, reasoning_tokens, cache_read_tokens, cache_write_tokens, + total_usd_micros, summary_json +) VALUES ( + ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ? +) +ON CONFLICT(id) DO UPDATE SET + source_last_seq = excluded.source_last_seq, + created_at_ms = excluded.created_at_ms, + started_at_ms = excluded.started_at_ms, + last_event_at_ms = excluded.last_event_at_ms, + completed_at_ms = excluded.completed_at_ms, + status = excluded.status, + archived_at_ms = excluded.archived_at_ms, + parent_id = excluded.parent_id, + title = excluded.title, + workflow_slug = excluded.workflow_slug, + workflow_name = excluded.workflow_name, + repository_name = excluded.repository_name, + automation_id = excluded.automation_id, + diff_files_changed = excluded.diff_files_changed, + diff_additions = excluded.diff_additions, + diff_deletions = excluded.diff_deletions, + input_tokens = excluded.input_tokens, + output_tokens = excluded.output_tokens, + reasoning_tokens = excluded.reasoning_tokens, + cache_read_tokens = excluded.cache_read_tokens, + cache_write_tokens = excluded.cache_write_tokens, + total_usd_micros = excluded.total_usd_micros, + summary_json = excluded.summary_json +WHERE excluded.source_last_seq > runs.source_last_seq +"; + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub enum RunSummarySort { + #[default] + CreatedAt, + UpdatedAt, + Status, + Elapsed, + Repository, + Title, + Workflow, + Changes, + Size, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub enum RunSummarySortDirection { + Asc, + #[default] + Desc, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum RunSummaryVisibility { + All, + Default { + include_archived: bool, + }, + Selected { + statuses: Vec, + archived: bool, + }, +} + +impl Default for RunSummaryVisibility { + fn default() -> Self { + Self::Default { + include_archived: false, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RunSummaryListQuery { + pub parent_id: Option, + pub automation_id: Option, + pub visibility: RunSummaryVisibility, + pub sort: RunSummarySort, + pub direction: RunSummarySortDirection, + pub limit: u32, + pub offset: u32, +} + +impl Default for RunSummaryListQuery { + fn default() -> Self { + Self { + parent_id: None, + automation_id: None, + visibility: RunSummaryVisibility::default(), + sort: RunSummarySort::default(), + direction: RunSummarySortDirection::default(), + limit: 100, + offset: 0, + } + } +} + +#[derive(Debug, Clone, PartialEq)] +pub struct RunSummaryPage { + pub data: Vec, + pub total: u64, + pub has_more: bool, +} + +#[derive(Clone)] +pub struct RunSummaryStore { + pool: SqlitePool, +} + +impl std::fmt::Debug for RunSummaryStore { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("RunSummaryStore").finish_non_exhaustive() + } +} + +impl RunSummaryStore { + #[must_use] + pub fn new(pool: SqlitePool) -> Self { + Self { pool } + } + + pub(crate) async fn upsert_projection(&self, entry: &CachedRunProjection) -> Result<()> { + let record = ProjectedRunSummary::from_entry(entry); + let mut connection = self.pool.acquire().await?; + upsert_run(&mut connection, &record).await?; + Ok(()) + } + + pub(crate) async fn reconcile(&self, entries: &[CachedRunProjection]) -> Result<()> { + let authoritative_ids = entries + .iter() + .map(|entry| entry.run_id.to_string()) + .collect::>(); + let mut transaction = self.pool.begin().await?; + for entry in entries { + let record = ProjectedRunSummary::from_entry(entry); + upsert_run(&mut transaction, &record).await?; + } + + let stored_ids = sqlx::query_scalar::<_, String>("SELECT id FROM runs") + .fetch_all(&mut *transaction) + .await?; + for stored_id in stored_ids { + if !authoritative_ids.contains(&stored_id) { + sqlx::query("DELETE FROM runs WHERE id = ?") + .bind(stored_id) + .execute(&mut *transaction) + .await?; + } + } + transaction.commit().await?; + Ok(()) + } + + pub async fn get(&self, run_id: &RunId, now: DateTime) -> Result> { + let row = sqlx::query( + r" +SELECT runs.id, runs.summary_json, + (SELECT COUNT(*) FROM runs AS child WHERE child.parent_id = runs.id) AS children_count +FROM runs +WHERE runs.id = ? +", + ) + .bind(run_id.to_string()) + .fetch_optional(&self.pool) + .await?; + row.map(|row| decode_run_row(&row, now)).transpose() + } + + pub async fn list( + &self, + query: &RunSummaryListQuery, + now: DateTime, + ) -> Result { + let mut transaction = self.pool.begin().await?; + + let mut count_query = QueryBuilder::::new("SELECT COUNT(*) FROM runs"); + push_filters(&mut count_query, query); + let total: i64 = count_query + .build_query_scalar() + .fetch_one(&mut *transaction) + .await?; + + let mut rows_query = QueryBuilder::::new( + r" +SELECT runs.id, runs.summary_json, + (SELECT COUNT(*) FROM runs AS child WHERE child.parent_id = runs.id) AS children_count +FROM runs", + ); + push_filters(&mut rows_query, query); + push_order(&mut rows_query, query.sort, query.direction, now); + rows_query.push(" LIMIT ").push_bind(i64::from(query.limit)); + rows_query + .push(" OFFSET ") + .push_bind(i64::from(query.offset)); + let rows = rows_query.build().fetch_all(&mut *transaction).await?; + transaction.commit().await?; + + let data = rows + .iter() + .map(|row| decode_run_row(row, now)) + .collect::>>()?; + let total = u64::try_from(total).map_err(|_| Error::RunSummaryMismatch { + run_id: "".to_string(), + field: "negative total", + })?; + let consumed = u64::from(query.offset).saturating_add(data.len() as u64); + Ok(RunSummaryPage { + data, + total, + has_more: consumed < total, + }) + } + + pub async fn delete(&self, run_id: &RunId) -> Result<()> { + sqlx::query("DELETE FROM runs WHERE id = ?") + .bind(run_id.to_string()) + .execute(&self.pool) + .await?; + Ok(()) + } +} + +#[derive(Debug)] +struct ProjectedRunSummary { + run: Run, + last_seq: u32, + workflow_name: Option, + repository_name: Option, + input_tokens: i64, + output_tokens: i64, + reasoning_tokens: i64, + cache_read_tokens: i64, + cache_write_tokens: i64, + total_usd_micros: Option, +} + +impl ProjectedRunSummary { + fn from_entry(entry: &CachedRunProjection) -> Self { + let mut run = entry.summary.clone(); + if run.timing.is_none() { + let at = run + .timestamps + .last_event_at + .unwrap_or(run.timestamps.created_at); + run.timing = entry.projection.live_run_timing(at); + } + let billing = projected_billing(&entry.projection); + let workflow_name = run + .workflow + .name + .clone() + .or_else(|| run.workflow.graph_name.clone()) + .or_else(|| run.workflow.slug.clone()); + let repository_name = run + .repository + .as_ref() + .map(|repository| repository.name.clone()); + + Self { + run, + last_seq: entry.last_seq, + workflow_name, + repository_name, + input_tokens: billing.input_tokens, + output_tokens: billing.output_tokens, + reasoning_tokens: billing.reasoning_tokens, + cache_read_tokens: billing.cache_read_tokens, + cache_write_tokens: billing.cache_write_tokens, + total_usd_micros: billing.total_usd_micros, + } + } +} + +async fn upsert_run(connection: &mut SqliteConnection, record: &ProjectedRunSummary) -> Result<()> { + let run = &record.run; + let diff = run.diff.unwrap_or_default(); + let summary_json = serde_json::to_string(run)?; + sqlx::query(UPSERT_RUN_SQL) + .bind(run.id.to_string()) + .bind(i64::from(record.last_seq)) + .bind(run.timestamps.created_at.timestamp_millis()) + .bind( + run.timestamps + .started_at + .map(|value| value.timestamp_millis()), + ) + .bind( + run.timestamps + .last_event_at + .unwrap_or(run.timestamps.created_at) + .timestamp_millis(), + ) + .bind( + run.timestamps + .completed_at + .map(|value| value.timestamp_millis()), + ) + .bind(run.lifecycle.status.kind().to_string()) + .bind( + run.lifecycle + .archived_at + .map(|value| value.timestamp_millis()), + ) + .bind(run.parent_id.map(|value| value.to_string())) + .bind(&run.title) + .bind(&run.workflow.slug) + .bind(&record.workflow_name) + .bind(&record.repository_name) + .bind(run.automation.as_ref().map(|automation| &automation.id)) + .bind(diff.files_changed) + .bind(diff.additions) + .bind(diff.deletions) + .bind(record.input_tokens) + .bind(record.output_tokens) + .bind(record.reasoning_tokens) + .bind(record.cache_read_tokens) + .bind(record.cache_write_tokens) + .bind(record.total_usd_micros) + .bind(summary_json) + .execute(connection) + .await?; + Ok(()) +} + +fn push_filters(builder: &mut QueryBuilder, query: &RunSummaryListQuery) { + builder.push(" WHERE 1 = 1"); + if let Some(parent_id) = query.parent_id { + builder + .push(" AND parent_id = ") + .push_bind(parent_id.to_string()); + } + if let Some(automation_id) = &query.automation_id { + builder + .push(" AND automation_id = ") + .push_bind(automation_id.clone()); + } + + match &query.visibility { + RunSummaryVisibility::All => {} + RunSummaryVisibility::Default { include_archived } => { + if *include_archived { + builder.push( + " AND (archived_at_ms IS NOT NULL OR (archived_at_ms IS NULL AND status <> 'removing'))", + ); + } else { + builder.push(" AND archived_at_ms IS NULL AND status <> 'removing'"); + } + } + RunSummaryVisibility::Selected { statuses, archived } => { + builder.push(" AND ("); + let mut has_condition = false; + if *archived { + builder.push("archived_at_ms IS NOT NULL"); + has_condition = true; + } + if !statuses.is_empty() { + if has_condition { + builder.push(" OR "); + } + builder.push("(archived_at_ms IS NULL AND status IN ("); + let mut separated = builder.separated(", "); + for status in statuses { + separated.push_bind(status.to_string()); + } + separated.push_unseparated("))"); + has_condition = true; + } + if !has_condition { + builder.push("0"); + } + builder.push(")"); + } + } +} + +fn push_order( + builder: &mut QueryBuilder, + sort: RunSummarySort, + direction: RunSummarySortDirection, + now: DateTime, +) { + builder.push(" ORDER BY "); + match sort { + RunSummarySort::CreatedAt => builder.push("created_at_ms"), + RunSummarySort::UpdatedAt => builder.push("last_event_at_ms"), + RunSummarySort::Status => builder.push( + r"CASE + WHEN archived_at_ms IS NOT NULL THEN 7 + WHEN status IN ('submitted', 'pending') THEN 0 + WHEN status = 'runnable' THEN 1 + WHEN status = 'starting' THEN 2 + WHEN status IN ('running', 'paused') THEN 3 + WHEN status = 'blocked' THEN 4 + WHEN status = 'succeeded' THEN 5 + WHEN status IN ('failed', 'dead') THEN 6 + WHEN status = 'removing' THEN 8 + ELSE 9 + END", + ), + RunSummarySort::Elapsed => builder + .push("(COALESCE(completed_at_ms, ") + .push_bind(now.timestamp_millis()) + .push(") - COALESCE(started_at_ms, created_at_ms))"), + RunSummarySort::Repository => builder.push("COALESCE(repository_name, '') COLLATE NOCASE"), + RunSummarySort::Title => builder.push("TRIM(title) COLLATE NOCASE"), + RunSummarySort::Workflow => builder.push("COALESCE(workflow_name, '') COLLATE NOCASE"), + RunSummarySort::Changes => builder.push("(diff_additions + diff_deletions)"), + RunSummarySort::Size => builder.push( + r"CASE + WHEN COALESCE(total_usd_micros, 0) <= 20000000 THEN 0 + WHEN total_usd_micros <= 50000000 THEN 1 + WHEN total_usd_micros <= 100000000 THEN 2 + WHEN total_usd_micros <= 200000000 THEN 3 + ELSE 4 + END", + ), + }; + match direction { + RunSummarySortDirection::Asc => builder.push(" ASC"), + RunSummarySortDirection::Desc => builder.push(" DESC"), + }; + builder.push(", id DESC"); +} + +fn decode_run_row(row: &SqliteRow, now: DateTime) -> Result { + let stored_id: String = row.try_get("id")?; + let summary_json: String = row.try_get("summary_json")?; + let children_count: i64 = row.try_get("children_count")?; + let mut run: Run = serde_json::from_str(&summary_json)?; + if stored_id != run.id.to_string() { + return Err(Error::RunSummaryMismatch { + run_id: stored_id, + field: "id", + }); + } + run.children_count = u64::try_from(children_count).map_err(|_| Error::RunSummaryMismatch { + run_id: run.id.to_string(), + field: "children_count", + })?; + apply_read_overlays(&mut run, now); + Ok(run) +} + +fn apply_read_overlays(run: &mut Run, now: DateTime) { + if run.timestamps.completed_at.is_some() { + return; + } + let Some(started_at) = run.timestamps.started_at else { + return; + }; + let wall_time_ms = u64::try_from( + now.signed_duration_since(started_at) + .num_milliseconds() + .max(0), + ) + .expect("non-negative milliseconds fit in u64"); + run.timing = Some( + run.timing + .unwrap_or_else(|| RunTiming::wall_only(wall_time_ms)) + .with_wall_time(wall_time_ms), + ); +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use chrono::{DateTime, Utc}; + use fabro_types::{ + AutomationRef, BilledTokenCounts, Conclusion, DiffSummary, Graph, RunDiff, RunId, + RunProjection, RunSize, RunSpec, RunStatus, RunTiming, StageOutcome, SuccessReason, + WorkflowSettings, test_support, + }; + use ulid::Ulid; + + use super::{ + RunSummaryListQuery, RunSummarySort, RunSummarySortDirection, RunSummaryStore, + RunSummaryVisibility, + }; + use crate::slate::CachedRunProjection; + + fn dt(value: &str) -> DateTime { + value.parse().unwrap() + } + + fn run_id(timestamp_ms: u64, random: u128) -> RunId { + RunId::from(Ulid::from_parts(timestamp_ms, random)) + } + + fn projection(run_id: RunId, title: &str, created_at: DateTime) -> RunProjection { + RunProjection::new( + title.to_string(), + RunSpec { + run_id, + settings: WorkflowSettings::default(), + graph: Graph::new("test"), + graph_source: None, + workflow_slug: Some("test-workflow".to_string()), + automation: None, + source_directory: None, + labels: HashMap::new(), + provenance: test_support::test_run_provenance(), + manifest_blob: None, + definition_blob: None, + git: None, + fork_source_ref: None, + }, + created_at, + ) + } + + fn entry(projection: RunProjection, last_seq: u32) -> CachedRunProjection { + CachedRunProjection::from_projection(projection.spec.run_id, projection, last_seq) + } + + async fn store() -> (tempfile::TempDir, RunSummaryStore) { + let directory = tempfile::tempdir().unwrap(); + let database = fabro_db::Database::connect(directory.path().join("fabro.sqlite3")) + .await + .unwrap(); + database.migrate().await.unwrap(); + (directory, RunSummaryStore::new(database.clone_pool())) + } + + #[tokio::test] + async fn upsert_is_monotonic_and_get_applies_children_count() { + let (_directory, store) = store().await; + let created_at = dt("2026-07-11T12:00:00Z"); + let parent_id = run_id(created_at.timestamp_millis().cast_unsigned(), 1); + let child_id = run_id(created_at.timestamp_millis().cast_unsigned() + 1, 2); + + let parent = entry(projection(parent_id, "parent", created_at), 1); + store.upsert_projection(&parent).await.unwrap(); + + let mut child_projection = projection(child_id, "new title", created_at); + child_projection.parent_id = Some(parent_id); + child_projection.last_event_at = created_at + chrono::Duration::seconds(2); + store + .upsert_projection(&entry(child_projection, 2)) + .await + .unwrap(); + + let mut stale = projection(child_id, "stale title", created_at); + stale.parent_id = Some(parent_id); + store.upsert_projection(&entry(stale, 1)).await.unwrap(); + + let parent = store.get(&parent_id, created_at).await.unwrap().unwrap(); + let child = store.get(&child_id, created_at).await.unwrap().unwrap(); + assert_eq!(parent.children_count, 1); + assert_eq!(child.title, "new title"); + } + + #[tokio::test] + async fn list_filters_sorts_and_paginates_in_sqlite() { + let (_directory, store) = store().await; + let created_at = dt("2026-07-11T12:00:00Z"); + let first_id = run_id(created_at.timestamp_millis().cast_unsigned(), 1); + let second_id = run_id(created_at.timestamp_millis().cast_unsigned() + 1, 2); + let archived_id = run_id(created_at.timestamp_millis().cast_unsigned() + 2, 3); + + let mut first = projection(first_id, "bravo", created_at); + first.spec.automation = Some(AutomationRef { + id: "nightly".to_string(), + name: None, + trigger_id: None, + }); + let mut second = projection(second_id, "alpha", created_at); + second.spec.automation = Some(AutomationRef { + id: "nightly".to_string(), + name: None, + trigger_id: None, + }); + let mut archived = projection(archived_id, "charlie", created_at); + archived.archived_at = Some(created_at); + for projected in [first, second, archived] { + store.upsert_projection(&entry(projected, 1)).await.unwrap(); + } + + let page = store + .list( + &RunSummaryListQuery { + automation_id: Some("nightly".to_string()), + sort: RunSummarySort::Title, + direction: RunSummarySortDirection::Asc, + limit: 1, + ..RunSummaryListQuery::default() + }, + created_at, + ) + .await + .unwrap(); + assert_eq!(page.total, 2); + assert!(page.has_more); + assert_eq!(page.data[0].title, "alpha"); + + let archived = store + .list( + &RunSummaryListQuery { + visibility: RunSummaryVisibility::Selected { + statuses: Vec::new(), + archived: true, + }, + ..RunSummaryListQuery::default() + }, + created_at, + ) + .await + .unwrap(); + assert_eq!(archived.data.len(), 1); + assert_eq!(archived.data[0].id, archived_id); + } + + #[tokio::test] + async fn projection_persists_billing_diff_and_derived_size() { + let (_directory, store) = store().await; + let created_at = dt("2026-07-11T12:00:00Z"); + let run_id = run_id(created_at.timestamp_millis().cast_unsigned(), 1); + let mut projection = projection(run_id, "billed", created_at); + projection.spec.automation = Some(AutomationRef { + id: "nightly".to_string(), + name: None, + trigger_id: None, + }); + projection.status = RunStatus::Succeeded { + reason: SuccessReason::Completed, + }; + projection.last_event_at = created_at + chrono::Duration::minutes(1); + projection.conclusion = Some(Conclusion { + timestamp: projection.last_event_at, + status: StageOutcome::Succeeded, + timing: RunTiming::wall_only(60_000), + failure: None, + final_git_commit_sha: None, + stages: Vec::new(), + billing: Some(BilledTokenCounts { + input_tokens: 100, + output_tokens: 20, + total_tokens: 135, + reasoning_tokens: 5, + cache_read_tokens: 10, + cache_write_tokens: 0, + total_usd_micros: Some(21_000_000), + }), + total_retries: 0, + diff: RunDiff { + patch: None, + summary: Some(DiffSummary { + files_changed: 2, + additions: 10, + deletions: 3, + }), + }, + }); + store + .upsert_projection(&entry(projection, 4)) + .await + .unwrap(); + + let row = sqlx::query( + "SELECT source_last_seq, created_at_ms, last_event_at_ms, status, title, workflow_slug, \ + automation_id, input_tokens, reasoning_tokens, cache_read_tokens, total_usd_micros, \ + diff_files_changed, diff_additions, diff_deletions FROM runs WHERE id = ?", + ) + .bind(run_id.to_string()) + .fetch_one(&store.pool) + .await + .unwrap(); + assert_eq!(sqlx::Row::get::(&row, "source_last_seq"), 4); + assert_eq!( + sqlx::Row::get::(&row, "created_at_ms"), + created_at.timestamp_millis() + ); + assert_eq!( + sqlx::Row::get::(&row, "last_event_at_ms"), + (created_at + chrono::Duration::minutes(1)).timestamp_millis() + ); + assert_eq!(sqlx::Row::get::(&row, "status"), "succeeded"); + assert_eq!(sqlx::Row::get::(&row, "title"), "billed"); + assert_eq!( + sqlx::Row::get::(&row, "workflow_slug"), + "test-workflow" + ); + assert_eq!( + sqlx::Row::get::(&row, "automation_id"), + "nightly" + ); + assert_eq!(sqlx::Row::get::(&row, "input_tokens"), 100); + assert_eq!(sqlx::Row::get::(&row, "reasoning_tokens"), 5); + assert_eq!(sqlx::Row::get::(&row, "cache_read_tokens"), 10); + assert_eq!( + sqlx::Row::get::(&row, "total_usd_micros"), + 21_000_000 + ); + assert_eq!(sqlx::Row::get::(&row, "diff_files_changed"), 2); + assert_eq!(sqlx::Row::get::(&row, "diff_additions"), 10); + assert_eq!(sqlx::Row::get::(&row, "diff_deletions"), 3); + + let run = store.get(&run_id, created_at).await.unwrap().unwrap(); + assert_eq!(run.size, RunSize::S); + } + + #[tokio::test] + async fn reconcile_removes_rows_absent_from_authoritative_entries() { + let (_directory, store) = store().await; + let created_at = dt("2026-07-11T12:00:00Z"); + let kept_id = run_id(created_at.timestamp_millis().cast_unsigned(), 1); + let removed_id = run_id(created_at.timestamp_millis().cast_unsigned() + 1, 2); + let kept = entry(projection(kept_id, "kept", created_at), 1); + let removed = entry(projection(removed_id, "removed", created_at), 1); + store.upsert_projection(&kept).await.unwrap(); + store.upsert_projection(&removed).await.unwrap(); + + store.reconcile(std::slice::from_ref(&kept)).await.unwrap(); + + assert!(store.get(&kept_id, created_at).await.unwrap().is_some()); + assert!(store.get(&removed_id, created_at).await.unwrap().is_none()); + } + + #[tokio::test] + async fn failed_reconcile_rolls_back_and_can_be_retried() { + let (_directory, store) = store().await; + let created_at = dt("2026-07-11T12:00:00Z"); + let stale_id = run_id(created_at.timestamp_millis().cast_unsigned(), 1); + let good_id = run_id(created_at.timestamp_millis().cast_unsigned() + 1, 2); + let recovered_id = run_id(created_at.timestamp_millis().cast_unsigned() + 2, 3); + store + .upsert_projection(&entry(projection(stale_id, "stale", created_at), 1)) + .await + .unwrap(); + + let good = entry(projection(good_id, "good", created_at), 1); + let mut invalid_projection = projection(recovered_id, "recovered", created_at); + invalid_projection.conclusion = Some(Conclusion { + timestamp: created_at, + status: StageOutcome::Succeeded, + timing: RunTiming::default(), + failure: None, + final_git_commit_sha: None, + stages: Vec::new(), + billing: Some(BilledTokenCounts { + input_tokens: -1, + ..BilledTokenCounts::default() + }), + total_retries: 0, + diff: RunDiff::default(), + }); + let invalid = entry(invalid_projection.clone(), 1); + + assert!(store.reconcile(&[good.clone(), invalid]).await.is_err()); + assert!(store.get(&stale_id, created_at).await.unwrap().is_some()); + assert!(store.get(&good_id, created_at).await.unwrap().is_none()); + + invalid_projection.conclusion.as_mut().unwrap().billing = Some(BilledTokenCounts { + input_tokens: 1, + total_tokens: 1, + ..BilledTokenCounts::default() + }); + let recovered = entry(invalid_projection, 1); + store.reconcile(&[good, recovered]).await.unwrap(); + + assert!(store.get(&stale_id, created_at).await.unwrap().is_none()); + assert!(store.get(&good_id, created_at).await.unwrap().is_some()); + assert!( + store + .get(&recovered_id, created_at) + .await + .unwrap() + .is_some() + ); + } +} diff --git a/lib/crates/fabro-store/src/slate/mod.rs b/lib/crates/fabro-store/src/slate/mod.rs index 8cc1bd41d..a38e623b4 100644 --- a/lib/crates/fabro-store/src/slate/mod.rs +++ b/lib/crates/fabro-store/src/slate/mod.rs @@ -7,7 +7,7 @@ mod run_store; use std::collections::HashMap; use std::path::PathBuf; -use std::sync::Arc; +use std::sync::{Arc, OnceLock}; use std::time::Duration; pub use auth_codes::{AuthCode, AuthCodeStore}; @@ -25,7 +25,7 @@ use slatedb::config::{CompressionCodec, Settings}; use tokio::sync::{Mutex, OnceCell}; use tracing::warn; -use crate::{Error, ListRunsQuery, Result, RunProjection, keys}; +use crate::{Error, ListRunsQuery, Result, RunProjection, RunSummaryStore, keys}; #[derive(Debug, Clone, PartialEq, Eq)] pub struct UnreadableRun { @@ -53,6 +53,7 @@ pub struct Database { refresh_tokens: Arc>>, projection_cache: Arc, projection_cache_warmed: Arc>, + run_summary_store: Arc>>, } impl std::fmt::Debug for Database { @@ -85,9 +86,18 @@ impl Database { refresh_tokens: Arc::new(OnceCell::new()), projection_cache: Arc::new(RunProjectionCache::default()), projection_cache_warmed: Arc::new(OnceCell::new()), + run_summary_store: Arc::new(OnceLock::new()), } } + pub fn attach_run_summary_store(&self, store: Arc) -> Arc { + Arc::clone(self.run_summary_store.get_or_init(|| store)) + } + + fn run_summary_store(&self) -> Option> { + self.run_summary_store.get().cloned() + } + fn shared_db_prefix(&self) -> String { self.base_prefix.clone() } @@ -154,8 +164,13 @@ impl Database { } self.catalog_index().await?.add(run_id).await?; - let run_store = - RunDatabase::open_writer(*run_id, db, Arc::clone(&self.projection_cache)).await?; + let run_store = RunDatabase::open_writer( + *run_id, + db, + Arc::clone(&self.projection_cache), + self.run_summary_store(), + ) + .await?; Self::cache_active_run(&mut active_runs, &run_store); Ok(run_store) } @@ -178,8 +193,13 @@ impl Database { if !RunDatabase::has_any_events(&db, run_id).await? { return Err(Error::RunNotFound(run_id.to_string())); } - let run_store = - RunDatabase::open_writer(*run_id, db, Arc::clone(&self.projection_cache)).await?; + let run_store = RunDatabase::open_writer( + *run_id, + db, + Arc::clone(&self.projection_cache), + self.run_summary_store(), + ) + .await?; Self::cache_active_run(&mut active_runs, &run_store); Ok(run_store) } @@ -197,7 +217,13 @@ impl Database { if !RunDatabase::has_any_events(&db, run_id).await? { return Err(Error::RunNotFound(run_id.to_string())); } - RunDatabase::open_reader(*run_id, db, Arc::clone(&self.projection_cache)).await + RunDatabase::open_reader( + *run_id, + db, + Arc::clone(&self.projection_cache), + self.run_summary_store(), + ) + .await } pub async fn list_runs(&self, query: &ListRunsQuery, now: DateTime) -> Result> { @@ -245,6 +271,9 @@ impl Database { } } } + if let Some(store) = self.run_summary_store() { + store.reconcile(&entries).await?; + } self.projection_cache.replace_all(entries).await; Ok::<_, Error>(()) }) @@ -356,6 +385,9 @@ impl Database { self.delete_session_indexes_for_run(run_id).await?; self.catalog_index().await?.remove(run_id).await?; self.remove_cached_run(run_id).await; + if let Some(store) = self.run_summary_store() { + store.delete(run_id).await?; + } Ok(()) } @@ -527,6 +559,18 @@ mod tests { (object_store, store) } + async fn make_summary_store() -> (tempfile::TempDir, Arc) { + let directory = tempfile::tempdir().unwrap(); + let database = fabro_db::Database::connect(directory.path().join("fabro.sqlite3")) + .await + .unwrap(); + database.migrate().await.unwrap(); + ( + directory, + Arc::new(RunSummaryStore::new(database.clone_pool())), + ) + } + fn sample_run_spec(label: &str) -> RunSpec { let mut graph = Graph::new("night-sky"); graph.attrs.insert( @@ -1325,6 +1369,8 @@ mod tests { #[tokio::test] async fn append_event_refreshes_projection_cache_and_delete_removes_it() { let (_object_store, store) = make_store(); + let (_directory, summaries) = make_summary_store().await; + store.attach_run_summary_store(Arc::clone(&summaries)); let run = store.create_run(&test_run_id("run-1")).await.unwrap(); append_created(&run, "run-1", dt("2026-03-27T12:00:00Z")).await; store.warm_projection_cache().await.unwrap(); @@ -1436,7 +1482,7 @@ mod tests { Some("abc123") ); - let summaries = store + let cached_summaries = store .list_runs(&ListRunsQuery::default(), Utc::now()) .await .unwrap(); @@ -1444,7 +1490,7 @@ mod tests { .list_runs_with_projection(&ListRunsQuery::default(), Utc::now()) .await .unwrap(); - assert_eq!(summaries, vec![cached.summary.clone()]); + assert_eq!(cached_summaries, vec![cached.summary.clone()]); assert_eq!(projected[0].0, cached.summary); assert_eq!( projected[0] @@ -1460,6 +1506,18 @@ mod tests { .git_commit_sha .as_deref() ); + let comparison_time = dt("2026-03-27T12:00:10Z"); + let cache_summary = store + .get_cached_summary(&test_run_id("run-1"), comparison_time) + .await + .unwrap() + .unwrap(); + let sql_summary = summaries + .get(&test_run_id("run-1"), comparison_time) + .await + .unwrap() + .unwrap(); + assert_eq!(sql_summary, cache_summary); store.delete_run(&test_run_id("run-1")).await.unwrap(); assert!( @@ -1476,6 +1534,34 @@ mod tests { .unwrap() .is_empty() ); + assert!( + summaries + .get(&test_run_id("run-1"), Utc::now()) + .await + .unwrap() + .is_none() + ); + } + + #[tokio::test] + async fn projection_cache_warmup_backfills_sqlite_run_summaries() { + let (object_store, store) = make_store(); + let run = store.create_run(&test_run_id("run-1")).await.unwrap(); + append_completed(&run, "run-1", dt("2026-03-27T12:00:00Z")).await; + + let reopened = Database::new(object_store, "runs", Duration::from_millis(1), None); + let (_directory, summaries) = make_summary_store().await; + reopened.attach_run_summary_store(Arc::clone(&summaries)); + reopened.warm_projection_cache().await.unwrap(); + + let summary = summaries + .get(&test_run_id("run-1"), Utc::now()) + .await + .unwrap() + .unwrap(); + assert_eq!(summary.lifecycle.status, RunStatus::Succeeded { + reason: SuccessReason::Completed, + }); } #[tokio::test] diff --git a/lib/crates/fabro-store/src/slate/projection_cache.rs b/lib/crates/fabro-store/src/slate/projection_cache.rs index aa4974c2f..73f8aaf2e 100644 --- a/lib/crates/fabro-store/src/slate/projection_cache.rs +++ b/lib/crates/fabro-store/src/slate/projection_cache.rs @@ -184,25 +184,27 @@ impl RunProjectionCache { Some(entry.summary) } - pub(crate) async fn apply_event(&self, run_id: &RunId, event: &EventEnvelope) -> Result<()> { + pub(crate) async fn apply_event( + &self, + run_id: &RunId, + event: &EventEnvelope, + ) -> Result { let mut state = self.state.lock().await; let Some(entry) = state.entries.get(run_id).cloned() else { if event.seq == 1 { let projection = RunProjection::apply_events(std::slice::from_ref(event))?; - state.insert(CachedRunProjection::from_projection( - *run_id, projection, event.seq, - )); - } else { - return Err(Error::InvalidEvent(format!( - "projection cache cannot initialize run {run_id} from event seq {}", - event.seq - ))); + let entry = CachedRunProjection::from_projection(*run_id, projection, event.seq); + state.insert(entry.clone()); + return Ok(entry); } - return Ok(()); + return Err(Error::InvalidEvent(format!( + "projection cache cannot initialize run {run_id} from event seq {}", + event.seq + ))); }; if event.seq <= entry.last_seq { - return Ok(()); + return Ok(entry); } if event.seq != entry.last_seq.saturating_add(1) { return Err(Error::Other(format!( @@ -213,10 +215,9 @@ impl RunProjectionCache { let mut projection = (*entry.projection).clone(); projection.apply_event(event)?; - state.insert(CachedRunProjection::from_projection( - *run_id, projection, event.seq, - )); - Ok(()) + let entry = CachedRunProjection::from_projection(*run_id, projection, event.seq); + state.insert(entry.clone()); + Ok(entry) } pub(crate) async fn remove(&self, run_id: &RunId) { diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index 94116fa94..3e2cb4db5 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -9,12 +9,14 @@ use futures::Stream; use slatedb::{Db, DbRead}; use tokio::sync::{Mutex, broadcast, mpsc}; use tokio_stream::wrappers::UnboundedReceiverStream; -use tracing::warn; +use tracing::{error, warn}; use super::blob_store::BlobStore; use super::projection_cache::{CachedRunProjection, RunProjectionCache}; use crate::run_state::{EventProjectionCache, RunProjectionReducer}; -use crate::{Error, EventEnvelope, EventPayload, Result, RunProjection, StageId, keys}; +use crate::{ + Error, EventEnvelope, EventPayload, Result, RunProjection, RunSummaryStore, StageId, keys, +}; const DEFAULT_EVENT_TAIL_LIMIT: usize = 1024; #[derive(Clone)] @@ -41,6 +43,7 @@ pub(crate) struct RunDatabaseInner { state_lock: Mutex<()>, projection_cache: Mutex, shared_projection_cache: Arc, + run_summary_store: Option>, recent_events: Mutex>, recent_event_limit: usize, event_tx: broadcast::Sender, @@ -51,16 +54,25 @@ impl RunDatabase { run_id: RunId, db: Db, shared_projection_cache: Arc, + run_summary_store: Option>, ) -> Result { - Self::build(run_id, db, false, shared_projection_cache).await + Self::build( + run_id, + db, + false, + shared_projection_cache, + run_summary_store, + ) + .await } pub(crate) async fn open_reader( run_id: RunId, db: Db, shared_projection_cache: Arc, + run_summary_store: Option>, ) -> Result { - Self::build(run_id, db, true, shared_projection_cache).await + Self::build(run_id, db, true, shared_projection_cache, run_summary_store).await } async fn build( @@ -68,6 +80,7 @@ impl RunDatabase { db: Db, read_only: bool, shared_projection_cache: Arc, + run_summary_store: Option>, ) -> Result { let event_seq = recover_next_seq(&db, keys::run_events_prefix(&run_id), keys::parse_event_seq).await?; @@ -83,6 +96,7 @@ impl RunDatabase { state_lock: Mutex::new(()), projection_cache: Mutex::new(EventProjectionCache::default()), shared_projection_cache, + run_summary_store, recent_events: Mutex::new(VecDeque::with_capacity(DEFAULT_EVENT_TAIL_LIMIT)), recent_event_limit: DEFAULT_EVENT_TAIL_LIMIT, event_tx, @@ -257,43 +271,74 @@ impl RunDatabase { ) .await?; self.cache_event(&event).await?; - if let Err(err) = self + Box::pin(self.update_summary_projection_after_append(&event)).await?; + Ok(event) + } + + async fn update_summary_projection_after_append(&self, event: &EventEnvelope) -> Result<()> { + let cached = match self .inner .shared_projection_cache - .apply_event(&self.inner.run_id, &event) + .apply_event(&self.inner.run_id, event) .await { - match Self::build_cached_projection(&self.inner.db, &self.inner.run_id).await { - Ok(Some(entry)) => { - self.inner.shared_projection_cache.replace(entry).await; - return Ok(event); - } - Ok(None) => { - self.inner - .shared_projection_cache - .remove(&self.inner.run_id) - .await; - } - Err(rebuild_err) => { - self.inner - .shared_projection_cache - .remove(&self.inner.run_id) - .await; - warn!( - run_id = %self.inner.run_id, - error = %rebuild_err, - "Failed to rebuild run projection cache after append" - ); + Ok(entry) => entry, + Err(err) => { + match Self::build_cached_projection(&self.inner.db, &self.inner.run_id).await { + Ok(Some(entry)) => { + self.inner + .shared_projection_cache + .replace(entry.clone()) + .await; + entry + } + Ok(None) => { + self.inner + .shared_projection_cache + .remove(&self.inner.run_id) + .await; + warn!( + run_id = %self.inner.run_id, + error = ?err, + "Failed to update run projection cache after append" + ); + return Err(err); + } + Err(rebuild_err) => { + self.inner + .shared_projection_cache + .remove(&self.inner.run_id) + .await; + warn!( + run_id = %self.inner.run_id, + error = %rebuild_err, + "Failed to rebuild run projection cache after append" + ); + warn!( + run_id = %self.inner.run_id, + error = ?err, + "Failed to update run projection cache after append" + ); + return Err(err); + } } } - warn!( - run_id = %self.inner.run_id, - error = %err, - "Failed to update run projection cache after append" - ); - return Err(err); + }; + if let Some(store) = &self.inner.run_summary_store { + let source_last_seq = cached.last_seq; + let store = Arc::clone(store); + let upsert = Box::pin(async move { store.upsert_projection(&cached).await }); + if let Err(err) = upsert.await { + error!( + run_id = %self.inner.run_id, + source_last_seq, + error = %err, + "Failed to update SQLite run summary after append" + ); + return Err(err); + } } - Ok(event) + Ok(()) } pub async fn list_events(&self) -> Result> { From 2843b33d92fcdcdb45cc02a55166d5d0c4290ba1 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 11 Jul 2026 15:46:51 -0400 Subject: [PATCH 2/2] Simplify runs read model: single-source mappings, leaner queries Consolidate duplicated logic from the SQLite runs read model review: - Derive the status sort CASE and board-column filter from a new RunStatusKind::board_rank(), replacing three hand-maintained copies of the status/column mapping; add a test upserting every status variant so the migration CHECK can't silently drift - Share RunSize bucket thresholds between from_total_usd_micros and the generated size-sort CASE via RunSize::BUCKET_MAX_USD_MICROS - Resolve run selectors from a lean identity query instead of decoding every stored summary per request - Delete the RunsSortKey/RunsSortDirection adapter enums; the store sort enums now carry the wire serde names - Consolidate the workflow display-name fallback chain into WorkflowRef::display_name() (store, CLI, run lookup) - Share pagination clamping and the paginated list envelope across handlers - Reconcile now skips rows whose source seq is unchanged and batch-deletes stale rows; drop the two indexes no query can use - Hold the summary store OnceLock cell in RunDatabaseInner instead of a snapshot so late attachment reaches already-open writers - Misc: expect() on COUNT(*) sign, %err logging, shared wall-time helper, shared SQLite test fixture, dead billing fallback removed Co-Authored-By: Claude Fable 5 --- Cargo.lock | 1 + lib/crates/fabro-cli/src/server_runs.rs | 6 +- .../fabro-db/migrations/2026071104_runs.sql | 2 - lib/crates/fabro-db/tests/sqlite.rs | 2 +- lib/crates/fabro-server/src/server.rs | 25 +- .../src/server/handler/automations.rs | 28 +- .../fabro-server/src/server/handler/runs.rs | 198 +++++-------- lib/crates/fabro-store/Cargo.toml | 1 + lib/crates/fabro-store/src/lib.rs | 6 +- lib/crates/fabro-store/src/run_state.rs | 6 +- .../fabro-store/src/run_summary_store.rs | 266 ++++++++++++------ lib/crates/fabro-store/src/slate/mod.rs | 19 +- lib/crates/fabro-store/src/slate/run_store.rs | 51 ++-- lib/crates/fabro-store/src/test_util.rs | 10 + lib/crates/fabro-types/src/run_projection.rs | 7 +- lib/crates/fabro-types/src/run_summary.rs | 34 ++- lib/crates/fabro-types/src/status.rs | 24 +- lib/crates/fabro-types/src/timing.rs | 8 + lib/crates/fabro-workflow/src/run_lookup.rs | 8 +- 19 files changed, 376 insertions(+), 326 deletions(-) create mode 100644 lib/crates/fabro-store/src/test_util.rs diff --git a/Cargo.lock b/Cargo.lock index 2198d929b..4329fef8b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3126,6 +3126,7 @@ dependencies = [ "serde_json", "slatedb", "sqlx", + "strum 0.28.0", "tempfile", "thiserror 2.0.18", "tokio", diff --git a/lib/crates/fabro-cli/src/server_runs.rs b/lib/crates/fabro-cli/src/server_runs.rs index adb2f9566..a2fe5c88f 100644 --- a/lib/crates/fabro-cli/src/server_runs.rs +++ b/lib/crates/fabro-cli/src/server_runs.rs @@ -38,11 +38,7 @@ impl ServerRunInfo { } pub(crate) fn workflow_display_name(&self) -> String { - self.workflow_name() - .or_else(|| self.workflow_graph_name()) - .or_else(|| self.workflow_slug()) - .unwrap_or("-") - .to_string() + self.run.workflow.display_name().unwrap_or("-").to_string() } pub(crate) fn workflow_matches(&self, pattern: &str) -> bool { diff --git a/lib/crates/fabro-db/migrations/2026071104_runs.sql b/lib/crates/fabro-db/migrations/2026071104_runs.sql index 42d2b7a7d..1dbba7da6 100644 --- a/lib/crates/fabro-db/migrations/2026071104_runs.sql +++ b/lib/crates/fabro-db/migrations/2026071104_runs.sql @@ -53,6 +53,4 @@ CREATE INDEX runs_by_created_at ON runs(created_at_ms DESC, id DESC); CREATE INDEX runs_by_updated_at ON runs(last_event_at_ms DESC, id DESC); CREATE INDEX runs_by_status ON runs(archived_at_ms, status, last_event_at_ms DESC, id DESC); CREATE INDEX runs_by_parent ON runs(parent_id, created_at_ms DESC, id DESC); -CREATE INDEX runs_by_workflow ON runs(workflow_slug, created_at_ms DESC, id DESC); -CREATE INDEX runs_by_repository ON runs(repository_name, created_at_ms DESC, id DESC); CREATE INDEX runs_by_automation ON runs(automation_id, created_at_ms DESC, id DESC); diff --git a/lib/crates/fabro-db/tests/sqlite.rs b/lib/crates/fabro-db/tests/sqlite.rs index e875209b8..bc4981b06 100644 --- a/lib/crates/fabro-db/tests/sqlite.rs +++ b/lib/crates/fabro-db/tests/sqlite.rs @@ -349,7 +349,7 @@ async fn runs_schema_creates_indexes_and_rejects_invalid_rows() -> anyhow::Resul ) .fetch_one(database.pool()) .await?; - assert_eq!(index_count, 7); + assert_eq!(index_count, 5); insert_minimal_run(database.pool(), "submitted", 0, r#"{"id":"run"}"#).await?; for (status, input_tokens, summary_json) in [ diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 1b37763d6..ab4c0a726 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -204,9 +204,17 @@ pub struct PaginationParams { pub offset: u32, } +pub(crate) fn clamp_page_limit(limit: u32) -> u32 { + limit.clamp(1, 100) +} + +pub(crate) fn clamp_page_offset(offset: u32) -> u32 { + offset.min(MAX_PAGE_OFFSET) +} + pub(crate) fn paginate_items(items: Vec, pagination: &PaginationParams) -> (Vec, bool) { - let limit = pagination.limit.clamp(1, 100) as usize; - let offset = pagination.offset.min(MAX_PAGE_OFFSET) as usize; + let limit = clamp_page_limit(pagination.limit) as usize; + let offset = clamp_page_offset(pagination.offset) as usize; let mut data: Vec<_> = items.into_iter().skip(offset).take(limit + 1).collect(); let has_more = data.len() > limit; data.truncate(limit); @@ -219,7 +227,7 @@ pub(crate) struct DfParams { pub(crate) verbose: bool, } -/// Non-paginated list response wrapper with `has_more: false`. +/// List response envelope with pagination metadata. #[derive(serde::Serialize)] pub struct ListResponse { data: T, @@ -227,6 +235,7 @@ pub struct ListResponse { } impl ListResponse { + /// Non-paginated response with `has_more: false`. pub fn new(data: T) -> Self { Self { data, @@ -236,6 +245,16 @@ impl ListResponse { }, } } + + pub fn paginated(data: T, has_more: bool, total: u64) -> Self { + Self { + data, + meta: PaginationMeta { + has_more, + total: i64::try_from(total).ok(), + }, + } + } } /// Snapshot of a managed run. diff --git a/lib/crates/fabro-server/src/server/handler/automations.rs b/lib/crates/fabro-server/src/server/handler/automations.rs index 8d9848678..15fcff772 100644 --- a/lib/crates/fabro-server/src/server/handler/automations.rs +++ b/lib/crates/fabro-server/src/server/handler/automations.rs @@ -2,7 +2,6 @@ use std::sync::Arc; use axum::http::HeaderMap; use axum_extra::extract::Query as ExtraQuery; -use chrono::Utc; use fabro_automation::{ Automation, AutomationDraft, AutomationId, AutomationReplace, AutomationStoreError, }; @@ -11,8 +10,8 @@ use fabro_types::{AutomationRef, RunId}; use serde::Serialize; use super::super::{ - ApiError, AppState, IntoResponse, Json, MAX_PAGE_OFFSET, PaginationParams, Path, RequiredUser, - Response, Router, State, StatusCode, get, + ApiError, AppState, IntoResponse, Json, PaginationParams, Path, RequiredUser, Response, Router, + State, StatusCode, clamp_page_limit, clamp_page_offset, get, }; use super::{json_with_etag_response, lifecycle, parse_required_if_match, runs}; use crate::automation_materializer::AutomationRunMaterializeInput; @@ -84,28 +83,11 @@ async fn list_automation_runs( let query = RunSummaryListQuery { automation_id: Some(id.to_string()), visibility: RunSummaryVisibility::All, - limit: pagination.limit.clamp(1, 100), - offset: pagination.offset.min(MAX_PAGE_OFFSET), + limit: clamp_page_limit(pagination.limit), + offset: clamp_page_offset(pagination.offset), ..RunSummaryListQuery::default() }; - let page = match state.stores.run_summaries.list(&query, Utc::now()).await { - Ok(page) => page, - Err(err) => { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); - } - }; - - let data = state.decorate_run_summaries(page.data).await; - - ( - StatusCode::OK, - Json(serde_json::json!({ - "data": data, - "meta": { "has_more": page.has_more, "total": page.total } - })), - ) - .into_response() + runs::run_summary_page_response(&state, &query).await } async fn create_automation_run( diff --git a/lib/crates/fabro-server/src/server/handler/runs.rs b/lib/crates/fabro-server/src/server/handler/runs.rs index cfdfbc2b0..6ea48942f 100644 --- a/lib/crates/fabro-server/src/server/handler/runs.rs +++ b/lib/crates/fabro-server/src/server/handler/runs.rs @@ -24,20 +24,21 @@ use fabro_store::{ use fabro_types::settings::ResolveError; use fabro_types::{ AutomationRef, Principal, RunClientProvenance, RunId, RunProvenance, RunServerProvenance, - StageContextWindow, StageContextWindowStaleness, StageContextWindowUnavailableReason, - StageHandler, StageModelUsage, StageProjection, SystemActorKind, WorkflowSettings, - parse_blob_ref, + RunStatusKind, StageContextWindow, StageContextWindowStaleness, + StageContextWindowUnavailableReason, StageHandler, StageModelUsage, StageProjection, + SystemActorKind, WorkflowSettings, parse_blob_ref, }; use fabro_util::version::FABRO_VERSION; use fabro_workflow::command_log::{command_log_path, read_json_string_blob, read_log_slice}; use fabro_workflow::run_status::RunStatus; use fabro_workflow::{Error as WorkflowError, operations}; +use strum::VariantArray as _; use tokio::fs; use tracing::info; use super::super::{ - AppState, DeleteRunOutcome, ListResponse, MAX_PAGE_OFFSET, RunExecutionMode, VariableError, - answer_from_request, api_question_from_pending_interview, default_page_limit, + AppState, DeleteRunOutcome, ListResponse, RunExecutionMode, VariableError, answer_from_request, + api_question_from_pending_interview, clamp_page_limit, clamp_page_offset, default_page_limit, delete_run_internal, load_pending_interview, managed_run, parse_run_id_path, parse_stage_id_path, reject_if_archived, submit_pending_interview_answer, workflow_event, }; @@ -88,29 +89,6 @@ pub(super) fn routes() -> Router> { .merge(manifest_routes()) } -#[derive(Debug, Clone, Copy, Default, serde::Deserialize)] -#[serde(rename_all = "snake_case")] -enum RunsSortKey { - #[default] - CreatedAt, - UpdatedAt, - Status, - Elapsed, - Repo, - Title, - Workflow, - Changes, - Size, -} - -#[derive(Debug, Clone, Copy, Default, serde::Deserialize)] -#[serde(rename_all = "lowercase")] -enum RunsSortDirection { - Asc, - #[default] - Desc, -} - #[derive(serde::Deserialize)] struct ListRunsParams { #[serde(rename = "page[limit]", default = "default_page_limit")] @@ -124,9 +102,9 @@ struct ListRunsParams { #[serde(default)] status: Vec, #[serde(default)] - sort: RunsSortKey, + sort: RunSummarySort, #[serde(default)] - direction: RunsSortDirection, + direction: RunSummarySortDirection, } impl ListRunsParams { @@ -134,10 +112,10 @@ impl ListRunsParams { RunSummaryListQuery { parent_id: self.parent_id, visibility: summary_visibility(&self.status, self.include_archived), - sort: summary_sort(self.sort), - direction: summary_sort_direction(self.direction), - limit: self.limit.clamp(1, 100), - offset: self.offset.min(MAX_PAGE_OFFSET), + sort: self.sort, + direction: self.direction, + limit: clamp_page_limit(self.limit), + offset: clamp_page_offset(self.offset), ..RunSummaryListQuery::default() } } @@ -151,35 +129,14 @@ fn summary_visibility(selected: &[BoardColumn], include_archived: bool) -> RunSu let mut statuses = HashSet::new(); let mut archived = false; for column in selected { - match column { - BoardColumn::Pending => { - statuses.insert(fabro_types::RunStatusKind::Submitted); - statuses.insert(fabro_types::RunStatusKind::Pending); - } - BoardColumn::Runnable => { - statuses.insert(fabro_types::RunStatusKind::Runnable); - } - BoardColumn::Initializing => { - statuses.insert(fabro_types::RunStatusKind::Starting); - } - BoardColumn::Running => { - statuses.insert(fabro_types::RunStatusKind::Running); - statuses.insert(fabro_types::RunStatusKind::Paused); - } - BoardColumn::Blocked => { - statuses.insert(fabro_types::RunStatusKind::Blocked); - } - BoardColumn::Succeeded => { - statuses.insert(fabro_types::RunStatusKind::Succeeded); - } - BoardColumn::Failed => { - statuses.insert(fabro_types::RunStatusKind::Failed); - statuses.insert(fabro_types::RunStatusKind::Dead); - } - BoardColumn::Archived => archived = true, - BoardColumn::Removing => { - statuses.insert(fabro_types::RunStatusKind::Removing); - } + match board_column_rank(*column) { + None => archived = true, + Some(rank) => statuses.extend( + RunStatusKind::VARIANTS + .iter() + .copied() + .filter(|kind| kind.board_rank() == rank), + ), } } RunSummaryVisibility::Selected { @@ -188,24 +145,20 @@ fn summary_visibility(selected: &[BoardColumn], include_archived: bool) -> RunSu } } -fn summary_sort(sort: RunsSortKey) -> RunSummarySort { - match sort { - RunsSortKey::CreatedAt => RunSummarySort::CreatedAt, - RunsSortKey::UpdatedAt => RunSummarySort::UpdatedAt, - RunsSortKey::Status => RunSummarySort::Status, - RunsSortKey::Elapsed => RunSummarySort::Elapsed, - RunsSortKey::Repo => RunSummarySort::Repository, - RunsSortKey::Title => RunSummarySort::Title, - RunsSortKey::Workflow => RunSummarySort::Workflow, - RunsSortKey::Changes => RunSummarySort::Changes, - RunsSortKey::Size => RunSummarySort::Size, - } -} - -fn summary_sort_direction(direction: RunsSortDirection) -> RunSummarySortDirection { - match direction { - RunsSortDirection::Asc => RunSummarySortDirection::Asc, - RunsSortDirection::Desc => RunSummarySortDirection::Desc, +/// Rank of each board column, mirroring the `BoardColumn` enum order. +/// Statuses map to columns through [`RunStatusKind::board_rank`]; `archived` +/// has no rank because it selects on the archival overlay, not a status. +fn board_column_rank(column: BoardColumn) -> Option { + match column { + BoardColumn::Pending => Some(0), + BoardColumn::Runnable => Some(1), + BoardColumn::Initializing => Some(2), + BoardColumn::Running => Some(3), + BoardColumn::Blocked => Some(4), + BoardColumn::Succeeded => Some(5), + BoardColumn::Failed => Some(6), + BoardColumn::Archived => None, + BoardColumn::Removing => Some(8), } } @@ -361,29 +314,28 @@ async fn list_runs( State(state): State>, ExtraQuery(params): ExtraQuery, ) -> Response { - let page = match state - .stores - .run_summaries - .list(¶ms.summary_query(), Utc::now()) - .await - { - Ok(page) => page, - Err(err) => { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); + run_summary_page_response(&state, ¶ms.summary_query()).await +} + +/// List run summaries matching `query`, decorate them, and wrap them in the +/// paginated list envelope. Shared by the runs and automation-runs lists. +pub(super) async fn run_summary_page_response( + state: &AppState, + query: &RunSummaryListQuery, +) -> Response { + match state.stores.run_summaries.list(query, Utc::now()).await { + Ok(page) => { + let data = state.decorate_run_summaries(page.data).await; + ( + StatusCode::OK, + Json(ListResponse::paginated(data, page.has_more, page.total)), + ) + .into_response() } - }; - - let data = state.decorate_run_summaries(page.data).await; - - ( - StatusCode::OK, - Json(serde_json::json!({ - "data": data, - "meta": { "has_more": page.has_more, "total": page.total } - })), - ) - .into_response() + Err(err) => { + ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() + } + } } #[derive(Debug, serde::Deserialize)] @@ -425,49 +377,33 @@ async fn resolve_run( State(state): State>, Query(query): Query, ) -> Response { - let summary_query = RunSummaryListQuery { - visibility: RunSummaryVisibility::All, - limit: u32::MAX, - ..RunSummaryListQuery::default() - }; - let runs = match state - .stores - .run_summaries - .list(&summary_query, Utc::now()) - .await - { - Ok(page) => page.data, + let identities = match state.stores.run_summaries.list_identities().await { + Ok(identities) => identities, Err(err) => { return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) .into_response(); } }; - match resolve_run_by_selector( - &runs, + let resolved_id = match resolve_run_by_selector( + &identities, &query.selector, |run| run.id.to_string(), - |run| run.workflow.slug.clone(), - |run| run.workflow.name.clone(), + |run| run.workflow_slug.clone(), + |run| run.workflow_name.clone(), |run| run.id.created_at(), |run| run.id.created_at().to_rfc3339(), - |run| { - run.repository - .as_ref() - .and_then(|repository| repository.origin_url.clone()) - }, + |run| run.repository_origin_url.clone(), ) { - Ok(run) => { - let run = state.decorate_run_summary(run.clone()).await; - (StatusCode::OK, Json(run)).into_response() - } + Ok(identity) => identity.id, Err(err @ (ResolveRunError::InvalidSelector | ResolveRunError::AmbiguousPrefix { .. })) => { - ApiError::bad_request(err.to_string()).into_response() + return ApiError::bad_request(err.to_string()).into_response(); } Err(err @ ResolveRunError::NotFound { .. }) => { - ApiError::not_found(err.to_string()).into_response() + return ApiError::not_found(err.to_string()).into_response(); } - } + }; + updated_run_response(&state, &resolved_id).await } async fn delete_run( diff --git a/lib/crates/fabro-store/Cargo.toml b/lib/crates/fabro-store/Cargo.toml index 1fd4c10dd..234da7ad4 100644 --- a/lib/crates/fabro-store/Cargo.toml +++ b/lib/crates/fabro-store/Cargo.toml @@ -25,6 +25,7 @@ dashmap.workspace = true serde.workspace = true serde_json.workspace = true sqlx.workspace = true +strum.workspace = true chrono = { workspace = true, features = ["serde"] } bytes.workspace = true thiserror.workspace = true diff --git a/lib/crates/fabro-store/src/lib.rs b/lib/crates/fabro-store/src/lib.rs index d33d8aa91..0eb246d04 100644 --- a/lib/crates/fabro-store/src/lib.rs +++ b/lib/crates/fabro-store/src/lib.rs @@ -10,6 +10,8 @@ mod run_state; mod run_summary_store; mod serializable_projection; mod slate; +#[cfg(test)] +mod test_util; mod types; pub use artifact_store::{ @@ -27,8 +29,8 @@ pub use run_sessions::{ }; pub use run_state::RunProjectionReducer; pub use run_summary_store::{ - RunSummaryListQuery, RunSummaryPage, RunSummarySort, RunSummarySortDirection, RunSummaryStore, - RunSummaryVisibility, + RunSummaryIdentity, RunSummaryListQuery, RunSummaryPage, RunSummarySort, + RunSummarySortDirection, RunSummaryStore, RunSummaryVisibility, }; pub use serializable_projection::SerializableProjection; pub use slate::{ diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index e0275a795..ec67aa4f0 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -935,7 +935,7 @@ pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> Run { .as_ref() .map(|conclusion| conclusion.timing); let terminal_total = terminal_total_usd_micros(state); - let current_total = terminal_total.or_else(|| projected_total_usd_micros(state)); + let current_total = projected_billing(state).total_usd_micros; Run { id: *run_id, @@ -1003,10 +1003,6 @@ fn terminal_total_usd_micros(state: &RunProjection) -> Option { .and_then(|billing| billing.total_usd_micros) } -fn projected_total_usd_micros(state: &RunProjection) -> Option { - projected_billing(state).total_usd_micros -} - pub(crate) fn projected_billing(state: &RunProjection) -> BilledTokenCounts { if let Some(billing) = state .conclusion diff --git a/lib/crates/fabro-store/src/run_summary_store.rs b/lib/crates/fabro-store/src/run_summary_store.rs index 289530d1c..15a54adc8 100644 --- a/lib/crates/fabro-store/src/run_summary_store.rs +++ b/lib/crates/fabro-store/src/run_summary_store.rs @@ -1,9 +1,12 @@ -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; +use std::fmt::Write as _; +use std::sync::LazyLock; use chrono::{DateTime, Utc}; -use fabro_types::{Run, RunId, RunStatusKind, RunTiming}; +use fabro_types::{Run, RunId, RunSize, RunStatusKind, RunTiming}; use sqlx::sqlite::{SqliteConnection, SqliteRow}; use sqlx::{QueryBuilder, Row as _, Sqlite, SqlitePool}; +use strum::VariantArray as _; use crate::run_state::projected_billing; use crate::slate::CachedRunProjection; @@ -46,13 +49,20 @@ ON CONFLICT(id) DO UPDATE SET WHERE excluded.source_last_seq > runs.source_last_seq "; -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +const SELECT_RUN_SUMMARIES_SQL: &str = r" +SELECT runs.id, runs.summary_json, + (SELECT COUNT(*) FROM runs AS child WHERE child.parent_id = runs.id) AS children_count +FROM runs"; + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Deserialize)] +#[serde(rename_all = "snake_case")] pub enum RunSummarySort { #[default] CreatedAt, UpdatedAt, Status, Elapsed, + #[serde(rename = "repo")] Repository, Title, Workflow, @@ -60,7 +70,8 @@ pub enum RunSummarySort { Size, } -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Deserialize)] +#[serde(rename_all = "lowercase")] pub enum RunSummarySortDirection { Asc, #[default] @@ -144,46 +155,84 @@ impl RunSummaryStore { } pub(crate) async fn reconcile(&self, entries: &[CachedRunProjection]) -> Result<()> { - let authoritative_ids = entries - .iter() - .map(|entry| entry.run_id.to_string()) - .collect::>(); let mut transaction = self.pool.begin().await?; + let stored_seqs: HashMap = + sqlx::query_as::<_, (String, i64)>("SELECT id, source_last_seq FROM runs") + .fetch_all(&mut *transaction) + .await? + .into_iter() + .collect(); + + let mut authoritative_ids = HashSet::new(); for entry in entries { - let record = ProjectedRunSummary::from_entry(entry); - upsert_run(&mut transaction, &record).await?; + let run_id = entry.run_id.to_string(); + let up_to_date = stored_seqs + .get(&run_id) + .is_some_and(|stored_seq| *stored_seq >= i64::from(entry.last_seq)); + authoritative_ids.insert(run_id); + if up_to_date { + continue; + } + upsert_run(&mut transaction, &ProjectedRunSummary::from_entry(entry)).await?; } - let stored_ids = sqlx::query_scalar::<_, String>("SELECT id FROM runs") - .fetch_all(&mut *transaction) - .await?; - for stored_id in stored_ids { - if !authoritative_ids.contains(&stored_id) { - sqlx::query("DELETE FROM runs WHERE id = ?") - .bind(stored_id) - .execute(&mut *transaction) - .await?; + let stale_ids = stored_seqs + .keys() + .filter(|stored_id| !authoritative_ids.contains(stored_id.as_str())) + .collect::>(); + for chunk in stale_ids.chunks(500) { + let mut delete = QueryBuilder::::new("DELETE FROM runs WHERE id IN ("); + let mut separated = delete.separated(", "); + for stale_id in chunk { + separated.push_bind(stale_id.as_str()); } + delete.push(")"); + delete.build().execute(&mut *transaction).await?; } transaction.commit().await?; Ok(()) } pub async fn get(&self, run_id: &RunId, now: DateTime) -> Result> { - let row = sqlx::query( - r" -SELECT runs.id, runs.summary_json, - (SELECT COUNT(*) FROM runs AS child WHERE child.parent_id = runs.id) AS children_count -FROM runs -WHERE runs.id = ? -", - ) - .bind(run_id.to_string()) - .fetch_optional(&self.pool) - .await?; + let mut query = QueryBuilder::::new(SELECT_RUN_SUMMARIES_SQL); + query + .push(" WHERE runs.id = ") + .push_bind(run_id.to_string()); + let row = query.build().fetch_optional(&self.pool).await?; row.map(|row| decode_run_row(&row, now)).transpose() } + /// Identity fields for every stored run, for selector resolution without + /// decoding full summaries. + pub async fn list_identities(&self) -> Result> { + let rows = sqlx::query( + r" +SELECT id, workflow_slug, + json_extract(summary_json, '$.workflow.name') AS workflow_name, + json_extract(summary_json, '$.repository.origin_url') AS repository_origin_url +FROM runs", + ) + .fetch_all(&self.pool) + .await?; + rows.iter() + .map(|row| { + let stored_id: String = row.try_get("id")?; + let id = stored_id + .parse::() + .map_err(|_| Error::RunSummaryMismatch { + run_id: stored_id, + field: "id", + })?; + Ok(RunSummaryIdentity { + id, + workflow_slug: row.try_get("workflow_slug")?, + workflow_name: row.try_get("workflow_name")?, + repository_origin_url: row.try_get("repository_origin_url")?, + }) + }) + .collect() + } + pub async fn list( &self, query: &RunSummaryListQuery, @@ -198,12 +247,7 @@ WHERE runs.id = ? .fetch_one(&mut *transaction) .await?; - let mut rows_query = QueryBuilder::::new( - r" -SELECT runs.id, runs.summary_json, - (SELECT COUNT(*) FROM runs AS child WHERE child.parent_id = runs.id) AS children_count -FROM runs", - ); + let mut rows_query = QueryBuilder::::new(SELECT_RUN_SUMMARIES_SQL); push_filters(&mut rows_query, query); push_order(&mut rows_query, query.sort, query.direction, now); rows_query.push(" LIMIT ").push_bind(i64::from(query.limit)); @@ -217,10 +261,7 @@ FROM runs", .iter() .map(|row| decode_run_row(row, now)) .collect::>>()?; - let total = u64::try_from(total).map_err(|_| Error::RunSummaryMismatch { - run_id: "".to_string(), - field: "negative total", - })?; + let total = u64::try_from(total).expect("COUNT(*) is non-negative"); let consumed = u64::from(query.offset).saturating_add(data.len() as u64); Ok(RunSummaryPage { data, @@ -238,6 +279,16 @@ FROM runs", } } +/// Identity fields of a stored run summary, cheap to list for selector +/// resolution. +#[derive(Debug, Clone)] +pub struct RunSummaryIdentity { + pub id: RunId, + pub workflow_slug: Option, + pub workflow_name: Option, + pub repository_origin_url: Option, +} + #[derive(Debug)] struct ProjectedRunSummary { run: Run, @@ -263,12 +314,7 @@ impl ProjectedRunSummary { run.timing = entry.projection.live_run_timing(at); } let billing = projected_billing(&entry.projection); - let workflow_name = run - .workflow - .name - .clone() - .or_else(|| run.workflow.graph_name.clone()) - .or_else(|| run.workflow.slug.clone()); + let workflow_name = run.workflow.display_name().map(str::to_string); let repository_name = run .repository .as_ref() @@ -356,12 +402,13 @@ fn push_filters(builder: &mut QueryBuilder, query: &RunSummaryListQuery) match &query.visibility { RunSummaryVisibility::All => {} RunSummaryVisibility::Default { include_archived } => { + let not_removing = format!("status <> '{}'", RunStatusKind::Removing); if *include_archived { - builder.push( - " AND (archived_at_ms IS NOT NULL OR (archived_at_ms IS NULL AND status <> 'removing'))", - ); + builder.push(format!( + " AND (archived_at_ms IS NOT NULL OR {not_removing})" + )); } else { - builder.push(" AND archived_at_ms IS NULL AND status <> 'removing'"); + builder.push(format!(" AND archived_at_ms IS NULL AND {not_removing}")); } } RunSummaryVisibility::Selected { statuses, archived } => { @@ -391,6 +438,32 @@ fn push_filters(builder: &mut QueryBuilder, query: &RunSummaryListQuery) } } +/// Status sort rank derived from [`RunStatusKind::board_rank`], so the SQL +/// order and the board column order share one source. Archived runs rank 7, +/// matching the `archived` board column. +static STATUS_RANK_CASE_SQL: LazyLock = LazyLock::new(|| { + let mut case = String::from("CASE WHEN archived_at_ms IS NOT NULL THEN 7"); + for kind in RunStatusKind::VARIANTS { + let _ = write!(case, " WHEN status = '{kind}' THEN {}", kind.board_rank()); + } + case.push_str(" ELSE 9 END"); + case +}); + +/// Size sort rank derived from [`RunSize::BUCKET_MAX_USD_MICROS`], so the SQL +/// order and the displayed size buckets share one source. +static SIZE_RANK_CASE_SQL: LazyLock = LazyLock::new(|| { + let mut case = String::from("CASE"); + for (rank, (_, max_usd_micros)) in RunSize::BUCKET_MAX_USD_MICROS.iter().enumerate() { + let _ = write!( + case, + " WHEN COALESCE(total_usd_micros, 0) <= {max_usd_micros} THEN {rank}" + ); + } + let _ = write!(case, " ELSE {} END", RunSize::BUCKET_MAX_USD_MICROS.len()); + case +}); + fn push_order( builder: &mut QueryBuilder, sort: RunSummarySort, @@ -401,20 +474,7 @@ fn push_order( match sort { RunSummarySort::CreatedAt => builder.push("created_at_ms"), RunSummarySort::UpdatedAt => builder.push("last_event_at_ms"), - RunSummarySort::Status => builder.push( - r"CASE - WHEN archived_at_ms IS NOT NULL THEN 7 - WHEN status IN ('submitted', 'pending') THEN 0 - WHEN status = 'runnable' THEN 1 - WHEN status = 'starting' THEN 2 - WHEN status IN ('running', 'paused') THEN 3 - WHEN status = 'blocked' THEN 4 - WHEN status = 'succeeded' THEN 5 - WHEN status IN ('failed', 'dead') THEN 6 - WHEN status = 'removing' THEN 8 - ELSE 9 - END", - ), + RunSummarySort::Status => builder.push(STATUS_RANK_CASE_SQL.as_str()), RunSummarySort::Elapsed => builder .push("(COALESCE(completed_at_ms, ") .push_bind(now.timestamp_millis()) @@ -423,15 +483,7 @@ fn push_order( RunSummarySort::Title => builder.push("TRIM(title) COLLATE NOCASE"), RunSummarySort::Workflow => builder.push("COALESCE(workflow_name, '') COLLATE NOCASE"), RunSummarySort::Changes => builder.push("(diff_additions + diff_deletions)"), - RunSummarySort::Size => builder.push( - r"CASE - WHEN COALESCE(total_usd_micros, 0) <= 20000000 THEN 0 - WHEN total_usd_micros <= 50000000 THEN 1 - WHEN total_usd_micros <= 100000000 THEN 2 - WHEN total_usd_micros <= 200000000 THEN 3 - ELSE 4 - END", - ), + RunSummarySort::Size => builder.push(SIZE_RANK_CASE_SQL.as_str()), }; match direction { RunSummarySortDirection::Asc => builder.push(" ASC"), @@ -455,23 +507,18 @@ fn decode_run_row(row: &SqliteRow, now: DateTime) -> Result { run_id: run.id.to_string(), field: "children_count", })?; - apply_read_overlays(&mut run, now); + overlay_live_wall_time(&mut run, now); Ok(run) } -fn apply_read_overlays(run: &mut Run, now: DateTime) { +fn overlay_live_wall_time(run: &mut Run, now: DateTime) { if run.timestamps.completed_at.is_some() { return; } let Some(started_at) = run.timestamps.started_at else { return; }; - let wall_time_ms = u64::try_from( - now.signed_duration_since(started_at) - .num_milliseconds() - .max(0), - ) - .expect("non-negative milliseconds fit in u64"); + let wall_time_ms = RunTiming::wall_time_ms_since(started_at, now); run.timing = Some( run.timing .unwrap_or_else(|| RunTiming::wall_only(wall_time_ms)) @@ -485,10 +532,11 @@ mod tests { use chrono::{DateTime, Utc}; use fabro_types::{ - AutomationRef, BilledTokenCounts, Conclusion, DiffSummary, Graph, RunDiff, RunId, - RunProjection, RunSize, RunSpec, RunStatus, RunTiming, StageOutcome, SuccessReason, - WorkflowSettings, test_support, + AutomationRef, BilledTokenCounts, BlockedReason, Conclusion, DiffSummary, FailureReason, + Graph, PendingReason, RunDiff, RunId, RunProjection, RunSize, RunSpec, RunStatus, + RunStatusKind, RunTiming, StageOutcome, SuccessReason, WorkflowSettings, test_support, }; + use strum::VariantArray as _; use ulid::Ulid; use super::{ @@ -496,6 +544,7 @@ mod tests { RunSummaryVisibility, }; use crate::slate::CachedRunProjection; + use crate::test_util; fn dt(value: &str) -> DateTime { value.parse().unwrap() @@ -532,12 +581,49 @@ mod tests { } async fn store() -> (tempfile::TempDir, RunSummaryStore) { - let directory = tempfile::tempdir().unwrap(); - let database = fabro_db::Database::connect(directory.path().join("fabro.sqlite3")) - .await - .unwrap(); - database.migrate().await.unwrap(); - (directory, RunSummaryStore::new(database.clone_pool())) + test_util::sqlite_summary_store().await + } + + fn sample_status(kind: RunStatusKind) -> RunStatus { + match kind { + RunStatusKind::Submitted => RunStatus::Submitted, + RunStatusKind::Pending => RunStatus::Pending { + reason: PendingReason::ApprovalRequired, + }, + RunStatusKind::Runnable => RunStatus::Runnable, + RunStatusKind::Starting => RunStatus::Starting, + RunStatusKind::Running => RunStatus::Running, + RunStatusKind::Blocked => RunStatus::Blocked { + blocked_reason: BlockedReason::HumanInputRequired, + }, + RunStatusKind::Paused => RunStatus::Paused { prior_block: None }, + RunStatusKind::Removing => RunStatus::Removing, + RunStatusKind::Succeeded => RunStatus::Succeeded { + reason: SuccessReason::Completed, + }, + RunStatusKind::Failed => RunStatus::Failed { + reason: FailureReason::WorkflowError, + }, + RunStatusKind::Dead => RunStatus::Dead, + } + } + + /// The migration's `CHECK (status IN (...))` freezes the status strings; + /// prove every `RunStatusKind` variant passes it so an enum change that + /// forgets a follow-up migration fails in CI instead of at runtime. + #[tokio::test] + async fn every_status_kind_upserts_within_schema_check() { + let (_directory, store) = store().await; + let created_at = dt("2026-07-11T12:00:00Z"); + for (index, kind) in RunStatusKind::VARIANTS.iter().enumerate() { + let id = run_id( + created_at.timestamp_millis().cast_unsigned(), + u128::try_from(index).unwrap() + 1, + ); + let mut projected = projection(id, "status", created_at); + projected.status = sample_status(*kind); + store.upsert_projection(&entry(projected, 1)).await.unwrap(); + } } #[tokio::test] diff --git a/lib/crates/fabro-store/src/slate/mod.rs b/lib/crates/fabro-store/src/slate/mod.rs index a38e623b4..60c40d1c4 100644 --- a/lib/crates/fabro-store/src/slate/mod.rs +++ b/lib/crates/fabro-store/src/slate/mod.rs @@ -168,7 +168,7 @@ impl Database { *run_id, db, Arc::clone(&self.projection_cache), - self.run_summary_store(), + Arc::clone(&self.run_summary_store), ) .await?; Self::cache_active_run(&mut active_runs, &run_store); @@ -197,7 +197,7 @@ impl Database { *run_id, db, Arc::clone(&self.projection_cache), - self.run_summary_store(), + Arc::clone(&self.run_summary_store), ) .await?; Self::cache_active_run(&mut active_runs, &run_store); @@ -221,7 +221,7 @@ impl Database { *run_id, db, Arc::clone(&self.projection_cache), - self.run_summary_store(), + Arc::clone(&self.run_summary_store), ) .await } @@ -511,7 +511,7 @@ mod tests { use object_store::path::Path; use super::*; - use crate::{EventPayload, keys}; + use crate::{EventPayload, keys, test_util}; fn dt(value: &str) -> DateTime { value.parse().unwrap() @@ -560,15 +560,8 @@ mod tests { } async fn make_summary_store() -> (tempfile::TempDir, Arc) { - let directory = tempfile::tempdir().unwrap(); - let database = fabro_db::Database::connect(directory.path().join("fabro.sqlite3")) - .await - .unwrap(); - database.migrate().await.unwrap(); - ( - directory, - Arc::new(RunSummaryStore::new(database.clone_pool())), - ) + let (directory, store) = test_util::sqlite_summary_store().await; + (directory, Arc::new(store)) } fn sample_run_spec(label: &str) -> RunSpec { diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index 3e2cb4db5..2b67a066a 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -1,6 +1,6 @@ use std::collections::VecDeque; -use std::sync::Arc; use std::sync::atomic::{AtomicU32, Ordering}; +use std::sync::{Arc, OnceLock}; use bytes::Bytes; use chrono::Utc; @@ -43,7 +43,9 @@ pub(crate) struct RunDatabaseInner { state_lock: Mutex<()>, projection_cache: Mutex, shared_projection_cache: Arc, - run_summary_store: Option>, + // Shared cell rather than a snapshot so a summary store attached after + // this writer opened is still picked up by later appends. + run_summary_store: Arc>>, recent_events: Mutex>, recent_event_limit: usize, event_tx: broadcast::Sender, @@ -54,7 +56,7 @@ impl RunDatabase { run_id: RunId, db: Db, shared_projection_cache: Arc, - run_summary_store: Option>, + run_summary_store: Arc>>, ) -> Result { Self::build( run_id, @@ -70,7 +72,7 @@ impl RunDatabase { run_id: RunId, db: Db, shared_projection_cache: Arc, - run_summary_store: Option>, + run_summary_store: Arc>>, ) -> Result { Self::build(run_id, db, true, shared_projection_cache, run_summary_store).await } @@ -80,7 +82,7 @@ impl RunDatabase { db: Db, read_only: bool, shared_projection_cache: Arc, - run_summary_store: Option>, + run_summary_store: Arc>>, ) -> Result { let event_seq = recover_next_seq(&db, keys::run_events_prefix(&run_id), keys::parse_event_seq).await?; @@ -271,6 +273,8 @@ impl RunDatabase { ) .await?; self.cache_event(&event).await?; + // Box::pin keeps append_event_envelope's future small enough for the + // clippy::large_futures budget of its many callers. Box::pin(self.update_summary_projection_after_append(&event)).await?; Ok(event) } @@ -292,31 +296,21 @@ impl RunDatabase { .await; entry } - Ok(None) => { + rebuild => { self.inner .shared_projection_cache .remove(&self.inner.run_id) .await; + if let Err(rebuild_err) = rebuild { + warn!( + run_id = %self.inner.run_id, + error = %rebuild_err, + "Failed to rebuild run projection cache after append" + ); + } warn!( run_id = %self.inner.run_id, - error = ?err, - "Failed to update run projection cache after append" - ); - return Err(err); - } - Err(rebuild_err) => { - self.inner - .shared_projection_cache - .remove(&self.inner.run_id) - .await; - warn!( - run_id = %self.inner.run_id, - error = %rebuild_err, - "Failed to rebuild run projection cache after append" - ); - warn!( - run_id = %self.inner.run_id, - error = ?err, + error = %err, "Failed to update run projection cache after append" ); return Err(err); @@ -324,14 +318,11 @@ impl RunDatabase { } } }; - if let Some(store) = &self.inner.run_summary_store { - let source_last_seq = cached.last_seq; - let store = Arc::clone(store); - let upsert = Box::pin(async move { store.upsert_projection(&cached).await }); - if let Err(err) = upsert.await { + if let Some(store) = self.inner.run_summary_store.get() { + if let Err(err) = store.upsert_projection(&cached).await { error!( run_id = %self.inner.run_id, - source_last_seq, + source_last_seq = cached.last_seq, error = %err, "Failed to update SQLite run summary after append" ); diff --git a/lib/crates/fabro-store/src/test_util.rs b/lib/crates/fabro-store/src/test_util.rs new file mode 100644 index 000000000..bbc0b8717 --- /dev/null +++ b/lib/crates/fabro-store/src/test_util.rs @@ -0,0 +1,10 @@ +use crate::RunSummaryStore; + +pub(crate) async fn sqlite_summary_store() -> (tempfile::TempDir, RunSummaryStore) { + let directory = tempfile::tempdir().unwrap(); + let database = fabro_db::Database::connect(directory.path().join("fabro.sqlite3")) + .await + .unwrap(); + database.migrate().await.unwrap(); + (directory, RunSummaryStore::new(database.clone_pool())) +} diff --git a/lib/crates/fabro-types/src/run_projection.rs b/lib/crates/fabro-types/src/run_projection.rs index 40130ed9c..6b67282ee 100644 --- a/lib/crates/fabro-types/src/run_projection.rs +++ b/lib/crates/fabro-types/src/run_projection.rs @@ -616,12 +616,7 @@ impl RunProjection { #[must_use] pub fn live_run_timing(&self, now: DateTime) -> Option { let start = self.start.as_ref()?; - let wall_time_ms = u64::try_from( - now.signed_duration_since(start.start_time) - .num_milliseconds() - .max(0), - ) - .expect("non-negative milliseconds fit in u64"); + let wall_time_ms = RunTiming::wall_time_ms_since(start.start_time, now); let active = self .stages .values() diff --git a/lib/crates/fabro-types/src/run_summary.rs b/lib/crates/fabro-types/src/run_summary.rs index 3c4321052..0966df360 100644 --- a/lib/crates/fabro-types/src/run_summary.rs +++ b/lib/crates/fabro-types/src/run_summary.rs @@ -101,6 +101,18 @@ pub struct WorkflowRef { pub edge_count: i64, } +impl WorkflowRef { + /// Best available human-facing workflow name: explicit name, then graph + /// name, then slug. + #[must_use] + pub fn display_name(&self) -> Option<&str> { + self.name + .as_deref() + .or(self.graph_name.as_deref()) + .or(self.slug.as_deref()) + } +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct AutomationRef { pub id: String, @@ -233,15 +245,23 @@ pub enum RunSize { } impl RunSize { + /// Inclusive upper bounds in USD micros for each bucket below [`Self::Xl`], + /// ordered smallest to largest. Shared with the SQLite size sort so both + /// stay in step. + pub const BUCKET_MAX_USD_MICROS: [(Self, i64); 4] = [ + (Self::Xs, 20_000_000), + (Self::S, 50_000_000), + (Self::M, 100_000_000), + (Self::L, 200_000_000), + ]; + #[must_use] pub fn from_total_usd_micros(total_usd_micros: Option) -> Self { - match total_usd_micros.unwrap_or(0) { - ..=20_000_000 => Self::Xs, - 20_000_001..=50_000_000 => Self::S, - 50_000_001..=100_000_000 => Self::M, - 100_000_001..=200_000_000 => Self::L, - _ => Self::Xl, - } + let total = total_usd_micros.unwrap_or(0); + Self::BUCKET_MAX_USD_MICROS + .iter() + .find(|(_, max)| total <= *max) + .map_or(Self::Xl, |(size, _)| *size) } } diff --git a/lib/crates/fabro-types/src/status.rs b/lib/crates/fabro-types/src/status.rs index 3ca84723f..770f6aa5d 100644 --- a/lib/crates/fabro-types/src/status.rs +++ b/lib/crates/fabro-types/src/status.rs @@ -1,7 +1,7 @@ use std::fmt; use serde::{Deserialize, Serialize}; -use strum::{Display, EnumString, IntoStaticStr}; +use strum::{Display, EnumString, IntoStaticStr, VariantArray}; #[derive( Debug, @@ -15,6 +15,7 @@ use strum::{Display, EnumString, IntoStaticStr}; Display, EnumString, IntoStaticStr, + VariantArray, )] #[serde(rename_all = "snake_case")] #[strum(serialize_all = "snake_case")] @@ -32,6 +33,27 @@ pub enum RunStatusKind { Dead, } +impl RunStatusKind { + /// Position of this status in the run-board column order. Ranks mirror + /// the API `BoardColumn` enum: pending 0, runnable 1, initializing 2, + /// running 3, blocked 4, succeeded 5, failed 6, archived 7, removing 8. + /// Rank 7 is reserved for archived runs, which is an overlay flag rather + /// than a status. + #[must_use] + pub fn board_rank(self) -> u8 { + match self { + Self::Submitted | Self::Pending => 0, + Self::Runnable => 1, + Self::Starting => 2, + Self::Running | Self::Paused => 3, + Self::Blocked => 4, + Self::Succeeded => 5, + Self::Failed | Self::Dead => 6, + Self::Removing => 8, + } + } +} + #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum RunStatus { diff --git a/lib/crates/fabro-types/src/timing.rs b/lib/crates/fabro-types/src/timing.rs index 49fa57b06..0d1e555e8 100644 --- a/lib/crates/fabro-types/src/timing.rs +++ b/lib/crates/fabro-types/src/timing.rs @@ -15,6 +15,7 @@ //! child branches carry their own work timing; run-level active time sums work //! across stage visits and can exceed run wall time when work runs in parallel. +use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; /// Timing breakdown for one stage visit. @@ -135,6 +136,13 @@ impl RunTiming { ..self } } + + /// Milliseconds elapsed from `start` to `now`, clamped at zero. + #[must_use] + pub fn wall_time_ms_since(start: DateTime, now: DateTime) -> u64 { + u64::try_from(now.signed_duration_since(start).num_milliseconds().max(0)) + .expect("non-negative milliseconds fit in u64") + } } impl From for RunTiming { diff --git a/lib/crates/fabro-workflow/src/run_lookup.rs b/lib/crates/fabro-workflow/src/run_lookup.rs index 01788b89e..167b234cb 100644 --- a/lib/crates/fabro-workflow/src/run_lookup.rs +++ b/lib/crates/fabro-workflow/src/run_lookup.rs @@ -83,13 +83,7 @@ impl RunInfo { pub fn workflow_display_name(&self) -> String { self.summary.as_ref().map_or_else( || "[no run spec]".to_string(), - |_| { - self.workflow_name() - .or_else(|| self.workflow_graph_name()) - .or_else(|| self.workflow_slug()) - .unwrap_or("-") - .to_string() - }, + |summary| summary.workflow.display_name().unwrap_or("-").to_string(), ) }