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};
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
use corium_log::{NativeLogStorage, NativeVersionedLog};
use corium_store::{BlobId, BlobIdStream, BlobStore, FsStore, MemoryStore, RootStore, StoreError};
#[cfg(feature = "postgres")]
use corium_store::PostgresBlobStore;
#[cfg(feature = "s3")]
use corium_store::S3BlobStore;
#[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,
},
#[cfg(feature = "s3")]
S3 {
bucket: String,
prefix: 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(),
#[cfg(feature = "s3")]
Self::S3 { bucket, prefix } => formatter
.debug_struct("S3")
.field("bucket", bucket)
.field("prefix", prefix)
.finish(),
}
}
}
pub enum NodeStore {
Mem(MemoryStore),
Fs(FsStore),
#[cfg(feature = "postgres")]
Postgres(PostgresBlobStore),
#[cfg(feature = "turso")]
Turso(TursoBlobStore),
#[cfg(feature = "s3")]
S3(S3BlobStore),
}
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?)),
#[cfg(feature = "s3")]
StoreSpec::S3 { bucket, prefix } => {
Ok(Self::S3(S3BlobStore::connect(bucket, prefix).await?))
}
}
}
#[allow(clippy::unused_async)]
pub async fn open_existing(
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_existing(connection_string).await?,
)),
#[cfg(feature = "turso")]
StoreSpec::Turso { path } => {
Ok(Self::Turso(TursoBlobStore::open_existing(path).await?))
}
#[cfg(feature = "s3")]
StoreSpec::S3 { bucket, prefix } => {
Ok(Self::S3(S3BlobStore::connect(bucket, prefix).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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(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,
#[cfg(feature = "s3")]
Self::S3(store) => store.list_roots(prefix).await,
}
}
}
pub enum LogBackend {
Fs(PathBuf),
Mem(MemLogRegistry),
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
Native(Arc<dyn NativeLogStorage>),
}
impl LogBackend {
#[must_use]
#[allow(clippy::needless_pass_by_value)]
pub fn for_spec(
spec: &StoreSpec,
data_dir: &std::path::Path,
#[cfg_attr(
not(any(feature = "postgres", feature = "turso", feature = "s3")),
allow(unused_variables)
)]
store: Arc<NodeStore>,
) -> Self {
match spec {
StoreSpec::Memory => Self::Mem(MemLogRegistry::new()),
StoreSpec::Fs => Self::Fs(data_dir.join("logs")),
#[cfg(feature = "postgres")]
StoreSpec::Postgres { .. } => Self::Native(Arc::new(NativeRootLogStore::new(store))),
#[cfg(feature = "turso")]
StoreSpec::Turso { .. } => Self::Native(Arc::new(NativeRootLogStore::new(store))),
#[cfg(feature = "s3")]
StoreSpec::S3 { .. } => Self::Native(Arc::new(NativeRootLogStore::new(store))),
}
}
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))),
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
Self::Native(storage) => Ok(Arc::new(NativeVersionedLog::open(
Arc::clone(storage),
name,
write_version,
)?)),
}
}
#[must_use]
pub fn exists(&self, name: &str) -> bool {
match self {
Self::Fs(dir) => VersionedLog::exists(dir, name),
Self::Mem(registry) => registry.exists(name),
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
Self::Native(storage) => storage
.list_chunks(name)
.is_ok_and(|chunks| !chunks.is_empty()),
}
}
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(())
}
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
Self::Native(storage) => storage.delete_all(name),
}
}
}
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
struct NativeRootLogStore {
store: Arc<NodeStore>,
}
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
impl NativeRootLogStore {
fn new(store: Arc<NodeStore>) -> Self {
Self { store }
}
}
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
impl NativeRootLogStore {
fn key(name: &str, version: u64, chunk: u64) -> String {
if chunk == 0 {
format!("log:{name}:v{version:020}")
} else {
format!("log:{name}:v{version:020}:c{chunk:020}")
}
}
fn prefix(name: &str) -> String {
format!("log:{name}:v")
}
fn parse_key(prefix: &str, key: &str) -> Option<(u64, u64)> {
let rest = key.strip_prefix(prefix)?;
match rest.split_once(":c") {
Some((version, chunk)) => Some((version.parse().ok()?, chunk.parse().ok()?)),
None => Some((rest.parse().ok()?, 0)),
}
}
fn block_on<T>(
&self,
future: impl std::future::Future<Output = Result<T, StoreError>>,
) -> Result<T, LogError> {
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current()
.block_on(future)
.map_err(|error| LogError::Native(error.to_string()))
})
}
}
#[cfg(any(feature = "postgres", feature = "turso", feature = "s3"))]
impl NativeLogStorage for NativeRootLogStore {
fn read_chunk(
&self,
name: &str,
version: u64,
chunk: u64,
) -> Result<Option<Vec<u8>>, LogError> {
self.block_on(self.store.get_root(&Self::key(name, version, chunk)))
}
fn cas_chunk(
&self,
name: &str,
version: u64,
chunk: u64,
expected: Option<&[u8]>,
new: &[u8],
) -> Result<(), LogError> {
self.block_on(
self.store
.cas_root(&Self::key(name, version, chunk), expected, new),
)
}
fn list_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError> {
let prefix = Self::prefix(name);
let names = self.block_on(self.store.list_roots(&prefix))?;
names
.into_iter()
.map(|key| Self::parse_key(&prefix, &key).ok_or(LogError::Corrupt))
.collect()
}
fn delete_all(&self, name: &str) -> Result<(), LogError> {
for (version, chunk) in self.list_chunks(name)? {
self.block_on(self.store.delete_root(&Self::key(name, version, chunk)))?;
}
Ok(())
}
}