use std::path::PathBuf;
use std::sync::Arc;
use async_trait::async_trait;
use meerkat_core::{ArtifactStore, BlobStore, DurabilityDeclaration, RealmLocator, StorageLayout};
use meerkat_jobs::DetachedJobStore;
use meerkat_runtime::RuntimeStore;
use meerkat_schedule::ScheduleStore;
use meerkat_workgraph::WorkGraphStore;
use crate::SessionStore;
use crate::persistence::PersistenceError;
use meerkat_store::{RealmManifestPin, RealmPaths};
#[derive(Clone)]
pub struct RealmOpenContext {
pub locator: RealmLocator,
pub manifest: RealmManifestPin,
pub paths: RealmPaths,
pub layout: Option<StorageLayout>,
}
pub struct RealmStoreSet {
pub session_store: Arc<dyn SessionStore>,
pub runtime_store: Arc<dyn RuntimeStore>,
pub schedule_store: Arc<dyn ScheduleStore>,
pub workgraph_store: Arc<dyn WorkGraphStore>,
pub job_store: Arc<dyn DetachedJobStore>,
pub blob_store: Arc<dyn BlobStore>,
pub artifact_store: Arc<dyn ArtifactStore>,
pub store_path: PathBuf,
pub projection_root: Option<PathBuf>,
pub durability: Vec<DurabilityDeclaration>,
}
#[async_trait]
pub trait RealmStorageProvider: Send + Sync {
fn name(&self) -> &str;
async fn open(&self, ctx: &RealmOpenContext) -> Result<RealmStoreSet, PersistenceError>;
fn migrator(&self) -> Option<&dyn meerkat_core::StorageMigrator> {
None
}
}
pub const REQUIRED_DURABILITY_DOMAINS: [&str; 7] = [
"sessions",
"runtime",
"schedule",
"workgraph",
"jobs",
"blobs",
"artifacts",
];
pub fn enforce_fail_closed_durability(
set: &RealmStoreSet,
ephemeral_domains: &[String],
) -> Result<(), PersistenceError> {
for required in REQUIRED_DURABILITY_DOMAINS {
let count = set
.durability
.iter()
.filter(|declaration| declaration.domain == required)
.count();
if count != 1 {
return Err(PersistenceError::DurabilityViolation {
domain: format!(
"{required} (provider supplied {count} durability declarations for this \
slot; exactly one is required)"
),
});
}
}
for declaration in &set.durability {
if declaration.is_undeclared_nonpersistent_durable()
&& !ephemeral_domains
.iter()
.any(|domain| domain == &declaration.domain)
{
return Err(PersistenceError::DurabilityViolation {
domain: declaration.domain.clone(),
});
}
}
Ok(())
}
pub async fn open_realm_persistence_with_layout(
layout: StorageLayout,
realm_id: &str,
backend_hint: Option<meerkat_store::RealmBackend>,
origin_hint: Option<meerkat_store::RealmOrigin>,
) -> Result<(meerkat_store::RealmManifest, crate::PersistenceBundle), PersistenceError> {
crate::persistence::open_realm_persistence_builtin_with_layout(
layout,
realm_id,
backend_hint,
origin_hint,
)
.await
}
#[derive(Debug, Clone, Copy, Default)]
pub struct DiskStorageProvider;
#[async_trait]
impl RealmStorageProvider for DiskStorageProvider {
fn name(&self) -> &'static str {
"disk"
}
async fn open(&self, ctx: &RealmOpenContext) -> Result<RealmStoreSet, PersistenceError> {
crate::persistence::open_disk_store_set(ctx)
}
fn migrator(&self) -> Option<&dyn meerkat_core::StorageMigrator> {
static DISK_MIGRATOR: meerkat_store::doctor::DiskStorageMigrator =
meerkat_store::doctor::DiskStorageMigrator;
Some(&DISK_MIGRATOR)
}
}