kmp_adapter_embedded/adapter/
migration.rs1use std::fs;
10use std::path::Path;
11
12use kmp_domain::{ContextUpdatedEvent, PortError, ProjectionMutation};
13use serde::{Deserialize, Serialize};
14
15use super::engine::{Key, Table};
16use super::format_version::{self, StorageEngine};
17use super::store::EmbeddedKernelStore;
18
19#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
21pub struct StoreMigrationReceipt {
22 pub source_format: u32,
23 pub source_sha256: String,
24 pub destination_format: u32,
25 pub events_migrated: u64,
26 pub mutations_applied: u64,
27 pub kernel_version: String,
28}
29
30impl StoreMigrationReceipt {
31 pub const MIGRATION_ID: &'static str = "store-format-migration";
33}
34
35impl EmbeddedKernelStore {
36 pub async fn migrate_data_dir<F>(
38 source_dir: &Path,
39 destination_dir: &Path,
40 derive: F,
41 ) -> Result<(Self, StoreMigrationReceipt), PortError>
42 where
43 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
44 {
45 Self::migrate_data_dir_to(source_dir, destination_dir, StorageEngine::Sqlite, derive).await
46 }
47
48 pub async fn migrate_data_dir_to<F>(
50 source_dir: &Path,
51 destination_dir: &Path,
52 destination_engine: StorageEngine,
53 _derive: F,
54 ) -> Result<(Self, StoreMigrationReceipt), PortError>
55 where
56 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
57 {
58 if same_file(source_dir, destination_dir) {
59 return Err(PortError::InvalidState(
60 "migration source and destination are the same data directory".to_string(),
61 ));
62 }
63 let source_format = format_version::read_stamped_version(source_dir)?;
64 if source_format > StorageEngine::NEWEST_KNOWN_FORMAT_VERSION {
65 return Err(PortError::InvalidState(format!(
66 "migration source `{}` uses format version {source_format}, newer than this \
67 binary supports ({}); upgrade the binary",
68 source_dir.display(),
69 StorageEngine::NEWEST_KNOWN_FORMAT_VERSION
70 )));
71 }
72 if source_format == StorageEngine::Sqlite.format_version() {
73 return Err(PortError::Unavailable(format!(
74 "migration from a SQLite format-2 store is unnecessary and unsupported; the \
75 source at `{}` is left untouched",
76 source_dir.display()
77 )));
78 }
79 let _ = destination_engine;
80 Err(PortError::InvalidState(format!(
81 "migration source `{}` uses unsupported format version {source_format}; current \
82 KMP left it untouched. Preserve the source, use an explicitly archived compatible \
83 exporter to create `.kmp/memory.jsonl`, then import that bundle into an empty \
84 current store",
85 source_dir.display()
86 )))
87 }
88
89 pub async fn open_or_migrate_data_dir<F>(
91 source_dir: &Path,
92 destination_dir: &Path,
93 derive: F,
94 ) -> Result<(Self, Option<StoreMigrationReceipt>), PortError>
95 where
96 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
97 {
98 Self::open_or_migrate_data_dir_to(
99 source_dir,
100 destination_dir,
101 StorageEngine::Sqlite,
102 derive,
103 )
104 .await
105 }
106
107 pub async fn open_or_migrate_data_dir_to<F>(
109 source_dir: &Path,
110 destination_dir: &Path,
111 destination_engine: StorageEngine,
112 derive: F,
113 ) -> Result<(Self, Option<StoreMigrationReceipt>), PortError>
114 where
115 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
116 {
117 if format_version::existing_store_file(destination_dir).is_some() {
118 let store = Self::open(destination_dir)?;
119 let receipt = store.migration_receipt().await?;
120 return Ok((store, receipt));
121 }
122 let (store, receipt) =
123 Self::migrate_data_dir_to(source_dir, destination_dir, destination_engine, derive)
124 .await?;
125 Ok((store, Some(receipt)))
126 }
127
128 pub async fn migration_receipt(&self) -> Result<Option<StoreMigrationReceipt>, PortError> {
130 self.run(|store| {
131 let tx = store.begin_read()?;
132 let Some(raw) = tx.get(
133 Table::Migrations,
134 Key::Str(StoreMigrationReceipt::MIGRATION_ID),
135 )?
136 else {
137 return Ok(None);
138 };
139 let receipt = serde_json::from_slice(&raw).map_err(|error| {
140 PortError::InvalidState(format!("migration receipt is unreadable: {error}"))
141 })?;
142 Ok(Some(receipt))
143 })
144 .await
145 }
146}
147
148fn same_file(left: &Path, right: &Path) -> bool {
149 match (fs::canonicalize(left), fs::canonicalize(right)) {
150 (Ok(left), Ok(right)) => left == right,
151 _ => left == right,
152 }
153}