use std::fmt;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::SystemTime;
use async_trait::async_trait;
use corium_log::{LogError, MemLogRegistry, TransactionLog, VersionedLog};
use corium_store::{BlobId, BlobIdStream, BlobStore, FsStore, MemoryStore, RootStore, StoreError};
#[cfg(feature = "postgres")]
use corium_store::PostgresBlobStore;
#[cfg(feature = "turso")]
use corium_store::TursoBlobStore;
#[derive(Clone, Default)]
pub enum StoreSpec {
Memory,
#[default]
Fs,
#[cfg(feature = "postgres")]
Postgres {
connection_string: String,
},
#[cfg(feature = "turso")]
Turso {
path: String,
},
}
impl fmt::Debug for StoreSpec {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Memory => formatter.write_str("Memory"),
Self::Fs => formatter.write_str("Fs"),
#[cfg(feature = "postgres")]
Self::Postgres { .. } => formatter
.debug_struct("Postgres")
.field("connection_string", &"[REDACTED]")
.finish(),
#[cfg(feature = "turso")]
Self::Turso { path } => formatter.debug_struct("Turso").field("path", path).finish(),
}
}
}
pub enum NodeStore {
Mem(MemoryStore),
Fs(FsStore),
#[cfg(feature = "postgres")]
Postgres(PostgresBlobStore),
#[cfg(feature = "turso")]
Turso(TursoBlobStore),
}
impl NodeStore {
#[allow(clippy::unused_async)]
pub async fn open(spec: &StoreSpec, data_dir: &std::path::Path) -> Result<Self, StoreError> {
match spec {
StoreSpec::Memory => Ok(Self::Mem(MemoryStore::default())),
StoreSpec::Fs => Ok(Self::Fs(FsStore::open(data_dir.join("store"))?)),
#[cfg(feature = "postgres")]
StoreSpec::Postgres { connection_string } => Ok(Self::Postgres(
PostgresBlobStore::connect(connection_string).await?,
)),
#[cfg(feature = "turso")]
StoreSpec::Turso { path } => Ok(Self::Turso(TursoBlobStore::open(path).await?)),
}
}
}
#[async_trait]
impl BlobStore for NodeStore {
async fn put(&self, bytes: &[u8]) -> Result<BlobId, StoreError> {
match self {
Self::Mem(store) => store.put(bytes).await,
Self::Fs(store) => store.put(bytes).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.put(bytes).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.put(bytes).await,
}
}
async fn get(&self, id: &BlobId) -> Result<Option<Vec<u8>>, StoreError> {
match self {
Self::Mem(store) => store.get(id).await,
Self::Fs(store) => store.get(id).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.get(id).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.get(id).await,
}
}
async fn contains(&self, id: &BlobId) -> Result<bool, StoreError> {
match self {
Self::Mem(store) => store.contains(id).await,
Self::Fs(store) => store.contains(id).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.contains(id).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.contains(id).await,
}
}
async fn delete(&self, id: &BlobId) -> Result<(), StoreError> {
match self {
Self::Mem(store) => store.delete(id).await,
Self::Fs(store) => store.delete(id).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.delete(id).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.delete(id).await,
}
}
async fn list(&self) -> Result<BlobIdStream, StoreError> {
match self {
Self::Mem(store) => store.list().await,
Self::Fs(store) => store.list().await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.list().await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.list().await,
}
}
async fn modified_at(&self, id: &BlobId) -> Result<Option<SystemTime>, StoreError> {
match self {
Self::Mem(store) => store.modified_at(id).await,
Self::Fs(store) => store.modified_at(id).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.modified_at(id).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.modified_at(id).await,
}
}
}
#[async_trait]
impl RootStore for NodeStore {
async fn get_root(&self, name: &str) -> Result<Option<Vec<u8>>, StoreError> {
match self {
Self::Mem(store) => store.get_root(name).await,
Self::Fs(store) => store.get_root(name).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.get_root(name).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.get_root(name).await,
}
}
async fn cas_root(
&self,
name: &str,
expected: Option<&[u8]>,
new: &[u8],
) -> Result<(), StoreError> {
match self {
Self::Mem(store) => store.cas_root(name, expected, new).await,
Self::Fs(store) => store.cas_root(name, expected, new).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.cas_root(name, expected, new).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.cas_root(name, expected, new).await,
}
}
async fn delete_root(&self, name: &str) -> Result<(), StoreError> {
match self {
Self::Mem(store) => store.delete_root(name).await,
Self::Fs(store) => store.delete_root(name).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.delete_root(name).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.delete_root(name).await,
}
}
async fn list_roots(&self, prefix: &str) -> Result<Vec<String>, StoreError> {
match self {
Self::Mem(store) => store.list_roots(prefix).await,
Self::Fs(store) => store.list_roots(prefix).await,
#[cfg(feature = "postgres")]
Self::Postgres(store) => store.list_roots(prefix).await,
#[cfg(feature = "turso")]
Self::Turso(store) => store.list_roots(prefix).await,
}
}
}
pub enum LogBackend {
Fs(PathBuf),
Mem(MemLogRegistry),
}
impl LogBackend {
#[must_use]
pub fn for_spec(spec: &StoreSpec, data_dir: &std::path::Path) -> Self {
match spec {
StoreSpec::Memory => Self::Mem(MemLogRegistry::new()),
StoreSpec::Fs => Self::Fs(data_dir.join("logs")),
#[cfg(feature = "postgres")]
StoreSpec::Postgres { .. } => Self::Fs(data_dir.join("logs")),
#[cfg(feature = "turso")]
StoreSpec::Turso { .. } => Self::Fs(data_dir.join("logs")),
}
}
pub fn open(
&self,
name: &str,
write_version: u64,
) -> Result<Arc<dyn TransactionLog>, LogError> {
match self {
Self::Fs(dir) => Ok(Arc::new(VersionedLog::open(dir, name, write_version)?)),
Self::Mem(registry) => Ok(Arc::new(registry.open(name, write_version))),
}
}
pub fn delete_all(&self, name: &str) -> Result<(), LogError> {
match self {
Self::Fs(dir) => VersionedLog::delete_all(dir, name),
Self::Mem(registry) => {
registry.delete_all(name);
Ok(())
}
}
}
}