feat(slatedb): add disk_cache setting for S3-backed deployments

When `disk_cache = true` in `[server.slatedb]`, Fabro enables SlateDB's
object-store cache at `<storage_root>/cache/slatedb`, caching raw S3
bytes on local disk to reduce read latency. All cache parameters use
SlateDB defaults (16 GB max, 4 MB parts). A warning is emitted if
enabled with `provider = "local"` since the cache adds overhead when
the object store is already on the local filesystem.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-04-15 08:02:46 -04:00
parent cfd0005319
commit a147fecc00
No known key found for this signature in database
36 changed files with 174 additions and 18 deletions

View file

@ -58,6 +58,14 @@ strategy = "tailscale_funnel"
[server.storage]
root = "/var/lib/fabro"
[server.slatedb]
provider = "s3"
disk_cache = true
[server.slatedb.s3]
bucket = "my-fabro-data"
region = "us-east-1"
[server.scheduler]
max_concurrent_runs = 8
@ -147,6 +155,33 @@ GitHub-specific auth policy.
The GitHub OAuth client ID still lives under `[server.integrations.github].client_id`.
### `[server.slatedb]` section
Configure the embedded SlateDB key-value store used for run event storage.
| Key | Description | Default |
|---|---|---|
| `provider` | Object store backend: `local` or `s3` | `"local"` |
| `prefix` | Key prefix within the object store | `""` |
| `flush_interval` | How often to flush the write-ahead log | `"1ms"` |
| `disk_cache` | Enable a local disk cache for object store reads | `false` |
When `disk_cache = true`, Fabro creates a cache directory at `<storage_root>/cache/slatedb` and
configures SlateDB to cache object store bytes on local disk (16 GB max, 4 MB parts). This
significantly reduces read latency and costs for S3-backed deployments. A warning is emitted if
enabled with `provider = "local"` since the disk cache adds overhead when the object store is
already local.
```toml title="settings.toml"
[server.slatedb]
provider = "s3"
disk_cache = true
[server.slatedb.s3]
bucket = "{{ env.SLATEDB_BUCKET }}"
region = "us-east-1"
```
### Run defaults
The `[run.*]` sections in `settings.toml` act as defaults for every run.

View file

@ -332,6 +332,7 @@ mod tests {
Arc::clone(&object_store),
"",
Duration::from_millis(1),
None,
));
let artifact_store = ArtifactStore::new(object_store, "artifacts");
(store, artifact_store)

View file

@ -13,6 +13,7 @@ pub(crate) async fn rebuild_run_store(
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
));
let run_store = store.create_run(run_id).await?;
for event in events {

View file

@ -52,6 +52,7 @@ prefix = ""
provider = "local"
prefix = ""
flush_interval = "1ms"
disk_cache = false
[features]
session_sandboxes = false

View file

@ -467,6 +467,7 @@ fn combine_server_slatedb(
flush_interval: higher.flush_interval.or(lower.flush_interval),
local: higher.local.or(lower.local),
s3: higher.s3.or(lower.s3),
disk_cache: higher.disk_cache.or(lower.disk_cache),
}
}

View file

