From 50aff2787b718424571cb9e0a079eaeec267c3f5 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 5 Apr 2026 02:53:17 -0400 Subject: [PATCH] refactor(store): route test helpers through server-owned runs --- lib/crates/fabro-cli/tests/it/cmd/support.rs | 179 +- lib/crates/fabro-server/src/server.rs | 10 +- lib/crates/fabro-store/src/keys.rs | 109 +- lib/crates/fabro-store/src/slate/catalog.rs | 169 +- lib/crates/fabro-store/src/slate/mod.rs | 2108 ++--------------- lib/crates/fabro-store/src/slate/run_store.rs | 350 ++- 6 files changed, 625 insertions(+), 2300 deletions(-) diff --git a/lib/crates/fabro-cli/tests/it/cmd/support.rs b/lib/crates/fabro-cli/tests/it/cmd/support.rs index 93e50d38d..655e1a2d6 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/support.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/support.rs @@ -1,18 +1,96 @@ use std::collections::BTreeSet; use std::path::{Path, PathBuf}; use std::process::Output; -use std::sync::Arc; use std::time::{Duration, Instant}; -use fabro_store::{EventEnvelope, RunProjection, SlateRunStore, SlateStore}; +use fabro_store::EventEnvelope; use fabro_test::TestContext; -use fabro_types::RunId; -use object_store::local::LocalFileSystem; +use fabro_types::{ + Checkpoint, Conclusion, NodeStatusRecord, PullRequestRecord, Retro, RunRecord, + RunStatusRecord, SandboxRecord, StageId, StartRecord, +}; use serde_json::Value; use shlex::try_quote; const COMMAND_TIMEOUT: Duration = Duration::from_secs(30); +#[allow(dead_code)] +#[derive(Debug, Clone, Default, serde::Deserialize)] +pub(crate) struct RunProjection { + #[serde(default)] + pub run: Option, + #[serde(default)] + pub graph_source: Option, + #[serde(default)] + pub start: Option, + #[serde(default)] + pub status: Option, + #[serde(default)] + pub checkpoint: Option, + #[serde(default)] + pub checkpoints: Vec<(u32, Checkpoint)>, + #[serde(default)] + pub conclusion: Option, + #[serde(default)] + pub retro: Option, + #[serde(default)] + pub retro_prompt: Option, + #[serde(default)] + pub retro_response: Option, + #[serde(default)] + pub sandbox: Option, + #[serde(default)] + pub final_patch: Option, + #[serde(default)] + pub pull_request: Option, + #[serde(default)] + pub nodes: std::collections::HashMap, +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Default, serde::Deserialize)] +pub(crate) struct NodeState { + #[serde(default)] + pub prompt: Option, + #[serde(default)] + pub response: Option, + #[serde(default)] + pub status: Option, + #[serde(default)] + pub provider_used: Option, + #[serde(default)] + pub diff: Option, + #[serde(default)] + pub script_invocation: Option, + #[serde(default)] + pub script_timing: Option, + #[serde(default)] + pub parallel_results: Option, + #[serde(default)] + pub stdout: Option, + #[serde(default)] + pub stderr: Option, +} + +#[derive(Debug, Clone, Default, serde::Deserialize)] +struct RunSummaryRecord { + run_id: String, + #[serde(default)] + labels: std::collections::HashMap, +} + +impl RunProjection { + pub(crate) fn iter_nodes(&self) -> impl Iterator { + self.nodes + .iter() + .filter_map(|(stage_id, state)| stage_id.parse::().ok().map(|id| (id, state))) + } + + pub(crate) fn is_empty(&self) -> bool { + self.nodes.is_empty() + } +} + pub(crate) struct RunSetup { pub(crate) run_id: String, pub(crate) run_dir: PathBuf, @@ -209,12 +287,7 @@ pub(crate) fn setup_detached_dry_run(context: &TestContext) -> RunSetup { .to_string(); let run = resolve_run(context, &run_id); let deadline = Instant::now() + COMMAND_TIMEOUT; - while { - let store = run_store(&run.run_dir); - block_on(store.list_events()) - .ok() - .is_none_or(|events| events.is_empty()) - } { + while run_events(&run.run_dir).is_empty() { assert!( Instant::now() < deadline, "timed out waiting for store events for {run_id}" @@ -332,13 +405,7 @@ worktree_mode = "never" ); let run = run_local_workflow(context, &workspace_dir, "run.toml"); - let store = run_store(&run.run_dir); - assert!( - block_on(store.state()) - .ok() - .and_then(|state| state.sandbox) - .is_some() - ); + assert!(run_state(&run.run_dir).sandbox.is_some()); WorkspaceRunSetup { run, workspace_dir } } @@ -423,9 +490,9 @@ pub(crate) fn write_gated_workflow(path: &Path, name: &str, goal: &str) -> Workf pub(crate) fn wait_for_status(run_dir: &Path, expected: &[&str]) -> String { let deadline = Instant::now() + COMMAND_TIMEOUT; loop { - if let Some(status) = block_on(run_store(run_dir).state()) - .ok() - .and_then(|state| state.status.map(|record| record.status.to_string())) + if let Some(status) = run_state(run_dir) + .status + .map(|record| record.status.to_string()) { if expected.iter().any(|candidate| *candidate == status) { return status; @@ -461,23 +528,15 @@ pub(crate) fn run_count_for_test_case(context: &TestContext) -> usize { } fn run_dirs_for_test_case(context: &TestContext) -> Vec { - let runs_dir = context.storage_dir.join("runs"); - let entries = match std::fs::read_dir(&runs_dir) { - Ok(entries) => entries, - Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Vec::new(), - Err(err) => panic!("failed to read {}: {err}", runs_dir.display()), - }; - entries - .filter_map(Result::ok) - .map(|entry| entry.path()) - .filter(|path| path.is_dir()) - .filter(|path| { - std::panic::catch_unwind(|| run_state(path)).ok().and_then(|state| state.run).is_some_and(|run| { - run.labels - .get("fabro_test_case") - .is_some_and(|value| value == context.test_case_id()) - }) + let runs: Vec = + block_on(get_server_json_for_storage(&context.storage_dir, "/api/v1/runs")); + runs.into_iter() + .filter(|run| { + run.labels + .get("fabro_test_case") + .is_some_and(|value| value == context.test_case_id()) }) + .filter_map(|run| find_run_dir(&context.storage_dir, &run.run_id)) .collect() } @@ -550,28 +609,50 @@ fn block_on(future: impl std::future::Future) -> T { .block_on(future) } -fn run_store(run_dir: &Path) -> SlateRunStore { +fn server_http_client(storage_dir: &Path) -> reqwest::Client { + reqwest::ClientBuilder::new() + .unix_socket(storage_dir.join("fabro.sock")) + .no_proxy() + .build() + .expect("test HTTP client should build") +} + +async fn get_server_json(run_dir: &Path, path: &str) -> T { let runs_dir = run_dir.parent().expect("run dir should have parent"); let storage_dir = runs_dir.parent().expect("runs dir should have parent"); - let run_id: RunId = infer_run_id(run_dir).parse().expect("run id should parse"); - let object_store = Arc::new( - LocalFileSystem::new_with_prefix(storage_dir.join("store")) - .expect("test store path should be accessible"), + get_server_json_for_storage(storage_dir, path).await +} + +async fn get_server_json_for_storage( + storage_dir: &Path, + path: &str, +) -> T { + let response = server_http_client(storage_dir) + .get(format!("http://fabro{path}")) + .send() + .await + .expect("server request should succeed"); + assert!( + response.status().is_success(), + "server request failed for {path}: {}", + response.status() ); - let store = Arc::new(SlateStore::new(object_store, "", Duration::from_millis(1))); - block_on(store.open_run_reader(&run_id)).expect("run store should exist") + response + .json::() + .await + .expect("server response should parse") } pub(crate) fn run_state(run_dir: &Path) -> RunProjection { - let store = run_store(run_dir); - block_on(store.state()).expect("run store state should exist") + let run_id = infer_run_id(run_dir); + block_on(get_server_json(run_dir, &format!("/api/v1/runs/{run_id}/state"))) } pub(crate) fn run_events(run_dir: &Path) -> Vec { - let store = run_store(run_dir); - block_on(store.list_events()) - .ok() - .expect("run store events should exist") + let run_id = infer_run_id(run_dir); + let response: serde_json::Value = + block_on(get_server_json(run_dir, &format!("/api/v1/runs/{run_id}/events"))); + serde_json::from_value(response["data"].clone()).expect("event list should parse") } pub(crate) fn git_stdout(repo_dir: &Path, args: &[&str]) -> String { diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 69bb228a7..a9f78628d 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -2,7 +2,7 @@ use std::collections::HashMap; use std::str::FromStr; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, RwLock}; -use std::time::Duration; +use std::time::{Duration, Instant}; #[cfg(test)] use axum::body::to_bytes; @@ -127,6 +127,7 @@ struct ManagedRun { status: RunStatus, error: Option, created_at: chrono::DateTime, + enqueued_at: Instant, // Populated when running: interviewer: Option>, event_tx: Option>, @@ -182,6 +183,7 @@ impl AppState { } } + /// Build the axum Router with all run endpoints and embedded static assets. pub fn build_router(state: Arc, auth_mode: AuthMode) -> Router { let middleware_state = Arc::clone(&state); @@ -687,6 +689,7 @@ fn managed_run( status, error: None, created_at, + enqueued_at: Instant::now(), interviewer: None, event_tx: None, checkpoint: None, @@ -912,7 +915,7 @@ async fn start_run( /// Execute a single run: transitions queued → starting → running → completed/failed/cancelled. async fn execute_run(state: Arc, run_id: RunId) { // Transition to Starting and set up cancel infrastructure - let (cancel_rx, run_dir, event_tx, cancel_token, execution_mode) = { + let (cancel_rx, run_dir, event_tx, cancel_token, execution_mode, queued_for) = { let mut runs = state.runs.lock().expect("runs lock poisoned"); let managed_run = match runs.get_mut(&run_id) { Some(r) if r.status == RunStatus::Queued => r, @@ -937,8 +940,10 @@ async fn execute_run(state: Arc, run_id: RunId) { managed_run.event_tx.clone(), cancel_token, managed_run.execution_mode, + managed_run.enqueued_at.elapsed(), ) }; + let _ = queued_for; // Create interviewer and event plumbing (this is the "provisioning" phase) let interviewer = Arc::new(WebInterviewer::new()); @@ -1351,7 +1356,6 @@ async fn list_run_events( }; let since_seq = params.since_seq(); let limit = params.limit(); - match state.store.open_run_reader(&id).await { Ok(run_store) => match run_store.list_events_from_with_limit(since_seq, limit).await { Ok(mut events) => { diff --git a/lib/crates/fabro-store/src/keys.rs b/lib/crates/fabro-store/src/keys.rs index c598bedea..f44242205 100644 --- a/lib/crates/fabro-store/src/keys.rs +++ b/lib/crates/fabro-store/src/keys.rs @@ -1,45 +1,83 @@ use crate::StageId; -use fabro_types::RunBlobId; +use fabro_types::{RunBlobId, RunId}; +pub(crate) const RUNS_PREFIX: &str = "runs/"; +pub(crate) const CATALOG_BY_ID_PREFIX: &str = "_catalog/by-id/"; +pub(crate) const CATALOG_BY_START_PREFIX: &str = "_catalog/by-start/"; pub(crate) const INIT_KEY: &str = "_init.json"; pub(crate) const EVENTS_PREFIX: &str = "events#"; pub(crate) const BLOBS_PREFIX: &str = "blobs#"; pub(crate) const ARTIFACT_NODES_PREFIX: &str = "artifacts#nodes#"; -pub(crate) fn init() -> &'static str { - INIT_KEY +pub(crate) fn run_prefix(run_id: &RunId) -> String { + format!("{RUNS_PREFIX}{run_id}/") } -pub(crate) fn event_key(seq: u32, epoch_ms: i64) -> String { - format!("{EVENTS_PREFIX}{seq:06}-{epoch_ms}.json") +pub(crate) fn init_key(run_id: &RunId) -> String { + format!("{}{INIT_KEY}", run_prefix(run_id)) } -pub(crate) fn blob_key(id: &RunBlobId) -> String { - format!("{BLOBS_PREFIX}{id}") +pub(crate) fn events_prefix(run_id: &RunId) -> String { + format!("{}{EVENTS_PREFIX}", run_prefix(run_id)) } -pub(crate) fn node_artifact_prefix(node: &StageId) -> String { +pub(crate) fn event_key(run_id: &RunId, seq: u32, epoch_ms: i64) -> String { + format!("{}{seq:06}-{epoch_ms}.json", events_prefix(run_id)) +} + +pub(crate) fn blobs_prefix(run_id: &RunId) -> String { + format!("{}{BLOBS_PREFIX}", run_prefix(run_id)) +} + +pub(crate) fn blob_key(run_id: &RunId, id: &RunBlobId) -> String { + format!("{}{id}", blobs_prefix(run_id)) +} + +pub(crate) fn node_artifact_prefix(run_id: &RunId, node: &StageId) -> String { format!( - "{ARTIFACT_NODES_PREFIX}{}#visit-{}", + "{}{ARTIFACT_NODES_PREFIX}{}#visit-{}", + run_prefix(run_id), node.node_id(), node.visit() ) } -pub(crate) fn node_artifact(node: &StageId, filename: &str) -> String { - format!("{}#{filename}", node_artifact_prefix(node)) +pub(crate) fn node_artifact(run_id: &RunId, node: &StageId, filename: &str) -> String { + format!("{}#{filename}", node_artifact_prefix(run_id, node)) +} + +pub(crate) fn catalog_by_id_key(run_id: &RunId) -> String { + format!("{CATALOG_BY_ID_PREFIX}{run_id}.json") +} + +pub(crate) fn catalog_by_start_prefix() -> &'static str { + CATALOG_BY_START_PREFIX +} + +pub(crate) fn catalog_by_start_key(run_id: &RunId) -> String { + format!( + "{CATALOG_BY_START_PREFIX}{}/{run_id}.json", + run_id.created_at().format("%Y-%m-%d-%H-%M") + ) } pub(crate) fn parse_event_seq(key: &str) -> Option { - parse_seq(key, EVENTS_PREFIX) + parse_seq(key.rsplit('/').next()?, EVENTS_PREFIX) } pub(crate) fn parse_blob_id(key: &str) -> Option { - key.strip_prefix(BLOBS_PREFIX)?.parse().ok() + key.rsplit('/').next()?.strip_prefix(BLOBS_PREFIX)?.parse().ok() } pub(crate) fn parse_node_artifact_key(key: &str) -> Option<(StageId, String)> { - parse_visit_scoped_key(key, ARTIFACT_NODES_PREFIX) + let artifact_start = key.find(ARTIFACT_NODES_PREFIX)?; + parse_visit_scoped_key(&key[artifact_start..], ARTIFACT_NODES_PREFIX) +} + +pub(crate) fn parse_run_id_from_catalog_key(key: &str) -> Option { + let filename = key.rsplit('/').next()?; + let run_id = filename.strip_suffix(".json").unwrap_or(filename); + run_id.parse().ok() } fn parse_seq(key: &str, prefix: &str) -> Option { @@ -59,33 +97,50 @@ mod tests { #[test] fn top_level_keys_match_spec() { - assert_eq!(init(), "_init.json"); - assert_eq!(event_key(7, 123), "events#000007-123.json"); + let run_id = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + assert_eq!(INIT_KEY, "_init.json"); + assert_eq!( + event_key(&run_id, 7, 123), + "runs/01JT56VE4Z5NZ814GZN2JZD65A/events#000007-123.json" + ); } #[test] fn sequence_keys_are_zero_padded() { - assert_eq!(event_key(7, 123), "events#000007-123.json"); + let run_id = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + assert_eq!( + event_key(&run_id, 7, 123), + "runs/01JT56VE4Z5NZ814GZN2JZD65A/events#000007-123.json" + ); } #[test] fn artifact_keys_match_spec() { let node = StageId::new("code", 2); - let blob_id = RunBlobId::new(&"01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(), b"summary"); - assert_eq!(blob_key(&blob_id), format!("blobs#{blob_id}")); + let run_id = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let blob_id = RunBlobId::new(&run_id, b"summary"); + assert_eq!(blob_key(&run_id, &blob_id), format!("runs/{run_id}/blobs#{blob_id}")); assert_eq!( - node_artifact(&node, "src/main.rs"), - "artifacts#nodes#code#visit-2#src/main.rs" + node_artifact(&run_id, &node, "src/main.rs"), + "runs/01JT56VE4Z5NZ814GZN2JZD65A/artifacts#nodes#code#visit-2#src/main.rs" ); } #[test] fn parse_helpers_extract_sequences_and_node_visits() { - assert_eq!(parse_event_seq("events#000007-123.json"), Some(7)); - let blob_id = RunBlobId::new(&"01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(), b"summary"); - assert_eq!(parse_blob_id(&format!("blobs#{blob_id}")), Some(blob_id)); assert_eq!( - parse_node_artifact_key("artifacts#nodes#code#visit-2#src/main.rs"), + parse_event_seq("runs/01JT56VE4Z5NZ814GZN2JZD65A/events#000007-123.json"), + Some(7) + ); + let blob_id = RunBlobId::new(&"01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(), b"summary"); + assert_eq!( + parse_blob_id(&format!("runs/01JT56VE4Z5NZ814GZN2JZD65A/blobs#{blob_id}")), + Some(blob_id) + ); + assert_eq!( + parse_node_artifact_key( + "runs/01JT56VE4Z5NZ814GZN2JZD65A/artifacts#nodes#code#visit-2#src/main.rs" + ), Some((StageId::new("code", 2), "src/main.rs".to_string())) ); } @@ -103,7 +158,9 @@ mod tests { #[test] fn asset_filename_with_slashes_parses_correctly() { assert_eq!( - parse_node_artifact_key("artifacts#nodes#build#visit-1#deep/nested/path/file.rs"), + parse_node_artifact_key( + "runs/01JT56VE4Z5NZ814GZN2JZD65A/artifacts#nodes#build#visit-1#deep/nested/path/file.rs" + ), Some(( StageId::new("build", 1), "deep/nested/path/file.rs".to_string() diff --git a/lib/crates/fabro-store/src/slate/catalog.rs b/lib/crates/fabro-store/src/slate/catalog.rs index 8a7e0d070..fa26d8253 100644 --- a/lib/crates/fabro-store/src/slate/catalog.rs +++ b/lib/crates/fabro-store/src/slate/catalog.rs @@ -1,56 +1,35 @@ -use std::collections::HashSet; -use std::sync::Arc; - -use bytes::Bytes; -use futures::TryStreamExt; -use object_store::ObjectStore; -use object_store::path::Path; +use chrono::{Datelike, Timelike}; +use slatedb::Db; +use crate::keys; use crate::{ListRunsQuery, Result}; use fabro_types::RunId; -pub(crate) async fn write_catalog( - store: Arc, - base_prefix: &str, - run_id: &RunId, -) -> Result<()> { - store - .put(&by_id_path(base_prefix, run_id), Bytes::new().into()) - .await?; - store - .put(&by_start_path(base_prefix, run_id), Bytes::new().into()) - .await?; +pub(crate) async fn write_catalog(db: &Db, run_id: &RunId) -> Result<()> { + db.put(keys::catalog_by_id_key(run_id), []).await?; + db.put(keys::catalog_by_start_key(run_id), []).await?; Ok(()) } -pub(crate) async fn read_locator( - store: Arc, - base_prefix: &str, - run_id: &RunId, -) -> Result { - match store.head(&by_id_path(base_prefix, run_id)).await { - Ok(_) => Ok(true), - Err(object_store::Error::NotFound { .. }) => Ok(false), - Err(err) => Err(err.into()), - } +pub(crate) async fn read_locator(db: &Db, run_id: &RunId) -> Result { + Ok(db.get(keys::catalog_by_id_key(run_id)).await?.is_some()) } -pub(crate) async fn list_run_ids( - store: Arc, - base_prefix: &str, - query: &ListRunsQuery, -) -> Result> { - let prefix = Path::from(format!("{base_prefix}by-start")); - let metas = store.list(Some(&prefix)).try_collect::>().await?; +pub(crate) async fn delete_catalog(db: &Db, run_id: &RunId) -> Result<()> { + db.delete(keys::catalog_by_id_key(run_id)).await?; + db.delete(keys::catalog_by_start_key(run_id)).await?; + Ok(()) +} + +pub(crate) async fn list_run_ids(db: &Db, query: &ListRunsQuery) -> Result> { + let mut iter = db.scan_prefix(keys::catalog_by_start_prefix()).await?; let mut run_ids = Vec::new(); - let mut seen = HashSet::new(); - for meta in metas { - let Some(run_id) = parse_run_id_from_path(&meta.location) else { + while let Some(entry) = iter.next().await? { + let key = String::from_utf8(entry.key.to_vec()) + .map_err(|err| crate::StoreError::Other(format!("stored key is not valid UTF-8: {err}")))?; + let Some(run_id) = keys::parse_run_id_from_catalog_key(&key) else { continue; }; - if !seen.insert(run_id) { - continue; - } let created_at = run_id.created_at(); if let Some(start) = query.start { if created_at < start { @@ -64,102 +43,16 @@ pub(crate) async fn list_run_ids( } run_ids.push(run_id); } + run_ids.sort_by_key(|run_id| { + let created_at = run_id.created_at(); + ( + created_at.year(), + created_at.month(), + created_at.day(), + created_at.hour(), + created_at.minute(), + *run_id, + ) + }); Ok(run_ids) } - -pub(crate) fn parse_run_id_from_path(path: &Path) -> Option { - let filename = path.filename()?; - let run_id = filename.strip_suffix(".json").unwrap_or(filename); - run_id.parse().ok() -} - -pub(crate) fn db_prefix(base_prefix: &str, run_id: &RunId) -> String { - format!( - "{base_prefix}db/{}/{run_id}/", - run_id.created_at().format("%Y-%m-%d-%H-%M-%S-%3f") - ) -} - -pub(crate) fn by_id_path(base_prefix: &str, run_id: &RunId) -> Path { - Path::from(format!("{base_prefix}by-id/{run_id}.json")) -} - -pub(crate) fn by_start_path(base_prefix: &str, run_id: &RunId) -> Path { - Path::from(format!( - "{base_prefix}by-start/{}/{run_id}.json", - run_id.created_at().format("%Y-%m-%d-%H-%M") - )) -} - -#[cfg(test)] -pub(super) mod test_support { - use super::*; - - pub(crate) async fn repair_catalog( - store: Arc, - base_prefix: &str, - ) -> Result<()> { - let by_id_prefix = Path::from(format!("{base_prefix}by-id")); - let by_start_prefix = Path::from(format!("{base_prefix}by-start")); - - let by_id_metas = store - .list(Some(&by_id_prefix)) - .try_collect::>() - .await?; - let run_ids = by_id_metas - .iter() - .filter_map(|meta| parse_run_id_from_path(&meta.location)) - .collect::>(); - - for run_id in &run_ids { - let path = by_start_path(base_prefix, run_id); - if !object_exists(store.clone(), &path).await? { - store.put(&path, Bytes::new().into()).await?; - } - } - - let by_start_metas = store - .list(Some(&by_start_prefix)) - .try_collect::>() - .await?; - let canonical = run_ids.into_iter().collect::>(); - let mut seen = HashSet::new(); - for meta in by_start_metas { - let location = meta.location.clone(); - let Some(run_id) = parse_run_id_from_path(&location) else { - delete_if_exists(store.clone(), &location).await?; - continue; - }; - let expected = by_start_path(base_prefix, &run_id); - if canonical.contains(&run_id) && expected == location { - seen.insert(run_id); - continue; - } - delete_if_exists(store.clone(), &location).await?; - } - - for run_id in canonical { - if !seen.contains(&run_id) { - store - .put(&by_start_path(base_prefix, &run_id), Bytes::new().into()) - .await?; - } - } - Ok(()) - } - - async fn object_exists(store: Arc, path: &Path) -> Result { - match store.head(path).await { - Ok(_) => Ok(true), - Err(object_store::Error::NotFound { .. }) => Ok(false), - Err(err) => Err(err.into()), - } - } - - async fn delete_if_exists(store: Arc, path: &Path) -> Result<()> { - match store.delete(path).await { - Ok(()) | Err(object_store::Error::NotFound { .. }) => Ok(()), - Err(err) => Err(err.into()), - } - } -} diff --git a/lib/crates/fabro-store/src/slate/mod.rs b/lib/crates/fabro-store/src/slate/mod.rs index 8e1b66ebd..76bb590fd 100644 --- a/lib/crates/fabro-store/src/slate/mod.rs +++ b/lib/crates/fabro-store/src/slate/mod.rs @@ -5,12 +5,9 @@ use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; -use futures::TryStreamExt; use object_store::ObjectStore; -use object_store::path::Path; -use slatedb::DbReader; -use slatedb::config::{DbReaderOptions, Settings}; -use tokio::sync::Mutex; +use slatedb::config::Settings; +use tokio::sync::{Mutex, OnceCell}; use crate::keys; use crate::{ListRunsQuery, Result, RunSummary, StoreError}; @@ -23,7 +20,8 @@ pub struct SlateStore { object_store: Arc, base_prefix: String, flush_interval: Duration, - active_runs: Arc>>>, + db: Arc>, + active_runs: Arc>>>, } impl std::fmt::Debug for SlateStore { @@ -45,221 +43,131 @@ impl SlateStore { object_store, base_prefix: normalize_base_prefix(base_prefix.into()), flush_interval, + db: Arc::new(OnceCell::new()), active_runs: Arc::new(Mutex::new(HashMap::new())), } } - async fn open_db(&self, db_prefix: &str) -> Result { - Ok( - slatedb::Db::builder(db_prefix.to_string(), self.object_store.clone()) - .with_settings(Settings { - flush_interval: Some(self.flush_interval), - ..Settings::default() - }) - .build() - .await?, - ) + fn shared_db_prefix(&self) -> String { + format!("{}db", self.base_prefix) } - async fn open_reader(&self, db_prefix: &str) -> Result { - Ok(DbReader::open( - db_prefix.to_string(), - self.object_store.clone(), - None, - DbReaderOptions { - manifest_poll_interval: Duration::from_millis(5), - ..DbReaderOptions::default() - }, - ) - .await?) - } - - async fn db_prefix_has_objects(&self, db_prefix: &str) -> Result { - let prefix = Path::from(db_prefix.to_string()); - let mut items = self.object_store.list(Some(&prefix)); - Ok(items.try_next().await?.is_some()) + async fn open_db(&self) -> Result { + let db = self + .db + .get_or_try_init(|| async { + slatedb::Db::builder(self.shared_db_prefix(), self.object_store.clone()) + .with_settings(Settings { + flush_interval: Some(self.flush_interval), + ..Settings::default() + }) + .build() + .await + }) + .await?; + Ok(db.clone()) } async fn get_active_run(&self, run_id: &RunId) -> Option { - let mut active_runs = self.active_runs.lock().await; - let weak = active_runs.get(run_id).cloned()?; - if let Some(inner) = weak.upgrade() { - Some(SlateRunStore::from_inner(inner)) - } else { - active_runs.remove(run_id); - None - } + let active_runs = self.active_runs.lock().await; + active_runs + .get(run_id) + .cloned() + .map(SlateRunStore::from_inner) } async fn cache_active_run(&self, run_store: &SlateRunStore) { self.active_runs .lock() .await - .insert(run_store.run_id(), run_store.downgrade()); + .insert(run_store.run_id(), run_store.inner_arc()); } async fn remove_active_run(&self, run_id: &RunId) -> Option { - let weak = self.active_runs.lock().await.remove(run_id)?; - weak.upgrade().map(SlateRunStore::from_inner) - } - - async fn open_run_store( - &self, - run_id: &RunId, - db_prefix: &str, - ) -> Result> { - if let Some(active) = self.get_active_run(run_id).await { - if active.matches_run(run_id, db_prefix) { - return Ok(Some(active)); - } - return Err(StoreError::Other(format!( - "active run cache mismatch for run_id {run_id:?}" - ))); - } - if !self.db_prefix_has_objects(db_prefix).await? { - return Ok(None); - } - let db = self.open_db(db_prefix).await?; - let has_init = match SlateRunStore::validate_init(&db, run_id).await { - Ok(has_init) => has_init, - Err(err) => { - let _ = db.close().await; - return Err(err); - } - }; - if !has_init { - let _ = db.close().await; - return Ok(None); - } - let run_store = SlateRunStore::open_writer(*run_id, db_prefix.to_string(), db).await?; - self.cache_active_run(&run_store).await; - Ok(Some(run_store)) - } - - async fn open_run_reader_store( - &self, - run_id: &RunId, - db_prefix: &str, - ) -> Result> { - if !self.db_prefix_has_objects(db_prefix).await? { - return Ok(None); - } - let reader = self.open_reader(db_prefix).await?; - let has_init = match SlateRunStore::validate_init(&reader, run_id).await { - Ok(has_init) => has_init, - Err(err) => { - let _ = reader.close().await; - return Err(err); - } - }; - if !has_init { - let _ = reader.close().await; - return Ok(None); - } - SlateRunStore::open_reader(*run_id, db_prefix.to_string(), reader) + self.active_runs + .lock() .await - .map(Some) + .remove(run_id) + .map(SlateRunStore::from_inner) } - async fn delete_db_prefix(&self, db_prefix: &str) -> Result<()> { - let prefix = Path::from(db_prefix.to_string()); - let metas = self - .object_store - .list(Some(&prefix)) - .try_collect::>() - .await?; - for meta in metas { - delete_path(self.object_store.clone(), &meta.location).await?; - } - Ok(()) - } -} - -impl SlateStore { pub async fn create_run(&self, run_id: &RunId) -> Result { - let locator_exists = - catalog::read_locator(self.object_store.clone(), &self.base_prefix, run_id).await?; - let db_prefix = catalog::db_prefix(&self.base_prefix, run_id); + let db = self.open_db().await?; + let locator_exists = catalog::read_locator(&db, run_id).await?; if let Some(active) = self.get_active_run(run_id).await { - if locator_exists && !active.matches_run(run_id, &db_prefix) { + if locator_exists && !active.matches_run(run_id) { return Err(StoreError::RunAlreadyExists(run_id.to_string())); } - catalog::write_catalog(self.object_store.clone(), &self.base_prefix, run_id).await?; + catalog::write_catalog(&db, run_id).await?; return Ok(active); } - if locator_exists && self.db_prefix_has_objects(&db_prefix).await? { + if locator_exists { return Err(StoreError::RunAlreadyExists(run_id.to_string())); } - let db = self.open_db(&db_prefix).await?; SlateRunStore::validate_init(&db, run_id).await?; - db.put(keys::init(), serde_json::to_vec(run_id)?).await?; - let run_store = SlateRunStore::open_writer(*run_id, db_prefix.clone(), db).await?; + db.put(keys::init_key(run_id), serde_json::to_vec(run_id)?).await?; + catalog::write_catalog(&db, run_id).await?; + let run_store = SlateRunStore::open_writer(*run_id, db).await?; self.cache_active_run(&run_store).await; - catalog::write_catalog(self.object_store.clone(), &self.base_prefix, run_id).await?; Ok(run_store) } pub async fn open_run(&self, run_id: &RunId) -> Result { - let exists = - catalog::read_locator(self.object_store.clone(), &self.base_prefix, run_id).await?; - if !exists { + let db = self.open_db().await?; + if !catalog::read_locator(&db, run_id).await? { return Err(StoreError::RunNotFound(run_id.to_string())); } - let db_prefix = catalog::db_prefix(&self.base_prefix, run_id); - - let run_store = self - .open_run_store(run_id, &db_prefix) - .await? - .ok_or_else(|| StoreError::RunNotFound(run_id.to_string()))?; + if let Some(active) = self.get_active_run(run_id).await { + if !active.matches_run(run_id) { + return Err(StoreError::Other(format!( + "active run cache mismatch for run_id {run_id:?}" + ))); + } + return Ok(active); + } + if !SlateRunStore::validate_init(&db, run_id).await? { + return Err(StoreError::RunNotFound(run_id.to_string())); + } + let run_store = SlateRunStore::open_writer(*run_id, db).await?; + self.cache_active_run(&run_store).await; Ok(run_store) } pub async fn open_run_reader(&self, run_id: &RunId) -> Result { - let exists = - catalog::read_locator(self.object_store.clone(), &self.base_prefix, run_id).await?; - if !exists { + let db = self.open_db().await?; + if !catalog::read_locator(&db, run_id).await? { return Err(StoreError::RunNotFound(run_id.to_string())); } - let db_prefix = catalog::db_prefix(&self.base_prefix, run_id); - - let run_store = self - .open_run_reader_store(run_id, &db_prefix) - .await? - .ok_or_else(|| StoreError::RunNotFound(run_id.to_string()))?; - Ok(run_store) + if let Some(active) = self.get_active_run(run_id).await { + if !active.matches_run(run_id) { + return Err(StoreError::Other(format!( + "active run cache mismatch for run_id {run_id:?}" + ))); + } + return Ok(active.into_read_only()); + } + if !SlateRunStore::validate_init(&db, run_id).await? { + return Err(StoreError::RunNotFound(run_id.to_string())); + } + SlateRunStore::open_reader(*run_id, db).await } pub async fn list_runs(&self, query: &ListRunsQuery) -> Result> { - let run_ids = - catalog::list_run_ids(self.object_store.clone(), &self.base_prefix, query).await?; + let db = self.open_db().await?; + let run_ids = catalog::list_run_ids(&db, query).await?; let mut summaries = Vec::new(); for run_id in run_ids { - let db_prefix = catalog::db_prefix(&self.base_prefix, &run_id); if let Some(active) = self.get_active_run(&run_id).await { - if !active.matches_run(&run_id, &db_prefix) { - return Err(StoreError::Other(format!( - "active run cache mismatch for run_id {run_id:?}" - ))); - } - let snapshot = active.snapshot().await?; - summaries.push(SlateRunStore::build_summary(snapshot.as_ref(), &run_id).await?); + summaries.push(active.state().await?.build_summary(&run_id)); continue; } - if !self.db_prefix_has_objects(&db_prefix).await? { + if !SlateRunStore::validate_init(&db, &run_id).await? { continue; } - let reader = self.open_reader(&db_prefix).await?; - if !SlateRunStore::validate_init(&reader, &run_id).await? { - let _ = reader.close().await; - continue; - } - let summary = SlateRunStore::build_summary(&reader, &run_id).await; - let _ = reader.close().await; - let summary = summary?; - summaries.push(summary); + summaries.push(SlateRunStore::build_summary(&db, &run_id).await?); } summaries.sort_by(|a, b| b.run_id.created_at().cmp(&a.run_id.created_at())); Ok(summaries) @@ -271,63 +179,23 @@ impl SlateStore { active.close().await?; } - let db_prefix = catalog::db_prefix(&self.base_prefix, run_id); - - if catalog::read_locator(self.object_store.clone(), &self.base_prefix, run_id).await? { - delete_path( - self.object_store.clone(), - &catalog::by_start_path(&self.base_prefix, run_id), - ) - .await?; - self.delete_db_prefix(&db_prefix).await?; - delete_path( - self.object_store.clone(), - &catalog::by_id_path(&self.base_prefix, run_id), - ) - .await?; - return Ok(()); + let db = self.open_db().await?; + let prefix = keys::run_prefix(run_id); + let mut iter = db.scan_prefix(prefix.as_bytes()).await?; + let mut keys_to_delete = Vec::new(); + while let Some(entry) = iter.next().await? { + keys_to_delete.push(String::from_utf8(entry.key.to_vec()).map_err(|err| { + StoreError::Other(format!("stored key is not valid UTF-8: {err}")) + })?); } - - if active.is_some() { - delete_path( - self.object_store.clone(), - &catalog::by_start_path(&self.base_prefix, run_id), - ) - .await?; - self.delete_db_prefix(&db_prefix).await?; - delete_path( - self.object_store.clone(), - &catalog::by_id_path(&self.base_prefix, run_id), - ) - .await?; - return Ok(()); + for key in keys_to_delete { + db.delete(key).await?; } - - let by_start_prefix = Path::from(format!("{}by-start", self.base_prefix)); - let metas = self - .object_store - .list(Some(&by_start_prefix)) - .try_collect::>() - .await?; - let expected_name = format!("{run_id}.json"); - for meta in metas { - if meta.location.filename() != Some(expected_name.as_str()) { - continue; - } - delete_path(self.object_store.clone(), &meta.location).await?; - } - self.delete_db_prefix(&db_prefix).await?; + catalog::delete_catalog(&db, run_id).await?; Ok(()) } } -async fn delete_path(store: Arc, path: &Path) -> Result<()> { - match store.delete(path).await { - Ok(()) | Err(object_store::Error::NotFound { .. }) => Ok(()), - Err(err) => Err(err.into()), - } -} - pub(crate) fn normalize_base_prefix(prefix: String) -> String { if prefix.is_empty() { return String::new(); @@ -343,36 +211,26 @@ pub(crate) fn normalize_base_prefix(prefix: String) -> String { mod tests { use super::*; - use std::collections::HashMap; - use std::path::PathBuf; - use std::time::Duration; - - use bytes::Bytes; - use chrono::{DateTime, Duration as ChronoDuration, Utc}; - use fabro_types::{ - AttrValue, Checkpoint, Conclusion, Graph, PullRequestRecord, Retro, RunId, RunRecord, - RunStatus, RunStatusRecord, SandboxRecord, Settings, StageStatus, StartRecord, - StatusReason, - }; + use chrono::{DateTime, Utc}; + use futures::TryStreamExt; + use object_store::path::Path; use object_store::memory::InMemory; - use slatedb::config::Settings as SlateSettings; - use slatedb::{CloseReason, ErrorKind}; - use tokio::time::timeout; + use fabro_types::{AttrValue, Graph, RunRecord, RunStatus, Settings, StatusReason}; + use std::path::PathBuf; - use crate::{EventPayload, StageId}; + use crate::EventPayload; - #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] - struct CatalogRecord { - run_id: RunId, - created_at: DateTime, - db_prefix: String, - run_dir: Option, + fn dt(value: &str) -> DateTime { + value.parse().unwrap() } - fn dt(rfc3339: &str) -> DateTime { - DateTime::parse_from_rfc3339(rfc3339) - .unwrap() - .with_timezone(&Utc) + fn test_run_id(label: &str) -> RunId { + let (timestamp_ms, random) = match label { + "run-1" => (dt("2026-03-27T12:00:00Z").timestamp_millis() as u64, 1), + "run-2" => (dt("2026-03-27T12:00:10Z").timestamp_millis() as u64, 2), + _ => panic!("unknown test run id: {label}"), + }; + RunId::from(ulid::Ulid::from_parts(timestamp_ms, random)) } fn make_store() -> (Arc, SlateStore) { @@ -381,33 +239,18 @@ mod tests { (object_store, store) } - async fn repair_catalog_for_tests(store: &SlateStore) -> Result<()> { - catalog::test_support::repair_catalog(store.object_store.clone(), &store.base_prefix).await - } - - fn test_run_id(label: &str) -> RunId { - let (timestamp_ms, random) = match label { - "run-1" => (dt("2026-03-27T12:00:00Z").timestamp_millis() as u64, 1), - "other-run" => (dt("2026-03-27T12:00:00Z").timestamp_millis() as u64, 2), - "run-early" => (dt("2026-03-27T10:00:00Z").timestamp_millis() as u64, 3), - "run-late" => (dt("2026-03-27T12:00:00Z").timestamp_millis() as u64, 4), - _ => panic!("unknown test run id: {label}"), - }; - RunId::from(ulid::Ulid::from_parts(timestamp_ms, random)) - } - - fn sample_run_record(run_id: &str, _created_at: DateTime) -> RunRecord { + fn sample_run_record(label: &str) -> RunRecord { let mut graph = Graph::new("night-sky"); graph.attrs.insert( "goal".to_string(), AttrValue::String("map the constellations".to_string()), ); RunRecord { - run_id: test_run_id(run_id), + run_id: test_run_id(label), settings: Settings::default(), graph, workflow_slug: Some("night-sky".to_string()), - working_directory: PathBuf::from("/tmp/night-sky"), + working_directory: PathBuf::from(format!("/tmp/{label}")), host_repo_path: Some("github.com/fabro-sh/fabro".to_string()), base_branch: Some("main".to_string()), labels: std::collections::HashMap::from([("team".to_string(), "infra".to_string())]), @@ -418,234 +261,58 @@ mod tests { run_id: &str, ts: &str, event: &str, - node_id: Option<&str>, properties: serde_json::Value, ) -> EventPayload { - let properties = normalize_test_event_properties(run_id, event, properties); - let mut value = serde_json::json!({ - "id": format!("evt-{run_id}-{event}"), - "ts": ts, - "run_id": test_run_id(run_id).to_string(), - "event": event, - "properties": properties, - }); - if let Some(node_id) = node_id { - value["node_id"] = serde_json::Value::String(node_id.to_string()); - } - EventPayload::new(value, &test_run_id(run_id)).unwrap() + EventPayload::new( + serde_json::json!({ + "id": format!("evt-{run_id}-{event}"), + "ts": ts, + "run_id": test_run_id(run_id).to_string(), + "event": event, + "properties": properties, + }), + &test_run_id(run_id), + ) + .unwrap() } - fn normalize_test_event_properties( - run_id: &str, - event: &str, - properties: serde_json::Value, - ) -> serde_json::Value { - let serde_json::Value::Object(mut props) = properties else { - return properties; - }; - - match event { - "run.created" => { - props - .entry("run_dir") - .or_insert_with(|| serde_json::Value::String(format!("/tmp/{run_id}"))); - } - "run.started" => { - props - .entry("name") - .or_insert_with(|| serde_json::Value::String("night-sky".to_string())); - } - "sandbox.initialized" => { - props - .entry("provider") - .or_insert_with(|| serde_json::Value::String("local".to_string())); - } - "checkpoint.completed" => { - props - .entry("completed_nodes") - .or_insert_with(|| serde_json::Value::Array(Vec::new())); - props - .entry("node_retries") - .or_insert_with(|| serde_json::json!({})); - props - .entry("context_values") - .or_insert_with(|| serde_json::json!({})); - props - .entry("node_outcomes") - .or_insert_with(|| serde_json::json!({})); - props - .entry("node_visits") - .or_insert_with(|| serde_json::json!({})); - } - "parallel.completed" => { - props - .entry("visit") - .or_insert_with(|| serde_json::Value::from(1)); - props - .entry("duration_ms") - .or_insert_with(|| serde_json::Value::from(0)); - props - .entry("success_count") - .or_insert_with(|| serde_json::Value::from(0)); - props - .entry("failure_count") - .or_insert_with(|| serde_json::Value::from(0)); - } - "stage.completed" => { - if let Some(visit) = props.get("visit").cloned() { - props.entry("index").or_insert(visit); - } - props - .entry("index") - .or_insert_with(|| serde_json::Value::from(0)); - props - .entry("duration_ms") - .or_insert_with(|| serde_json::Value::from(0)); - props - .entry("attempt") - .or_insert_with(|| serde_json::Value::from(1)); - props - .entry("max_attempts") - .or_insert_with(|| serde_json::Value::from(1)); - } - "command.started" => { - let default_script = props - .get("command") - .cloned() - .unwrap_or_else(|| serde_json::Value::String(String::new())); - props.entry("script").or_insert(default_script); - props - .entry("language") - .or_insert_with(|| serde_json::Value::String("shell".to_string())); - } - "command.completed" => { - props - .entry("duration_ms") - .or_insert_with(|| serde_json::Value::from(0)); - props - .entry("timed_out") - .or_insert_with(|| serde_json::Value::Bool(false)); - } - "run.completed" => { - props - .entry("artifact_count") - .or_insert_with(|| serde_json::Value::from(0)); - } - "pull_request.created" => { - props - .entry("draft") - .or_insert_with(|| serde_json::Value::Bool(false)); - } - "retro.completed" => { - props - .entry("duration_ms") - .or_insert_with(|| serde_json::Value::from(0)); - } - _ => {} - } - - serde_json::Value::Object(props) + async fn append_created(run: &SlateRunStore, label: &str, created_at: DateTime) { + let run_record = sample_run_record(label); + run.append_event(&event_payload( + label, + &created_at.to_rfc3339(), + "run.created", + serde_json::json!({ + "settings": run_record.settings, + "graph": run_record.graph, + "workflow_slug": run_record.workflow_slug, + "working_directory": run_record.working_directory, + "run_dir": format!("/tmp/{label}"), + "host_repo_path": run_record.host_repo_path, + "base_branch": run_record.base_branch, + "labels": run_record.labels, + }), + )) + .await + .unwrap(); } - fn sample_start_record(run_id: &str, created_at: DateTime) -> StartRecord { - StartRecord { - run_id: test_run_id(run_id), - start_time: created_at + ChronoDuration::seconds(5), - run_branch: Some("fabro/run/demo".to_string()), - base_sha: Some("abc123".to_string()), - } - } - - fn sample_status(status: RunStatus, reason: Option) -> RunStatusRecord { - RunStatusRecord { - status, - reason, - updated_at: dt("2026-03-27T12:05:00Z"), - } - } - - fn sample_checkpoint() -> Checkpoint { - Checkpoint { - timestamp: dt("2026-03-27T12:10:00Z"), - current_node: "code".to_string(), - completed_nodes: vec!["plan".to_string()], - node_retries: HashMap::from([("code".to_string(), 1)]), - context_values: HashMap::from([( - "artifact".to_string(), - serde_json::json!({"kind": "summary"}), - )]), - node_outcomes: HashMap::new(), - next_node_id: Some("review".to_string()), - git_commit_sha: Some("def456".to_string()), - loop_failure_signatures: HashMap::new(), - restart_failure_signatures: HashMap::new(), - node_visits: HashMap::from([("code".to_string(), 2)]), - } - } - - fn sample_conclusion() -> Conclusion { - Conclusion { - timestamp: dt("2026-03-27T12:15:00Z"), - status: StageStatus::Success, - duration_ms: 3210, - failure_reason: None, - final_git_commit_sha: Some("feedbeef".to_string()), - stages: Vec::new(), - total_cost: Some(1.25), - total_retries: 2, - total_input_tokens: 10, - total_output_tokens: 20, - total_cache_read_tokens: 30, - total_cache_write_tokens: 40, - total_reasoning_tokens: 50, - has_pricing: true, - } - } - - fn sample_retro(run_id: &str) -> Retro { - Retro { - run_id: test_run_id(run_id), - workflow_name: "night-sky".to_string(), - goal: "map the constellations".to_string(), - timestamp: dt("2026-03-27T12:20:00Z"), - smoothness: None, - stages: Vec::new(), - stats: fabro_types::AggregateStats { - total_duration_ms: 3210, - total_cost: Some(1.25), - total_retries: 2, - files_touched: vec!["src/lib.rs".to_string()], - stages_completed: 3, - stages_failed: 0, - }, - intent: Some("ship the fix".to_string()), - outcome: Some("done".to_string()), - learnings: None, - friction_points: None, - open_items: None, - } - } - - fn sample_sandbox() -> SandboxRecord { - SandboxRecord { - provider: "local".to_string(), - working_directory: "/tmp/night-sky".to_string(), - identifier: Some("sandbox-1".to_string()), - host_working_directory: Some("/tmp/night-sky".to_string()), - container_mount_point: None, - } - } - - fn sample_pull_request() -> PullRequestRecord { - PullRequestRecord { - html_url: "https://github.com/fabro-sh/fabro/pull/123".to_string(), - number: 123, - owner: "fabro-sh".to_string(), - repo: "fabro".to_string(), - base_branch: "main".to_string(), - head_branch: "fabro/run/demo".to_string(), - title: "Map the constellations".to_string(), - } + async fn append_completed(run: &SlateRunStore, label: &str, created_at: DateTime) { + append_created(run, label, created_at).await; + run.append_event(&event_payload( + label, + "2026-03-27T12:00:02Z", + "run.completed", + serde_json::json!({ + "duration_ms": 3210, + "artifact_count": 1, + "status": "success", + "reason": "completed", + "total_cost": 1.25, + }), + )) + .await + .unwrap(); } async fn list_paths(store: Arc, prefix: &str) -> Vec { @@ -659,72 +326,68 @@ mod tests { items } - async fn object_exists(store: Arc, path: &Path) -> bool { - store.head(path).await.is_ok() - } + #[tokio::test] + async fn create_open_list_and_delete_full_lifecycle_in_shared_db() { + let (object_store, store) = make_store(); + let run_1 = store.create_run(&test_run_id("run-1")).await.unwrap(); + let run_2 = store.create_run(&test_run_id("run-2")).await.unwrap(); + append_completed(&run_1, "run-1", dt("2026-03-27T12:00:00Z")).await; + append_created(&run_2, "run-2", dt("2026-03-27T12:00:10Z")).await; - async fn seed_db( - object_store: Arc, - record: &CatalogRecord, - include_init: bool, - ) -> slatedb::Db { - let db = slatedb::Db::builder(record.db_prefix.clone(), object_store) - .with_settings(SlateSettings { - flush_interval: Some(Duration::from_millis(1)), - ..SlateSettings::default() - }) - .build() - .await - .unwrap(); - if include_init { - db.put(keys::init(), serde_json::to_vec(&record.run_id).unwrap()) - .await - .unwrap(); - } - db + let summary = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + assert_eq!(summary.len(), 2); + assert_eq!(summary[0].run_id, test_run_id("run-2")); + assert_eq!(summary[1].run_id, test_run_id("run-1")); + assert_eq!(summary[1].workflow_name, Some("night-sky".to_string())); + assert_eq!(summary[1].goal, Some("map the constellations".to_string())); + assert_eq!(summary[1].status, Some(RunStatus::Succeeded)); + assert_eq!(summary[1].status_reason, Some(StatusReason::Completed)); + + let reopened = store.open_run(&test_run_id("run-1")).await.unwrap(); + let stored = reopened.state().await.unwrap().run.unwrap(); + assert_eq!(stored.run_id, test_run_id("run-1")); + + store.delete_run(&test_run_id("run-1")).await.unwrap(); + assert!(store.open_run(&test_run_id("run-1")).await.is_err()); + let remaining = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + assert_eq!(remaining.len(), 1); + assert_eq!(remaining[0].run_id, test_run_id("run-2")); + assert!(!list_paths(object_store, "runs/db").await.is_empty()); } #[tokio::test] - async fn create_open_list_and_delete_full_lifecycle() { - let (object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); + async fn open_run_reader_is_read_only() { + let (_object_store, store) = make_store(); let run = store.create_run(&test_run_id("run-1")).await.unwrap(); + append_created(&run, "run-1", dt("2026-03-27T12:00:00Z")).await; + + let reader = store.open_run_reader(&test_run_id("run-1")).await.unwrap(); + let err = reader + .append_event(&event_payload( + "run-1", + "2026-03-27T12:00:01Z", + "run.completed", + serde_json::json!({ "reason": "completed" }), + )) + .await + .unwrap_err(); + assert!(matches!(err, StoreError::ReadOnly)); + } + + #[tokio::test] + async fn reader_sees_cached_projection_and_recent_events_for_active_run() { + let (_object_store, store) = make_store(); + let run = store.create_run(&test_run_id("run-1")).await.unwrap(); + append_created(&run, "run-1", dt("2026-03-27T12:00:00Z")).await; + + let reader = store.open_run_reader(&test_run_id("run-1")).await.unwrap(); + let state = reader.state().await.unwrap(); + assert_eq!(state.run.unwrap().run_id, test_run_id("run-1")); - let run_record = sample_run_record("run-1", created_at); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": run_record.settings, - "graph": run_record.graph, - "workflow_slug": run_record.workflow_slug, - "working_directory": run_record.working_directory, - "host_repo_path": run_record.host_repo_path, - "base_branch": run_record.base_branch, - "labels": run_record.labels, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:01Z", - "run.started", - None, - serde_json::json!({ - "run_branch": "fabro/run/demo", - "base_sha": "abc123", - }), - )) - .await - .unwrap(); run.append_event(&event_payload( "run-1", "2026-03-27T12:00:02Z", "run.completed", - None, serde_json::json!({ "duration_ms": 3210, "artifact_count": 1, @@ -736,1368 +399,21 @@ mod tests { .await .unwrap(); - let by_id = catalog::by_id_path("runs/", &test_run_id("run-1")); - let by_start = catalog::by_start_path("runs/", &test_run_id("run-1")); - assert!(object_exists(object_store.clone(), &by_id).await); - assert!(object_exists(object_store.clone(), &by_start).await); + let recent = reader.list_events_from_with_limit(2, 10).await.unwrap(); + assert_eq!(recent.len(), 1); + assert_eq!(recent[0].seq, 2); + } - let summary = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + #[tokio::test] + async fn reopening_store_rebuilds_from_shared_db() { + 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 = SlateStore::new(object_store, "runs", Duration::from_millis(1)); + let summary = reopened.list_runs(&ListRunsQuery::default()).await.unwrap(); assert_eq!(summary.len(), 1); assert_eq!(summary[0].run_id, test_run_id("run-1")); - assert_eq!(summary[0].workflow_name, Some("night-sky".to_string())); - assert_eq!(summary[0].goal, Some("map the constellations".to_string())); assert_eq!(summary[0].status, Some(RunStatus::Succeeded)); - assert_eq!(summary[0].status_reason, Some(StatusReason::Completed)); - - let reopened = store.open_run(&test_run_id("run-1")).await.unwrap(); - let stored = reopened.state().await.unwrap().run.unwrap(); - assert_eq!(stored.run_id, test_run_id("run-1")); - - store.delete_run(&test_run_id("run-1")).await.unwrap(); - assert!(store.open_run(&test_run_id("run-1")).await.is_err()); - assert!(!object_exists(object_store.clone(), &by_id).await); - assert!(!object_exists(object_store.clone(), &by_start).await); - assert!(list_paths(object_store, "runs/db").await.is_empty()); - } - - #[tokio::test] - async fn by_id_without_by_start_opens_but_is_omitted_from_list_and_repair_restores_index() { - let (object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); - let record = CatalogRecord { - run_id: test_run_id("run-1"), - created_at, - db_prefix: catalog::db_prefix("runs/", &test_run_id("run-1")), - run_dir: None, - }; - - let db = seed_db(object_store.clone(), &record, true).await; - db.put( - keys::event_key(1, created_at.timestamp_millis()), - serde_json::to_vec(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": sample_run_record("run-1", created_at).settings, - "graph": sample_run_record("run-1", created_at).graph, - "workflow_slug": sample_run_record("run-1", created_at).workflow_slug, - "working_directory": sample_run_record("run-1", created_at).working_directory, - "host_repo_path": sample_run_record("run-1", created_at).host_repo_path, - "base_branch": sample_run_record("run-1", created_at).base_branch, - "labels": sample_run_record("run-1", created_at).labels, - }), - )) - .unwrap(), - ) - .await - .unwrap(); - db.close().await.unwrap(); - - object_store - .put( - &catalog::by_id_path("runs/", &test_run_id("run-1")), - serde_json::to_vec(&record).unwrap().into(), - ) - .await - .unwrap(); - - assert!(store.open_run(&test_run_id("run-1")).await.is_ok()); - assert!( - store - .list_runs(&ListRunsQuery::default()) - .await - .unwrap() - .is_empty() - ); - - repair_catalog_for_tests(&store).await.unwrap(); - let listed = store.list_runs(&ListRunsQuery::default()).await.unwrap(); - assert_eq!(listed.len(), 1); - assert!( - object_exists( - object_store, - &catalog::by_start_path("runs/", &test_run_id("run-1")) - ) - .await - ); - } - - #[tokio::test] - async fn reopen_recovers_event_sequences() { - let (_object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let run_record = sample_run_record("run-1", created_at); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": run_record.settings, - "graph": run_record.graph, - "workflow_slug": run_record.workflow_slug, - "working_directory": run_record.working_directory, - "host_repo_path": run_record.host_repo_path, - "base_branch": run_record.base_branch, - "labels": run_record.labels, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:01Z", - "Started", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:02Z", - "Next", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - drop(run); - - let reopened = store.open_run(&test_run_id("run-1")).await.unwrap(); - let next_event = reopened - .append_event(&event_payload( - "run-1", - "2026-03-27T12:00:02Z", - "AfterReopen", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - assert_eq!(next_event, 4); - } - - #[tokio::test] - async fn open_run_and_list_runs_skip_empty_databases_without_init() { - let (object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); - let record = CatalogRecord { - run_id: test_run_id("run-1"), - created_at, - db_prefix: catalog::db_prefix("runs/", &test_run_id("run-1")), - run_dir: None, - }; - - let db = seed_db(object_store.clone(), &record, false).await; - db.close().await.unwrap(); - object_store - .put( - &catalog::by_id_path("runs/", &test_run_id("run-1")), - serde_json::to_vec(&record).unwrap().into(), - ) - .await - .unwrap(); - object_store - .put( - &catalog::by_start_path("runs/", &test_run_id("run-1")), - serde_json::to_vec(&record).unwrap().into(), - ) - .await - .unwrap(); - - assert!(store.open_run(&test_run_id("run-1")).await.is_err()); - assert!( - store - .list_runs(&ListRunsQuery::default()) - .await - .unwrap() - .is_empty() - ); - } - - #[tokio::test] - async fn create_run_allows_idempotent_retry_and_rejects_conflict() { - let (_object_store, store) = make_store(); - let _created_at = dt("2026-03-27T12:00:00Z"); - // Hold the first run store so the active cache Weak ref stays alive. - let _run = store.create_run(&test_run_id("run-1")).await.unwrap(); - store.create_run(&test_run_id("run-1")).await.unwrap(); - - let conflict = store.create_run(&test_run_id("other-run")).await; - // Different run_id should work fine - assert!(conflict.is_ok()); - - // Dropping and re-creating with locator already written should reject - drop(_run); - let conflict = store.create_run(&test_run_id("run-1")).await; - assert!(matches!(conflict, Err(StoreError::RunAlreadyExists(_)))); - } - - #[tokio::test] - async fn list_runs_and_open_run_reuse_active_handle_without_fencing() { - let (_object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let run_record = sample_run_record("run-1", created_at); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": run_record.settings, - "graph": run_record.graph, - "workflow_slug": run_record.workflow_slug, - "working_directory": run_record.working_directory, - "host_repo_path": run_record.host_repo_path, - "base_branch": run_record.base_branch, - "labels": run_record.labels, - }), - )) - .await - .unwrap(); - - let listed = store.list_runs(&ListRunsQuery::default()).await.unwrap(); - assert_eq!(listed.len(), 1); - - let reopened = store.open_run(&test_run_id("run-1")).await.unwrap(); - let first_event = run - .append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "Started", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - let second_event = reopened - .append_event(&event_payload( - "run-1", - "2026-03-27T12:00:01Z", - "Continued", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - assert_eq!(first_event, 2); - assert_eq!(second_event, 3); - } - - #[tokio::test] - async fn watch_events_from_polls_new_events() { - let (_object_store, store) = make_store(); - let _created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let mut stream = run.watch_events_from(1).unwrap(); - - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "Started", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - - let event = timeout( - Duration::from_secs(2), - futures::StreamExt::next(&mut stream), - ) - .await - .unwrap() - .unwrap() - .unwrap(); - assert_eq!(event.seq, 1); - } - - #[tokio::test] - async fn delete_run_is_idempotent_and_fallback_cleans_by_start_orphans() { - let (object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); - let record = CatalogRecord { - run_id: test_run_id("run-1"), - created_at, - db_prefix: catalog::db_prefix("runs/", &test_run_id("run-1")), - run_dir: None, - }; - let db = seed_db(object_store.clone(), &record, true).await; - db.put( - keys::event_key(1, created_at.timestamp_millis()), - serde_json::to_vec(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": sample_run_record("run-1", created_at).settings, - "graph": sample_run_record("run-1", created_at).graph, - "workflow_slug": sample_run_record("run-1", created_at).workflow_slug, - "working_directory": sample_run_record("run-1", created_at).working_directory, - "host_repo_path": sample_run_record("run-1", created_at).host_repo_path, - "base_branch": sample_run_record("run-1", created_at).base_branch, - "labels": sample_run_record("run-1", created_at).labels, - }), - )) - .unwrap(), - ) - .await - .unwrap(); - db.close().await.unwrap(); - object_store - .put( - &catalog::by_start_path("runs/", &test_run_id("run-1")), - serde_json::to_vec(&record).unwrap().into(), - ) - .await - .unwrap(); - - store.delete_run(&test_run_id("run-1")).await.unwrap(); - store.delete_run(&test_run_id("run-1")).await.unwrap(); - assert!(list_paths(object_store, "runs").await.is_empty()); - } - - #[tokio::test] - async fn delete_run_closes_active_handles() { - let (object_store, store) = make_store(); - let _created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let blob_id = run.write_blob(br#"{"done":true}"#).await.unwrap(); - assert_eq!( - run.read_blob(&blob_id).await.unwrap(), - Some(Bytes::from_static(br#"{"done":true}"#)) - ); - - store.delete_run(&test_run_id("run-1")).await.unwrap(); - - let err = run.write_blob(br#"{"done":false}"#).await.unwrap_err(); - assert!(matches!( - err, - StoreError::Slate(err) if matches!(err.kind(), ErrorKind::Closed(CloseReason::Clean)) - )); - assert!(store.open_run(&test_run_id("run-1")).await.is_err()); - assert!(list_paths(object_store, "runs").await.is_empty()); - } - - #[tokio::test] - async fn repair_catalog_removes_stale_wrong_time_prefixes() { - let (object_store, store) = make_store(); - let _created_at = dt("2026-03-27T12:00:00Z"); - let _wrong_time = dt("2026-03-27T11:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - run.write_blob(br#"{"done":true}"#).await.unwrap(); - - let _locator = catalog::read_locator(object_store.clone(), "runs/", &test_run_id("run-1")) - .await - .unwrap(); - object_store - .put( - &catalog::by_start_path("runs/", &test_run_id("run-1")), - Bytes::new().into(), - ) - .await - .unwrap(); - - repair_catalog_for_tests(&store).await.unwrap(); - let paths = list_paths(object_store, "runs/by-start").await; - assert_eq!(paths.len(), 1); - assert!(paths[0].contains(&format!("2026-03-27-12-00/{}.json", test_run_id("run-1")))); - } - - #[tokio::test] - async fn create_run_uses_distinct_db_prefix_for_same_minute_orphan() { - let (object_store, store) = make_store(); - let old_created_at = dt("2026-03-27T12:00:00Z"); - let _new_created_at = dt("2026-03-27T12:00:30Z"); - let orphan = CatalogRecord { - run_id: test_run_id("run-1"), - created_at: old_created_at, - db_prefix: catalog::db_prefix("runs/", &test_run_id("run-1")), - run_dir: None, - }; - let new_prefix = catalog::db_prefix("runs/", &test_run_id("run-1")); - assert_eq!(orphan.db_prefix, new_prefix); - - let db = seed_db(object_store.clone(), &orphan, true).await; - db.close().await.unwrap(); - - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - assert!(run.state().await.unwrap().graph_source.is_none()); - - let locator = catalog::read_locator(object_store, "runs/", &test_run_id("run-1")) - .await - .unwrap(); - assert!(locator); - } - - #[tokio::test] - async fn create_run_rejects_mismatched_init_for_existing_prefix() { - let (object_store, store) = make_store(); - let _created_at = dt("2026-03-27T12:00:00Z"); - let db_prefix = catalog::db_prefix("runs/", &test_run_id("run-1")); - let db = slatedb::Db::builder(db_prefix.clone(), object_store) - .with_settings(SlateSettings { - flush_interval: Some(Duration::from_millis(1)), - ..SlateSettings::default() - }) - .build() - .await - .unwrap(); - db.put( - keys::init(), - serde_json::to_vec(&test_run_id("other-run")).unwrap(), - ) - .await - .unwrap(); - db.close().await.unwrap(); - - let Err(err) = store.create_run(&test_run_id("run-1")).await else { - panic!("expected create_run to reject mismatched _init.json"); - }; - assert!(matches!( - err, - StoreError::Other(message) if message.contains("_init.json") - )); - } - - #[tokio::test] - async fn slate_run_store_round_trips_assets_and_projects_events() { - let (_object_store, store) = make_store(); - let _created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let node = StageId::new("code", 2); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:01:00Z", - "stage.prompt", - Some("code"), - serde_json::json!({"text": "Plan", "visit": 2}), - )) - .await - .unwrap(); - run.put_artifact(&node, "src/lib.rs", b"fn main() {}") - .await - .unwrap(); - - let state = run.state().await.unwrap(); - let node_state = state.node(&node).unwrap(); - assert_eq!(node_state.prompt, Some("Plan".to_string())); - assert_eq!( - run.get_artifact(&node, "src/lib.rs").await.unwrap(), - Some(Bytes::from_static(b"fn main() {}")) - ); - } - - #[tokio::test] - async fn slate_run_store_lists_blobs_and_assets() { - let (_object_store, store) = make_store(); - let _created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let summary_blob = run.write_blob(br#"{"done":true}"#).await.unwrap(); - let plan_blob = run.write_blob(br#"{"steps":3}"#).await.unwrap(); - - let snapshot_node = StageId::new("code", 2); - run.put_artifact(&snapshot_node, "src/lib.rs", b"fn main() {}") - .await - .unwrap(); - - let artifact_only_node = StageId::new("artifact-only", 7); - run.put_artifact(&artifact_only_node, "logs/output.txt", b"hello") - .await - .unwrap(); - - assert_eq!( - run.list_blobs().await.unwrap(), - vec![plan_blob, summary_blob] - ); - assert_eq!( - run.list_all_artifacts().await.unwrap(), - vec![ - crate::slate::NodeArtifact { - node: crate::StageId::new("artifact-only", 7), - filename: "logs/output.txt".to_string(), - }, - crate::slate::NodeArtifact { - node: crate::StageId::new("code", 2), - filename: "src/lib.rs".to_string(), - } - ] - ); - } - - #[tokio::test] - async fn slate_run_store_lists_events_with_limit() { - let (_object_store, store) = make_store(); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - - for (idx, ts) in [ - "2026-03-27T12:00:00Z", - "2026-03-27T12:00:01Z", - "2026-03-27T12:00:02Z", - ] - .into_iter() - .enumerate() - { - run.append_event(&event_payload( - "run-1", - ts, - "run.submitted", - None, - serde_json::json!({"index": idx}), - )) - .await - .unwrap(); - } - - let events = run.list_events_from_with_limit(2, 1).await.unwrap(); - assert_eq!(events.iter().map(|event| event.seq).collect::>(), vec![2, 3]); - } - - #[tokio::test] - async fn slate_run_store_lists_artifacts_for_stage_only() { - let (_object_store, store) = make_store(); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let code_stage = StageId::new("code", 2); - let build_stage = StageId::new("build", 1); - - run.put_artifact(&code_stage, "src/lib.rs", b"fn main() {}") - .await - .unwrap(); - run.put_artifact(&code_stage, "src/main.rs", b"fn main() {}") - .await - .unwrap(); - run.put_artifact(&build_stage, "target/output.txt", b"ok") - .await - .unwrap(); - - assert_eq!( - run.list_artifacts_for_stage(&code_stage).await.unwrap(), - vec!["src/lib.rs".to_string(), "src/main.rs".to_string()] - ); - } - - #[tokio::test] - async fn create_run_state_and_node_storage_round_trip() { - let (_object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - - let run_record = sample_run_record("run-1", created_at); - let start_record = sample_start_record("run-1", created_at); - let status_record = - sample_status(RunStatus::Running, Some(StatusReason::SandboxInitializing)); - let checkpoint = sample_checkpoint(); - let conclusion = sample_conclusion(); - let retro = sample_retro("run-1"); - let sandbox = sample_sandbox(); - let node = StageId::new("code", 2); - let pull_request = sample_pull_request(); - - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": run_record.settings, - "graph": run_record.graph, - "workflow_source": "digraph night_sky {}", - "workflow_slug": run_record.workflow_slug, - "working_directory": run_record.working_directory, - "host_repo_path": run_record.host_repo_path, - "base_branch": run_record.base_branch, - "labels": run_record.labels, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:05Z", - "run.started", - None, - serde_json::json!({ - "run_branch": start_record.run_branch, - "base_sha": start_record.base_sha, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:06Z", - "run.running", - None, - serde_json::json!({ - "reason": status_record.reason, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:07Z", - "checkpoint.completed", - Some("code"), - serde_json::json!({ - "status": "success", - "current_node": checkpoint.current_node, - "completed_nodes": checkpoint.completed_nodes, - "node_retries": checkpoint.node_retries, - "context_values": checkpoint.context_values, - "node_outcomes": checkpoint.node_outcomes, - "next_node_id": checkpoint.next_node_id, - "git_commit_sha": checkpoint.git_commit_sha, - "loop_failure_signatures": serde_json::json!({}), - "restart_failure_signatures": serde_json::json!({}), - "node_visits": checkpoint.node_visits, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08Z", - "sandbox.initialized", - None, - serde_json::json!({ - "provider": sandbox.provider, - "working_directory": sandbox.working_directory, - "identifier": sandbox.identifier, - "host_working_directory": sandbox.host_working_directory, - "container_mount_point": sandbox.container_mount_point, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08.1Z", - "stage.prompt", - Some("code"), - serde_json::json!({ - "visit": 2, - "text": "Plan the fix", - "mode": "prompt", - "provider": "openai", - "model": "gpt-5.4" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08.2Z", - "command.started", - Some("code"), - serde_json::json!({ - "visit": 2, - "command": "cargo test" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08.3Z", - "command.completed", - Some("code"), - serde_json::json!({ - "visit": 2, - "stdout": "ok", - "stderr": "", - "exit_code": 0 - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08.4Z", - "checkpoint.completed", - Some("code"), - serde_json::json!({ - "status": "success", - "ordinal": 2, - "current_node": checkpoint.current_node, - "completed_nodes": checkpoint.completed_nodes, - "node_retries": checkpoint.node_retries, - "context_values": checkpoint.context_values, - "node_outcomes": checkpoint.node_outcomes, - "next_node_id": checkpoint.next_node_id, - "git_commit_sha": checkpoint.git_commit_sha, - "node_visits": checkpoint.node_visits, - "diff": "diff --git a/src/lib.rs b/src/lib.rs" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08.5Z", - "parallel.completed", - Some("code"), - serde_json::json!({ - "visit": 2, - "results": [{"node_id": "lint", "status": "success"}] - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08.6Z", - "stage.completed", - Some("code"), - serde_json::json!({ - "visit": 2, - "status": "success", - "notes": "all good", - "response": "Implemented", - "files_touched": ["src/lib.rs"] - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08.7Z", - "retro.started", - None, - serde_json::json!({ - "prompt": "How did it go?" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:09Z", - "retro.completed", - None, - serde_json::json!({ - "response": "Smooth enough", - "retro": retro, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:10Z", - "run.completed", - None, - serde_json::json!({ - "status": conclusion.status, - "duration_ms": conclusion.duration_ms, - "total_cost": conclusion.total_cost, - "final_git_commit_sha": conclusion.final_git_commit_sha, - "final_patch": "diff --git a/src/lib.rs b/src/lib.rs\n", - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:11Z", - "pull_request.created", - None, - serde_json::json!({ - "pr_url": pull_request.html_url, - "pr_number": pull_request.number, - "owner": pull_request.owner, - "repo": pull_request.repo, - "base_branch": pull_request.base_branch, - "head_branch": pull_request.head_branch, - "title": pull_request.title, - }), - )) - .await - .unwrap(); - let summary_blob = run.write_blob(br#"{"done":true}"#).await.unwrap(); - run.put_artifact(&node, "src/lib.rs", b"fn main() {}") - .await - .unwrap(); - - let state = run.state().await.unwrap(); - let stored_run = state.run.as_ref().unwrap(); - assert_eq!(stored_run.run_id, run_record.run_id); - assert_eq!( - stored_run.run_id.created_at(), - run_record.run_id.created_at() - ); - assert_eq!(stored_run.workflow_slug, run_record.workflow_slug); - assert_eq!(stored_run.graph.name, run_record.graph.name); - assert_eq!(state.graph_source.as_deref(), Some("digraph night_sky {}")); - - let stored_start = state.start.as_ref().unwrap(); - assert_eq!(stored_start.run_id, start_record.run_id); - assert_eq!(stored_start.start_time, start_record.start_time); - - let stored_status = state.status.as_ref().unwrap(); - assert_eq!(stored_status.status, RunStatus::Succeeded); - assert_eq!(stored_status.reason, None); - - let stored_checkpoint = state.checkpoint.as_ref().unwrap(); - assert_eq!(stored_checkpoint.current_node, checkpoint.current_node); - assert_eq!(stored_checkpoint.next_node_id, checkpoint.next_node_id); - - let stored_conclusion = state.conclusion.as_ref().unwrap(); - assert_eq!(stored_conclusion.status, conclusion.status); - assert_eq!(stored_conclusion.duration_ms, conclusion.duration_ms); - assert_eq!(stored_conclusion.total_cost, conclusion.total_cost); - - let stored_retro = state.retro.as_ref().unwrap(); - assert_eq!(stored_retro.run_id, retro.run_id); - assert_eq!(stored_retro.intent, retro.intent); - let stored_sandbox = state.sandbox.as_ref().unwrap(); - assert_eq!(stored_sandbox.provider, sandbox.provider); - assert_eq!(stored_sandbox.working_directory, sandbox.working_directory); - assert_eq!(state.retro_prompt.as_deref(), Some("How did it go?")); - assert_eq!(state.retro_response.as_deref(), Some("Smooth enough")); - assert_eq!( - run.read_blob(&summary_blob).await.unwrap(), - Some(Bytes::from_static(br#"{"done":true}"#)) - ); - assert_eq!( - run.get_artifact(&node, "src/lib.rs").await.unwrap(), - Some(Bytes::from_static(b"fn main() {}")) - ); - assert_eq!( - state.final_patch.as_deref(), - Some("diff --git a/src/lib.rs b/src/lib.rs\n") - ); - assert_eq!(state.pull_request, Some(pull_request.clone())); - assert!(state.iter_nodes().any(|(node, _)| node.node_id() == "code")); - let node_state = state - .node(&node) - .expect("node state should exist for code:2"); - assert_eq!(node_state.prompt.as_deref(), Some("Plan the fix")); - assert_eq!(node_state.response.as_deref(), Some("Implemented")); - assert_eq!(node_state.stdout.as_deref(), Some("ok")); - assert_eq!(node_state.stderr.as_deref(), Some("")); - assert_eq!( - node_state.diff.as_deref(), - Some("diff --git a/src/lib.rs b/src/lib.rs") - ); - assert_eq!( - node_state - .provider_used - .as_ref() - .and_then(|v| v.get("provider")) - .and_then(|v| v.as_str()), - Some("openai") - ); - } - - #[tokio::test] - async fn state_projects_event_stream() { - let (_object_store, store) = make_store(); - let created_at = dt("2026-03-27T12:00:00Z"); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let run_record = sample_run_record("run-1", created_at); - let retro = sample_retro("run-1"); - - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": run_record.settings, - "graph": run_record.graph, - "workflow_source": "digraph night_sky {}", - "workflow_slug": run_record.workflow_slug, - "working_directory": run_record.working_directory, - "host_repo_path": run_record.host_repo_path, - "base_branch": run_record.base_branch, - "labels": run_record.labels, - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:05Z", - "run.started", - None, - serde_json::json!({ - "run_branch": "fabro/run/demo", - "base_sha": "abc123" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:06Z", - "run.running", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:07Z", - "stage.prompt", - Some("code"), - serde_json::json!({ - "visit": 2, - "text": "Plan the fix", - "mode": "prompt", - "provider": "openai", - "model": "gpt-5.4" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:08Z", - "stage.completed", - Some("code"), - serde_json::json!({ - "status": "success", - "notes": "all good", - "response": "Implemented", - "files_touched": ["src/lib.rs"], - "node_visits": {"code": 2} - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:09Z", - "checkpoint.completed", - Some("code"), - serde_json::json!({ - "status": "success", - "current_node": "code", - "completed_nodes": ["plan"], - "context_values": {"artifact": {"kind": "summary"}}, - "next_node_id": "review", - "git_commit_sha": "def456", - "node_visits": {"code": 2}, - "diff": "diff --git a/src/lib.rs b/src/lib.rs" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:10Z", - "sandbox.initialized", - None, - serde_json::json!({ - "provider": "local", - "working_directory": "/tmp/night-sky", - "identifier": "sandbox-1", - "host_working_directory": "/tmp/night-sky" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:11Z", - "retro.started", - None, - serde_json::json!({ - "prompt": "How did it go?" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:12Z", - "retro.completed", - None, - serde_json::json!({ - "response": "Smooth enough", - "retro": retro - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:13Z", - "pull_request.created", - None, - serde_json::json!({ - "pr_url": "https://github.com/fabro-sh/fabro/pull/123", - "pr_number": 123, - "owner": "fabro-sh", - "repo": "fabro", - "base_branch": "main", - "head_branch": "fabro/run/demo", - "title": "Map the constellations", - "draft": false - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:15Z", - "run.completed", - None, - serde_json::json!({ - "duration_ms": 3210, - "artifact_count": 1, - "status": "success", - "total_cost": 1.25, - "final_git_commit_sha": "feedbeef", - "final_patch": "diff --git a/src/lib.rs b/src/lib.rs\n" - }), - )) - .await - .unwrap(); - - let state = run.state().await.unwrap(); - assert_eq!( - state.run.as_ref().map(|run| run.run_id), - Some(test_run_id("run-1")) - ); - assert_eq!(state.graph_source.as_deref(), Some("digraph night_sky {}")); - assert_eq!( - state - .start - .as_ref() - .and_then(|start| start.run_branch.as_deref()), - Some("fabro/run/demo") - ); - assert_eq!( - state.status.as_ref().map(|status| status.status), - Some(RunStatus::Succeeded) - ); - assert_eq!( - state - .checkpoint - .as_ref() - .map(|checkpoint| checkpoint.current_node.as_str()), - Some("code") - ); - assert_eq!(state.checkpoints.len(), 1); - assert_eq!( - state.final_patch.as_deref(), - Some("diff --git a/src/lib.rs b/src/lib.rs\n") - ); - assert_eq!(state.retro_prompt.as_deref(), Some("How did it go?")); - assert_eq!(state.retro_response.as_deref(), Some("Smooth enough")); - assert_eq!(state.pull_request.as_ref().map(|pr| pr.number), Some(123)); - assert_eq!( - state - .sandbox - .as_ref() - .map(|sandbox| sandbox.provider.as_str()), - Some("local") - ); - assert_eq!(state.list_node_visits("code"), vec![2]); - let node = state.node(&StageId::new("code", 2)).unwrap(); - assert_eq!(node.prompt.as_deref(), Some("Plan the fix")); - assert_eq!(node.response.as_deref(), Some("Implemented")); - assert_eq!( - node.diff.as_deref(), - Some("diff --git a/src/lib.rs b/src/lib.rs") - ); - assert_eq!( - node.provider_used - .as_ref() - .and_then(|value| value.get("provider")) - .and_then(|value| value.as_str()), - Some("openai") - ); - } - - #[tokio::test] - async fn state_rewind_keeps_active_projection_only() { - let (_object_store, store) = make_store(); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": Settings::default(), - "graph": Graph::new("night-sky"), - "working_directory": "/tmp/night-sky", - "labels": {} - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:01Z", - "stage.prompt", - Some("code"), - serde_json::json!({ - "visit": 1, - "text": "before rewind" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:02Z", - "pull_request.created", - None, - serde_json::json!({ - "pr_url": "https://github.com/fabro-sh/fabro/pull/123", - "pr_number": 123, - "owner": "fabro-sh", - "repo": "fabro", - "base_branch": "main", - "head_branch": "fabro/run/demo", - "title": "Map the constellations", - "draft": false - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:03Z", - "run.completed", - None, - serde_json::json!({ - "duration_ms": 10, - "artifact_count": 0, - "status": "success", - "final_patch": "old patch" - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:04Z", - "run.rewound", - None, - serde_json::json!({ - "target_checkpoint_ordinal": 1, - "target_node_id": "plan", - "target_visit": 1 - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:05Z", - "checkpoint.completed", - Some("plan"), - serde_json::json!({ - "status": "success", - "current_node": "plan", - "completed_nodes": [], - "node_visits": {"plan": 1} - }), - )) - .await - .unwrap(); - run.append_event(&event_payload( - "run-1", - "2026-03-27T12:00:06Z", - "run.submitted", - None, - serde_json::json!({}), - )) - .await - .unwrap(); - - let state = run.state().await.unwrap(); - assert_eq!( - state.status.as_ref().map(|status| status.status), - Some(RunStatus::Submitted) - ); - assert!(state.conclusion.is_none()); - assert!(state.final_patch.is_none()); - assert!(state.pull_request.is_none()); - assert_eq!(state.checkpoints.len(), 1); - assert_eq!( - state - .checkpoint - .as_ref() - .map(|checkpoint| checkpoint.current_node.as_str()), - Some("plan") - ); - assert!(state.is_empty()); - } - - #[tokio::test] - async fn append_event_validates_payload_shape_and_run_id() { - let (_object_store, store) = make_store(); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - - let invalid_missing: EventPayload = serde_json::from_value(serde_json::json!({ - "run_id": "run-1" - })) - .unwrap(); - let err = run.append_event(&invalid_missing).await.unwrap_err(); - assert!(matches!(err, StoreError::InvalidEvent(_))); - - let invalid_run_id: EventPayload = serde_json::from_value(serde_json::json!({ - "id": "evt-invalid-run", - "ts": "2026-03-27T12:00:00Z", - "run_id": "other-run", - "event": "StageStarted" - })) - .unwrap(); - let err = run.append_event(&invalid_run_id).await.unwrap_err(); - assert!(matches!(err, StoreError::InvalidEvent(_))); - } - - #[tokio::test] - async fn state_retains_checkpoint_history_by_event_sequence() { - let (_object_store, store) = make_store(); - let run = store.create_run(&test_run_id("run-1")).await.unwrap(); - let checkpoint = sample_checkpoint(); - let seq = run - .append_event(&event_payload( - "run-1", - "2026-03-27T12:00:00Z", - "checkpoint.completed", - Some(&checkpoint.current_node), - serde_json::json!({ - "status": "success", - "current_node": checkpoint.current_node, - "completed_nodes": checkpoint.completed_nodes, - "node_retries": checkpoint.node_retries, - "context_values": checkpoint.context_values, - "node_outcomes": checkpoint.node_outcomes, - "next_node_id": checkpoint.next_node_id, - "git_commit_sha": checkpoint.git_commit_sha, - "loop_failure_signatures": serde_json::json!({}), - "restart_failure_signatures": serde_json::json!({}), - "node_visits": checkpoint.node_visits, - }), - )) - .await - .unwrap(); - let state = run.state().await.unwrap(); - assert_eq!(seq, 1); - assert_eq!(state.checkpoints.len(), 1); - assert_eq!(state.checkpoints[0].0, 1); - assert_eq!(state.checkpoints[0].1.current_node, checkpoint.current_node); - assert_eq!( - state.checkpoint.as_ref().unwrap().current_node, - checkpoint.current_node - ); - } - - #[tokio::test] - async fn list_runs_filters_dates_and_tolerates_missing_status() { - let (_object_store, store) = make_store(); - let early = dt("2026-03-27T10:00:00Z"); - let late = dt("2026-03-27T12:00:00Z"); - - let early_run = store.create_run(&test_run_id("run-early")).await.unwrap(); - let early_record = sample_run_record("run-early", early); - early_run - .append_event(&event_payload( - "run-early", - "2026-03-27T10:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": early_record.settings, - "graph": early_record.graph, - "workflow_slug": early_record.workflow_slug, - "working_directory": early_record.working_directory, - "host_repo_path": early_record.host_repo_path, - "base_branch": early_record.base_branch, - "labels": early_record.labels, - }), - )) - .await - .unwrap(); - - let late_run = store.create_run(&test_run_id("run-late")).await.unwrap(); - let late_record = sample_run_record("run-late", late); - late_run - .append_event(&event_payload( - "run-late", - "2026-03-27T12:00:00Z", - "run.created", - None, - serde_json::json!({ - "settings": late_record.settings, - "graph": late_record.graph, - "workflow_slug": late_record.workflow_slug, - "working_directory": late_record.working_directory, - "host_repo_path": late_record.host_repo_path, - "base_branch": late_record.base_branch, - "labels": late_record.labels, - }), - )) - .await - .unwrap(); - late_run - .append_event(&event_payload( - "run-late", - "2026-03-27T12:00:01Z", - "run.started", - None, - serde_json::json!({ - "run_branch": "fabro/run/demo", - "base_sha": "abc123", - }), - )) - .await - .unwrap(); - late_run - .append_event(&event_payload( - "run-late", - "2026-03-27T12:00:02Z", - "run.completed", - None, - serde_json::json!({ - "duration_ms": 3210, - "artifact_count": 1, - "status": "success", - "reason": "completed", - "total_cost": 1.25, - }), - )) - .await - .unwrap(); - - let all = store.list_runs(&ListRunsQuery::default()).await.unwrap(); - assert_eq!(all.len(), 2); - assert_eq!(all[0].run_id, test_run_id("run-late")); - assert_eq!(all[0].workflow_name, Some("night-sky".to_string())); - assert_eq!(all[0].goal, Some("map the constellations".to_string())); - assert_eq!( - all[0].host_repo_path, - Some("github.com/fabro-sh/fabro".to_string()) - ); - assert_eq!(all[0].duration_ms, Some(3210)); - assert_eq!(all[0].total_cost, Some(1.25)); - assert_eq!(all[0].status_reason, Some(StatusReason::Completed)); - assert_eq!(all[1].status, None); - - let filtered = store - .list_runs(&ListRunsQuery { - start: Some(dt("2026-03-27T11:00:00Z")), - end: Some(dt("2026-03-27T13:00:00Z")), - }) - .await - .unwrap(); - assert_eq!(filtered.len(), 1); - assert_eq!(filtered[0].run_id, test_run_id("run-late")); } } diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index 9d70a78b4..467636135 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -1,15 +1,13 @@ +use std::collections::VecDeque; use std::sync::atomic::{AtomicU32, Ordering}; -use std::sync::{Arc, Weak}; -use std::time::Duration; +use std::sync::Arc; use bytes::Bytes; use chrono::Utc; use futures::Stream; -use serde::Serialize; use serde::de::DeserializeOwned; -use slatedb::{CloseReason, DbRead, DbReader, ErrorKind}; -use tokio::sync::{Mutex, mpsc}; -use tokio::time; +use slatedb::{CloseReason, Db, DbRead, ErrorKind}; +use tokio::sync::{Mutex, broadcast, mpsc}; use tokio_stream::wrappers::UnboundedReceiverStream; use crate::keys; @@ -17,6 +15,8 @@ use crate::run_state::EventProjectionCache; use crate::{EventEnvelope, EventPayload, Result, RunProjection, RunSummary, StageId, StoreError}; use fabro_types::{RunBlobId, RunId}; +const DEFAULT_EVENT_TAIL_LIMIT: usize = 1024; + #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] pub struct NodeArtifact { pub node: StageId, @@ -26,108 +26,113 @@ pub struct NodeArtifact { #[derive(Clone)] pub struct SlateRunStore { inner: Arc, + read_only: bool, } impl std::fmt::Debug for SlateRunStore { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("SlateRunStore") .field("run_id", &self.inner.run_id) - .field("db_prefix", &self.inner.db_prefix) + .field("read_only", &self.read_only) .finish_non_exhaustive() } } pub(crate) struct SlateRunStoreInner { run_id: RunId, - db_prefix: String, - db: SlateRunDb, + db: Db, event_seq: AtomicU32, close_lock: Mutex<()>, projection_cache: Mutex, -} - -enum SlateRunDb { - Writer(slatedb::Db), - Reader(Box), + recent_events: Mutex>, + recent_event_limit: usize, + event_tx: broadcast::Sender, } impl SlateRunStore { - pub(crate) async fn open_writer( - run_id: RunId, - db_prefix: String, - db: slatedb::Db, - ) -> Result { - let event_seq = recover_next_seq(&db, keys::EVENTS_PREFIX, keys::parse_event_seq).await?; + pub(crate) async fn open_writer(run_id: RunId, db: Db) -> Result { + let event_seq = recover_next_seq(&db, &keys::events_prefix(&run_id), keys::parse_event_seq).await?; + let (event_tx, _) = broadcast::channel(DEFAULT_EVENT_TAIL_LIMIT.max(16)); Ok(Self { inner: Arc::new(SlateRunStoreInner { run_id, - db_prefix, - db: SlateRunDb::Writer(db), + db, event_seq: AtomicU32::new(event_seq), close_lock: Mutex::new(()), projection_cache: Mutex::new(EventProjectionCache::default()), + recent_events: Mutex::new(VecDeque::with_capacity(DEFAULT_EVENT_TAIL_LIMIT)), + recent_event_limit: DEFAULT_EVENT_TAIL_LIMIT, + event_tx, }), + read_only: false, }) } - pub(crate) async fn open_reader( - run_id: RunId, - db_prefix: String, - db: DbReader, - ) -> Result { - let event_seq = recover_next_seq(&db, keys::EVENTS_PREFIX, keys::parse_event_seq).await?; + pub(crate) async fn open_reader(run_id: RunId, db: Db) -> Result { + let event_seq = recover_next_seq(&db, &keys::events_prefix(&run_id), keys::parse_event_seq).await?; + let (event_tx, _) = broadcast::channel(DEFAULT_EVENT_TAIL_LIMIT.max(16)); Ok(Self { inner: Arc::new(SlateRunStoreInner { run_id, - db_prefix, - db: SlateRunDb::Reader(Box::new(db)), + db, event_seq: AtomicU32::new(event_seq), close_lock: Mutex::new(()), projection_cache: Mutex::new(EventProjectionCache::default()), + recent_events: Mutex::new(VecDeque::with_capacity(DEFAULT_EVENT_TAIL_LIMIT)), + recent_event_limit: DEFAULT_EVENT_TAIL_LIMIT, + event_tx, }), + read_only: true, }) } pub(crate) fn from_inner(inner: Arc) -> Self { - Self { inner } + Self { + inner, + read_only: false, + } } - pub(crate) fn downgrade(&self) -> Weak { - Arc::downgrade(&self.inner) + pub(crate) fn into_read_only(&self) -> Self { + Self { + inner: Arc::clone(&self.inner), + read_only: true, + } + } + + pub(crate) fn inner_arc(&self) -> Arc { + Arc::clone(&self.inner) } pub(crate) fn run_id(&self) -> RunId { self.inner.run_id } - pub(crate) fn matches_run(&self, run_id: &RunId, db_prefix: &str) -> bool { - self.inner.run_id == *run_id && self.inner.db_prefix == db_prefix + pub(crate) fn matches_run(&self, run_id: &RunId) -> bool { + self.inner.run_id == *run_id } pub(crate) async fn close(&self) -> Result<()> { let _guard = self.inner.close_lock.lock().await; - match self.inner.db.close().await { - Ok(()) => Ok(()), - Err(err) if matches!(err.kind(), ErrorKind::Closed(CloseReason::Clean)) => Ok(()), - Err(err) => Err(err.into()), + if Arc::strong_count(&self.inner) <= 1 { + match self.inner.db.close().await { + Ok(()) => Ok(()), + Err(err) if matches!(err.kind(), ErrorKind::Closed(CloseReason::Clean)) => Ok(()), + Err(err) => Err(err.into()), + } + } else { + Ok(()) } } - pub(crate) async fn snapshot(&self) -> Result> { - match &self.inner.db { - SlateRunDb::Writer(db) => Ok(db.snapshot().await?), - SlateRunDb::Reader(_) => Err(StoreError::ReadOnly), - } - } - - pub(crate) async fn validate_init(db: &R, expected: &RunId) -> Result + pub(crate) async fn validate_init(db: &R, run_id: &RunId) -> Result where R: DbRead + Sync, { - match get_json::(db, keys::init()).await? { - Some(existing) if existing == *expected => Ok(true), + match get_json::(db, &keys::init_key(run_id)).await? { + Some(existing) if existing == *run_id => Ok(true), Some(existing) => Err(StoreError::Other(format!( - "existing _init.json {existing:?} does not match requested run_id {expected:?}" + "existing init record {existing:?} does not match requested run_id {run_id:?}" ))), None => Ok(false), } @@ -137,7 +142,7 @@ impl SlateRunStore { where R: DbRead + Sync, { - let events = list_events_from(db, 1).await?; + let events = list_events_from(db, run_id, 1).await?; let state = RunProjection::apply_events(&events)?; Ok(state.build_summary(run_id)) } @@ -147,7 +152,7 @@ impl SlateRunStore { let cache = self.inner.projection_cache.lock().await; cache.last_seq.saturating_add(1) }; - let events = self.inner.db.list_events_from(next_seq).await?; + let events = list_events_from(&self.inner.db, &self.inner.run_id, next_seq).await?; let mut cache = self.inner.projection_cache.lock().await; for event in &events { cache.state.apply_event(event)?; @@ -155,24 +160,62 @@ impl SlateRunStore { } Ok(cache.state.clone()) } + + async fn cache_event(&self, event: &EventEnvelope) -> Result<()> { + { + let mut projection_cache = self.inner.projection_cache.lock().await; + projection_cache.state.apply_event(event)?; + projection_cache.last_seq = event.seq; + } + let mut recent_events = self.inner.recent_events.lock().await; + recent_events.push_back(event.clone()); + while recent_events.len() > self.inner.recent_event_limit { + recent_events.pop_front(); + } + let _ = self.inner.event_tx.send(event.clone()); + Ok(()) + } + + async fn cached_events_from(&self, start_seq: u32, limit: usize) -> Option> { + let recent_events = self.inner.recent_events.lock().await; + let oldest_seq = recent_events.front().map(|event| event.seq)?; + if start_seq < oldest_seq { + return None; + } + let mut events = recent_events + .iter() + .filter(|event| event.seq >= start_seq) + .take(limit.saturating_add(1)) + .cloned() + .collect::>(); + if events.is_empty() && start_seq <= self.inner.event_seq.load(Ordering::SeqCst) { + events = Vec::new(); + } + Some(events) + } } impl SlateRunStore { pub async fn append_event(&self, payload: &EventPayload) -> Result { + if self.read_only { + return Err(StoreError::ReadOnly); + } payload.validate(&self.inner.run_id)?; let seq = self.inner.event_seq.fetch_add(1, Ordering::SeqCst); - self.inner - .db - .put_json( - &keys::event_key(seq, Utc::now().timestamp_millis()), - payload, - ) - .await?; + let event = EventEnvelope { + seq, + payload: payload.clone(), + }; + self.inner.db.put( + keys::event_key(&self.inner.run_id, seq, Utc::now().timestamp_millis()), + serde_json::to_vec(payload)?, + ).await?; + self.cache_event(&event).await?; Ok(seq) } pub async fn list_events(&self) -> Result> { - self.inner.db.list_events_from(1).await + self.list_events_from_with_limit(1, usize::MAX / 2).await } pub async fn list_events_from_with_limit( @@ -180,10 +223,10 @@ impl SlateRunStore { start_seq: u32, limit: usize, ) -> Result> { - self.inner - .db - .list_events_from_with_limit(start_seq, limit) - .await + if let Some(events) = self.cached_events_from(start_seq, limit).await { + return Ok(events); + } + list_events_from_with_limit(&self.inner.db, &self.inner.run_id, start_seq, limit).await } pub fn watch_events_from( @@ -192,72 +235,86 @@ impl SlateRunStore { ) -> Result> + Send>>> { let inner = Arc::clone(&self.inner); let (sender, receiver) = mpsc::unbounded_channel(); - tokio::spawn(async move { + let cached = { + let recent_events = inner.recent_events.lock().await; + recent_events + .iter() + .filter(|event| event.seq >= seq) + .cloned() + .collect::>() + }; let mut next_seq = seq; - loop { - if sender.is_closed() { + for event in cached { + next_seq = event.seq.saturating_add(1); + if sender.send(Ok(event)).is_err() { return; } + } - match inner.db.list_events_from(next_seq).await { - Ok(events) => { - if events.is_empty() { - time::sleep(Duration::from_millis(100)).await; - continue; - } - for event in events { - next_seq = event.seq.saturating_add(1); - if sender.send(Ok(event)).is_err() { - return; - } - } - } - Err(err) => { - let _ = sender.send(Err(err)); - return; - } + let mut rx = inner.event_tx.subscribe(); + while let Ok(event) = rx.recv().await { + if event.seq < next_seq { + continue; + } + next_seq = event.seq.saturating_add(1); + if sender.send(Ok(event)).is_err() { + return; } } }); - Ok(Box::pin(UnboundedReceiverStream::new(receiver))) } pub async fn write_blob(&self, data: &[u8]) -> Result { + if self.read_only { + return Err(StoreError::ReadOnly); + } let id = RunBlobId::new(&self.inner.run_id, data); - self.inner.db.put_bytes(&keys::blob_key(&id), data).await?; + self.inner + .db + .put(keys::blob_key(&self.inner.run_id, &id), data) + .await?; Ok(id) } pub async fn read_blob(&self, id: &RunBlobId) -> Result> { - self.inner.db.get_bytes(&keys::blob_key(id)).await + Ok(self + .inner + .db + .get(keys::blob_key(&self.inner.run_id, id)) + .await?) } pub async fn list_blobs(&self) -> Result> { - self.inner.db.list_blobs().await + list_blobs(&self.inner.db, &self.inner.run_id).await } pub async fn put_artifact(&self, node: &StageId, filename: &str, data: &[u8]) -> Result<()> { + if self.read_only { + return Err(StoreError::ReadOnly); + } self.inner .db - .put_bytes(&keys::node_artifact(node, filename), data) - .await + .put(keys::node_artifact(&self.inner.run_id, node, filename), data) + .await?; + Ok(()) } pub async fn get_artifact(&self, node: &StageId, filename: &str) -> Result> { - self.inner + Ok(self + .inner .db - .get_bytes(&keys::node_artifact(node, filename)) - .await + .get(keys::node_artifact(&self.inner.run_id, node, filename)) + .await?) } pub async fn list_all_artifacts(&self) -> Result> { - self.inner.db.list_all_artifacts().await + list_all_artifacts(&self.inner.db, &self.inner.run_id).await } pub async fn list_artifacts_for_stage(&self, stage_id: &StageId) -> Result> { - self.inner.db.list_artifacts_for_stage(stage_id).await + list_artifacts_for_stage(&self.inner.db, &self.inner.run_id, stage_id).await } pub async fn state(&self) -> Result { @@ -265,81 +322,6 @@ impl SlateRunStore { } } -impl SlateRunDb { - fn writer(&self) -> Result<&slatedb::Db> { - match self { - Self::Writer(db) => Ok(db), - Self::Reader(_) => Err(StoreError::ReadOnly), - } - } - - async fn close(&self) -> std::result::Result<(), slatedb::Error> { - match self { - Self::Writer(db) => db.close().await, - Self::Reader(db) => db.close().await, - } - } - - async fn put_json(&self, key: &str, value: &T) -> Result<()> { - put_json(self.writer()?, key, value).await - } - - async fn get_bytes(&self, key: &str) -> Result> { - match self { - Self::Writer(db) => get_bytes(db, key).await, - Self::Reader(db) => db.get(key).await.map_err(Into::into), - } - } - - async fn put_bytes(&self, key: &str, value: &[u8]) -> Result<()> { - put_bytes(self.writer()?, key, value).await - } - - async fn list_events_from(&self, start_seq: u32) -> Result> { - match self { - Self::Writer(db) => list_events_from(db, start_seq).await, - Self::Reader(db) => list_events_from(db.as_ref(), start_seq).await, - } - } - - async fn list_events_from_with_limit( - &self, - start_seq: u32, - limit: usize, - ) -> Result> { - match self { - Self::Writer(db) => list_events_from_with_limit(db, start_seq, limit).await, - Self::Reader(db) => list_events_from_with_limit(db.as_ref(), start_seq, limit).await, - } - } - - async fn list_blobs(&self) -> Result> { - match self { - Self::Writer(db) => list_blobs(db).await, - Self::Reader(db) => list_blobs(db.as_ref()).await, - } - } - - async fn list_all_artifacts(&self) -> Result> { - match self { - Self::Writer(db) => list_all_artifacts(db).await, - Self::Reader(db) => list_all_artifacts(db.as_ref()).await, - } - } - - async fn list_artifacts_for_stage(&self, stage_id: &StageId) -> Result> { - match self { - Self::Writer(db) => list_artifacts_for_stage(db, stage_id).await, - Self::Reader(db) => list_artifacts_for_stage(db.as_ref(), stage_id).await, - } - } -} - -async fn put_json(db: &slatedb::Db, key: &str, value: &T) -> Result<()> { - db.put(key, serde_json::to_vec(value)?).await?; - Ok(()) -} - async fn get_json(db: &R, key: &str) -> Result> where R: DbRead + Sync, @@ -352,15 +334,6 @@ where .map_err(Into::into) } -async fn put_bytes(db: &slatedb::Db, key: &str, value: &[u8]) -> Result<()> { - db.put(key, value).await?; - Ok(()) -} - -async fn get_bytes(db: &slatedb::Db, key: &str) -> Result> { - Ok(db.get(key).await?) -} - async fn recover_next_seq(db: &R, prefix: &str, parse: fn(&str) -> Option) -> Result where R: DbRead + Sync, @@ -376,11 +349,11 @@ where Ok(max_seq.saturating_add(1).max(1)) } -async fn list_events_from(db: &R, start_seq: u32) -> Result> +async fn list_events_from(db: &R, run_id: &RunId, start_seq: u32) -> Result> where R: DbRead + Sync, { - let mut iter = db.scan_prefix(keys::EVENTS_PREFIX.as_bytes()).await?; + let mut iter = db.scan_prefix(keys::events_prefix(run_id).as_bytes()).await?; let mut events = Vec::new(); while let Some(entry) = iter.next().await? { let key = key_to_string(&entry.key)?; @@ -401,22 +374,23 @@ where async fn list_events_from_with_limit( db: &R, + run_id: &RunId, start_seq: u32, limit: usize, ) -> Result> where R: DbRead + Sync, { - let mut events = list_events_from(db, start_seq).await?; + let mut events = list_events_from(db, run_id, start_seq).await?; events.truncate(limit.saturating_add(1)); Ok(events) } -async fn list_blobs(db: &R) -> Result> +async fn list_blobs(db: &R, run_id: &RunId) -> Result> where R: DbRead + Sync, { - let mut iter = db.scan_prefix(keys::BLOBS_PREFIX.as_bytes()).await?; + let mut iter = db.scan_prefix(keys::blobs_prefix(run_id).as_bytes()).await?; let mut blob_ids = Vec::new(); while let Some(entry) = iter.next().await? { let key = key_to_string(&entry.key)?; @@ -429,12 +403,12 @@ where Ok(blob_ids) } -async fn list_all_artifacts(db: &R) -> Result> +async fn list_all_artifacts(db: &R, run_id: &RunId) -> Result> where R: DbRead + Sync, { let mut iter = db - .scan_prefix(keys::ARTIFACT_NODES_PREFIX.as_bytes()) + .scan_prefix(keys::run_prefix(run_id).as_bytes()) .await?; let mut assets = Vec::new(); while let Some(entry) = iter.next().await? { @@ -448,11 +422,11 @@ where Ok(assets) } -async fn list_artifacts_for_stage(db: &R, stage_id: &StageId) -> Result> +async fn list_artifacts_for_stage(db: &R, run_id: &RunId, stage_id: &StageId) -> Result> where R: DbRead + Sync, { - let prefix = keys::node_artifact_prefix(stage_id); + let prefix = keys::node_artifact_prefix(run_id, stage_id); let mut iter = db.scan_prefix(prefix.as_bytes()).await?; let mut filenames = Vec::new(); while let Some(entry) = iter.next().await? {