mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
Merge pull request #572 from fabro-sh/sqlite-runs-read-model
Add SQLite runs read model
This commit is contained in:
commit
8e8aef854c
23 changed files with 1490 additions and 370 deletions
3
Cargo.lock
generated
3
Cargo.lock
generated
|
|
@ -3114,6 +3114,7 @@ dependencies = [
|
|||
"bytes",
|
||||
"chrono",
|
||||
"dashmap",
|
||||
"fabro-db",
|
||||
"fabro-types",
|
||||
"fabro-util",
|
||||
"futures",
|
||||
|
|
@ -3124,6 +3125,8 @@ dependencies = [
|
|||
"serde",
|
||||
"serde_json",
|
||||
"slatedb",
|
||||
"sqlx",
|
||||
"strum 0.28.0",
|
||||
"tempfile",
|
||||
"thiserror 2.0.18",
|
||||
"tokio",
|
||||
|
|
|
|||
68
docs/plans/2026-07-11-sqlite-runs-read-model-plan.md
Normal file
68
docs/plans/2026-07-11-sqlite-runs-read-model-plan.md
Normal file
|
|
@ -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?
|
||||
|
|
@ -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 {
|
||||
|
|
|
|||
56
lib/crates/fabro-db/migrations/2026071104_runs.sql
Normal file
56
lib/crates/fabro-db/migrations/2026071104_runs.sql
Normal file
|
|
@ -0,0 +1,56 @@
|
|||
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_automation ON runs(automation_id, created_at_ms DESC, id DESC);
|
||||
|
|
@ -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, 5);
|
||||
|
||||
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()?;
|
||||
|
|
|
|||
|
|
@ -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!(
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
@ -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<T>(items: Vec<T>, pagination: &PaginationParams) -> (Vec<T>, 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<T: serde::Serialize> {
|
||||
data: T,
|
||||
|
|
@ -227,6 +235,7 @@ pub struct ListResponse<T: serde::Serialize> {
|
|||
}
|
||||
|
||||
impl<T: serde::Serialize> ListResponse<T> {
|
||||
/// Non-paginated response with `has_more: false`.
|
||||
pub fn new(data: T) -> Self {
|
||||
Self {
|
||||
data,
|
||||
|
|
@ -236,6 +245,16 @@ impl<T: serde::Serialize> ListResponse<T> {
|
|||
},
|
||||
}
|
||||
}
|
||||
|
||||
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.
|
||||
|
|
@ -1107,12 +1126,13 @@ pub struct AppState {
|
|||
}
|
||||
|
||||
pub(crate) struct AppStores {
|
||||
pub(crate) runs: Arc<Database>,
|
||||
pub(crate) automations: Arc<AutomationStore>,
|
||||
pub(crate) environments: Arc<EnvironmentStore>,
|
||||
pub(crate) mcp_servers: Arc<McpServerStore>,
|
||||
pub(crate) vault: Arc<SecretStore>,
|
||||
pub(crate) variables: Arc<VariableStore>,
|
||||
pub(crate) runs: Arc<Database>,
|
||||
pub(crate) run_summaries: Arc<RunSummaryStore>,
|
||||
pub(crate) automations: Arc<AutomationStore>,
|
||||
pub(crate) environments: Arc<EnvironmentStore>,
|
||||
pub(crate) mcp_servers: Arc<McpServerStore>,
|
||||
pub(crate) vault: Arc<SecretStore>,
|
||||
pub(crate) variables: Arc<VariableStore>,
|
||||
}
|
||||
|
||||
type PullRequestCreateLocks = Arc<Mutex<HashMap<RunId, Arc<AsyncMutex<()>>>>>;
|
||||
|
|
@ -2386,6 +2406,8 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
|
|||
})
|
||||
.context("load environments")?,
|
||||
);
|
||||
let run_summaries =
|
||||
store.attach_run_summary_store(Arc::new(RunSummaryStore::new(db_pool.clone())));
|
||||
let mcp_server_dir = mcp_server_dir_for_active_config(&active_config_path);
|
||||
let mcp_server_pool = db_pool.clone();
|
||||
let mcp_server_store = Arc::new(
|
||||
|
|
@ -2491,6 +2513,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
|
|||
aggregate_billing: Mutex::new(BillingAccumulator::default()),
|
||||
stores: AppStores {
|
||||
runs: store,
|
||||
run_summaries,
|
||||
automations: automation_store,
|
||||
environments: environment_store,
|
||||
mcp_servers: mcp_server_store,
|
||||
|
|
|
|||
|
|
@ -2,16 +2,16 @@ 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,
|
||||
};
|
||||
use fabro_store::{RunSummaryListQuery, RunSummaryVisibility};
|
||||
use fabro_types::{AutomationRef, RunId};
|
||||
use serde::Serialize;
|
||||
|
||||
use super::super::{
|
||||
ApiError, AppState, IntoResponse, Json, PaginationParams, Path, RequiredUser, Response, Router,
|
||||
State, StatusCode, get, paginate_items,
|
||||
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;
|
||||
|
|
@ -80,47 +80,14 @@ async fn list_automation_runs(
|
|||
Err(err) => 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,
|
||||
Err(err) => {
|
||||
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
||||
.into_response();
|
||||
}
|
||||
let query = RunSummaryListQuery {
|
||||
automation_id: Some(id.to_string()),
|
||||
visibility: RunSummaryVisibility::All,
|
||||
limit: clamp_page_limit(pagination.limit),
|
||||
offset: clamp_page_offset(pagination.offset),
|
||||
..RunSummaryListQuery::default()
|
||||
};
|
||||
|
||||
let mut runs: Vec<fabro_types::Run> = 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;
|
||||
|
||||
(
|
||||
StatusCode::OK,
|
||||
Json(serde_json::json!({
|
||||
"data": data,
|
||||
"meta": { "has_more": has_more, "total": total }
|
||||
})),
|
||||
)
|
||||
.into_response()
|
||||
runs::run_summary_page_response(&state, &query).await
|
||||
}
|
||||
|
||||
async fn create_automation_run(
|
||||
|
|
|
|||
|
|
@ -11,31 +11,35 @@ 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,
|
||||
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, PaginationParams, 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,
|
||||
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,
|
||||
};
|
||||
use crate::error::ApiError;
|
||||
|
|
@ -85,29 +89,6 @@ pub(super) fn routes() -> Router<Arc<AppState>> {
|
|||
.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")]
|
||||
|
|
@ -121,113 +102,64 @@ struct ListRunsParams {
|
|||
#[serde(default)]
|
||||
status: Vec<BoardColumn>,
|
||||
#[serde(default)]
|
||||
sort: RunsSortKey,
|
||||
sort: RunSummarySort,
|
||||
#[serde(default)]
|
||||
direction: RunsSortDirection,
|
||||
direction: RunSummarySortDirection,
|
||||
}
|
||||
|
||||
impl ListRunsParams {
|
||||
fn pagination(&self) -> PaginationParams {
|
||||
PaginationParams {
|
||||
limit: self.limit,
|
||||
offset: self.offset,
|
||||
}
|
||||
}
|
||||
|
||||
fn status_filter(&self) -> Option<HashSet<BoardColumn>> {
|
||||
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: self.sort,
|
||||
direction: self.direction,
|
||||
limit: clamp_page_limit(self.limit),
|
||||
offset: clamp_page_offset(self.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 board_column_rank(*column) {
|
||||
None => archived = true,
|
||||
Some(rank) => statuses.extend(
|
||||
RunStatusKind::VARIANTS
|
||||
.iter()
|
||||
.copied()
|
||||
.filter(|kind| kind.board_rank() == rank),
|
||||
),
|
||||
}
|
||||
}
|
||||
RunSummaryVisibility::Selected {
|
||||
statuses: statuses.into_iter().collect(),
|
||||
archived,
|
||||
}
|
||||
}
|
||||
|
||||
fn run_elapsed_ms(run: &fabro_types::Run, now: DateTime<Utc>) -> 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 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)
|
||||
/// 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<u8> {
|
||||
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),
|
||||
}
|
||||
}
|
||||
|
||||
async fn link_run_parent(
|
||||
|
|
@ -364,12 +296,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,56 +314,28 @@ async fn list_runs(
|
|||
State(state): State<Arc<AppState>>,
|
||||
ExtraQuery(params): ExtraQuery<ListRunsParams>,
|
||||
) -> Response {
|
||||
let entries = match state
|
||||
.stores
|
||||
.runs
|
||||
.list_cached_runs(
|
||||
&fabro_store::ListRunsQuery {
|
||||
parent_id: params.parent_id,
|
||||
..fabro_store::ListRunsQuery::default()
|
||||
},
|
||||
Utc::now(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(entries) => entries,
|
||||
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 status_filter = params.status_filter();
|
||||
let include_archived = params.include_archived;
|
||||
|
||||
let filtered: Vec<fabro_types::Run> = 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());
|
||||
|
||||
(
|
||||
StatusCode::OK,
|
||||
Json(serde_json::json!({
|
||||
"data": data,
|
||||
"meta": { "has_more": has_more, "total": total }
|
||||
})),
|
||||
)
|
||||
.into_response()
|
||||
Err(err) => {
|
||||
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, serde::Deserialize)]
|
||||
|
|
@ -478,44 +377,33 @@ async fn resolve_run(
|
|||
State(state): State<Arc<AppState>>,
|
||||
Query(query): Query<ResolveRunQuery>,
|
||||
) -> Response {
|
||||
let runs = match state
|
||||
.stores
|
||||
.runs
|
||||
.list_runs(&fabro_store::ListRunsQuery::default(), Utc::now())
|
||||
.await
|
||||
{
|
||||
Ok(runs) => runs,
|
||||
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(
|
||||
|
|
@ -1031,7 +919,7 @@ async fn get_run_status(
|
|||
RequireRunManagementTarget(id, _actor): RequireRunManagementTarget,
|
||||
State(state): State<Arc<AppState>>,
|
||||
) -> 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()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -24,6 +24,8 @@ tokio-stream.workspace = true
|
|||
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
|
||||
|
|
@ -37,3 +39,4 @@ tokio = { workspace = true, features = ["test-util", "macros"] }
|
|||
tempfile = "3"
|
||||
ulid.workspace = true
|
||||
insta = { workspace = true }
|
||||
fabro-db = { path = "../fabro-db" }
|
||||
|
|
|
|||
|
|
@ -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}")]
|
||||
|
|
|
|||
|
|
@ -7,8 +7,11 @@ mod keys;
|
|||
mod record;
|
||||
mod run_sessions;
|
||||
mod run_state;
|
||||
mod run_summary_store;
|
||||
mod serializable_projection;
|
||||
mod slate;
|
||||
#[cfg(test)]
|
||||
mod test_util;
|
||||
mod types;
|
||||
|
||||
pub use artifact_store::{
|
||||
|
|
@ -25,6 +28,10 @@ pub use run_sessions::{
|
|||
project_run_sessions,
|
||||
};
|
||||
pub use run_state::RunProjectionReducer;
|
||||
pub use run_summary_store::{
|
||||
RunSummaryIdentity, RunSummaryListQuery, RunSummaryPage, RunSummarySort,
|
||||
RunSummarySortDirection, RunSummaryStore, RunSummaryVisibility,
|
||||
};
|
||||
pub use serializable_projection::SerializableProjection;
|
||||
pub use slate::{
|
||||
AuthCode, AuthCodeStore, Blob, BlobStore, CachedRunProjection, ConsumeOutcome, Database,
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
|
|
@ -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,21 +1003,22 @@ fn terminal_total_usd_micros(state: &RunProjection) -> Option<i64> {
|
|||
.and_then(|billing| billing.total_usd_micros)
|
||||
}
|
||||
|
||||
fn projected_total_usd_micros(state: &RunProjection) -> Option<i64> {
|
||||
let mut total_usd_micros = 0_i64;
|
||||
let mut has_total = false;
|
||||
|
||||
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 {
|
||||
|
|
|
|||
875
lib/crates/fabro-store/src/run_summary_store.rs
Normal file
875
lib/crates/fabro-store/src/run_summary_store.rs
Normal file
|
|
@ -0,0 +1,875 @@
|
|||
use std::collections::{HashMap, HashSet};
|
||||
use std::fmt::Write as _;
|
||||
use std::sync::LazyLock;
|
||||
|
||||
use chrono::{DateTime, Utc};
|
||||
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;
|
||||
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
|
||||
";
|
||||
|
||||
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,
|
||||
Changes,
|
||||
Size,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Deserialize)]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
pub enum RunSummarySortDirection {
|
||||
Asc,
|
||||
#[default]
|
||||
Desc,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum RunSummaryVisibility {
|
||||
All,
|
||||
Default {
|
||||
include_archived: bool,
|
||||
},
|
||||
Selected {
|
||||
statuses: Vec<RunStatusKind>,
|
||||
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<RunId>,
|
||||
pub automation_id: Option<String>,
|
||||
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<Run>,
|
||||
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 mut transaction = self.pool.begin().await?;
|
||||
let stored_seqs: HashMap<String, i64> =
|
||||
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 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 stale_ids = stored_seqs
|
||||
.keys()
|
||||
.filter(|stored_id| !authoritative_ids.contains(stored_id.as_str()))
|
||||
.collect::<Vec<_>>();
|
||||
for chunk in stale_ids.chunks(500) {
|
||||
let mut delete = QueryBuilder::<Sqlite>::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<Utc>) -> Result<Option<Run>> {
|
||||
let mut query = QueryBuilder::<Sqlite>::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<Vec<RunSummaryIdentity>> {
|
||||
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::<RunId>()
|
||||
.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,
|
||||
now: DateTime<Utc>,
|
||||
) -> Result<RunSummaryPage> {
|
||||
let mut transaction = self.pool.begin().await?;
|
||||
|
||||
let mut count_query = QueryBuilder::<Sqlite>::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::<Sqlite>::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));
|
||||
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::<Result<Vec<_>>>()?;
|
||||
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,
|
||||
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(())
|
||||
}
|
||||
}
|
||||
|
||||
/// 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<String>,
|
||||
pub workflow_name: Option<String>,
|
||||
pub repository_origin_url: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct ProjectedRunSummary {
|
||||
run: Run,
|
||||
last_seq: u32,
|
||||
workflow_name: Option<String>,
|
||||
repository_name: Option<String>,
|
||||
input_tokens: i64,
|
||||
output_tokens: i64,
|
||||
reasoning_tokens: i64,
|
||||
cache_read_tokens: i64,
|
||||
cache_write_tokens: i64,
|
||||
total_usd_micros: Option<i64>,
|
||||
}
|
||||
|
||||
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.display_name().map(str::to_string);
|
||||
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<Sqlite>, 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 } => {
|
||||
let not_removing = format!("status <> '{}'", RunStatusKind::Removing);
|
||||
if *include_archived {
|
||||
builder.push(format!(
|
||||
" AND (archived_at_ms IS NOT NULL OR {not_removing})"
|
||||
));
|
||||
} else {
|
||||
builder.push(format!(" AND archived_at_ms IS NULL AND {not_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(")");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// 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<String> = 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<String> = 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<Sqlite>,
|
||||
sort: RunSummarySort,
|
||||
direction: RunSummarySortDirection,
|
||||
now: DateTime<Utc>,
|
||||
) {
|
||||
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(STATUS_RANK_CASE_SQL.as_str()),
|
||||
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(SIZE_RANK_CASE_SQL.as_str()),
|
||||
};
|
||||
match direction {
|
||||
RunSummarySortDirection::Asc => builder.push(" ASC"),
|
||||
RunSummarySortDirection::Desc => builder.push(" DESC"),
|
||||
};
|
||||
builder.push(", id DESC");
|
||||
}
|
||||
|
||||
fn decode_run_row(row: &SqliteRow, now: DateTime<Utc>) -> Result<Run> {
|
||||
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",
|
||||
})?;
|
||||
overlay_live_wall_time(&mut run, now);
|
||||
Ok(run)
|
||||
}
|
||||
|
||||
fn overlay_live_wall_time(run: &mut Run, now: DateTime<Utc>) {
|
||||
if run.timestamps.completed_at.is_some() {
|
||||
return;
|
||||
}
|
||||
let Some(started_at) = run.timestamps.started_at else {
|
||||
return;
|
||||
};
|
||||
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))
|
||||
.with_wall_time(wall_time_ms),
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::collections::HashMap;
|
||||
|
||||
use chrono::{DateTime, Utc};
|
||||
use fabro_types::{
|
||||
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::{
|
||||
RunSummaryListQuery, RunSummarySort, RunSummarySortDirection, RunSummaryStore,
|
||||
RunSummaryVisibility,
|
||||
};
|
||||
use crate::slate::CachedRunProjection;
|
||||
use crate::test_util;
|
||||
|
||||
fn dt(value: &str) -> DateTime<Utc> {
|
||||
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<Utc>) -> 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) {
|
||||
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]
|
||||
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::<i64, _>(&row, "source_last_seq"), 4);
|
||||
assert_eq!(
|
||||
sqlx::Row::get::<i64, _>(&row, "created_at_ms"),
|
||||
created_at.timestamp_millis()
|
||||
);
|
||||
assert_eq!(
|
||||
sqlx::Row::get::<i64, _>(&row, "last_event_at_ms"),
|
||||
(created_at + chrono::Duration::minutes(1)).timestamp_millis()
|
||||
);
|
||||
assert_eq!(sqlx::Row::get::<String, _>(&row, "status"), "succeeded");
|
||||
assert_eq!(sqlx::Row::get::<String, _>(&row, "title"), "billed");
|
||||
assert_eq!(
|
||||
sqlx::Row::get::<String, _>(&row, "workflow_slug"),
|
||||
"test-workflow"
|
||||
);
|
||||
assert_eq!(
|
||||
sqlx::Row::get::<String, _>(&row, "automation_id"),
|
||||
"nightly"
|
||||
);
|
||||
assert_eq!(sqlx::Row::get::<i64, _>(&row, "input_tokens"), 100);
|
||||
assert_eq!(sqlx::Row::get::<i64, _>(&row, "reasoning_tokens"), 5);
|
||||
assert_eq!(sqlx::Row::get::<i64, _>(&row, "cache_read_tokens"), 10);
|
||||
assert_eq!(
|
||||
sqlx::Row::get::<i64, _>(&row, "total_usd_micros"),
|
||||
21_000_000
|
||||
);
|
||||
assert_eq!(sqlx::Row::get::<i64, _>(&row, "diff_files_changed"), 2);
|
||||
assert_eq!(sqlx::Row::get::<i64, _>(&row, "diff_additions"), 10);
|
||||
assert_eq!(sqlx::Row::get::<i64, _>(&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()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<OnceCell<Arc<RefreshTokenStore>>>,
|
||||
projection_cache: Arc<RunProjectionCache>,
|
||||
projection_cache_warmed: Arc<OnceCell<()>>,
|
||||
run_summary_store: Arc<OnceLock<Arc<RunSummaryStore>>>,
|
||||
}
|
||||
|
||||
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<RunSummaryStore>) -> Arc<RunSummaryStore> {
|
||||
Arc::clone(self.run_summary_store.get_or_init(|| store))
|
||||
}
|
||||
|
||||
fn run_summary_store(&self) -> Option<Arc<RunSummaryStore>> {
|
||||
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),
|
||||
Arc::clone(&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),
|
||||
Arc::clone(&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),
|
||||
Arc::clone(&self.run_summary_store),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn list_runs(&self, query: &ListRunsQuery, now: DateTime<Utc>) -> Result<Vec<Run>> {
|
||||
|
|
@ -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(())
|
||||
}
|
||||
|
||||
|
|
@ -479,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<Utc> {
|
||||
value.parse().unwrap()
|
||||
|
|
@ -527,6 +559,11 @@ mod tests {
|
|||
(object_store, store)
|
||||
}
|
||||
|
||||
async fn make_summary_store() -> (tempfile::TempDir, Arc<RunSummaryStore>) {
|
||||
let (directory, store) = test_util::sqlite_summary_store().await;
|
||||
(directory, Arc::new(store))
|
||||
}
|
||||
|
||||
fn sample_run_spec(label: &str) -> RunSpec {
|
||||
let mut graph = Graph::new("night-sky");
|
||||
graph.attrs.insert(
|
||||
|
|
@ -1325,6 +1362,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 +1475,7 @@ mod tests {
|
|||
Some("abc123")
|
||||
);
|
||||
|
||||
let summaries = store
|
||||
let cached_summaries = store
|
||||
.list_runs(&ListRunsQuery::default(), Utc::now())
|
||||
.await
|
||||
.unwrap();
|
||||
|
|
@ -1444,7 +1483,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 +1499,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 +1527,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]
|
||||
|
|
|
|||
|
|
@ -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<CachedRunProjection> {
|
||||
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) {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
@ -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,9 @@ pub(crate) struct RunDatabaseInner {
|
|||
state_lock: Mutex<()>,
|
||||
projection_cache: Mutex<EventProjectionCache>,
|
||||
shared_projection_cache: Arc<RunProjectionCache>,
|
||||
// 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<OnceLock<Arc<RunSummaryStore>>>,
|
||||
recent_events: Mutex<VecDeque<EventEnvelope>>,
|
||||
recent_event_limit: usize,
|
||||
event_tx: broadcast::Sender<EventEnvelope>,
|
||||
|
|
@ -51,16 +56,25 @@ impl RunDatabase {
|
|||
run_id: RunId,
|
||||
db: Db,
|
||||
shared_projection_cache: Arc<RunProjectionCache>,
|
||||
run_summary_store: Arc<OnceLock<Arc<RunSummaryStore>>>,
|
||||
) -> Result<Self> {
|
||||
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<RunProjectionCache>,
|
||||
run_summary_store: Arc<OnceLock<Arc<RunSummaryStore>>>,
|
||||
) -> Result<Self> {
|
||||
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 +82,7 @@ impl RunDatabase {
|
|||
db: Db,
|
||||
read_only: bool,
|
||||
shared_projection_cache: Arc<RunProjectionCache>,
|
||||
run_summary_store: Arc<OnceLock<Arc<RunSummaryStore>>>,
|
||||
) -> Result<Self> {
|
||||
let event_seq =
|
||||
recover_next_seq(&db, keys::run_events_prefix(&run_id), keys::parse_event_seq).await?;
|
||||
|
|
@ -83,6 +98,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 +273,63 @@ impl RunDatabase {
|
|||
)
|
||||
.await?;
|
||||
self.cache_event(&event).await?;
|
||||
if let Err(err) = self
|
||||
// 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)
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
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.get() {
|
||||
if let Err(err) = store.upsert_projection(&cached).await {
|
||||
error!(
|
||||
run_id = %self.inner.run_id,
|
||||
source_last_seq = cached.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<Vec<EventEnvelope>> {
|
||||
|
|
|
|||
10
lib/crates/fabro-store/src/test_util.rs
Normal file
10
lib/crates/fabro-store/src/test_util.rs
Normal file
|
|
@ -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()))
|
||||
}
|
||||
|
|
@ -616,12 +616,7 @@ impl RunProjection {
|
|||
#[must_use]
|
||||
pub fn live_run_timing(&self, now: DateTime<Utc>) -> Option<RunTiming> {
|
||||
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()
|
||||
|
|
|
|||
|
|
@ -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<i64>) -> 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)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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<Utc>, now: DateTime<Utc>) -> u64 {
|
||||
u64::try_from(now.signed_duration_since(start).num_milliseconds().max(0))
|
||||
.expect("non-negative milliseconds fit in u64")
|
||||
}
|
||||
}
|
||||
|
||||
impl From<StageTiming> for RunTiming {
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
)
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue