Add fabro-store with in-memory and SlateDB backends

This commit is contained in:
Bryan Helmkamp 2026-03-28 00:13:25 -04:00
parent 71e151c699
commit c6c227684c
No known key found for this signature in database
14 changed files with 3965 additions and 175 deletions

1264
Cargo.lock generated

File diff suppressed because it is too large Load diff

View file

@ -62,6 +62,8 @@ daytona-api-client = { git = "https://github.com/brynary/daytona-sdk-rust", rev
sentry = { version = "0.35", default-features = false, features = ["backtrace", "contexts", "ureq", "rustls"] }
fork = "0.2"
exec = "0.3"
slatedb = "0.11.2"
object_store = "0.12.5"
[profile.release]
lto = "thin"

View file

@ -0,0 +1,26 @@
[package]
name = "fabro-store"
edition.workspace = true
version.workspace = true
license.workspace = true
[lib]
doctest = false
[dependencies]
fabro-types = { path = "../fabro-types" }
slatedb.workspace = true
object_store.workspace = true
async-trait.workspace = true
tokio = { workspace = true, features = ["full"] }
tokio-stream.workspace = true
serde.workspace = true
serde_json.workspace = true
chrono = { workspace = true, features = ["serde"] }
bytes.workspace = true
thiserror.workspace = true
tracing.workspace = true
futures.workspace = true
[dev-dependencies]
tokio = { workspace = true, features = ["test-util", "macros"] }

View file

@ -0,0 +1,19 @@
pub type Result<T, E = StoreError> = std::result::Result<T, E>;
#[derive(Debug, thiserror::Error)]
pub enum StoreError {
#[error("SlateDB error: {0}")]
Slate(#[from] slatedb::Error),
#[error("Object store error: {0}")]
ObjectStore(#[from] object_store::Error),
#[error("Serialization error: {0}")]
Serde(#[from] serde_json::Error),
#[error("Invalid event payload: {0}")]
InvalidEvent(String),
#[error("Run not found: {0}")]
RunNotFound(String),
#[error("Run already exists: {0}")]
RunAlreadyExists(String),
#[error("{0}")]
Other(String),
}

View file

@ -0,0 +1,208 @@
use crate::NodeVisitRef;
pub const INIT_KEY: &str = "_init.json";
pub const RUN_KEY: &str = "run.json";
pub const START_KEY: &str = "start.json";
pub const STATUS_KEY: &str = "status.json";
pub const CHECKPOINT_KEY: &str = "checkpoint.json";
pub const CONCLUSION_KEY: &str = "conclusion.json";
pub const RETRO_KEY: &str = "retro.json";
pub const GRAPH_KEY: &str = "graph.fabro";
pub const SANDBOX_KEY: &str = "sandbox.json";
pub const RETRO_PROMPT_KEY: &str = "retro/prompt.md";
pub const RETRO_RESPONSE_KEY: &str = "retro/response.md";
pub const EVENTS_PREFIX: &str = "events/";
pub const CHECKPOINTS_PREFIX: &str = "checkpoints/";
pub const ARTIFACT_VALUES_PREFIX: &str = "artifacts/values/";
pub const ARTIFACT_NODES_PREFIX: &str = "artifacts/nodes/";
pub fn init() -> &'static str {
INIT_KEY
}
pub fn run() -> &'static str {
RUN_KEY
}
pub fn start() -> &'static str {
START_KEY
}
pub fn status() -> &'static str {
STATUS_KEY
}
pub fn checkpoint() -> &'static str {
CHECKPOINT_KEY
}
pub fn conclusion() -> &'static str {
CONCLUSION_KEY
}
pub fn retro() -> &'static str {
RETRO_KEY
}
pub fn graph() -> &'static str {
GRAPH_KEY
}
pub fn sandbox() -> &'static str {
SANDBOX_KEY
}
pub fn node_visit_prefix(node: &NodeVisitRef<'_>) -> String {
format!("nodes/{}/visit-{}", node.node_id, node.visit)
}
pub fn node_prompt(node: &NodeVisitRef<'_>) -> String {
format!("{}/prompt.md", node_visit_prefix(node))
}
pub fn node_response(node: &NodeVisitRef<'_>) -> String {
format!("{}/response.md", node_visit_prefix(node))
}
pub fn node_status(node: &NodeVisitRef<'_>) -> String {
format!("{}/status.json", node_visit_prefix(node))
}
pub fn node_stdout(node: &NodeVisitRef<'_>) -> String {
format!("{}/stdout.log", node_visit_prefix(node))
}
pub fn node_stderr(node: &NodeVisitRef<'_>) -> String {
format!("{}/stderr.log", node_visit_prefix(node))
}
pub fn retro_prompt() -> &'static str {
RETRO_PROMPT_KEY
}
pub fn retro_response() -> &'static str {
RETRO_RESPONSE_KEY
}
pub fn event_key(seq: u32, epoch_ms: i64) -> String {
format!("{EVENTS_PREFIX}{seq:06}-{epoch_ms}.json")
}
pub fn checkpoint_history_key(seq: u32, epoch_ms: i64) -> String {
format!("{CHECKPOINTS_PREFIX}{seq:04}-{epoch_ms}.json")
}
pub fn artifact_value(artifact_id: &str) -> String {
format!("{ARTIFACT_VALUES_PREFIX}{artifact_id}.json")
}
pub fn node_asset_prefix(node: &NodeVisitRef<'_>) -> String {
format!(
"{ARTIFACT_NODES_PREFIX}{}/visit-{}",
node.node_id, node.visit
)
}
pub fn node_asset(node: &NodeVisitRef<'_>, filename: &str) -> String {
format!("{}/{filename}", node_asset_prefix(node))
}
pub fn parse_event_seq(key: &str) -> Option<u32> {
parse_seq(key, EVENTS_PREFIX)
}
pub fn parse_checkpoint_seq(key: &str) -> Option<u32> {
parse_seq(key, CHECKPOINTS_PREFIX)
}
pub fn parse_node_key(key: &str) -> Option<(String, u32, String)> {
parse_visit_scoped_key(key, "nodes/")
}
#[cfg(test)]
pub fn parse_node_asset_key(key: &str) -> Option<(String, u32, String)> {
parse_visit_scoped_key(key, ARTIFACT_NODES_PREFIX)
}
fn parse_seq(key: &str, prefix: &str) -> Option<u32> {
key.strip_prefix(prefix)?.split_once('-')?.0.parse().ok()
}
fn parse_visit_scoped_key(key: &str, prefix: &str) -> Option<(String, u32, String)> {
let rest = key.strip_prefix(prefix)?;
let (node_id, rest) = rest.split_once("/visit-")?;
let (visit, file) = rest.split_once('/')?;
Some((node_id.to_string(), visit.parse().ok()?, file.to_string()))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn top_level_keys_match_spec() {
assert_eq!(init(), "_init.json");
assert_eq!(run(), "run.json");
assert_eq!(graph(), "graph.fabro");
assert_eq!(retro_prompt(), "retro/prompt.md");
assert_eq!(retro_response(), "retro/response.md");
}
#[test]
fn node_keys_match_spec() {
let node = NodeVisitRef {
node_id: "plan",
visit: 3,
};
assert_eq!(node_visit_prefix(&node), "nodes/plan/visit-3");
assert_eq!(node_prompt(&node), "nodes/plan/visit-3/prompt.md");
assert_eq!(node_response(&node), "nodes/plan/visit-3/response.md");
assert_eq!(node_status(&node), "nodes/plan/visit-3/status.json");
assert_eq!(node_stdout(&node), "nodes/plan/visit-3/stdout.log");
assert_eq!(node_stderr(&node), "nodes/plan/visit-3/stderr.log");
}
#[test]
fn sequence_keys_are_zero_padded() {
assert_eq!(event_key(7, 123), "events/000007-123.json");
assert_eq!(checkpoint_history_key(42, 456), "checkpoints/0042-456.json");
}
#[test]
fn artifact_keys_match_spec() {
let node = NodeVisitRef {
node_id: "code",
visit: 2,
};
assert_eq!(artifact_value("summary"), "artifacts/values/summary.json");
assert_eq!(
node_asset(&node, "src/main.rs"),
"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));
assert_eq!(parse_checkpoint_seq("checkpoints/0042-456.json"), Some(42));
assert_eq!(
parse_node_key("nodes/plan/visit-3/status.json"),
Some(("plan".to_string(), 3, "status.json".to_string()))
);
assert_eq!(
parse_node_asset_key("artifacts/nodes/code/visit-2/src/main.rs"),
Some(("code".to_string(), 2, "src/main.rs".to_string()))
);
}
#[test]
fn parse_helpers_reject_invalid_keys() {
assert_eq!(parse_event_seq("events/not-a-seq.json"), None);
assert_eq!(parse_checkpoint_seq("checkpoints/oops.json"), None);
assert_eq!(parse_node_key("nodes/plan/status.json"), None);
assert_eq!(
parse_node_asset_key("artifacts/nodes/code/status.json"),
None
);
}
}

View file

@ -0,0 +1,106 @@
use std::pin::Pin;
use async_trait::async_trait;
use bytes::Bytes;
use chrono::{DateTime, Utc};
use futures::Stream;
mod error;
mod keys;
mod memory;
mod slate;
mod types;
pub use error::{Result, StoreError};
pub use memory::InMemoryStore;
pub use slate::SlateStore;
pub use types::{
CatalogRecord, EventEnvelope, EventPayload, NodeSnapshot, NodeVisitRef, RunSnapshot, RunSummary,
};
use fabro_types::{
Checkpoint, Conclusion, NodeStatusRecord, Retro, RunRecord, RunStatusRecord, SandboxRecord,
StartRecord,
};
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct ListRunsQuery {
pub start: Option<DateTime<Utc>>,
pub end: Option<DateTime<Utc>>,
}
#[async_trait]
pub trait Store: Send + Sync {
async fn create_run(
&self,
run_id: &str,
created_at: DateTime<Utc>,
) -> Result<Box<dyn RunStore>>;
async fn open_run(&self, run_id: &str) -> Result<Option<Box<dyn RunStore>>>;
async fn list_runs(&self, query: &ListRunsQuery) -> Result<Vec<RunSummary>>;
async fn delete_run(&self, run_id: &str) -> Result<()>;
}
#[async_trait]
pub trait RunStore: Send + Sync {
async fn put_run(&self, record: &RunRecord) -> Result<()>;
async fn get_run(&self) -> Result<Option<RunRecord>>;
async fn put_start(&self, record: &StartRecord) -> Result<()>;
async fn get_start(&self) -> Result<Option<StartRecord>>;
async fn put_status(&self, record: &RunStatusRecord) -> Result<()>;
async fn get_status(&self) -> Result<Option<RunStatusRecord>>;
async fn put_checkpoint(&self, record: &Checkpoint) -> Result<()>;
async fn get_checkpoint(&self) -> Result<Option<Checkpoint>>;
async fn append_checkpoint(&self, record: &Checkpoint) -> Result<u32>;
async fn list_checkpoints(&self) -> Result<Vec<(u32, Checkpoint)>>;
async fn put_conclusion(&self, record: &Conclusion) -> Result<()>;
async fn get_conclusion(&self) -> Result<Option<Conclusion>>;
async fn put_retro(&self, retro: &Retro) -> Result<()>;
async fn get_retro(&self) -> Result<Option<Retro>>;
async fn put_graph(&self, dot_source: &str) -> Result<()>;
async fn get_graph(&self) -> Result<Option<String>>;
async fn put_sandbox(&self, record: &SandboxRecord) -> Result<()>;
async fn get_sandbox(&self) -> Result<Option<SandboxRecord>>;
async fn put_node_prompt(&self, node: &NodeVisitRef<'_>, prompt: &str) -> Result<()>;
async fn put_node_response(&self, node: &NodeVisitRef<'_>, response: &str) -> Result<()>;
async fn put_node_status(
&self,
node: &NodeVisitRef<'_>,
status: &NodeStatusRecord,
) -> Result<()>;
async fn put_node_stdout(&self, node: &NodeVisitRef<'_>, log: &str) -> Result<()>;
async fn put_node_stderr(&self, node: &NodeVisitRef<'_>, log: &str) -> Result<()>;
async fn get_node(&self, node: &NodeVisitRef<'_>) -> Result<NodeSnapshot>;
async fn list_node_visits(&self, node_id: &str) -> Result<Vec<u32>>;
async fn append_event(&self, payload: &EventPayload) -> Result<u32>;
async fn list_events(&self) -> Result<Vec<EventEnvelope>>;
async fn list_events_from(&self, seq: u32) -> Result<Vec<EventEnvelope>>;
async fn watch_events_from(
&self,
seq: u32,
) -> Result<Pin<Box<dyn Stream<Item = Result<EventEnvelope>> + Send>>>;
async fn put_retro_prompt(&self, text: &str) -> Result<()>;
async fn get_retro_prompt(&self) -> Result<Option<String>>;
async fn put_retro_response(&self, text: &str) -> Result<()>;
async fn get_retro_response(&self) -> Result<Option<String>>;
async fn put_artifact_value(&self, artifact_id: &str, value: &serde_json::Value) -> Result<()>;
async fn get_artifact_value(&self, artifact_id: &str) -> Result<Option<serde_json::Value>>;
async fn put_asset(&self, node: &NodeVisitRef<'_>, filename: &str, data: &[u8]) -> Result<()>;
async fn get_asset(&self, node: &NodeVisitRef<'_>, filename: &str) -> Result<Option<Bytes>>;
async fn list_assets(&self, node: &NodeVisitRef<'_>) -> Result<Vec<String>>;
async fn get_snapshot(&self) -> Result<Option<RunSnapshot>>;
}

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,172 @@
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use chrono::{DateTime, Utc};
use futures::TryStreamExt;
use object_store::path::Path;
use object_store::ObjectStore;
use crate::{CatalogRecord, ListRunsQuery, Result};
pub(crate) async fn write_catalog(
store: Arc<dyn ObjectStore>,
base_prefix: &str,
run_id: &str,
created_at: DateTime<Utc>,
db_prefix: &str,
) -> Result<CatalogRecord> {
let record = CatalogRecord {
run_id: run_id.to_string(),
created_at,
db_prefix: db_prefix.to_string(),
};
let bytes = serde_json::to_vec(&record)?;
store
.put(&by_id_path(base_prefix, run_id), bytes.clone().into())
.await?;
store
.put(
&by_start_path(base_prefix, created_at, run_id),
bytes.into(),
)
.await?;
Ok(record)
}
pub(crate) async fn read_locator(
store: Arc<dyn ObjectStore>,
base_prefix: &str,
run_id: &str,
) -> Result<Option<CatalogRecord>> {
read_catalog_path(store, by_id_path(base_prefix, run_id)).await
}
pub(crate) async fn list_catalogs(
store: Arc<dyn ObjectStore>,
base_prefix: &str,
query: &ListRunsQuery,
) -> Result<Vec<CatalogRecord>> {
let prefix = Path::from(format!("{base_prefix}by-start"));
let metas = store.list(Some(&prefix)).try_collect::<Vec<_>>().await?;
let mut records = Vec::new();
for meta in metas {
let Some(record) = read_catalog_path(store.clone(), meta.location).await? else {
continue;
};
if let Some(start) = query.start {
if record.created_at < start {
continue;
}
}
if let Some(end) = query.end {
if record.created_at > end {
continue;
}
}
records.push(record);
}
Ok(records)
}
pub async fn repair_catalog(store: Arc<dyn ObjectStore>, 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::<Vec<_>>()
.await?;
let mut canonical = HashMap::new();
for meta in by_id_metas {
if let Some(record) = read_catalog_path(store.clone(), meta.location).await? {
canonical.insert(record.run_id.clone(), record);
}
}
for record in canonical.values() {
let path = by_start_path(base_prefix, record.created_at, &record.run_id);
if !object_exists(store.clone(), &path).await? {
store.put(&path, serde_json::to_vec(record)?.into()).await?;
}
}
let by_start_metas = store
.list(Some(&by_start_prefix))
.try_collect::<Vec<_>>()
.await?;
let mut seen = HashSet::new();
for meta in by_start_metas {
let location = meta.location.clone();
let Some(record) = read_catalog_path(store.clone(), location.clone()).await? else {
delete_if_exists(store.clone(), &location).await?;
continue;
};
let expected = canonical.get(&record.run_id).map(|canonical_record| {
by_start_path(base_prefix, canonical_record.created_at, &record.run_id)
});
match expected {
Some(expected) if expected == location => {
seen.insert(record.run_id);
}
_ => {
delete_if_exists(store.clone(), &location).await?;
}
}
}
for record in canonical.values() {
if !seen.contains(&record.run_id) {
store
.put(
&by_start_path(base_prefix, record.created_at, &record.run_id),
serde_json::to_vec(record)?.into(),
)
.await?;
}
}
Ok(())
}
pub(crate) fn db_prefix(base_prefix: &str, created_at: DateTime<Utc>, run_id: &str) -> String {
format!(
"{base_prefix}db/{}/{run_id}/",
created_at.format("%Y-%m-%d-%H-%M")
)
}
pub(crate) fn by_id_path(base_prefix: &str, run_id: &str) -> Path {
Path::from(format!("{base_prefix}by-id/{run_id}.json"))
}
pub(crate) fn by_start_path(base_prefix: &str, created_at: DateTime<Utc>, run_id: &str) -> Path {
Path::from(format!(
"{base_prefix}by-start/{}/{run_id}.json",
created_at.format("%Y-%m-%d-%H-%M")
))
}
pub(crate) async fn read_catalog_path(
store: Arc<dyn ObjectStore>,
path: Path,
) -> Result<Option<CatalogRecord>> {
match store.get(&path).await {
Ok(result) => Ok(Some(serde_json::from_slice(&result.bytes().await?)?)),
Err(object_store::Error::NotFound { .. }) => Ok(None),
Err(err) => Err(err.into()),
}
}
async fn object_exists(store: Arc<dyn ObjectStore>, path: &Path) -> Result<bool> {
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<dyn ObjectStore>, path: &Path) -> Result<()> {
match store.delete(path).await {
Ok(()) | Err(object_store::Error::NotFound { .. }) => Ok(()),
Err(err) => Err(err.into()),
}
}

View file

@ -0,0 +1,635 @@
mod catalog;
mod run_store;
use std::sync::Arc;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use futures::TryStreamExt;
use object_store::path::Path;
use object_store::ObjectStore;
use crate::keys;
use crate::{CatalogRecord, ListRunsQuery, Result, RunStore, RunSummary, Store, StoreError};
use run_store::SlateRunStore;
#[derive(Clone)]
pub struct SlateStore {
object_store: Arc<dyn ObjectStore>,
base_prefix: String,
}
impl std::fmt::Debug for SlateStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SlateStore")
.field("base_prefix", &self.base_prefix)
.finish()
}
}
impl SlateStore {
pub fn new(object_store: Arc<dyn ObjectStore>, base_prefix: impl Into<String>) -> Self {
Self {
object_store,
base_prefix: normalize_base_prefix(base_prefix.into()),
}
}
pub async fn repair_catalog(&self) -> Result<()> {
catalog::repair_catalog(self.object_store.clone(), &self.base_prefix).await
}
async fn open_db(&self, db_prefix: &str) -> Result<slatedb::Db> {
Ok(slatedb::Db::open(db_prefix.to_string(), self.object_store.clone()).await?)
}
async fn open_run_store(&self, record: &CatalogRecord) -> Result<Option<SlateRunStore>> {
let db = self.open_db(&record.db_prefix).await?;
if !SlateRunStore::has_init(&db).await? {
let _ = db.close().await;
return Ok(None);
}
Ok(Some(
SlateRunStore::open(record.run_id.clone(), record.created_at, db).await?,
))
}
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::<Vec<_>>()
.await?;
for meta in metas {
delete_path(self.object_store.clone(), &meta.location).await?;
}
Ok(())
}
}
#[async_trait]
impl Store for SlateStore {
async fn create_run(
&self,
run_id: &str,
created_at: DateTime<Utc>,
) -> Result<Box<dyn RunStore>> {
let locator =
catalog::read_locator(self.object_store.clone(), &self.base_prefix, run_id).await?;
let db_prefix = match locator {
Some(existing) if existing.created_at != created_at => {
return Err(StoreError::RunAlreadyExists(run_id.to_string()));
}
Some(existing) => existing.db_prefix,
None => catalog::db_prefix(&self.base_prefix, created_at, run_id),
};
let record = CatalogRecord {
run_id: run_id.to_string(),
created_at,
db_prefix: db_prefix.clone(),
};
let db = self.open_db(&db_prefix).await?;
db.put(keys::init(), serde_json::to_vec(&record)?).await?;
catalog::write_catalog(
self.object_store.clone(),
&self.base_prefix,
run_id,
created_at,
&db_prefix,
)
.await?;
Ok(Box::new(
SlateRunStore::open(run_id.to_string(), created_at, db).await?,
))
}
async fn open_run(&self, run_id: &str) -> Result<Option<Box<dyn RunStore>>> {
let Some(locator) =
catalog::read_locator(self.object_store.clone(), &self.base_prefix, run_id).await?
else {
return Ok(None);
};
let Some(run_store) = self.open_run_store(&locator).await? else {
return Ok(None);
};
Ok(Some(Box::new(run_store)))
}
async fn list_runs(&self, query: &ListRunsQuery) -> Result<Vec<RunSummary>> {
let catalogs =
catalog::list_catalogs(self.object_store.clone(), &self.base_prefix, query).await?;
let mut summaries = Vec::new();
for record in catalogs {
let db = self.open_db(&record.db_prefix).await?;
if !SlateRunStore::has_init(&db).await? {
let _ = db.close().await;
continue;
}
let summary = SlateRunStore::build_summary(&db, &record).await?;
let _ = db.close().await;
summaries.push(summary);
}
summaries.sort_by(|a, b| b.created_at.cmp(&a.created_at));
Ok(summaries)
}
async fn delete_run(&self, run_id: &str) -> Result<()> {
if let Some(locator) =
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, locator.created_at, run_id),
)
.await?;
self.delete_db_prefix(&locator.db_prefix).await?;
delete_path(
self.object_store.clone(),
&catalog::by_id_path(&self.base_prefix, run_id),
)
.await?;
return Ok(());
}
let by_start_prefix = Path::from(format!("{}by-start", self.base_prefix));
let metas = self
.object_store
.list(Some(&by_start_prefix))
.try_collect::<Vec<_>>()
.await?;
let expected_name = format!("{run_id}.json");
for meta in metas {
if meta.location.filename() != Some(expected_name.as_str()) {
continue;
}
let Some(record) =
catalog::read_catalog_path(self.object_store.clone(), meta.location.clone())
.await?
else {
delete_path(self.object_store.clone(), &meta.location).await?;
continue;
};
self.delete_db_prefix(&record.db_prefix).await?;
delete_path(self.object_store.clone(), &meta.location).await?;
}
Ok(())
}
}
async fn delete_path(store: Arc<dyn ObjectStore>, 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();
}
if prefix.ends_with('/') {
prefix
} else {
format!("{prefix}/")
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
use std::time::Duration;
use bytes::Bytes;
use fabro_types::{
AttrValue, Checkpoint, Conclusion, FabroSettings, Graph, NodeStatusRecord, RunRecord,
RunStatus, RunStatusRecord, StageStatus, StartRecord, StatusReason,
};
use object_store::memory::InMemory;
use crate::{EventPayload, NodeVisitRef};
fn dt(rfc3339: &str) -> DateTime<Utc> {
DateTime::parse_from_rfc3339(rfc3339)
.unwrap()
.with_timezone(&Utc)
}
fn make_store() -> (Arc<dyn ObjectStore>, SlateStore) {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let store = SlateStore::new(object_store.clone(), "runs/");
(object_store, store)
}
fn sample_run_record(run_id: &str, created_at: DateTime<Utc>) -> RunRecord {
let mut graph = Graph::new("night-sky");
graph.attrs.insert(
"goal".to_string(),
AttrValue::String("map the constellations".to_string()),
);
RunRecord {
run_id: run_id.to_string(),
created_at,
settings: FabroSettings::default(),
graph,
workflow_slug: Some("night-sky".to_string()),
working_directory: PathBuf::from("/tmp/night-sky"),
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())]),
}
}
fn sample_start_record(run_id: &str, created_at: DateTime<Utc>) -> StartRecord {
StartRecord {
run_id: run_id.to_string(),
start_time: created_at + chrono::Duration::seconds(5),
run_branch: Some("fabro/run/demo".to_string()),
base_sha: Some("abc123".to_string()),
}
}
fn sample_status(status: RunStatus, reason: Option<StatusReason>) -> 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: std::collections::HashMap::from([("code".to_string(), 1)]),
context_values: std::collections::HashMap::new(),
node_outcomes: std::collections::HashMap::new(),
next_node_id: Some("review".to_string()),
git_commit_sha: Some("def456".to_string()),
loop_failure_signatures: std::collections::HashMap::new(),
restart_failure_signatures: std::collections::HashMap::new(),
node_visits: std::collections::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_node_status() -> NodeStatusRecord {
NodeStatusRecord {
status: StageStatus::Success,
notes: Some("done".to_string()),
failure_reason: None,
timestamp: dt("2026-03-27T12:12:00Z"),
}
}
fn event_payload(run_id: &str, ts: &str, event: &str) -> EventPayload {
EventPayload::new(
serde_json::json!({
"ts": ts,
"run_id": run_id,
"event": event
}),
run_id,
)
.unwrap()
}
async fn list_paths(store: Arc<dyn ObjectStore>, prefix: &str) -> Vec<String> {
let mut items = store
.list(Some(&Path::from(prefix.to_string())))
.map_ok(|meta| meta.location.to_string())
.try_collect::<Vec<_>>()
.await
.unwrap();
items.sort();
items
}
async fn object_exists(store: Arc<dyn ObjectStore>, path: &Path) -> bool {
store.head(path).await.is_ok()
}
async fn seed_db(
object_store: Arc<dyn ObjectStore>,
record: &CatalogRecord,
include_init: bool,
) -> slatedb::Db {
let db = slatedb::Db::open(record.db_prefix.clone(), object_store)
.await
.unwrap();
if include_init {
db.put(keys::init(), serde_json::to_vec(record).unwrap())
.await
.unwrap();
}
db
}
#[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");
let run = store.create_run("run-1", created_at).await.unwrap();
run.put_run(&sample_run_record("run-1", created_at))
.await
.unwrap();
run.put_start(&sample_start_record("run-1", created_at))
.await
.unwrap();
run.put_status(&sample_status(
RunStatus::Succeeded,
Some(StatusReason::Completed),
))
.await
.unwrap();
run.put_conclusion(&sample_conclusion()).await.unwrap();
let by_id = catalog::by_id_path("runs/", "run-1");
let by_start = catalog::by_start_path("runs/", created_at, "run-1");
assert!(object_exists(object_store.clone(), &by_id).await);
assert!(object_exists(object_store.clone(), &by_start).await);
let summary = store.list_runs(&ListRunsQuery::default()).await.unwrap();
assert_eq!(summary.len(), 1);
assert_eq!(summary[0].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("run-1").await.unwrap().unwrap();
let stored = reopened.get_run().await.unwrap().unwrap();
assert_eq!(stored.run_id, "run-1");
store.delete_run("run-1").await.unwrap();
assert!(store.open_run("run-1").await.unwrap().is_none());
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: "run-1".to_string(),
created_at,
db_prefix: catalog::db_prefix("runs/", created_at, "run-1"),
};
let db = seed_db(object_store.clone(), &record, true).await;
db.put(
keys::run(),
serde_json::to_vec(&sample_run_record("run-1", created_at)).unwrap(),
)
.await
.unwrap();
db.close().await.unwrap();
object_store
.put(
&catalog::by_id_path("runs/", "run-1"),
serde_json::to_vec(&record).unwrap().into(),
)
.await
.unwrap();
assert!(store.open_run("run-1").await.unwrap().is_some());
assert!(store
.list_runs(&ListRunsQuery::default())
.await
.unwrap()
.is_empty());
store.repair_catalog().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/", created_at, "run-1")
)
.await
);
}
#[tokio::test]
async fn reopen_recovers_event_and_checkpoint_sequences() {
let (_object_store, store) = make_store();
let created_at = dt("2026-03-27T12:00:00Z");
let run = store.create_run("run-1", created_at).await.unwrap();
run.put_run(&sample_run_record("run-1", created_at))
.await
.unwrap();
run.append_event(&event_payload("run-1", "2026-03-27T12:00:00Z", "Started"))
.await
.unwrap();
run.append_event(&event_payload("run-1", "2026-03-27T12:00:01Z", "Next"))
.await
.unwrap();
run.append_checkpoint(&sample_checkpoint()).await.unwrap();
drop(run);
let reopened = store.open_run("run-1").await.unwrap().unwrap();
let next_event = reopened
.append_event(&event_payload(
"run-1",
"2026-03-27T12:00:02Z",
"AfterReopen",
))
.await
.unwrap();
let next_checkpoint = reopened
.append_checkpoint(&sample_checkpoint())
.await
.unwrap();
assert_eq!(next_event, 3);
assert_eq!(next_checkpoint, 2);
}
#[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: "run-1".to_string(),
created_at,
db_prefix: catalog::db_prefix("runs/", created_at, "run-1"),
};
let db = seed_db(object_store.clone(), &record, false).await;
db.close().await.unwrap();
object_store
.put(
&catalog::by_id_path("runs/", "run-1"),
serde_json::to_vec(&record).unwrap().into(),
)
.await
.unwrap();
object_store
.put(
&catalog::by_start_path("runs/", created_at, "run-1"),
serde_json::to_vec(&record).unwrap().into(),
)
.await
.unwrap();
assert!(store.open_run("run-1").await.unwrap().is_none());
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");
store.create_run("run-1", created_at).await.unwrap();
store.create_run("run-1", created_at).await.unwrap();
let conflict = store
.create_run("run-1", created_at + chrono::Duration::seconds(1))
.await;
assert!(matches!(conflict, Err(StoreError::RunAlreadyExists(_))));
}
#[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("run-1", created_at).await.unwrap();
let mut stream = run.watch_events_from(1).await.unwrap();
run.append_event(&event_payload("run-1", "2026-03-27T12:00:00Z", "Started"))
.await
.unwrap();
let event = tokio::time::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: "run-1".to_string(),
created_at,
db_prefix: catalog::db_prefix("runs/", created_at, "run-1"),
};
let db = seed_db(object_store.clone(), &record, true).await;
db.put(
keys::run(),
serde_json::to_vec(&sample_run_record("run-1", created_at)).unwrap(),
)
.await
.unwrap();
db.close().await.unwrap();
object_store
.put(
&catalog::by_start_path("runs/", created_at, "run-1"),
serde_json::to_vec(&record).unwrap().into(),
)
.await
.unwrap();
store.delete_run("run-1").await.unwrap();
store.delete_run("run-1").await.unwrap();
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("run-1", created_at).await.unwrap();
run.put_run(&sample_run_record("run-1", created_at))
.await
.unwrap();
let locator = catalog::read_locator(object_store.clone(), "runs/", "run-1")
.await
.unwrap()
.unwrap();
object_store
.put(
&catalog::by_start_path("runs/", wrong_time, "run-1"),
serde_json::to_vec(&locator).unwrap().into(),
)
.await
.unwrap();
store.repair_catalog().await.unwrap();
let paths = list_paths(object_store, "runs/by-start").await;
assert_eq!(paths.len(), 1);
assert!(paths[0].contains("2026-03-27-12-00/run-1.json"));
}
#[tokio::test]
async fn slate_run_store_round_trips_node_data_and_assets() {
let (_object_store, store) = make_store();
let created_at = dt("2026-03-27T12:00:00Z");
let run = store.create_run("run-1", created_at).await.unwrap();
run.put_run(&sample_run_record("run-1", created_at))
.await
.unwrap();
let node = NodeVisitRef {
node_id: "code",
visit: 2,
};
run.put_node_prompt(&node, "Plan").await.unwrap();
run.put_node_status(&node, &sample_node_status())
.await
.unwrap();
run.put_asset(&node, "src/lib.rs", b"fn main() {}")
.await
.unwrap();
let snapshot = run.get_node(&node).await.unwrap();
assert_eq!(snapshot.prompt, Some("Plan".to_string()));
assert_eq!(
run.get_asset(&node, "src/lib.rs").await.unwrap(),
Some(Bytes::from_static(b"fn main() {}"))
);
assert_eq!(
run.list_assets(&node).await.unwrap(),
vec!["src/lib.rs".to_string()]
);
}
}

