Skip to main content

kmp_adapter_embedded/adapter/
portability.rs

1//! Export/import (E6): the append-only event log is the portable form of an
2//! embedded store. Export dumps it in sequence order; import replays it into
3//! an empty store, reproducing identical revisions, idempotency outcomes and
4//! projections — temporal reads and relation proof survive the round trip by
5//! construction.
6
7use 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/// Inclusive positions in this bundle's event stream. Export is a portable
18/// replay, not a view over the store's internal sequence keys: full and
19/// filtered bundles both renumber their payload positions from one while
20/// preserving every event's aggregate revision (and therefore every ref).
21/// An empty snapshot has neither bound.
22#[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/// First line of a bundle file: identity and integrity metadata for fail-fast
29/// import. Fields added in bundle format 2 default only so format-1 bundles
30/// remain readable; every format-2 field is validated before replay.
31#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
32pub struct BundleHeader {
33    pub bundle_format: u32,
34    /// Format of the portable event payload, not the on-disk SQLite layout.
35    /// `store_format` is accepted from format-1 bundles because that
36    /// older name described the field ambiguously.
37    #[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/// Outcome of an import: events replayed and projections rebuilt.
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub struct ImportReport {
58    pub events_imported: u64,
59    pub rebuild: ProjectionRebuildReport,
60}
61
62impl EmbeddedKernelStore {
63    /// Synchronous export for already-blocking operational paths such as
64    /// Doctor. Async application paths should use [`Self::export_bundle`].
65    pub fn export_bundle_blocking(&self) -> Result<String, PortError> {
66        encode_bundle(&self.read_event_log()?, None)
67    }
68
69    /// Serializes the full event log as a JSON-Lines bundle: one header line
70    /// followed by one event per line, in sequence order.
71    pub async fn export_bundle(&self) -> Result<String, PortError> {
72        self.run(|store| store.export_bundle_blocking()).await
73    }
74
75    /// Serializes only events rooted at one of `requested_abouts`.
76    ///
77    /// Abouts are opaque routing identifiers: matching is exact, with no
78    /// trimming, case folding, prefix expansion or other normalisation. The
79    /// filtered stream keeps the store's event order and each aggregate's
80    /// recorded revisions, then receives bundle-local positions starting at
81    /// one. A missing requested about is refused before callers can create a
82    /// destination file.
83    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    /// Exports the same complete stream with a human-selected snapshot id.
93    /// The id is metadata, not a filename: callers may store the bundle in git,
94    /// an artifact store, or anywhere else without changing what it identifies.
95    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    /// Replays a bundle into this store. Fail-fast rules: the store must be
106    /// empty (no merge semantics in v1 — ADR-011 rationale applies), the
107    /// header must match supported formats, and every event must reproduce
108    /// exactly the revision it was exported with.
109    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        // Projection derivation is pure. Prove every payload is rebuildable
126        // before the first append so a malformed later line cannot leave a
127        // half-restored store behind.
128        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
142/// Validates a bundle without opening or mutating a store. This is the
143/// recovery check: identity, range, about coverage and digest are all proved
144/// before an operator trusts a saved copy.
145pub fn verify_bundle(bundle: &str) -> Result<BundleHeader, PortError> {
146    parse_bundle(bundle).map(|verified| verified.header)
147}
148
149/// Merges only histories that have a deterministic answer: identical streams
150/// or one exact prefix of the other. Two branches that both appended at the
151/// same position are a semantic conflict, so KMP refuses to invent an order.
152pub 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    // Content-addressed head exports must be byte-identical across storage
281    // layouts and repeated exports. Their creation coordinate is
282    // therefore the newest event time. A named recovery point records when
283    // the operator created that point.
284    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    /// Replays a history into this store, in order, checking that every
399    /// event lands on the revision it was recorded with.
400    ///
401    /// That check is the whole point: a replay that silently renumbers
402    /// history would produce a store that reads plausibly and cites
403    /// revisions that never existed. Shared by import and migration, which
404    /// are the same operation seen from two different distances.
405    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}