1use std::path::PathBuf;
27use std::sync::Arc;
28
29use async_trait::async_trait;
30use meerkat::storage_provider::RealmStorageProvider;
31use meerkat_core::{DurabilityDeclaration, DurabilityResolution};
32
33use crate::blob_store::{BinaryBlobStore, ObjectStoreBlobStore};
34use crate::console_aggregator::{ConsoleLogStore, SqliteConsoleLogStore};
35use crate::identity_first::contracts::{ContinuityStore, LeaseProvider};
36use crate::identity_first::{AgentMemoryProvider, LocalContinuityStore};
37use crate::memory::sqlite_store::SqliteAgentMemoryStore;
38use crate::runtime::{PersistentMetadataStore, SqliteMetadataStore};
39use crate::storage_doctor::MobKitStorageMigrator;
40use crate::storage_layout::MobKitStorageLayout;
41use crate::unified_runtime::EventLogStore;
42
43#[derive(Clone)]
45pub struct MobKitRealmOpenContext {
46 pub layout: MobKitStorageLayout,
49 pub state_dir: PathBuf,
51 pub declared_ephemeral_domains: Vec<String>,
55}
56
57impl MobKitRealmOpenContext {
58 pub fn for_state_dir(state_dir: impl Into<PathBuf>) -> Self {
60 let state_dir = state_dir.into();
61 Self {
62 layout: MobKitStorageLayout::with_injected_roots(state_dir.clone(), None),
63 state_dir,
64 declared_ephemeral_domains: Vec::new(),
65 }
66 }
67
68 pub fn is_declared_ephemeral(&self, domain: &str) -> bool {
70 self.declared_ephemeral_domains
71 .iter()
72 .any(|declared| declared == domain)
73 }
74}
75
76#[derive(Clone)]
81pub enum MobKitLeaseAuthority {
82 Provider(Arc<dyn LeaseProvider>),
83 FencingFloor(u64),
84}
85
86pub struct MobKitRealmStoreSet {
94 pub continuity_store: Arc<dyn ContinuityStore>,
95 pub lease_authority: MobKitLeaseAuthority,
96 pub event_log_store: Option<Box<dyn EventLogStore>>,
97 pub console_log_store: Arc<dyn ConsoleLogStore>,
98 pub metadata_store: Arc<dyn PersistentMetadataStore>,
99 pub blob_store: Arc<dyn BinaryBlobStore>,
100 pub agent_memory_provider: Option<Arc<dyn AgentMemoryProvider>>,
101 pub schedule_store: Arc<dyn meerkat::ScheduleStore>,
102 pub durability: Vec<DurabilityDeclaration>,
106}
107
108#[derive(Debug)]
110pub enum MobKitStorageProviderError {
111 Open { slot: String, message: String },
113 DurabilityViolation { domain: String },
116}
117
118impl std::fmt::Display for MobKitStorageProviderError {
119 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
120 match self {
121 Self::Open { slot, message } => {
122 write!(f, "failed to open the {slot} store: {message}")
123 }
124 Self::DurabilityViolation { domain } => write!(
125 f,
126 "durable storage slot '{domain}' resolved to a non-persistent \
127 backend without an explicit ephemeral declaration — declare \
128 the domain ephemeral or supply persistent storage \
129 (fail-closed durability, storage-unification principle 7)"
130 ),
131 }
132 }
133}
134
135impl std::error::Error for MobKitStorageProviderError {}
136
137#[async_trait]
141pub trait MobKitStorageProvider: Send + Sync {
142 fn name(&self) -> &str;
144
145 async fn open_realm(
150 &self,
151 ctx: &MobKitRealmOpenContext,
152 ) -> Result<MobKitRealmStoreSet, MobKitStorageProviderError>;
153
154 fn meerkat_provider(&self) -> &dyn RealmStorageProvider;
158
159 fn migrator(&self) -> Option<&dyn meerkat_core::StorageMigrator> {
163 None
164 }
165}
166
167pub const REQUIRED_MOBKIT_DURABILITY_DOMAINS: [&str; 7] = [
171 "continuity",
172 "event_log",
173 "console",
174 "metadata",
175 "blobs",
176 "agent_memory",
177 "schedule",
178];
179
180pub fn enforce_fail_closed_store_set(
190 set: &MobKitRealmStoreSet,
191 ctx: &MobKitRealmOpenContext,
192) -> Result<(), MobKitStorageProviderError> {
193 for required in REQUIRED_MOBKIT_DURABILITY_DOMAINS {
194 let count = set
195 .durability
196 .iter()
197 .filter(|declaration| declaration.domain == required)
198 .count();
199 if count != 1 {
200 return Err(MobKitStorageProviderError::DurabilityViolation {
201 domain: format!(
202 "{required} (provider supplied {count} durability declarations for this slot; exactly one is required)"
203 ),
204 });
205 }
206 }
207 for declaration in &set.durability {
208 if declaration.is_undeclared_nonpersistent_durable()
209 && !ctx.is_declared_ephemeral(&declaration.domain)
210 {
211 return Err(MobKitStorageProviderError::DurabilityViolation {
212 domain: declaration.domain.clone(),
213 });
214 }
215 }
216 Ok(())
217}
218
219pub const MEERKAT_LEVEL_REALM_ID: &str = "mobkit";
224
225#[derive(Clone)]
231pub(crate) struct ProviderMeerkatStores {
232 pub provider_name: String,
233 pub runtime_store: Arc<dyn meerkat_runtime::RuntimeStore>,
234 pub runtime_declaration: DurabilityDeclaration,
235 pub workgraph_store: Arc<dyn meerkat::WorkGraphStore>,
236 pub workgraph_declaration: DurabilityDeclaration,
237 pub job_store: Arc<dyn meerkat::DetachedJobStore>,
238 pub job_declaration: DurabilityDeclaration,
239}
240
241impl ProviderMeerkatStores {
242 pub(crate) fn runtime_slot_summary(&self) -> crate::storage_health::StorageSlotSummary {
243 self.slot_summary(&self.runtime_declaration)
244 }
245
246 pub(crate) fn workgraph_slot_summary(&self) -> crate::storage_health::StorageSlotSummary {
247 self.slot_summary(&self.workgraph_declaration)
248 }
249
250 pub(crate) fn job_slot_summary(&self) -> crate::storage_health::StorageSlotSummary {
251 self.slot_summary(&self.job_declaration)
252 }
253
254 fn slot_summary(
255 &self,
256 declaration: &DurabilityDeclaration,
257 ) -> crate::storage_health::StorageSlotSummary {
258 crate::storage_health::StorageSlotSummary {
259 declaration: declaration.clone(),
260 backend: format!("storage provider '{}'", self.provider_name),
261 detail: Some(
262 "meerkat-level slot from the composite provider's realm bundle".to_string(),
263 ),
264 degraded: false,
265 }
266 }
267}
268
269pub(crate) async fn open_provider_meerkat_stores(
275 provider: &dyn RealmStorageProvider,
276 ctx: &MobKitRealmOpenContext,
277) -> Result<ProviderMeerkatStores, MobKitStorageProviderError> {
278 let provider_name = provider.name().to_string();
279 let pin_name = (provider_name != "disk").then_some(provider_name.as_str());
280 let pin = meerkat_store::realm::ensure_realm_manifest_pin_with_candidates(
281 &ctx.state_dir,
282 &[],
283 MEERKAT_LEVEL_REALM_ID,
284 pin_name,
285 None,
286 None,
287 )
288 .await
289 .map_err(|error| MobKitStorageProviderError::Open {
290 slot: "meerkat-level realm manifest".to_string(),
291 message: error.to_string(),
292 })?;
293 let realm = meerkat_core::RealmId::parse(MEERKAT_LEVEL_REALM_ID).map_err(|error| {
294 MobKitStorageProviderError::Open {
295 slot: "meerkat-level realm".to_string(),
296 message: format!("invalid realm id '{MEERKAT_LEVEL_REALM_ID}': {error}"),
297 }
298 })?;
299 let open_ctx = meerkat::storage_provider::RealmOpenContext {
300 locator: meerkat_core::RealmLocator {
301 state_root: ctx.state_dir.clone(),
302 realm,
303 },
304 manifest: pin.clone(),
305 paths: meerkat_store::realm_paths_in(&ctx.state_dir, MEERKAT_LEVEL_REALM_ID),
306 layout: None,
307 };
308 let set = provider
309 .open(&open_ctx)
310 .await
311 .map_err(|error| MobKitStorageProviderError::Open {
312 slot: "meerkat-level realm store set".to_string(),
313 message: error.to_string(),
314 })?;
315 let mut ephemeral_domains = pin.ephemeral_domains().to_vec();
318 ephemeral_domains.extend(ctx.declared_ephemeral_domains.iter().cloned());
319 meerkat::storage_provider::enforce_fail_closed_durability(&set, &ephemeral_domains).map_err(
320 |error| match error {
321 meerkat::PersistenceError::DurabilityViolation { domain } => {
322 MobKitStorageProviderError::DurabilityViolation {
323 domain: format!("{domain} (meerkat level)"),
324 }
325 }
326 other => MobKitStorageProviderError::Open {
327 slot: "meerkat-level realm store set".to_string(),
328 message: other.to_string(),
329 },
330 },
331 )?;
332 let declaration = |domain: &str| {
333 set.durability
334 .iter()
335 .find(|declaration| declaration.domain == domain)
336 .cloned()
337 .ok_or_else(|| MobKitStorageProviderError::DurabilityViolation {
338 domain: format!("{domain} (meerkat level: declaration missing)"),
339 })
340 };
341 let runtime_declaration = declaration("runtime")?;
342 let workgraph_declaration = declaration("workgraph")?;
343 let job_declaration = declaration("jobs")?;
344 Ok(ProviderMeerkatStores {
345 provider_name,
346 runtime_store: set.runtime_store,
347 runtime_declaration,
348 workgraph_store: set.workgraph_store,
349 workgraph_declaration,
350 job_store: set.job_store,
351 job_declaration,
352 })
353}
354
355#[derive(Debug, Clone, Copy, Default)]
358pub struct DiskMobKitStorageProvider;
359
360impl DiskMobKitStorageProvider {
361 fn open_error(slot: &str, message: impl std::fmt::Display) -> MobKitStorageProviderError {
362 MobKitStorageProviderError::Open {
363 slot: slot.to_string(),
364 message: message.to_string(),
365 }
366 }
367}
368
369#[async_trait]
370impl MobKitStorageProvider for DiskMobKitStorageProvider {
371 fn name(&self) -> &'static str {
372 "disk"
373 }
374
375 async fn open_realm(
376 &self,
377 ctx: &MobKitRealmOpenContext,
378 ) -> Result<MobKitRealmStoreSet, MobKitStorageProviderError> {
379 std::fs::create_dir_all(&ctx.state_dir)
380 .map_err(|error| Self::open_error("state directory", error))?;
381 let layout = &ctx.layout;
382
383 let continuity_path = layout
384 .continuity_db()
385 .map_err(|error| Self::open_error("continuity", error))?
386 .path;
387 let (continuity_store, fencing_floor) =
388 LocalContinuityStore::open_with_fencing_floor(continuity_path)
389 .await
390 .map_err(|error| Self::open_error("continuity", error))?;
391
392 let console_path = layout
393 .console_db()
394 .map_err(|error| Self::open_error("console", error))?
395 .path;
396 let console_log_store = SqliteConsoleLogStore::open(&console_path)
397 .map_err(|error| Self::open_error("console", error))?;
398
399 let metadata_path = layout
400 .metadata_db()
401 .map_err(|error| Self::open_error("metadata", error))?
402 .path;
403 let metadata_store = SqliteMetadataStore::open(&metadata_path)
404 .map_err(|error| Self::open_error("metadata", error))?;
405
406 let (blob_store, blob_resolution): (Arc<dyn BinaryBlobStore>, DurabilityResolution) =
407 if ctx.is_declared_ephemeral("blobs") {
408 (
409 Arc::new(ObjectStoreBlobStore::memory()),
410 DurabilityResolution::DeclaredEphemeral,
411 )
412 } else {
413 (
414 Arc::new(
415 ObjectStoreBlobStore::local(layout.blob_root())
416 .map_err(|error| Self::open_error("blobs", error))?,
417 ),
418 DurabilityResolution::Persistent,
419 )
420 };
421
422 let schedule_store = meerkat::SqliteScheduleStore::open(layout.schedule_db())
423 .map_err(|error| Self::open_error("schedule", error))?;
424
425 let agent_memory_root = layout
426 .agent_memory_root()
427 .map_err(|error| Self::open_error("agent memory", error))?
428 .path;
429 let agent_memory_provider = SqliteAgentMemoryStore::open(&agent_memory_root)
430 .map_err(|error| Self::open_error("agent memory", error))?;
431
432 let set = MobKitRealmStoreSet {
433 continuity_store: Arc::new(continuity_store),
434 lease_authority: MobKitLeaseAuthority::FencingFloor(fencing_floor),
435 event_log_store: None,
440 console_log_store: Arc::new(console_log_store),
441 metadata_store: Arc::new(metadata_store),
442 blob_store,
443 agent_memory_provider: Some(Arc::new(agent_memory_provider)),
444 schedule_store: Arc::new(schedule_store),
445 durability: vec![
446 DurabilityDeclaration::durable("continuity", DurabilityResolution::Persistent),
447 DurabilityDeclaration::durable(
448 "event_log",
449 DurabilityResolution::DeclaredEphemeral,
450 ),
451 DurabilityDeclaration::durable("console", DurabilityResolution::Persistent),
452 DurabilityDeclaration::durable("metadata", DurabilityResolution::Persistent),
453 DurabilityDeclaration::durable("blobs", blob_resolution),
454 DurabilityDeclaration::durable("agent_memory", DurabilityResolution::Persistent),
455 DurabilityDeclaration::durable("schedule", DurabilityResolution::Persistent),
456 ],
457 };
458 enforce_fail_closed_store_set(&set, ctx)?;
459 Ok(set)
460 }
461
462 fn meerkat_provider(&self) -> &dyn RealmStorageProvider {
463 static DISK: meerkat::storage_provider::DiskStorageProvider =
464 meerkat::storage_provider::DiskStorageProvider;
465 &DISK
466 }
467
468 fn migrator(&self) -> Option<&dyn meerkat_core::StorageMigrator> {
469 static MIGRATOR: MobKitStorageMigrator = MobKitStorageMigrator;
470 Some(&MIGRATOR)
471 }
472}
473
474#[cfg(test)]
475#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
476mod tests {
477 use super::*;
478 use meerkat_core::DurabilityClass;
479
480 fn context(dir: &std::path::Path) -> MobKitRealmOpenContext {
481 MobKitRealmOpenContext::for_state_dir(dir.join("state"))
482 }
483
484 #[tokio::test]
488 async fn disk_provider_opens_the_canonical_layout_fail_closed() {
489 let dir = tempfile::tempdir().expect("tempdir");
490 let provider = DiskMobKitStorageProvider;
491 assert_eq!(provider.name(), "disk");
492 let set = provider
493 .open_realm(&context(dir.path()))
494 .await
495 .expect("disk realm opens");
496
497 assert!(matches!(
498 set.lease_authority,
499 MobKitLeaseAuthority::FencingFloor(0)
500 ));
501 assert!(set.event_log_store.is_none());
502 assert!(set.agent_memory_provider.is_some());
503 assert!(set.blob_store.is_persistent());
504 for declaration in &set.durability {
505 assert_eq!(declaration.class, DurabilityClass::Durable);
506 match declaration.domain.as_str() {
507 "event_log" => assert_eq!(
508 declaration.resolution,
509 DurabilityResolution::DeclaredEphemeral
510 ),
511 _ => assert_eq!(declaration.resolution, DurabilityResolution::Persistent),
512 }
513 }
514 assert!(provider.migrator().is_some());
515 assert_eq!(provider.meerkat_provider().name(), "disk");
516 }
517
518 #[tokio::test]
521 async fn disk_provider_honors_declared_ephemeral_blobs() {
522 let dir = tempfile::tempdir().expect("tempdir");
523 let mut ctx = context(dir.path());
524 ctx.declared_ephemeral_domains.push("blobs".to_string());
525 let set = DiskMobKitStorageProvider
526 .open_realm(&ctx)
527 .await
528 .expect("declared-ephemeral realm opens");
529 assert!(!set.blob_store.is_persistent());
530 let blob_declaration = set
531 .durability
532 .iter()
533 .find(|declaration| declaration.domain == "blobs")
534 .expect("blobs declaration");
535 assert_eq!(
536 blob_declaration.resolution,
537 DurabilityResolution::DeclaredEphemeral
538 );
539 }
540
541 #[tokio::test]
544 async fn enforce_refuses_undeclared_nonpersistent_durable_slots() {
545 let dir = tempfile::tempdir().expect("tempdir");
546 let ctx = context(dir.path());
547 let mut set = DiskMobKitStorageProvider
548 .open_realm(&ctx)
549 .await
550 .expect("disk realm opens");
551 for declaration in &mut set.durability {
552 if declaration.domain == "continuity" {
553 declaration.resolution = DurabilityResolution::NonPersistent;
554 }
555 }
556 let error = enforce_fail_closed_store_set(&set, &ctx)
557 .expect_err("undeclared non-persistent durable slot must refuse");
558 assert!(matches!(
559 error,
560 MobKitStorageProviderError::DurabilityViolation { ref domain } if domain == "continuity"
561 ));
562 assert!(error.to_string().contains("fail-closed"));
563
564 let mut declared = ctx.clone();
566 declared
567 .declared_ephemeral_domains
568 .push("continuity".to_string());
569 enforce_fail_closed_store_set(&set, &declared)
570 .expect("declared ephemeral domain must compose");
571 }
572
573 #[tokio::test]
576 async fn enforce_refuses_omitted_and_duplicate_declarations() {
577 let dir = tempfile::tempdir().expect("tempdir");
578 let ctx = context(dir.path());
579 let set = DiskMobKitStorageProvider
580 .open_realm(&ctx)
581 .await
582 .expect("disk realm opens");
583
584 let mut omitted = set.durability.clone();
585 omitted.retain(|declaration| declaration.domain != "metadata");
586 let mut incomplete = MobKitRealmStoreSet {
587 durability: omitted,
588 ..clone_stores(&set)
589 };
590 let error = enforce_fail_closed_store_set(&incomplete, &ctx)
591 .expect_err("an omitted durability declaration must refuse");
592 assert!(matches!(
593 error,
594 MobKitStorageProviderError::DurabilityViolation { ref domain }
595 if domain.starts_with("metadata") && domain.contains("0 durability declarations")
596 ));
597
598 incomplete.durability = set.durability;
599 incomplete.durability.push(DurabilityDeclaration::durable(
600 "metadata",
601 DurabilityResolution::Persistent,
602 ));
603 let error = enforce_fail_closed_store_set(&incomplete, &ctx)
604 .expect_err("a duplicated durability declaration must refuse");
605 assert!(matches!(
606 error,
607 MobKitStorageProviderError::DurabilityViolation { ref domain }
608 if domain.starts_with("metadata") && domain.contains("2 durability declarations")
609 ));
610 }
611
612 fn clone_stores(set: &MobKitRealmStoreSet) -> MobKitRealmStoreSet {
613 MobKitRealmStoreSet {
614 continuity_store: Arc::clone(&set.continuity_store),
615 lease_authority: set.lease_authority.clone(),
616 event_log_store: None,
617 console_log_store: Arc::clone(&set.console_log_store),
618 metadata_store: Arc::clone(&set.metadata_store),
619 blob_store: Arc::clone(&set.blob_store),
620 agent_memory_provider: set.agent_memory_provider.clone(),
621 schedule_store: Arc::clone(&set.schedule_store),
622 durability: set.durability.clone(),
623 }
624 }
625
626 struct StubMeerkatProvider {
630 opened: std::sync::atomic::AtomicBool,
631 runtime_resolution: DurabilityResolution,
632 }
633
634 impl StubMeerkatProvider {
635 fn declared_ephemeral() -> Self {
636 Self {
637 opened: std::sync::atomic::AtomicBool::new(false),
638 runtime_resolution: DurabilityResolution::DeclaredEphemeral,
639 }
640 }
641 }
642
643 #[async_trait]
644 impl RealmStorageProvider for StubMeerkatProvider {
645 fn name(&self) -> &'static str {
646 "stub-remote"
647 }
648
649 async fn open(
650 &self,
651 ctx: &meerkat::storage_provider::RealmOpenContext,
652 ) -> Result<meerkat::storage_provider::RealmStoreSet, meerkat::PersistenceError> {
653 self.opened.store(true, std::sync::atomic::Ordering::SeqCst);
654 Ok(meerkat::storage_provider::RealmStoreSet {
655 session_store: Arc::new(meerkat_store::MemoryStore::new()),
656 runtime_store: Arc::new(meerkat_runtime::InMemoryRuntimeStore::new()),
657 job_store: Arc::new(meerkat::MemoryDetachedJobStore::new()),
659 schedule_store: Arc::new(meerkat_schedule::MemoryScheduleStore::new()),
660 workgraph_store: Arc::new(meerkat::MemoryWorkGraphStore::new()),
661 blob_store: Arc::new(meerkat_store::MemoryBlobStore::new()),
662 artifact_store: Arc::new(meerkat_store::MemoryArtifactStore::new()),
663 store_path: ctx.paths.root.clone(),
664 projection_root: None,
665 durability: [
666 "sessions",
667 "jobs",
668 "schedule",
669 "workgraph",
670 "blobs",
671 "artifacts",
672 ]
673 .iter()
674 .map(|domain| {
675 DurabilityDeclaration::durable(domain, DurabilityResolution::DeclaredEphemeral)
676 })
677 .chain(std::iter::once(DurabilityDeclaration::durable(
678 "runtime",
679 self.runtime_resolution,
680 )))
681 .collect(),
682 })
683 }
684 }
685
686 #[tokio::test]
691 async fn meerkat_level_open_routes_the_provider_bundle_and_pins_external() {
692 let dir = tempfile::tempdir().expect("tempdir");
693 let ctx = context(dir.path());
694 let provider = StubMeerkatProvider::declared_ephemeral();
695
696 let stores = open_provider_meerkat_stores(&provider, &ctx)
697 .await
698 .expect("stub meerkat bundle opens");
699
700 assert!(provider.opened.load(std::sync::atomic::Ordering::SeqCst));
701 assert_eq!(stores.provider_name, "stub-remote");
702 assert_eq!(
703 stores.runtime_declaration.resolution,
704 DurabilityResolution::DeclaredEphemeral
705 );
706 assert_eq!(
707 stores.workgraph_declaration.resolution,
708 DurabilityResolution::DeclaredEphemeral
709 );
710 let runtime_slot = stores.runtime_slot_summary();
711 assert_eq!(runtime_slot.declaration.domain, "runtime");
712 assert!(runtime_slot.backend.contains("stub-remote"));
713 let manifest_path = ctx
714 .state_dir
715 .join(MEERKAT_LEVEL_REALM_ID)
716 .join("realm_manifest.json");
717 let manifest = std::fs::read_to_string(&manifest_path)
718 .expect("meerkat-level realm manifest must be pinned");
719 assert!(
720 manifest.contains("stub-remote"),
721 "the pin must carry the provider name, got: {manifest}"
722 );
723 }
724
725 #[tokio::test]
729 async fn meerkat_level_open_refuses_undeclared_nonpersistent_durables() {
730 let dir = tempfile::tempdir().expect("tempdir");
731 let ctx = context(dir.path());
732 let provider = StubMeerkatProvider {
733 opened: std::sync::atomic::AtomicBool::new(false),
734 runtime_resolution: DurabilityResolution::NonPersistent,
735 };
736
737 let error = match open_provider_meerkat_stores(&provider, &ctx).await {
738 Err(error) => error,
739 Ok(_) => panic!("an undeclared non-persistent runtime slot must refuse"),
740 };
741 assert!(matches!(
742 error,
743 MobKitStorageProviderError::DurabilityViolation { ref domain }
744 if domain.contains("runtime") && domain.contains("meerkat level")
745 ));
746
747 let mut declared = ctx.clone();
750 declared
751 .declared_ephemeral_domains
752 .push("runtime".to_string());
753 let provider = StubMeerkatProvider {
754 opened: std::sync::atomic::AtomicBool::new(false),
755 runtime_resolution: DurabilityResolution::NonPersistent,
756 };
757 open_provider_meerkat_stores(&provider, &declared)
758 .await
759 .expect("a declared-ephemeral runtime domain must compose");
760 }
761}