use std::path::Path;
use std::sync::Arc;
use boatramp_core::kv::{KvOpenPolicy, KvStore, MemoryKv};
use crate::error::Result;
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Deserialize)]
#[serde(rename_all = "lowercase")]
#[cfg_attr(feature = "clap", derive(clap::ValueEnum))]
pub enum BlobBackend {
Fs,
S3,
Gcs,
Azure,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[cfg_attr(feature = "clap", derive(clap::ValueEnum))]
pub enum KvBackend {
Slatedb,
Memory,
Cloudflare,
Sql,
}
pub const CONTROL_PLANE_FLUSH: std::time::Duration = std::time::Duration::from_millis(5);
#[derive(Debug, Clone)]
pub struct SlateKvS3 {
pub bucket: String,
pub endpoint: Option<String>,
pub region: Option<String>,
pub path_style: bool,
pub prefix: String,
}
pub async fn build_kv(
kv: KvBackend,
data_dir: &Path,
slate_s3: Option<&SlateKvS3>,
policy: KvOpenPolicy,
sql: Option<&crate::config::SqlKvConfig>,
) -> Result<Arc<dyn KvStore>> {
match kv {
KvBackend::Slatedb => build_slatedb_kv(data_dir, slate_s3, policy).await,
KvBackend::Memory => Ok(Arc::new(MemoryKv::new())),
KvBackend::Cloudflare => build_cloudflare_kv(),
KvBackend::Sql => build_sql_kv(sql).await,
}
}
#[cfg(feature = "sql")]
async fn build_sql_kv(sql: Option<&crate::config::SqlKvConfig>) -> Result<Arc<dyn KvStore>> {
use crate::error::Error;
let cfg = sql.ok_or_else(|| {
Error::SqlKvConfig(
"`--kv sql` needs a `[serve.kv.sql]` config block (or the `BOATRAMP_KV_SQL_*` env)"
.to_string(),
)
})?;
match cfg.kind.trim().to_ascii_lowercase().as_str() {
"" | "sqlite" | "sqlite3" | "libsql" => {
let path = cfg
.path
.as_deref()
.filter(|p| !p.is_empty())
.ok_or_else(|| {
Error::SqlKvConfig(
"`[serve.kv.sql] kind = sqlite` needs `path` (an on-disk file) — set it or \
`BOATRAMP_KV_SQL_PATH`"
.to_string(),
)
})?;
Ok(Arc::new(
boatramp_storage::SqlKv::open_sqlite_local(path).await?,
))
}
engine @ ("postgres" | "postgresql" | "pg") => build_pg_kv(cfg, engine).await,
engine @ ("mysql" | "mariadb") => build_mysql_kv(cfg, engine).await,
other => Err(Error::SqlKvConfig(format!(
"unknown `[serve.kv.sql] kind` {other:?}: expected sqlite | postgres | mysql"
))),
}
}
#[cfg(all(feature = "sql", any(feature = "sql-postgres", feature = "sql-mysql")))]
fn resolve_kv_url(cfg: &crate::config::SqlKvConfig, engine: &str) -> Result<String> {
use crate::error::Error;
let url_env = cfg
.url_env
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty())
.ok_or_else(|| {
Error::SqlKvConfig(format!(
"`[serve.kv.sql] kind = {engine}` needs `url_env` (the NAME of the env var holding \
the connection URL — never a raw URL in config) — set it or `BOATRAMP_KV_SQL_URL_ENV`"
))
})?;
std::env::var(url_env)
.ok()
.map(|v| v.trim().to_string())
.filter(|v| !v.is_empty())
.ok_or_else(|| {
Error::SqlKvConfig(format!(
"`[serve.kv.sql] url_env = {url_env:?}` names an unset or empty env var; set \
{url_env} to the {engine} connection URL"
))
})
}
#[cfg(all(feature = "sql", feature = "sql-postgres"))]
async fn build_pg_kv(cfg: &crate::config::SqlKvConfig, engine: &str) -> Result<Arc<dyn KvStore>> {
let url = resolve_kv_url(cfg, engine)?;
Ok(Arc::new(
boatramp_storage::SqlKv::open_postgres(url, cfg.pool_max).await?,
))
}
#[cfg(all(feature = "sql", not(feature = "sql-postgres")))]
async fn build_pg_kv(_cfg: &crate::config::SqlKvConfig, engine: &str) -> Result<Arc<dyn KvStore>> {
Err(crate::error::Error::SqlKvConfig(format!(
"the `{engine}` SQL KV backend needs the Postgres engine — rebuild with `--features sql-postgres`"
)))
}
#[cfg(all(feature = "sql", feature = "sql-mysql"))]
async fn build_mysql_kv(
cfg: &crate::config::SqlKvConfig,
engine: &str,
) -> Result<Arc<dyn KvStore>> {
let url = resolve_kv_url(cfg, engine)?;
Ok(Arc::new(
boatramp_storage::SqlKv::open_mysql(url, cfg.pool_max).await?,
))
}
#[cfg(all(feature = "sql", not(feature = "sql-mysql")))]
async fn build_mysql_kv(
_cfg: &crate::config::SqlKvConfig,
engine: &str,
) -> Result<Arc<dyn KvStore>> {
Err(crate::error::Error::SqlKvConfig(format!(
"the `{engine}` SQL KV backend needs the MySQL engine — rebuild with `--features sql-mysql`"
)))
}
#[cfg(not(feature = "sql"))]
async fn build_sql_kv(_sql: Option<&crate::config::SqlKvConfig>) -> Result<Arc<dyn KvStore>> {
Err(crate::error::Error::NoSqlSupport)
}
#[cfg(feature = "slatedb")]
async fn build_slatedb_kv(
data_dir: &Path,
slate_s3: Option<&SlateKvS3>,
policy: KvOpenPolicy,
) -> Result<Arc<dyn KvStore>> {
match slate_s3 {
Some(s3) => Ok(Arc::new(
boatramp_storage::SlateKv::open_s3_with_flush_policy(
&boatramp_storage::S3StoreConfig {
bucket: s3.bucket.clone(),
endpoint: s3.endpoint.clone(),
region: s3.region.clone(),
path_style: s3.path_style,
},
&s3.prefix,
CONTROL_PLANE_FLUSH,
policy,
)
.await?,
)),
None => Ok(Arc::new(
boatramp_storage::SlateKv::open_local_with_flush_policy(
data_dir.join("kv-slate"),
CONTROL_PLANE_FLUSH,
policy,
)
.await?,
)),
}
}
#[cfg(not(feature = "slatedb"))]
async fn build_slatedb_kv(
_data_dir: &Path,
_slate_s3: Option<&SlateKvS3>,
_policy: KvOpenPolicy,
) -> Result<Arc<dyn KvStore>> {
Err(crate::error::Error::NoSlatedbSupport)
}
#[cfg(feature = "cloudflare-kv")]
fn build_cloudflare_kv() -> Result<Arc<dyn KvStore>> {
Ok(Arc::new(boatramp_storage::CloudflareKv::from_env()?))
}
#[cfg(not(feature = "cloudflare-kv"))]
fn build_cloudflare_kv() -> Result<Arc<dyn KvStore>> {
Err(crate::error::Error::NoCloudflareKvSupport)
}