@ -182,11 +182,23 @@ fn resolve_slatedb(
.and_then(|slatedb| slatedb.provider)
.expect("defaults.toml should provide server.slatedb.provider");
let disk_cache = layer
.and_then(|slatedb| slatedb.disk_cache)
.expect("defaults.toml should provide server.slatedb.disk_cache");
if disk_cache && provider == ObjectStoreProvider::Local {
tracing::warn!(
"disk_cache enabled with local provider; \
disk cache is designed for S3-backed deployments \
and adds overhead on local filesystems"
);
}
ServerSlateDbSettings {
prefix: layer
prefix: layer
.and_then(|slatedb| slatedb.prefix.clone())
.expect("defaults.toml should provide server.slatedb.prefix"),
store: resolve_object_store(
store: resolve_object_store(
provider,
layer.and_then(|slatedb| slatedb.local.as_ref()),
layer.and_then(|slatedb| slatedb.s3.as_ref()),
@ -198,6 +210,7 @@ fn resolve_slatedb(
.and_then(|slatedb| slatedb.flush_interval)
.map(|duration| duration.as_std())
.expect("defaults.toml should provide server.slatedb.flush_interval"),
disk_cache,
}
}

View file

@ -60,6 +60,8 @@ fn resolves_server_defaults_from_empty_settings() {
}
ObjectStoreSettings::S3 { .. } => panic!("expected local slatedb store by default"),
}
assert!(!settings.slatedb.disk_cache);
}
#[test]
@ -195,3 +197,19 @@ enabled = true
GithubIntegrationStrategy::Token
);
}
#[test]
fn resolves_disk_cache_true_from_settings() {
let file = parse(
r"
_version = 1
[server.slatedb]
disk_cache = true
",
);
let settings = fabro_config::resolve_server_from_file(&file).expect("settings should resolve");
assert!(settings.slatedb.disk_cache);
}

View file

@ -263,10 +263,15 @@ pub fn build_artifact_object_store(
fn build_slatedb_store(
settings: &ResolvedServerSettings,
) -> anyhow::Result<(Arc<dyn ObjectStore>, String, Duration)> {
) -> anyhow::Result<(Arc<dyn ObjectStore>, String, Duration, bool)> {
let prefix = resolve_interp(&settings.slatedb.prefix)?;
let object_store = build_object_store_from_settings(&settings.slatedb.store)?;
Ok((object_store, prefix, settings.slatedb.flush_interval))
Ok((
object_store,
prefix,
settings.slatedb.flush_interval,
settings.slatedb.disk_cache,
))
}
/// Start the HTTP API server.
@ -315,12 +320,18 @@ where
};
let web_enabled = router_web_enabled(&resolved_server_settings);
let (object_store, slatedb_prefix, flush_interval) =
let (object_store, slatedb_prefix, flush_interval, disk_cache) =
build_slatedb_store(&resolved_server_settings)?;
let cache_path = if disk_cache {
Some(data_dir.join("cache").join("slatedb"))
} else {
None
};
let store = Arc::new(fabro_store::Database::new(
object_store,
slatedb_prefix,
flush_interval,
cache_path,
));
let (artifact_object_store, artifact_prefix) =
build_artifact_object_store(&resolved_server_settings)?;
@ -892,12 +903,31 @@ root = "{}"
));
let resolved = resolve_server_settings(&settings).expect("settings should resolve");
let (_object_store, prefix, flush_interval) =
let (_object_store, prefix, flush_interval, disk_cache) =
build_slatedb_store(&resolved).expect("slatedb store should build");
assert!(root.exists(), "configured SlateDB root should be created");
assert_eq!(prefix, "");
assert_eq!(flush_interval, Duration::from_millis(1));
assert!(!disk_cache);
}
#[test]
fn build_slatedb_store_returns_disk_cache_when_enabled() {
let settings = parse_settings(
r"
_version = 1
[server.slatedb]
disk_cache = true
",
);
let resolved = resolve_server_settings(&settings).expect("settings should resolve");
let (_object_store, _prefix, _flush_interval, disk_cache) =
build_slatedb_store(&resolved).expect("slatedb store should build");
assert!(disk_cache);
}
#[tokio::test]

View file

@ -2263,6 +2263,7 @@ fn test_store_bundle() -> (Arc<Database>, ArtifactStore) {
Arc::clone(&object_store),
"",
Duration::from_millis(1),
None,
));
let artifact_store = ArtifactStore::new(object_store, "artifacts");
(store, artifact_store)

View file