View file

@ -0,0 +1,481 @@
use std::collections::BTreeSet;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use async_trait::async_trait;
use bytes::Bytes;
use chrono::{DateTime, Utc};
use futures::Stream;
use serde::de::DeserializeOwned;
use serde::Serialize;
use tokio::sync::mpsc;
use tokio_stream::wrappers::UnboundedReceiverStream;
use crate::keys;
use crate::{
CatalogRecord, EventEnvelope, EventPayload, NodeSnapshot, NodeVisitRef, Result, RunSnapshot,
RunStore, RunSummary, StoreError,
};
use fabro_types::{
Checkpoint, Conclusion, NodeStatusRecord, Retro, RunRecord, RunStatusRecord, SandboxRecord,
StartRecord,
};
pub(crate) struct SlateRunStore {
run_id: String,
created_at: DateTime<Utc>,
db: slatedb::Db,
event_seq: AtomicU32,
checkpoint_seq: AtomicU32,
}
impl SlateRunStore {
pub(crate) async fn open(
run_id: String,
created_at: DateTime<Utc>,
db: slatedb::Db,
) -> Result<Self> {
let event_seq = recover_next_seq(&db, keys::EVENTS_PREFIX, keys::parse_event_seq).await?;
let checkpoint_seq =
recover_next_seq(&db, keys::CHECKPOINTS_PREFIX, keys::parse_checkpoint_seq).await?;
Ok(Self {
run_id,
created_at,
db,
event_seq: AtomicU32::new(event_seq),
checkpoint_seq: AtomicU32::new(checkpoint_seq),
})
}
pub(crate) async fn has_init(db: &slatedb::Db) -> Result<bool> {
Ok(db.get(keys::init()).await?.is_some())
}
pub(crate) async fn build_summary(
db: &slatedb::Db,
catalog: &CatalogRecord,
) -> Result<RunSummary> {
let run = get_json::<RunRecord>(db, keys::run()).await?;
let start = get_json::<StartRecord>(db, keys::start()).await?;
let status = get_json::<RunStatusRecord>(db, keys::status()).await?;
let conclusion = get_json::<Conclusion>(db, keys::conclusion()).await?;
let workflow_name = run.as_ref().map(|run| {
if run.graph.name.is_empty() {
"unnamed".to_string()
} else {
run.graph.name.clone()
}
});
let goal = run.as_ref().and_then(|run| {
let goal = run.graph.goal();
(!goal.is_empty()).then(|| goal.to_string())
});
Ok(RunSummary {
run_id: catalog.run_id.clone(),
created_at: catalog.created_at,
db_prefix: catalog.db_prefix.clone(),
workflow_name,
workflow_slug: run.as_ref().and_then(|run| run.workflow_slug.clone()),
goal,
labels: run
.as_ref()
.map(|run| run.labels.clone())
.unwrap_or_default(),
host_repo_path: run.as_ref().and_then(|run| run.host_repo_path.clone()),
start_time: start.map(|start| start.start_time),
status: status.as_ref().map(|status| status.status),
status_reason: status.and_then(|status| status.reason),
duration_ms: conclusion.as_ref().map(|conclusion| conclusion.duration_ms),
total_cost: conclusion.and_then(|conclusion| conclusion.total_cost),
})
}
fn validate_run_record(&self, record: &RunRecord) -> Result<()> {
if record.created_at != self.created_at {
return Err(StoreError::Other(format!(
"run record created_at {:?} does not match store created_at {:?}",
record.created_at, self.created_at
)));
}
if record.run_id != self.run_id {
return Err(StoreError::Other(format!(
"run record run_id {:?} does not match store run_id {:?}",
record.run_id, self.run_id
)));
}
Ok(())
}
async fn build_node_snapshot(&self, node: &NodeVisitRef<'_>) -> Result<NodeSnapshot> {
Ok(NodeSnapshot {
node_id: node.node_id.to_string(),
visit: node.visit,
prompt: get_text(&self.db, &keys::node_prompt(node)).await?,
response: get_text(&self.db, &keys::node_response(node)).await?,
status: get_json(&self.db, &keys::node_status(node)).await?,
stdout: get_text(&self.db, &keys::node_stdout(node)).await?,
stderr: get_text(&self.db, &keys::node_stderr(node)).await?,
})
}
}
#[async_trait]
impl RunStore for SlateRunStore {
async fn put_run(&self, record: &RunRecord) -> Result<()> {
self.validate_run_record(record)?;
put_json(&self.db, keys::run(), record).await
}
async fn get_run(&self) -> Result<Option<RunRecord>> {
get_json(&self.db, keys::run()).await
}
async fn put_start(&self, record: &StartRecord) -> Result<()> {
put_json(&self.db, keys::start(), record).await
}
async fn get_start(&self) -> Result<Option<StartRecord>> {
get_json(&self.db, keys::start()).await
}
async fn put_status(&self, record: &RunStatusRecord) -> Result<()> {
put_json(&self.db, keys::status(), record).await
}
async fn get_status(&self) -> Result<Option<RunStatusRecord>> {
get_json(&self.db, keys::status()).await
}
async fn put_checkpoint(&self, record: &Checkpoint) -> Result<()> {
put_json(&self.db, keys::checkpoint(), record).await
}
async fn get_checkpoint(&self) -> Result<Option<Checkpoint>> {
get_json(&self.db, keys::checkpoint()).await
}
async fn append_checkpoint(&self, record: &Checkpoint) -> Result<u32> {
let seq = self.checkpoint_seq.fetch_add(1, Ordering::SeqCst);
self.put_checkpoint(record).await?;
put_json(
&self.db,
&keys::checkpoint_history_key(seq, Utc::now().timestamp_millis()),
record,
)
.await?;
Ok(seq)
}
async fn list_checkpoints(&self) -> Result<Vec<(u32, Checkpoint)>> {
list_checkpoints(&self.db).await
}
async fn put_conclusion(&self, record: &Conclusion) -> Result<()> {
put_json(&self.db, keys::conclusion(), record).await
}
async fn get_conclusion(&self) -> Result<Option<Conclusion>> {
get_json(&self.db, keys::conclusion()).await
}
async fn put_retro(&self, retro: &Retro) -> Result<()> {
put_json(&self.db, keys::retro(), retro).await
}
async fn get_retro(&self) -> Result<Option<Retro>> {
get_json(&self.db, keys::retro()).await
}
async fn put_graph(&self, dot_source: &str) -> Result<()> {
put_text(&self.db, keys::graph(), dot_source).await
}
async fn get_graph(&self) -> Result<Option<String>> {
get_text(&self.db, keys::graph()).await
}
async fn put_sandbox(&self, record: &SandboxRecord) -> Result<()> {
put_json(&self.db, keys::sandbox(), record).await
}
async fn get_sandbox(&self) -> Result<Option<SandboxRecord>> {
get_json(&self.db, keys::sandbox()).await
}
async fn put_node_prompt(&self, node: &NodeVisitRef<'_>, prompt: &str) -> Result<()> {
put_text(&self.db, &keys::node_prompt(node), prompt).await
}
async fn put_node_response(&self, node: &NodeVisitRef<'_>, response: &str) -> Result<()> {
put_text(&self.db, &keys::node_response(node), response).await
}
async fn put_node_status(
&self,
node: &NodeVisitRef<'_>,
status: &NodeStatusRecord,
) -> Result<()> {
put_json(&self.db, &keys::node_status(node), status).await
}
async fn put_node_stdout(&self, node: &NodeVisitRef<'_>, log: &str) -> Result<()> {
put_text(&self.db, &keys::node_stdout(node), log).await
}
async fn put_node_stderr(&self, node: &NodeVisitRef<'_>, log: &str) -> Result<()> {
put_text(&self.db, &keys::node_stderr(node), log).await
}
async fn get_node(&self, node: &NodeVisitRef<'_>) -> Result<NodeSnapshot> {
self.build_node_snapshot(node).await
}
async fn list_node_visits(&self, node_id: &str) -> Result<Vec<u32>> {
let prefix = format!("nodes/{node_id}/visit-");
let mut iter = self.db.scan_prefix(prefix.as_bytes()).await?;
let mut visits = BTreeSet::new();
while let Some(entry) = iter.next().await? {
let key = key_to_string(entry.key)?;
if let Some((current_node_id, visit, _)) = keys::parse_node_key(&key) {
if current_node_id == node_id {
visits.insert(visit);
}
}
}
Ok(visits.into_iter().collect())
}
async fn append_event(&self, payload: &EventPayload) -> Result<u32> {
payload.validate(&self.run_id)?;
let seq = self.event_seq.fetch_add(1, Ordering::SeqCst);
put_json(
&self.db,
&keys::event_key(seq, Utc::now().timestamp_millis()),
payload,
)
.await?;
Ok(seq)
}
async fn list_events(&self) -> Result<Vec<EventEnvelope>> {
list_events_from(&self.db, 1).await
}
async fn list_events_from(&self, seq: u32) -> Result<Vec<EventEnvelope>> {
list_events_from(&self.db, seq).await
}
async fn watch_events_from(
&self,
seq: u32,
) -> Result<std::pin::Pin<Box<dyn Stream<Item = Result<EventEnvelope>> + Send>>> {
let db = self.db.clone();
let (sender, receiver) = mpsc::unbounded_channel();
tokio::spawn(async move {
let mut next_seq = seq;
loop {
if sender.is_closed() {
return;
}
match list_events_from(&db, next_seq).await {
Ok(events) => {
if events.is_empty() {
tokio::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;
}
}
}
});
Ok(Box::pin(UnboundedReceiverStream::new(receiver)))
}
async fn put_retro_prompt(&self, text: &str) -> Result<()> {
put_text(&self.db, keys::retro_prompt(), text).await
}
async fn get_retro_prompt(&self) -> Result<Option<String>> {
get_text(&self.db, keys::retro_prompt()).await
}
async fn put_retro_response(&self, text: &str) -> Result<()> {
put_text(&self.db, keys::retro_response(), text).await
}
async fn get_retro_response(&self) -> Result<Option<String>> {
get_text(&self.db, keys::retro_response()).await
}
async fn put_artifact_value(&self, artifact_id: &str, value: &serde_json::Value) -> Result<()> {
put_json(&self.db, &keys::artifact_value(artifact_id), value).await
}
async fn get_artifact_value(&self, artifact_id: &str) -> Result<Option<serde_json::Value>> {
get_json(&self.db, &keys::artifact_value(artifact_id)).await
}
async fn put_asset(&self, node: &NodeVisitRef<'_>, filename: &str, data: &[u8]) -> Result<()> {
put_bytes(&self.db, &keys::node_asset(node, filename), data).await
}
async fn get_asset(&self, node: &NodeVisitRef<'_>, filename: &str) -> Result<Option<Bytes>> {
get_bytes(&self.db, &keys::node_asset(node, filename)).await
}
async fn list_assets(&self, node: &NodeVisitRef<'_>) -> Result<Vec<String>> {
let prefix = format!("{}/", keys::node_asset_prefix(node));
let mut iter = self.db.scan_prefix(prefix.as_bytes()).await?;
let mut assets = Vec::new();
while let Some(entry) = iter.next().await? {
let key = key_to_string(entry.key)?;
if let Some(asset) = key.strip_prefix(&prefix) {
assets.push(asset.to_string());
}
}
assets.sort();
Ok(assets)
}
async fn get_snapshot(&self) -> Result<Option<RunSnapshot>> {
let Some(run) = self.get_run().await? else {
return Ok(None);
};
let mut iter = self.db.scan_prefix(b"nodes/").await?;
let mut visits = BTreeSet::new();
while let Some(entry) = iter.next().await? {
let key = key_to_string(entry.key)?;
if let Some((node_id, visit, _)) = keys::parse_node_key(&key) {
visits.insert((node_id, visit));
}
}
let mut nodes = Vec::new();
for (node_id, visit) in visits {
let node = NodeVisitRef {
node_id: &node_id,
visit,
};
nodes.push(self.build_node_snapshot(&node).await?);
}
Ok(Some(RunSnapshot {
run,
start: self.get_start().await?,
status: self.get_status().await?,
checkpoint: self.get_checkpoint().await?,
conclusion: self.get_conclusion().await?,
retro: self.get_retro().await?,
graph: self.get_graph().await?,
sandbox: self.get_sandbox().await?,
nodes,
}))
}
}
async fn put_json<T: Serialize>(db: &slatedb::Db, key: &str, value: &T) -> Result<()> {
db.put(key, serde_json::to_vec(value)?).await?;
Ok(())
}
async fn get_json<T: DeserializeOwned>(db: &slatedb::Db, key: &str) -> Result<Option<T>> {
db.get(key)
.await?
.map(|value| serde_json::from_slice(&value))
.transpose()
.map_err(Into::into)
}
async fn put_text(db: &slatedb::Db, key: &str, value: &str) -> Result<()> {
db.put(key, value.as_bytes()).await?;
Ok(())
}
async fn get_text(db: &slatedb::Db, key: &str) -> Result<Option<String>> {
db.get(key)
.await?
.map(|value| {
String::from_utf8(value.to_vec())
.map_err(|err| StoreError::Other(format!("stored text is not valid UTF-8: {err}")))
})
.transpose()
}
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<Option<Bytes>> {
Ok(db.get(key).await?)
}
async fn recover_next_seq(
db: &slatedb::Db,
prefix: &str,
parse: fn(&str) -> Option<u32>,
) -> Result<u32> {
let mut iter = db.scan_prefix(prefix.as_bytes()).await?;
let mut max_seq = 0;
while let Some(entry) = iter.next().await? {
let key = key_to_string(entry.key)?;
if let Some(seq) = parse(&key) {
max_seq = max_seq.max(seq);
}
}
Ok(max_seq.saturating_add(1).max(1))
}
async fn list_events_from(db: &slatedb::Db, start_seq: u32) -> Result<Vec<EventEnvelope>> {
let mut iter = db.scan_prefix(keys::EVENTS_PREFIX.as_bytes()).await?;
let mut events = Vec::new();
while let Some(entry) = iter.next().await? {
let key = key_to_string(entry.key)?;
let Some(seq) = keys::parse_event_seq(&key) else {
continue;
};
if seq < start_seq {
continue;
}
events.push(EventEnvelope {
seq,
payload: serde_json::from_slice(&entry.value)?,
});
}
events.sort_by_key(|event| event.seq);
Ok(events)
}
async fn list_checkpoints(db: &slatedb::Db) -> Result<Vec<(u32, Checkpoint)>> {
let mut iter = db.scan_prefix(keys::CHECKPOINTS_PREFIX.as_bytes()).await?;
let mut checkpoints = Vec::new();
while let Some(entry) = iter.next().await? {
let key = key_to_string(entry.key)?;
let Some(seq) = keys::parse_checkpoint_seq(&key) else {
continue;
};
checkpoints.push((seq, serde_json::from_slice(&entry.value)?));
}
checkpoints.sort_by_key(|(seq, _)| *seq);
Ok(checkpoints)
}
fn key_to_string(key: Bytes) -> Result<String> {
String::from_utf8(key.to_vec())
.map_err(|err| StoreError::Other(format!("stored key is not valid UTF-8: {err}")))
}

View file

@ -0,0 +1,117 @@
use std::collections::HashMap;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::{Result, StoreError};
use fabro_types::{
Checkpoint, Conclusion, NodeStatusRecord, Retro, RunRecord, RunStatus, RunStatusRecord,
SandboxRecord, StartRecord, StatusReason,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct NodeVisitRef<'a> {
pub node_id: &'a str,
pub visit: u32,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CatalogRecord {
pub run_id: String,
pub created_at: DateTime<Utc>,
pub db_prefix: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunSummary {
pub run_id: String,
pub created_at: DateTime<Utc>,
pub db_prefix: String,
pub workflow_name: Option<String>,
pub workflow_slug: Option<String>,
pub goal: Option<String>,
pub labels: HashMap<String, String>,
pub host_repo_path: Option<String>,
pub start_time: Option<DateTime<Utc>>,
pub status: Option<RunStatus>,
pub status_reason: Option<StatusReason>,
pub duration_ms: Option<u64>,
pub total_cost: Option<f64>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RunSnapshot {
pub run: RunRecord,
pub start: Option<StartRecord>,
pub status: Option<RunStatusRecord>,
pub checkpoint: Option<Checkpoint>,
pub conclusion: Option<Conclusion>,
pub retro: Option<Retro>,
pub graph: Option<String>,
pub sandbox: Option<SandboxRecord>,
pub nodes: Vec<NodeSnapshot>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NodeSnapshot {
pub node_id: String,
pub visit: u32,
pub prompt: Option<String>,
pub response: Option<String>,
pub status: Option<NodeStatusRecord>,
pub stdout: Option<String>,
pub stderr: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(transparent)]
pub struct EventPayload(serde_json::Value);
impl EventPayload {
pub fn new(value: serde_json::Value, expected_run_id: &str) -> Result<Self> {
let payload = Self(value);
payload.validate(expected_run_id)?;
Ok(payload)
}
pub fn validate(&self, expected_run_id: &str) -> Result<()> {
let obj = self.0.as_object().ok_or_else(|| {
StoreError::InvalidEvent("event payload must be a JSON object".into())
})?;
for field in ["ts", "run_id", "event"] {
match obj.get(field) {
Some(serde_json::Value::String(_)) => {}
_ => {
return Err(StoreError::InvalidEvent(format!(
"missing or non-string required field: {field}"
)));
}
}
}
match obj.get("run_id") {
Some(serde_json::Value::String(run_id)) if run_id == expected_run_id => Ok(()),
Some(serde_json::Value::String(run_id)) => Err(StoreError::InvalidEvent(format!(
"payload run_id {run_id:?} does not match store run_id {expected_run_id:?}"
))),
_ => Err(StoreError::InvalidEvent(
"missing or non-string required field: run_id".into(),
)),
}
}
pub fn into_inner(self) -> serde_json::Value {
self.0
}
pub fn as_value(&self) -> &serde_json::Value {
&self.0
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct EventEnvelope {
pub seq: u32,
pub payload: EventPayload,
}

View file

@ -5,6 +5,7 @@ pub mod combine;
pub mod conclusion;
pub mod failure_signature;
pub mod graph;
pub mod node_status;
pub mod outcome;
pub mod retro;
pub mod run;
@ -18,6 +19,7 @@ pub use checkpoint::Checkpoint;
pub use conclusion::{Conclusion, StageSummary};
pub use failure_signature::FailureSignature;
pub use graph::{is_llm_handler_type, shape_to_handler_type, AttrValue, Edge, Graph, Node};
pub use node_status::NodeStatusRecord;
pub use outcome::{FailureCategory, FailureDetail, NodeResult, Outcome, OutcomeMeta, StageStatus};
pub use retro::{
AggregateStats, FrictionKind, FrictionPoint, Learning, LearningCategory, OpenItem,

View file

@ -0,0 +1,14 @@
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::outcome::StageStatus;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NodeStatusRecord {
pub status: StageStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub notes: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub failure_reason: Option<String>,
pub timestamp: DateTime<Utc>,
}

View file

@ -1,6 +1,7 @@
use std::path::{Path, PathBuf};
use chrono::Utc;
use fabro_types::NodeStatusRecord;
use crate::context::Context;
use crate::outcome::{Outcome, OutcomeExt};
@ -50,12 +51,12 @@ pub fn visit_from_context(context: &Context) -> usize {
pub(crate) fn write_node_status(run_dir: &Path, node_id: &str, visit: usize, outcome: &Outcome) {
let node_dir = node_dir(run_dir, node_id, visit);
let _ = std::fs::create_dir_all(&node_dir);
let status = serde_json::json!({
"status": outcome.status.to_string(),
"notes": outcome.notes,
"failure_reason": outcome.failure_reason(),
"timestamp": Utc::now().to_rfc3339(),
});
let status = NodeStatusRecord {
status: outcome.status.clone(),
notes: outcome.notes.clone(),
failure_reason: outcome.failure_reason().map(ToOwned::to_owned),
timestamp: Utc::now(),
};
if let Ok(json) = serde_json::to_string_pretty(&status) {
let _ = std::fs::write(node_dir.join("status.json"), json);
}
@ -66,6 +67,9 @@ mod tests {
use super::*;
use std::path::Path;
use fabro_types::StageStatus;
use tempfile::TempDir;
use crate::context::Context;
#[test]
@ -107,4 +111,30 @@ mod tests {
root.join("nodes").join("work-visit_5")
);
}
#[test]
fn write_node_status_uses_typed_record_with_legacy_shape() {
let temp = TempDir::new().unwrap();
let outcome = Outcome {
status: StageStatus::Fail,
notes: Some("needs retry".to_string()),
failure: Some(crate::outcome::FailureDetail::new(
"boom",
crate::outcome::FailureCategory::Deterministic,
)),
..Outcome::default()
};
write_node_status(temp.path(), "work", 1, &outcome);
let data = std::fs::read_to_string(temp.path().join("nodes/work/status.json")).unwrap();
let value: serde_json::Value = serde_json::from_str(&data).unwrap();
assert_eq!(value.get("status"), Some(&serde_json::json!("fail")));
assert_eq!(value.get("notes"), Some(&serde_json::json!("needs retry")));
assert_eq!(
value.get("failure_reason"),
Some(&serde_json::json!("boom"))
);
assert!(value.get("timestamp").and_then(|v| v.as_str()).is_some());
}
}