1use std::collections::{BTreeMap, BTreeSet};
8use std::time::{Duration, SystemTime, UNIX_EPOCH};
9
10use kmp_domain::{ContextEventStore, ContextUpdatedEvent, PortError, ProjectionMutation};
11use serde::{Deserialize, Serialize};
12use sha2::{Digest, Sha256};
13
14use super::replay::ProjectionRebuildReport;
15use super::store::EmbeddedKernelStore;
16
17#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
20pub struct BundleEventRange {
21 pub first: Option<u64>,
22 pub last: Option<u64>,
23}
24
25#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
29pub struct BundleHeader {
30 pub bundle_format: u32,
31 #[serde(rename = "event_format", alias = "store_format")]
35 pub event_format: u32,
36 pub event_count: u64,
37 pub kernel_version: String,
38 #[serde(default)]
39 pub snapshot_id: String,
40 #[serde(default)]
41 pub created_at_unix_ms: u64,
42 #[serde(default)]
43 pub event_range: BundleEventRange,
44 #[serde(default)]
45 pub abouts: Vec<String>,
46 #[serde(default)]
47 pub content_digest: String,
48}
49
50pub const BUNDLE_FORMAT_VERSION: u32 = 2;
51
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub struct ImportReport {
55 pub events_imported: u64,
56 pub rebuild: ProjectionRebuildReport,
57}
58
59impl EmbeddedKernelStore {
60 pub async fn export_bundle(&self) -> Result<String, PortError> {
63 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
64 encode_bundle(&events, None)
65 }
66
67 pub async fn export_named_bundle(&self, snapshot_id: &str) -> Result<String, PortError> {
71 if snapshot_id.trim().is_empty() {
72 return Err(PortError::InvalidState(
73 "snapshot id must not be empty".to_string(),
74 ));
75 }
76 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
77 encode_bundle(&events, Some(snapshot_id))
78 }
79
80 pub async fn import_bundle<F>(&self, bundle: &str, derive: F) -> Result<ImportReport, PortError>
85 where
86 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
87 {
88 let (log_length, _) = self.event_log_stats().await?;
89 if log_length != 0 {
90 return Err(PortError::Conflict(format!(
91 "import requires an empty store; this store already holds {log_length} events \
92 (merging bundles is not supported)"
93 )));
94 }
95
96 let verified = parse_bundle(bundle)?;
97 let header = verified.header;
98 let events = verified.events;
99 validate_revisions(&events)?;
100 for event in &events {
104 derive(event)?;
105 }
106 let events_imported = self.replay_event_stream(events).await?;
107 debug_assert_eq!(events_imported, header.event_count);
108
109 let rebuild = self.rebuild_projections(derive).await?;
110 Ok(ImportReport {
111 events_imported,
112 rebuild,
113 })
114 }
115}
116
117pub fn verify_bundle(bundle: &str) -> Result<BundleHeader, PortError> {
121 parse_bundle(bundle).map(|verified| verified.header)
122}
123
124pub fn merge_bundles(left: &str, right: &str, snapshot_id: &str) -> Result<String, PortError> {
128 if snapshot_id.trim().is_empty() {
129 return Err(PortError::InvalidState(
130 "merged snapshot id must not be empty".to_string(),
131 ));
132 }
133 let left = parse_bundle(left)?;
134 let right = parse_bundle(right)?;
135 let shared = left.events.len().min(right.events.len());
136 if let Some(position) =
137 (0..shared).find(|position| left.events[*position] != right.events[*position])
138 {
139 return Err(PortError::Conflict(format!(
140 "bundle histories diverge at event position {}; KMP only fast-forwards an exact \
141 prefix and will not invent causal order for two branches",
142 position + 1
143 )));
144 }
145 let events = if left.events.len() >= right.events.len() {
146 left.events
147 } else {
148 right.events
149 };
150 encode_bundle(&events, Some(snapshot_id))
151}
152
153struct VerifiedBundle {
154 header: BundleHeader,
155 events: Vec<ContextUpdatedEvent>,
156}
157
158fn parse_bundle(bundle: &str) -> Result<VerifiedBundle, PortError> {
159 let mut lines = bundle.lines().filter(|line| !line.trim().is_empty());
160 let header: BundleHeader = decode_line(
161 "bundle header",
162 lines.next().ok_or_else(|| {
163 PortError::InvalidState("bundle is empty: missing header line".to_string())
164 })?,
165 )?;
166 if !matches!(header.bundle_format, 1 | BUNDLE_FORMAT_VERSION) {
167 return Err(PortError::InvalidState(format!(
168 "bundle format {} is not supported (this binary reads 1 and {})",
169 header.bundle_format, BUNDLE_FORMAT_VERSION
170 )));
171 }
172 if header.event_format != super::format_version::EVENT_FORMAT_VERSION {
173 return Err(PortError::InvalidState(format!(
174 "bundle carries event format {}, this binary supports {}",
175 header.event_format,
176 super::format_version::EVENT_FORMAT_VERSION
177 )));
178 }
179
180 let mut events = Vec::new();
181 let mut event_payload = String::new();
182 for line in lines {
183 events.push(decode_line::<ContextUpdatedEvent>("bundle event", line)?);
184 event_payload.push_str(line);
185 event_payload.push('\n');
186 }
187 if events.len() as u64 != header.event_count {
188 return Err(PortError::InvalidState(format!(
189 "bundle header declares {} events but {} were present",
190 header.event_count,
191 events.len()
192 )));
193 }
194
195 if header.bundle_format == BUNDLE_FORMAT_VERSION {
196 validate_v2_header(&header, &events, &event_payload)?;
197 }
198 Ok(VerifiedBundle { header, events })
199}
200
201fn validate_v2_header(
202 header: &BundleHeader,
203 events: &[ContextUpdatedEvent],
204 event_payload: &str,
205) -> Result<(), PortError> {
206 if header.snapshot_id.trim().is_empty() {
207 return Err(PortError::InvalidState(
208 "bundle format 2 requires snapshot_id".to_string(),
209 ));
210 }
211 if header.created_at_unix_ms == 0 {
212 return Err(PortError::InvalidState(
213 "bundle format 2 requires created_at_unix_ms".to_string(),
214 ));
215 }
216 let expected_range = event_range(events.len());
217 if header.event_range != expected_range {
218 return Err(PortError::InvalidState(format!(
219 "bundle event_range {:?} does not cover its {} events (expected {:?})",
220 header.event_range,
221 events.len(),
222 expected_range
223 )));
224 }
225 let expected_abouts = abouts(events);
226 if header.abouts != expected_abouts {
227 return Err(PortError::InvalidState(format!(
228 "bundle abouts do not match its events (expected {})",
229 expected_abouts.join(", ")
230 )));
231 }
232 let expected_digest = content_digest(event_payload.as_bytes());
233 if header.content_digest != expected_digest {
234 return Err(PortError::InvalidState(format!(
235 "bundle content digest mismatch: header says {}, events produce {expected_digest}",
236 header.content_digest
237 )));
238 }
239 Ok(())
240}
241
242fn encode_bundle(
243 events: &[ContextUpdatedEvent],
244 snapshot_id: Option<&str>,
245) -> Result<String, PortError> {
246 let mut event_payload = String::new();
247 for event in events {
248 event_payload.push_str(&encode_line("bundle event", event)?);
249 }
250 let digest = content_digest(event_payload.as_bytes());
251 let named = snapshot_id.is_some();
252 let snapshot_id = snapshot_id
253 .map(str::to_string)
254 .unwrap_or_else(|| format!("content-{}", &digest[7..23]));
255 let created_at = if named {
260 SystemTime::now()
261 } else {
262 events
263 .iter()
264 .map(|event| event.occurred_at)
265 .max()
266 .unwrap_or(UNIX_EPOCH + Duration::from_millis(1))
267 };
268 let created_at_unix_ms = created_at
269 .duration_since(UNIX_EPOCH)
270 .unwrap_or(Duration::ZERO)
271 .as_millis() as u64;
272 let header = BundleHeader {
273 bundle_format: BUNDLE_FORMAT_VERSION,
274 event_format: super::format_version::EVENT_FORMAT_VERSION,
275 event_count: events.len() as u64,
276 kernel_version: env!("CARGO_PKG_VERSION").to_string(),
277 snapshot_id,
278 created_at_unix_ms,
279 event_range: event_range(events.len()),
280 abouts: abouts(events),
281 content_digest: digest,
282 };
283 let mut out = encode_line("bundle header", &header)?;
284 out.push_str(&event_payload);
285 Ok(out)
286}
287
288fn event_range(event_count: usize) -> BundleEventRange {
289 if event_count == 0 {
290 BundleEventRange::default()
291 } else {
292 BundleEventRange {
293 first: Some(1),
294 last: Some(event_count as u64),
295 }
296 }
297}
298
299fn abouts(events: &[ContextUpdatedEvent]) -> Vec<String> {
300 events
301 .iter()
302 .map(|event| event.root_node_id.clone())
303 .collect::<BTreeSet<_>>()
304 .into_iter()
305 .collect()
306}
307
308fn content_digest(bytes: &[u8]) -> String {
309 format!("sha256:{:x}", Sha256::digest(bytes))
310}
311
312fn validate_revisions(events: &[ContextUpdatedEvent]) -> Result<(), PortError> {
313 let mut revisions: BTreeMap<(&str, &str), u64> = BTreeMap::new();
314 for (position, event) in events.iter().enumerate() {
315 let previous = revisions
316 .get(&(event.root_node_id.as_str(), event.role.as_str()))
317 .copied()
318 .unwrap_or(0);
319 let expected = previous + 1;
320 if event.revision != expected {
321 return Err(PortError::InvalidState(format!(
322 "bundle event position {} carries revision {} for ({}, {}), expected {}; no \
323 events were imported",
324 position + 1,
325 event.revision,
326 event.root_node_id,
327 event.role,
328 expected
329 )));
330 }
331 revisions.insert(
332 (event.root_node_id.as_str(), event.role.as_str()),
333 event.revision,
334 );
335 }
336 Ok(())
337}
338
339impl EmbeddedKernelStore {
340 pub(crate) async fn replay_event_stream<I>(&self, events: I) -> Result<u64, PortError>
348 where
349 I: IntoIterator<Item = ContextUpdatedEvent>,
350 {
351 let mut replayed = 0u64;
352 for event in events {
353 let recorded_revision = event.revision;
354 let expected_previous = recorded_revision.checked_sub(1).ok_or_else(|| {
355 PortError::InvalidState("event carries revision 0; the log is corrupt".to_string())
356 })?;
357 let assigned = self.append(event, expected_previous).await?;
358 if assigned != recorded_revision {
359 return Err(PortError::Conflict(format!(
360 "replay integrity violation: assigned revision {assigned}, \
361 history recorded {recorded_revision}"
362 )));
363 }
364 replayed += 1;
365 }
366 Ok(replayed)
367 }
368}
369
370fn encode_line<T: Serialize>(what: &str, value: &T) -> Result<String, PortError> {
371 let mut line = serde_json::to_string(value)
372 .map_err(|error| PortError::InvalidState(format!("could not encode {what}: {error}")))?;
373 line.push('\n');
374 Ok(line)
375}
376
377fn decode_line<T: for<'de> Deserialize<'de>>(what: &str, line: &str) -> Result<T, PortError> {
378 serde_json::from_str(line)
379 .map_err(|error| PortError::InvalidState(format!("could not decode {what}: {error}")))
380}
381
382#[cfg(test)]
383mod tests {
384 use super::*;
385
386 fn event(root: &str, revision: u64, content_hash: &str) -> ContextUpdatedEvent {
387 ContextUpdatedEvent {
388 root_node_id: root.to_string(),
389 role: "agent".to_string(),
390 revision,
391 content_hash: content_hash.to_string(),
392 changes: Vec::new(),
393 idempotency_key: Some(format!("{root}:{revision}")),
394 logical_digest: None,
395 requested_by: Some("portability-test".to_string()),
396 occurred_at: UNIX_EPOCH + Duration::from_secs(revision),
397 }
398 }
399
400 #[test]
401 fn format_two_identifies_and_covers_the_snapshot() {
402 let events = vec![event("project:b", 1, "b"), event("project:a", 1, "a")];
403 let bundle = encode_bundle(&events, Some("pre-release")).expect("bundle");
404 let header = verify_bundle(&bundle).expect("verified");
405
406 assert_eq!(header.bundle_format, BUNDLE_FORMAT_VERSION);
407 assert_eq!(
408 header.event_format,
409 super::super::format_version::EVENT_FORMAT_VERSION
410 );
411 assert_eq!(header.snapshot_id, "pre-release");
412 assert!(header.created_at_unix_ms > 0);
413 assert_eq!(
414 header.event_range,
415 BundleEventRange {
416 first: Some(1),
417 last: Some(2),
418 }
419 );
420 assert_eq!(header.abouts, ["project:a", "project:b"]);
421 assert!(header.content_digest.starts_with("sha256:"));
422 }
423
424 #[test]
425 fn tampering_is_rejected_before_a_bundle_can_be_replayed() {
426 let bundle =
427 encode_bundle(&[event("project:a", 1, "before")], Some("saved")).expect("bundle");
428 let tampered = bundle.replace("\"content_hash\":\"before\"", "\"content_hash\":\"after\"");
429 let error = verify_bundle(&tampered).expect_err("digest catches changed payload");
430 assert!(error.to_string().contains("content digest mismatch"));
431 }
432
433 #[test]
434 fn merge_fast_forwards_an_exact_prefix() {
435 let first = event("project:a", 1, "one");
436 let second = event("project:a", 2, "two");
437 let left = encode_bundle(std::slice::from_ref(&first), Some("left")).expect("left");
438 let right = encode_bundle(&[first, second], Some("right")).expect("right");
439
440 let merged = merge_bundles(&left, &right, "merged").expect("fast forward");
441 let header = verify_bundle(&merged).expect("verified merge");
442 assert_eq!(header.snapshot_id, "merged");
443 assert_eq!(header.event_count, 2);
444 }
445
446 #[test]
447 fn merge_refuses_two_histories_at_the_same_position() {
448 let left = encode_bundle(&[event("project:a", 1, "left")], Some("left")).expect("left");
449 let right = encode_bundle(&[event("project:a", 1, "right")], Some("right")).expect("right");
450
451 let error = merge_bundles(&left, &right, "invented").expect_err("must refuse");
452 assert!(error.to_string().contains("diverge at event position 1"));
453 assert!(error.to_string().contains("will not invent causal order"));
454 }
455
456 #[test]
457 fn legacy_format_one_remains_readable() {
458 let legacy =
459 r#"{"bundle_format":1,"store_format":1,"event_count":0,"kernel_version":"0.1.3"}"#;
460 let header = verify_bundle(legacy).expect("format one remains portable");
461 assert_eq!(header.event_format, 1);
462 assert!(header.snapshot_id.is_empty());
463 }
464
465 #[test]
466 fn invalid_later_revision_is_rejected_in_preflight() {
467 let events = [event("project:a", 1, "one"), event("project:a", 3, "three")];
468 let error = validate_revisions(&events).expect_err("revision gap");
469 assert!(error.to_string().contains("position 2"));
470 assert!(error.to_string().contains("no events were imported"));
471 }
472}