@ -2,6 +2,7 @@ mod catalog;
mod run_store;
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
@ -19,6 +20,7 @@ pub struct Database {
object_store: Arc<dyn ObjectStore>,
base_prefix: String,
flush_interval: Duration,
cache_path: Option<PathBuf>,
db: Arc<OnceCell<slatedb::Db>>,
active_runs: Arc<Mutex<HashMap<RunId, Arc<RunDatabaseInner>>>>,
}
@ -28,6 +30,7 @@ impl std::fmt::Debug for Database {
f.debug_struct("Database")
.field("base_prefix", &self.base_prefix)
.field("flush_interval", &self.flush_interval)
.field("cache_path", &self.cache_path)
.finish_non_exhaustive()
}
}
@ -37,11 +40,13 @@ impl Database {
object_store: Arc<dyn ObjectStore>,
base_prefix: impl Into<String>,
flush_interval: Duration,
cache_path: Option<PathBuf>,
) -> Self {
Self {
object_store,
base_prefix: normalize_base_prefix(base_prefix.into()),
flush_interval,
cache_path,
db: Arc::new(OnceCell::new()),
active_runs: Arc::new(Mutex::new(HashMap::new())),
}
@ -55,12 +60,16 @@ impl Database {
let db = self
.db
.get_or_try_init(|| async {
let mut settings = Settings {
flush_interval: Some(self.flush_interval),
compression_codec: Some(CompressionCodec::Zstd),
..Settings::default()
};
if let Some(ref cache_path) = self.cache_path {
settings.object_store_cache_options.root_folder = Some(cache_path.clone());
}
slatedb::Db::builder(self.shared_db_prefix(), self.object_store.clone())
.with_settings(Settings {
flush_interval: Some(self.flush_interval),
compression_codec: Some(CompressionCodec::Zstd),
..Settings::default()
})
.with_settings(settings)
.build()
.await
})
@ -267,7 +276,12 @@ mod tests {
fn make_store() -> (Arc<dyn ObjectStore>, Database) {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let store = Database::new(object_store.clone(), "runs/", Duration::from_millis(1));
let store = Database::new(
object_store.clone(),
"runs/",
Duration::from_millis(1),
None,
);
(object_store, store)
}
@ -563,7 +577,7 @@ mod tests {
let run = store.create_run(&test_run_id("run-1")).await.unwrap();
append_completed(&run, "run-1", dt("2026-03-27T12:00:00Z")).await;
let reopened = Database::new(object_store, "runs", Duration::from_millis(1));
let reopened = Database::new(object_store, "runs", Duration::from_millis(1), None);
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"));

View file

@ -395,7 +395,7 @@ mod tests {
#[tokio::test]
async fn list_blobs_reads_global_cas_namespace() {
let object_store = Arc::new(InMemory::new());
let store = Database::new(object_store, "", Duration::from_millis(1));
let store = Database::new(object_store, "", Duration::from_millis(1), None);
let run_id = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap();
let run = store.create_run(&run_id).await.unwrap();
let first_blob = br#"{"a":1}"#;

View file

@ -143,6 +143,7 @@ mod tests {
root: InterpString::parse("/srv/slatedb"),
},
flush_interval: StdDuration::from_secs(30),
disk_cache: false,
},
..ServerSettings::default()
},

View file

@ -146,6 +146,7 @@ pub struct ServerSlateDbSettings {
pub store: ObjectStoreSettings,
#[serde(serialize_with = "serialize_std_duration")]
pub flush_interval: StdDuration,
pub disk_cache: bool,
}
impl Default for ServerSlateDbSettings {
@ -154,6 +155,7 @@ impl Default for ServerSlateDbSettings {
prefix: InterpString::parse(""),
store: ObjectStoreSettings::default(),
flush_interval: StdDuration::ZERO,
disk_cache: false,
}
}
}
@ -373,6 +375,8 @@ pub struct ServerSlateDbLayer {
pub local: Option<ObjectStoreLocalLayer>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub s3: Option<ObjectStoreS3Layer>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub disk_cache: Option<bool>,
}
/// Closed enum of object-store providers. Unknown providers hard-fail

View file

@ -372,7 +372,7 @@ mod tests {
async fn make_run_store(label: &str) -> fabro_store::RunDatabase {
let object_store = Arc::new(InMemory::new());
let store = Database::new(object_store, "runs/", Duration::from_millis(1));
let store = Database::new(object_store, "runs/", Duration::from_millis(1), None);
store.create_run(&test_run_id(label)).await.unwrap()
}

View file

@ -3077,6 +3077,7 @@ mod tests {
std::sync::Arc::new(object_store::memory::InMemory::new()),
"",
std::time::Duration::from_millis(1),
None,
);
let run_store = store.create_run(&fixtures::RUN_7).await.unwrap();
let stored = to_run_event(&fixtures::RUN_7, &Event::RunNotice {

View file

@ -358,6 +358,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -418,6 +418,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -197,6 +197,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -239,6 +239,7 @@ impl Handler for SubWorkflowHandler {
object_store.clone(),
"",
Duration::from_millis(1),
None,
));
let run_store = store
.create_run(&child_run_options.run_id)

View file

@ -102,6 +102,7 @@ impl EngineServices {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
));
Self {
registry: Arc::new(HandlerRegistry::new(Box::new(start::StartHandler))),

View file

@ -654,6 +654,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -197,6 +197,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -438,6 +438,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}
@ -965,7 +966,12 @@ mod tests {
std::fs::create_dir_all(storage_dir.join("store")).unwrap();
let object_store =
Arc::new(LocalFileSystem::new_with_prefix(storage_dir.join("store")).unwrap());
let store = Arc::new(Database::new(object_store, "", Duration::from_millis(1)));
let store = Arc::new(Database::new(
object_store,
"",
Duration::from_millis(1),
None,
));
let created = create(store.as_ref(), CreateRunInput {
workflow: WorkflowInput::DotSource {
source: MINIMAL_DOT.to_string(),
@ -1001,7 +1007,12 @@ mod tests {
std::fs::create_dir_all(storage_dir.join("store")).unwrap();
let object_store =
Arc::new(LocalFileSystem::new_with_prefix(storage_dir.join("store")).unwrap());
let store = Arc::new(Database::new(object_store, "", Duration::from_millis(1)));
let store = Arc::new(Database::new(
object_store,
"",
Duration::from_millis(1),
None,
));
let created = create(store.as_ref(), CreateRunInput {
workflow: WorkflowInput::DotSource {
source: MINIMAL_DOT.to_string(),

View file

@ -366,6 +366,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -991,6 +991,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -166,6 +166,7 @@ async fn test_run_store(run_id: &RunId) -> fabro_store::RunDatabase {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
));
store.create_run(run_id).await.unwrap()
}

View file

@ -338,6 +338,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -782,6 +782,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -71,6 +71,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -675,6 +675,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -222,6 +222,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -405,6 +405,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}

View file

@ -125,6 +125,7 @@ mod tests {
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
));
store.create_run(&fixtures::RUN_1).await.unwrap()
}

View file

@ -65,6 +65,7 @@ async fn initialized(
),
"",
Duration::from_millis(1),
None,
));
let inner_store = store
.create_run(&run_options.run_id)

View file

@ -71,6 +71,7 @@ fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error:
object_store,
"",
std::time::Duration::from_millis(1),
None,
));
let state = if tokio::runtime::Handle::try_current().is_ok() {
std::thread::spawn(

View file

@ -102,7 +102,12 @@ fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error:
test_store_dir(&run_dir)
};
let object_store = Arc::new(LocalFileSystem::new_with_prefix(store_dir)?);
let store = Arc::new(Database::new(object_store, "", Duration::from_millis(1)));
let store = Arc::new(Database::new(
object_store,
"",
Duration::from_millis(1),
None,
));
let state = if tokio::runtime::Handle::try_current().is_ok() {
std::thread::spawn(
move || -> Result<_, Box<dyn std::error::Error + Send + Sync>> {