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)]
32pub struct BundleHeader {
33 pub bundle_format: u32,
34 #[serde(rename = "event_format", alias = "store_format")]
38 pub event_format: u32,
39 pub event_count: u64,
40 pub kernel_version: String,
41 #[serde(default)]
42 pub snapshot_id: String,
43 #[serde(default)]
44 pub created_at_unix_ms: u64,
45 #[serde(default)]
46 pub event_range: BundleEventRange,
47 #[serde(default)]
48 pub abouts: Vec<String>,
49 #[serde(default)]
50 pub content_digest: String,
51}
52
53pub const BUNDLE_FORMAT_VERSION: u32 = 2;
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub struct ImportReport {
58 pub events_imported: u64,
59 pub rebuild: ProjectionRebuildReport,
60}
61
62impl EmbeddedKernelStore {
63 pub fn export_bundle_blocking(&self) -> Result<String, PortError> {
66 encode_bundle(&self.read_event_log()?, None)
67 }
68
69 pub async fn export_bundle(&self) -> Result<String, PortError> {
72 self.run(|store| store.export_bundle_blocking()).await
73 }
74
75 pub async fn export_bundle_for_abouts(
84 &self,
85 requested_abouts: &[String],
86 ) -> Result<String, PortError> {
87 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
88 let events = filter_events_for_abouts(events, requested_abouts)?;
89 encode_bundle(&events, None)
90 }
91
92 pub async fn export_named_bundle(&self, snapshot_id: &str) -> Result<String, PortError> {
96 if snapshot_id.trim().is_empty() {
97 return Err(PortError::InvalidState(
98 "snapshot id must not be empty".to_string(),
99 ));
100 }
101 let events = self.run(EmbeddedKernelStore::read_event_log).await?;
102 encode_bundle(&events, Some(snapshot_id))
103 }
104
105 pub async fn import_bundle<F>(&self, bundle: &str, derive: F) -> Result<ImportReport, PortError>
110 where
111 F: Fn(&ContextUpdatedEvent) -> Result<Vec<ProjectionMutation>, PortError> + Send + 'static,
112 {
113 let (log_length, _) = self.event_log_stats().await?;
114 if log_length != 0 {
115 return Err(PortError::Conflict(format!(
116 "import requires an empty store; this store already holds {log_length} events \
117 (merging bundles is not supported)"
118 )));
119 }
120
121 let verified = parse_bundle(bundle)?;
122 let header = verified.header;
123 let events = verified.events;
124 validate_revisions(&events)?;
125 for event in &events {
129 derive(event)?;
130 }
131 let events_imported = self.replay_event_stream(events).await?;
132 debug_assert_eq!(events_imported, header.event_count);
133
134 let rebuild = self.rebuild_projections(derive).await?;
135 Ok(ImportReport {
136 events_imported,
137 rebuild,
138 })
139 }
140}
141
142pub fn verify_bundle(bundle: &str) -> Result<BundleHeader, PortError> {
146 parse_bundle(bundle).map(|verified| verified.header)
147}
148
149pub fn merge_bundles(left: &str, right: &str, snapshot_id: &str) -> Result<String, PortError> {
153 if snapshot_id.trim().is_empty() {
154 return Err(PortError::InvalidState(
155 "merged snapshot id must not be empty".to_string(),
156 ));
157 }
158 let left = parse_bundle(left)?;
159 let right = parse_bundle(right)?;
160 let shared = left.events.len().min(right.events.len());
161 if let Some(position) =
162 (0..shared).find(|position| left.events[*position] != right.events[*position])
163 {
164 return Err(PortError::Conflict(format!(
165 "bundle histories diverge at event position {}; KMP only fast-forwards an exact \
166 prefix and will not invent causal order for two branches",
167 position + 1
168 )));
169 }
170 let events = if left.events.len() >= right.events.len() {
171 left.events
172 } else {
173 right.events
174 };
175 encode_bundle(&events, Some(snapshot_id))
176}
177
178struct VerifiedBundle {
179 header: BundleHeader,
180 events: Vec<ContextUpdatedEvent>,
181}
182
183fn parse_bundle(bundle: &str) -> Result<VerifiedBundle, PortError> {
184 let mut lines = bundle.lines().filter(|line| !line.trim().is_empty());
185 let header: BundleHeader = decode_line(
186 "bundle header",
187 lines.next().ok_or_else(|| {
188 PortError::InvalidState("bundle is empty: missing header line".to_string())
189 })?,
190 )?;
191 if !matches!(header.bundle_format, 1 | BUNDLE_FORMAT_VERSION) {
192 return Err(PortError::InvalidState(format!(
193 "bundle format {} is not supported (this binary reads 1 and {})",
194 header.bundle_format, BUNDLE_FORMAT_VERSION
195 )));
196 }
197 if header.event_format != super::format_version::EVENT_FORMAT_VERSION {
198 return Err(PortError::InvalidState(format!(
199 "bundle carries event format {}, this binary supports {}",
200 header.event_format,
201 super::format_version::EVENT_FORMAT_VERSION
202 )));
203 }
204
205 let mut events = Vec::new();
206 let mut event_payload = String::new();
207 for line in lines {
208 events.push(decode_line::<ContextUpdatedEvent>("bundle event", line)?);
209 event_payload.push_str(line);
210 event_payload.push('\n');
211 }
212 if events.len() as u64 != header.event_count {
213 return Err(PortError::InvalidState(format!(
214 "bundle header declares {} events but {} were present",
215 header.event_count,
216 events.len()
217 )));
218 }
219
220 if header.bundle_format == BUNDLE_FORMAT_VERSION {
221 validate_v2_header(&header, &events, &event_payload)?;
222 }
223 Ok(VerifiedBundle { header, events })
224}
225
226fn validate_v2_header(
227 header: &BundleHeader,
228 events: &[ContextUpdatedEvent],
229 event_payload: &str,
230) -> Result<(), PortError> {
231 if header.snapshot_id.trim().is_empty() {
232 return Err(PortError::InvalidState(
233 "bundle format 2 requires snapshot_id".to_string(),
234 ));
235 }
236 if header.created_at_unix_ms == 0 {
237 return Err(PortError::InvalidState(
238 "bundle format 2 requires created_at_unix_ms".to_string(),
239 ));
240 }
241 let expected_range = event_range(events.len());
242 if header.event_range != expected_range {
243 return Err(PortError::InvalidState(format!(
244 "bundle event_range {:?} does not cover its {} events (expected {:?})",
245 header.event_range,
246 events.len(),
247 expected_range
248 )));
249 }
250 let expected_abouts = abouts(events);
251 if header.abouts != expected_abouts {
252 return Err(PortError::InvalidState(format!(
253 "bundle abouts do not match its events (expected {})",
254 expected_abouts.join(", ")
255 )));
256 }
257 let expected_digest = content_digest(event_payload.as_bytes());
258 if header.content_digest != expected_digest {
259 return Err(PortError::InvalidState(format!(
260 "bundle content digest mismatch: header says {}, events produce {expected_digest}",
261 header.content_digest
262 )));
263 }
264 Ok(())
265}
266
267fn encode_bundle(
268 events: &[ContextUpdatedEvent],
269 snapshot_id: Option<&str>,
270) -> Result<String, PortError> {
271 let mut event_payload = String::new();
272 for event in events {
273 event_payload.push_str(&encode_line("bundle event", event)?);
274 }
275 let digest = content_digest(event_payload.as_bytes());
276 let named = snapshot_id.is_some();
277 let snapshot_id = snapshot_id
278 .map(str::to_string)
279 .unwrap_or_else(|| format!("content-{}", &digest[7..23]));
280 let created_at = if named {
285 SystemTime::now()
286 } else {
287 events
288 .iter()
289 .map(|event| event.occurred_at)
290 .max()
291 .unwrap_or(UNIX_EPOCH + Duration::from_millis(1))
292 };
293 let created_at_unix_ms = created_at
294 .duration_since(UNIX_EPOCH)
295 .unwrap_or(Duration::ZERO)
296 .as_millis() as u64;
297 let header = BundleHeader {
298 bundle_format: BUNDLE_FORMAT_VERSION,
299 event_format: super::format_version::EVENT_FORMAT_VERSION,
300 event_count: events.len() as u64,
301 kernel_version: env!("CARGO_PKG_VERSION").to_string(),
302 snapshot_id,
303 created_at_unix_ms,
304 event_range: event_range(events.len()),
305 abouts: abouts(events),
306 content_digest: digest,
307 };
308 let mut out = encode_line("bundle header", &header)?;
309 out.push_str(&event_payload);
310 Ok(out)
311}
312
313fn event_range(event_count: usize) -> BundleEventRange {
314 if event_count == 0 {
315 BundleEventRange::default()
316 } else {
317 BundleEventRange {
318 first: Some(1),
319 last: Some(event_count as u64),
320 }
321 }
322}
323
324fn abouts(events: &[ContextUpdatedEvent]) -> Vec<String> {
325 events
326 .iter()
327 .map(|event| event.root_node_id.clone())
328 .collect::<BTreeSet<_>>()
329 .into_iter()
330 .collect()
331}
332
333fn filter_events_for_abouts(
334 events: Vec<ContextUpdatedEvent>,
335 requested_abouts: &[String],
336) -> Result<Vec<ContextUpdatedEvent>, PortError> {
337 if requested_abouts.is_empty() {
338 return Err(PortError::InvalidState(
339 "filtered export requires at least one about".to_string(),
340 ));
341 }
342 let requested = requested_abouts.iter().cloned().collect::<BTreeSet<_>>();
343 let found = events
344 .iter()
345 .filter(|event| requested.contains(&event.root_node_id))
346 .map(|event| event.root_node_id.clone())
347 .collect::<BTreeSet<_>>();
348 let missing = requested.difference(&found).cloned().collect::<Vec<_>>();
349 if !missing.is_empty() {
350 return Err(PortError::InvalidState(format!(
351 "cannot export missing about{}: {}",
352 if missing.len() == 1 { "" } else { "s" },
353 missing
354 .iter()
355 .map(|about| format!("`{about}`"))
356 .collect::<Vec<_>>()
357 .join(", ")
358 )));
359 }
360 Ok(events
361 .into_iter()
362 .filter(|event| requested.contains(&event.root_node_id))
363 .collect())
364}
365
366fn content_digest(bytes: &[u8]) -> String {
367 format!("sha256:{:x}", Sha256::digest(bytes))
368}
369
370fn validate_revisions(events: &[ContextUpdatedEvent]) -> Result<(), PortError> {
371 let mut revisions: BTreeMap<(&str, &str), u64> = BTreeMap::new();
372 for (position, event) in events.iter().enumerate() {
373 let previous = revisions
374 .get(&(event.root_node_id.as_str(), event.role.as_str()))
375 .copied()
376 .unwrap_or(0);
377 let expected = previous + 1;
378 if event.revision != expected {
379 return Err(PortError::InvalidState(format!(
380 "bundle event position {} carries revision {} for ({}, {}), expected {}; no \
381 events were imported",
382 position + 1,
383 event.revision,
384 event.root_node_id,
385 event.role,
386 expected
387 )));
388 }
389 revisions.insert(
390 (event.root_node_id.as_str(), event.role.as_str()),
391 event.revision,
392 );
393 }
394 Ok(())
395}
396
397impl EmbeddedKernelStore {
398 pub(crate) async fn replay_event_stream<I>(&self, events: I) -> Result<u64, PortError>
406 where
407 I: IntoIterator<Item = ContextUpdatedEvent>,
408 {
409 let mut replayed = 0u64;
410 for event in events {
411 let recorded_revision = event.revision;
412 let expected_previous = recorded_revision.checked_sub(1).ok_or_else(|| {
413 PortError::InvalidState("event carries revision 0; the log is corrupt".to_string())
414 })?;
415 let assigned = self.append(event, expected_previous).await?;
416 if assigned != recorded_revision {
417 return Err(PortError::Conflict(format!(
418 "replay integrity violation: assigned revision {assigned}, \
419 history recorded {recorded_revision}"
420 )));
421 }
422 replayed += 1;
423 }
424 Ok(replayed)
425 }
426}
427
428fn encode_line<T: Serialize>(what: &str, value: &T) -> Result<String, PortError> {
429 let mut line = serde_json::to_string(value)
430 .map_err(|error| PortError::InvalidState(format!("could not encode {what}: {error}")))?;
431 line.push('\n');
432 Ok(line)
433}
434
435fn decode_line<T: for<'de> Deserialize<'de>>(what: &str, line: &str) -> Result<T, PortError> {
436 serde_json::from_str(line)
437 .map_err(|error| PortError::InvalidState(format!("could not decode {what}: {error}")))
438}
439
440#[cfg(test)]
441mod tests {
442 use super::*;
443
444 fn event(root: &str, revision: u64, content_hash: &str) -> ContextUpdatedEvent {
445 ContextUpdatedEvent {
446 root_node_id: root.to_string(),
447 role: "agent".to_string(),
448 revision,
449 content_hash: content_hash.to_string(),
450 changes: Vec::new(),
451 idempotency_key: Some(format!("{root}:{revision}")),
452 logical_digest: None,
453 requested_by: Some("portability-test".to_string()),
454 occurred_at: UNIX_EPOCH + Duration::from_secs(revision),
455 }
456 }
457
458 #[test]
459 fn format_two_identifies_and_covers_the_snapshot() {
460 let events = vec![event("project:b", 1, "b"), event("project:a", 1, "a")];
461 let bundle = encode_bundle(&events, Some("pre-release")).expect("bundle");
462 let header = verify_bundle(&bundle).expect("verified");
463
464 assert_eq!(header.bundle_format, BUNDLE_FORMAT_VERSION);
465 assert_eq!(
466 header.event_format,
467 super::super::format_version::EVENT_FORMAT_VERSION
468 );
469 assert_eq!(header.snapshot_id, "pre-release");
470 assert!(header.created_at_unix_ms > 0);
471 assert_eq!(
472 header.event_range,
473 BundleEventRange {
474 first: Some(1),
475 last: Some(2),
476 }
477 );
478 assert_eq!(header.abouts, ["project:a", "project:b"]);
479 assert!(header.content_digest.starts_with("sha256:"));
480 }
481
482 #[test]
483 fn filtered_export_matches_opaque_abouts_exactly_and_renumbers_its_range() {
484 let events = vec![
485 event("project:a", 1, "a1"),
486 event("project:ab", 1, "ab1"),
487 event("project:a", 2, "a2"),
488 ];
489 let filtered = filter_events_for_abouts(events, &["project:a".to_string()])
490 .expect("exact about exists");
491 assert_eq!(filtered.len(), 2);
492 assert!(
493 filtered
494 .iter()
495 .all(|event| event.root_node_id == "project:a")
496 );
497
498 let bundle = encode_bundle(&filtered, None).expect("filtered bundle");
499 let header = verify_bundle(&bundle).expect("filtered bundle verifies");
500 assert_eq!(header.abouts, ["project:a"]);
501 assert_eq!(header.event_count, 2);
502 assert_eq!(
503 header.event_range,
504 BundleEventRange {
505 first: Some(1),
506 last: Some(2),
507 }
508 );
509 }
510
511 #[test]
512 fn filtered_export_names_every_requested_about_that_is_missing() {
513 let error = filter_events_for_abouts(
514 vec![event("project:a", 1, "a")],
515 &["project:a".into(), "project:none".into()],
516 )
517 .expect_err("missing about must fail");
518 assert!(error.to_string().contains("`project:none`"), "{error}");
519 }
520
521 #[test]
522 fn tampering_is_rejected_before_a_bundle_can_be_replayed() {
523 let bundle =
524 encode_bundle(&[event("project:a", 1, "before")], Some("saved")).expect("bundle");
525 let tampered = bundle.replace("\"content_hash\":\"before\"", "\"content_hash\":\"after\"");
526 let error = verify_bundle(&tampered).expect_err("digest catches changed payload");
527 assert!(error.to_string().contains("content digest mismatch"));
528 }
529
530 #[test]
531 fn merge_fast_forwards_an_exact_prefix() {
532 let first = event("project:a", 1, "one");
533 let second = event("project:a", 2, "two");
534 let left = encode_bundle(std::slice::from_ref(&first), Some("left")).expect("left");
535 let right = encode_bundle(&[first, second], Some("right")).expect("right");
536
537 let merged = merge_bundles(&left, &right, "merged").expect("fast forward");
538 let header = verify_bundle(&merged).expect("verified merge");
539 assert_eq!(header.snapshot_id, "merged");
540 assert_eq!(header.event_count, 2);
541 }
542
543 #[test]
544 fn merge_refuses_two_histories_at_the_same_position() {
545 let left = encode_bundle(&[event("project:a", 1, "left")], Some("left")).expect("left");
546 let right = encode_bundle(&[event("project:a", 1, "right")], Some("right")).expect("right");
547
548 let error = merge_bundles(&left, &right, "invented").expect_err("must refuse");
549 assert!(error.to_string().contains("diverge at event position 1"));
550 assert!(error.to_string().contains("will not invent causal order"));
551 }
552
553 #[test]
554 fn legacy_format_one_remains_readable() {
555 let legacy =
556 r#"{"bundle_format":1,"store_format":1,"event_count":0,"kernel_version":"0.1.3"}"#;
557 let header = verify_bundle(legacy).expect("format one remains portable");
558 assert_eq!(header.event_format, 1);
559 assert!(header.snapshot_id.is_empty());
560 }
561
562 #[test]
563 fn invalid_later_revision_is_rejected_in_preflight() {
564 let events = [event("project:a", 1, "one"), event("project:a", 3, "three")];
565 let error = validate_revisions(&events).expect_err("revision gap");
566 assert!(error.to_string().contains("position 2"));
567 assert!(error.to_string().contains("no events were imported"));
568 }
569}