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)]
23pub struct BundleEventRange {
24 pub first: Option<u64>,
25 pub last: Option<u64>,
26}
27
28#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
31pub struct BundleHeader {
32 pub bundle_format: u32,
33 pub event_format: u32,
35 pub event_count: u64,
36 pub kernel_version: String,
37 #[serde(default)]
38 pub snapshot_id: String,
39 #[serde(default)]
40 pub created_at_unix_ms: u64,
41 #[serde(default)]
42 pub event_range: BundleEventRange,
43 #[serde(default)]
44 pub abouts: Vec<String>,
45 #[serde(default)]
46 pub content_digest: String,
47}
48
49pub const BUNDLE_FORMAT_VERSION: u32 = 3;
50
51#[derive(Debug, Clone, Copy, PartialEq, Eq)]
53pub struct ImportReport {
54 pub events_imported: u64,
55 pub rebuild: ProjectionRebuildReport,
56}
57
58impl EmbeddedKernelStore {
59 pub fn export_bundle_blocking(&self) -> Result<String, PortError> {
62 encode_bundle(&self.read_event_log()?, None)
63 }
64
65 pub async fn export_bundle(&self) -> Result<String, PortError> {
68 self.run(|store| store.export_bundle_blocking()).await
69 }
70
71 pub async fn export_bundle_for_abouts(
80 &self,
81 requested_abouts: &[String],
82 ) -> Result<String, PortError> {
83 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
84 let events = filter_events_for_abouts(events, requested_abouts)?;
85 encode_bundle(&events, None)
86 }
87
88 pub fn export_bundle_excluding_abouts_blocking(
102 &self,
103 excluded_abouts: &[String],
104 ) -> Result<String, PortError> {
105 let events = self.read_event_log()?;
106 encode_bundle(
107 &filter_events_excluding_abouts(events, excluded_abouts),
108 None,
109 )
110 }
111
112 pub async fn export_bundle_excluding_abouts(
114 &self,
115 excluded_abouts: &[String],
116 ) -> Result<String, PortError> {
117 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
118 encode_bundle(
119 &filter_events_excluding_abouts(events, excluded_abouts),
120 None,
121 )
122 }
123
124 pub async fn export_named_bundle(&self, snapshot_id: &str) -> Result<String, PortError> {
128 if snapshot_id.trim().is_empty() {
129 return Err(PortError::InvalidState(
130 "snapshot id must not be empty".to_string(),
131 ));
132 }
133 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
134 encode_bundle(&events, Some(snapshot_id))
135 }
136
137 pub async fn import_bundle<F>(&self, bundle: &str, derive: F) -> Result<ImportReport, PortError>
142 where
143 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
144 {
145 let (log_length, _) = self.event_log_stats().await?;
146 if log_length != 0 {
147 return Err(PortError::Conflict(format!(
148 "import requires an empty store; this store already holds {log_length} events \
149 (merging bundles is not supported)"
150 )));
151 }
152
153 let verified = parse_bundle(bundle)?;
154 let header = verified.header;
155 let events = verified.events;
156 validate_revisions(&events)?;
157 for event in &events {
161 derive(event)?;
162 }
163 let events_imported = self.replay_event_stream(events).await?;
164 debug_assert_eq!(events_imported, header.event_count);
165
166 let rebuild = self.rebuild_projections(derive).await?;
167 Ok(ImportReport {
168 events_imported,
169 rebuild,
170 })
171 }
172}
173
174pub fn verify_bundle(bundle: &str) -> Result<BundleHeader, PortError> {
178 parse_bundle(bundle).map(|verified| verified.header)
179}
180
181pub fn bundle_excluding_abouts(
189 bundle: &str,
190 excluded_abouts: &[String],
191) -> Result<String, PortError> {
192 let verified = parse_bundle(bundle)?;
193 encode_bundle(
194 &filter_events_excluding_abouts(verified.events, excluded_abouts),
195 None,
196 )
197}
198
199pub fn merge_bundles(left: &str, right: &str, snapshot_id: &str) -> Result<String, PortError> {
203 if snapshot_id.trim().is_empty() {
204 return Err(PortError::InvalidState(
205 "merged snapshot id must not be empty".to_string(),
206 ));
207 }
208 let left = parse_bundle(left)?;
209 let right = parse_bundle(right)?;
210 let shared = left.events.len().min(right.events.len());
211 if let Some(position) =
212 (0..shared).find(|position| left.events[*position] != right.events[*position])
213 {
214 return Err(PortError::Conflict(format!(
215 "bundle histories diverge at event position {}; KMP only fast-forwards an exact \
216 prefix and will not invent causal order for two branches",
217 position + 1
218 )));
219 }
220 let events = if left.events.len() >= right.events.len() {
221 left.events
222 } else {
223 right.events
224 };
225 encode_bundle(&events, Some(snapshot_id))
226}
227
228struct VerifiedBundle {
229 header: BundleHeader,
230 events: Vec<ContextUpdatedEvent>,
231}
232
233fn parse_bundle(bundle: &str) -> Result<VerifiedBundle, PortError> {
234 let mut lines = bundle.lines().filter(|line| !line.trim().is_empty());
235 let header: BundleHeader = decode_line(
236 "bundle header",
237 lines.next().ok_or_else(|| {
238 PortError::InvalidState("bundle is empty: missing header line".to_string())
239 })?,
240 )?;
241 if header.bundle_format != BUNDLE_FORMAT_VERSION {
242 return Err(PortError::InvalidState(format!(
243 "bundle format {} is not supported (this binary reads {})",
244 header.bundle_format, BUNDLE_FORMAT_VERSION
245 )));
246 }
247 if !matches!(
248 header.event_format,
249 2 | super::format_version::EVENT_FORMAT_VERSION
250 ) {
251 return Err(PortError::InvalidState(format!(
252 "bundle carries event format {}, this binary supports {}",
253 header.event_format,
254 super::format_version::EVENT_FORMAT_VERSION
255 )));
256 }
257
258 let mut events = Vec::new();
259 let mut event_payload = String::new();
260 for line in lines {
261 events.push(decode_line::<ContextUpdatedEvent>("bundle event", line)?);
262 event_payload.push_str(line);
263 event_payload.push('\n');
264 }
265 if events.len() as u64 != header.event_count {
266 return Err(PortError::InvalidState(format!(
267 "bundle header declares {} events but {} were present",
268 header.event_count,
269 events.len()
270 )));
271 }
272
273 if header.event_format < 3
274 && events
275 .iter()
276 .any(|e| e.role == kmp_domain::NodeCardEvent::ROLE)
277 {
278 return Err(PortError::InvalidState(
279 "card history requires event format 3".into(),
280 ));
281 }
282 let mut card_heads: BTreeMap<(String, String, String), u64> = BTreeMap::new();
283 for event in &events {
284 if let Some(card) = kmp_domain::NodeCardEvent::card(event)? {
285 let identity = (event.root_node_id.clone(), card.node_id, card.language);
286 let previous = card_heads.get(&identity).copied();
287 let baseline = event.changes[0].operation == "BASELINE";
288 let valid = match previous {
289 None => baseline || card.card_revision == 1,
290 Some(revision) => !baseline && revision.checked_add(1) == Some(card.card_revision),
291 };
292 if !valid {
293 return Err(PortError::InvalidState(
294 "non-contiguous card history".into(),
295 ));
296 }
297 card_heads.insert(identity, card.card_revision);
298 }
299 }
300 validate_header(&header, &events, &event_payload)?;
301 Ok(VerifiedBundle { header, events })
302}
303
304fn validate_header(
305 header: &BundleHeader,
306 events: &[ContextUpdatedEvent],
307 event_payload: &str,
308) -> Result<(), PortError> {
309 if header.snapshot_id.trim().is_empty() {
310 return Err(PortError::InvalidState(
311 "bundle format 2 requires snapshot_id".to_string(),
312 ));
313 }
314 if header.created_at_unix_ms == 0 {
315 return Err(PortError::InvalidState(
316 "bundle format 2 requires created_at_unix_ms".to_string(),
317 ));
318 }
319 let expected_range = event_range(events.len());
320 if header.event_range != expected_range {
321 return Err(PortError::InvalidState(format!(
322 "bundle event_range {:?} does not cover its {} events (expected {:?})",
323 header.event_range,
324 events.len(),
325 expected_range
326 )));
327 }
328 let expected_abouts = abouts(events);
329 if header.abouts != expected_abouts {
330 return Err(PortError::InvalidState(format!(
331 "bundle abouts do not match its events (expected {})",
332 expected_abouts.join(", ")
333 )));
334 }
335 let expected_digest = content_digest(event_payload.as_bytes());
336 if header.content_digest != expected_digest {
337 return Err(PortError::InvalidState(format!(
338 "bundle content digest mismatch: header says {}, events produce {expected_digest}",
339 header.content_digest
340 )));
341 }
342 Ok(())
343}
344
345fn encode_bundle(
346 events: &[ContextUpdatedEvent],
347 snapshot_id: Option<&str>,
348) -> Result<String, PortError> {
349 let mut event_payload = String::new();
350 for event in events {
351 event_payload.push_str(&encode_line("bundle event", event)?);
352 }
353 let digest = content_digest(event_payload.as_bytes());
354 let named = snapshot_id.is_some();
355 let snapshot_id = snapshot_id
356 .map(str::to_string)
357 .unwrap_or_else(|| format!("content-{}", &digest[7..23]));
358 let created_at = if named {
363 SystemTime::now()
364 } else {
365 events
366 .iter()
367 .map(|event| event.occurred_at)
368 .max()
369 .unwrap_or(UNIX_EPOCH + Duration::from_millis(1))
370 };
371 let created_at_unix_ms = created_at
372 .duration_since(UNIX_EPOCH)
373 .unwrap_or(Duration::ZERO)
374 .as_millis() as u64;
375 let header = BundleHeader {
376 bundle_format: BUNDLE_FORMAT_VERSION,
377 event_format: if events
378 .iter()
379 .any(|event| event.role == kmp_domain::NodeCardEvent::ROLE)
380 {
381 super::format_version::EVENT_FORMAT_VERSION
382 } else {
383 2
384 },
385 event_count: events.len() as u64,
386 kernel_version: env!("CARGO_PKG_VERSION").to_string(),
387 snapshot_id,
388 created_at_unix_ms,
389 event_range: event_range(events.len()),
390 abouts: abouts(events),
391 content_digest: digest,
392 };
393 let mut out = encode_line("bundle header", &header)?;
394 out.push_str(&event_payload);
395 Ok(out)
396}
397
398fn event_range(event_count: usize) -> BundleEventRange {
399 if event_count == 0 {
400 BundleEventRange::default()
401 } else {
402 BundleEventRange {
403 first: Some(1),
404 last: Some(event_count as u64),
405 }
406 }
407}
408
409fn abouts(events: &[ContextUpdatedEvent]) -> Vec<String> {
410 events
411 .iter()
412 .map(|event| event.root_node_id.clone())
413 .collect::<BTreeSet<_>>()
414 .into_iter()
415 .collect()
416}
417
418fn filter_events_excluding_abouts(
419 events: Vec<ContextUpdatedEvent>,
420 excluded_abouts: &[String],
421) -> Vec<ContextUpdatedEvent> {
422 if excluded_abouts.is_empty() {
423 return events;
424 }
425 let excluded = excluded_abouts.iter().cloned().collect::<BTreeSet<_>>();
426 events
427 .into_iter()
428 .filter(|event| !excluded.contains(&event.root_node_id))
429 .collect()
430}
431
432fn filter_events_for_abouts(
433 events: Vec<ContextUpdatedEvent>,
434 requested_abouts: &[String],
435) -> Result<Vec<ContextUpdatedEvent>, PortError> {
436 if requested_abouts.is_empty() {
437 return Err(PortError::InvalidState(
438 "filtered export requires at least one about".to_string(),
439 ));
440 }
441 let requested = requested_abouts.iter().cloned().collect::<BTreeSet<_>>();
442 let found = events
443 .iter()
444 .filter(|event| requested.contains(&event.root_node_id))
445 .map(|event| event.root_node_id.clone())
446 .collect::<BTreeSet<_>>();
447 let missing = requested.difference(&found).cloned().collect::<Vec<_>>();
448 if !missing.is_empty() {
449 return Err(PortError::InvalidState(format!(
450 "cannot export missing about{}: {}",
451 if missing.len() == 1 { "" } else { "s" },
452 missing
453 .iter()
454 .map(|about| format!("`{about}`"))
455 .collect::<Vec<_>>()
456 .join(", ")
457 )));
458 }
459 Ok(events
460 .into_iter()
461 .filter(|event| requested.contains(&event.root_node_id))
462 .collect())
463}
464
465fn content_digest(bytes: &[u8]) -> String {
466 format!("sha256:{:x}", Sha256::digest(bytes))
467}
468
469fn validate_revisions(events: &[ContextUpdatedEvent]) -> Result<(), PortError> {
470 let mut revisions: BTreeMap<(&str, &str), u64> = BTreeMap::new();
471 for (position, event) in events.iter().enumerate() {
472 let previous = revisions
473 .get(&(event.root_node_id.as_str(), event.role.as_str()))
474 .copied()
475 .unwrap_or(0);
476 let expected = previous + 1;
477 if event.revision != expected {
478 return Err(PortError::InvalidState(format!(
479 "bundle event position {} carries revision {} for ({}, {}), expected {}; no \
480 events were imported",
481 position + 1,
482 event.revision,
483 event.root_node_id,
484 event.role,
485 expected
486 )));
487 }
488 revisions.insert(
489 (event.root_node_id.as_str(), event.role.as_str()),
490 event.revision,
491 );
492 }
493 Ok(())
494}
495
496impl EmbeddedKernelStore {
497 pub(crate) async fn replay_event_stream<I>(&self, events: I) -> Result<u64, PortError>
505 where
506 I: IntoIterator<Item = ContextUpdatedEvent>,
507 {
508 let mut replayed = 0u64;
509 for event in events {
510 let recorded_revision = event.revision;
511 let expected_previous = recorded_revision.checked_sub(1).ok_or_else(|| {
512 PortError::InvalidState("event carries revision 0; the log is corrupt".to_string())
513 })?;
514 let assigned = self.append(event, expected_previous).await?;
515 if assigned != recorded_revision {
516 return Err(PortError::Conflict(format!(
517 "replay integrity violation: assigned revision {assigned}, \
518 history recorded {recorded_revision}"
519 )));
520 }
521 replayed += 1;
522 }
523 Ok(replayed)
524 }
525}
526
527fn encode_line<T: Serialize>(what: &str, value: &T) -> Result<String, PortError> {
528 let mut line = serde_json::to_string(value)
529 .map_err(|error| PortError::InvalidState(format!("could not encode {what}: {error}")))?;
530 line.push('\n');
531 Ok(line)
532}
533
534fn decode_line<T: for<'de> Deserialize<'de>>(what: &str, line: &str) -> Result<T, PortError> {
535 serde_json::from_str(line)
536 .map_err(|error| PortError::InvalidState(format!("could not decode {what}: {error}")))
537}
538
539#[cfg(test)]
540mod tests {
541 use super::*;
542
543 fn event(root: &str, revision: u64, content_hash: &str) -> ContextUpdatedEvent {
544 ContextUpdatedEvent {
545 root_node_id: root.to_string(),
546 role: "agent".to_string(),
547 revision,
548 content_hash: content_hash.to_string(),
549 changes: Vec::new(),
550 idempotency_key: Some(format!("{root}:{revision}")),
551 logical_digest: None,
552 requested_by: Some("portability-test".to_string()),
553 occurred_at: UNIX_EPOCH + Duration::from_secs(revision),
554 }
555 }
556
557 #[test]
558 fn current_format_identifies_and_covers_the_snapshot() {
559 let events = vec![event("project:b", 1, "b"), event("project:a", 1, "a")];
560 let bundle = encode_bundle(&events, Some("pre-release")).expect("bundle");
561 let header = verify_bundle(&bundle).expect("verified");
562
563 assert_eq!(header.bundle_format, BUNDLE_FORMAT_VERSION);
564 assert_eq!(header.event_format, 2);
565 assert_eq!(header.snapshot_id, "pre-release");
566 assert!(header.created_at_unix_ms > 0);
567 assert_eq!(
568 header.event_range,
569 BundleEventRange {
570 first: Some(1),
571 last: Some(2),
572 }
573 );
574 assert_eq!(header.abouts, ["project:a", "project:b"]);
575 assert!(header.content_digest.starts_with("sha256:"));
576 }
577
578 #[test]
579 fn exclusion_keeps_everything_the_excluded_abouts_do_not_root() {
580 let events = vec![
581 event("project:a", 1, "a1"),
582 event("guide:kmp", 1, "g1"),
583 event("project:a", 2, "a2"),
584 event("guide:kmp-agent", 1, "g2"),
585 ];
586 let excluded = vec!["guide:kmp".to_string(), "guide:kmp-agent".to_string()];
587
588 let kept = filter_events_excluding_abouts(events, &excluded);
589
590 assert_eq!(
591 kept.iter()
592 .map(|event| (event.root_node_id.as_str(), event.revision))
593 .collect::<Vec<_>>(),
594 vec![("project:a", 1), ("project:a", 2)]
595 );
596 }
597
598 #[test]
599 fn exclusion_matches_exactly_because_abouts_are_opaque() {
600 let events = vec![
601 event("guide:kmp", 1, "g1"),
602 event("guide:kmp:extra", 1, "x1"),
603 event(" guide:kmp", 1, "s1"),
604 ];
605 let excluded = vec!["guide:kmp".to_string()];
606
607 let kept = filter_events_excluding_abouts(events, &excluded);
608
609 assert_eq!(
610 kept.iter()
611 .map(|event| event.root_node_id.as_str())
612 .collect::<Vec<_>>(),
613 vec!["guide:kmp:extra", " guide:kmp"]
614 );
615 }
616
617 #[test]
618 fn excluding_an_about_the_store_never_had_is_not_an_error() {
619 let events = vec![event("project:a", 1, "a1")];
623 let excluded = vec!["guide:kmp".to_string()];
624
625 let kept = filter_events_excluding_abouts(events, &excluded);
626
627 assert_eq!(kept.len(), 1);
628 assert_eq!(kept[0].root_node_id, "project:a");
629 }
630
631 #[test]
632 fn verified_bundle_exclusion_keeps_unexcluded_events() {
633 let bundle = encode_bundle(
634 &[
635 event("project:a", 1, "a1"),
636 event("guide:kmp-agent", 1, "g1"),
637 event("project:a", 2, "a2"),
638 ],
639 None,
640 )
641 .expect("bundle");
642
643 let filtered = bundle_excluding_abouts(&bundle, &["guide:kmp-agent".to_string()])
644 .expect("filtered verified bundle");
645 let verified = parse_bundle(&filtered).expect("verified filtered bundle");
646
647 assert_eq!(verified.header.event_count, 2);
648 assert_eq!(verified.header.abouts, ["project:a"]);
649 assert_eq!(
650 verified
651 .events
652 .iter()
653 .map(|event| event.content_hash.as_str())
654 .collect::<Vec<_>>(),
655 ["a1", "a2"]
656 );
657 }
658
659 #[test]
660 fn filtered_export_matches_opaque_abouts_exactly_and_renumbers_its_range() {
661 let events = vec![
662 event("project:a", 1, "a1"),
663 event("project:ab", 1, "ab1"),
664 event("project:a", 2, "a2"),
665 ];
666 let filtered = filter_events_for_abouts(events, &["project:a".to_string()])
667 .expect("exact about exists");
668 assert_eq!(filtered.len(), 2);
669 assert!(
670 filtered
671 .iter()
672 .all(|event| event.root_node_id == "project:a")
673 );
674
675 let bundle = encode_bundle(&filtered, None).expect("filtered bundle");
676 let header = verify_bundle(&bundle).expect("filtered bundle verifies");
677 assert_eq!(header.abouts, ["project:a"]);
678 assert_eq!(header.event_count, 2);
679 assert_eq!(
680 header.event_range,
681 BundleEventRange {
682 first: Some(1),
683 last: Some(2),
684 }
685 );
686 }
687
688 #[test]
689 fn filtered_export_names_every_requested_about_that_is_missing() {
690 let error = filter_events_for_abouts(
691 vec![event("project:a", 1, "a")],
692 &["project:a".into(), "project:none".into()],
693 )
694 .expect_err("missing about must fail");
695 assert!(error.to_string().contains("`project:none`"), "{error}");
696 }
697
698 #[test]
699 fn tampering_is_rejected_before_a_bundle_can_be_replayed() {
700 let bundle =
701 encode_bundle(&[event("project:a", 1, "before")], Some("saved")).expect("bundle");
702 let tampered = bundle.replace("\"content_hash\":\"before\"", "\"content_hash\":\"after\"");
703 let error = verify_bundle(&tampered).expect_err("digest catches changed payload");
704 assert!(error.to_string().contains("content digest mismatch"));
705 }
706
707 #[test]
708 fn merge_fast_forwards_an_exact_prefix() {
709 let first = event("project:a", 1, "one");
710 let second = event("project:a", 2, "two");
711 let left = encode_bundle(std::slice::from_ref(&first), Some("left")).expect("left");
712 let right = encode_bundle(&[first, second], Some("right")).expect("right");
713
714 let merged = merge_bundles(&left, &right, "merged").expect("fast forward");
715 let header = verify_bundle(&merged).expect("verified merge");
716 assert_eq!(header.snapshot_id, "merged");
717 assert_eq!(header.event_count, 2);
718 }
719
720 #[test]
721 fn merge_refuses_two_histories_at_the_same_position() {
722 let left = encode_bundle(&[event("project:a", 1, "left")], Some("left")).expect("left");
723 let right = encode_bundle(&[event("project:a", 1, "right")], Some("right")).expect("right");
724
725 let error = merge_bundles(&left, &right, "invented").expect_err("must refuse");
726 assert!(error.to_string().contains("diverge at event position 1"));
727 assert!(error.to_string().contains("will not invent causal order"));
728 }
729
730 #[test]
731 fn unsupported_bundle_and_event_formats_are_rejected() {
732 let bundle = encode_bundle(&[], None).expect("bundle");
733 for old in [1, 2] {
734 let legacy = bundle.replace("\"bundle_format\":3", &format!("\"bundle_format\":{old}"));
735 let error = verify_bundle(&legacy).expect_err("old bundle is unsupported");
736 assert!(error.to_string().contains("is not supported"), "{error}");
737 }
738 let legacy = bundle.replace("\"event_format\":2", "\"event_format\":1");
739 let error = verify_bundle(&legacy).expect_err("old events are unsupported");
740 assert!(
741 error.to_string().contains("bundle carries event format 1"),
742 "{error}"
743 );
744 }
745
746 #[test]
747 fn invalid_later_revision_is_rejected_in_preflight() {
748 let events = [event("project:a", 1, "one"), event("project:a", 3, "three")];
749 let error = validate_revisions(&events).expect_err("revision gap");
750 assert!(error.to_string().contains("position 2"));
751 assert!(error.to_string().contains("no events were imported"));
752 }
753}