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 the exported event stream. A full-store snapshot
18/// starts at one; an empty snapshot has neither bound.
19#[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/// First line of a bundle file: identity and integrity metadata for fail-fast
26/// import. Fields added in bundle format 2 default only so format-1 bundles
27/// remain readable; every format-2 field is validated before replay.
28#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
29pub struct BundleHeader {
30    pub bundle_format: u32,
31    /// Format of the portable event payload, not the on-disk redb/SQLite
32    /// layout. `store_format` is accepted from format-1 bundles because that
33    /// older name described the field ambiguously.
34    #[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/// Outcome of an import: events replayed and projections rebuilt.
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub struct ImportReport {
55    pub events_imported: u64,
56    pub rebuild: ProjectionRebuildReport,
57}
58
59impl EmbeddedKernelStore {
60    /// Serializes the full event log as a JSON-Lines bundle: one header line
61    /// followed by one event per line, in sequence order.
62    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    /// Exports the same complete stream with a human-selected snapshot id.
68    /// The id is metadata, not a filename: callers may store the bundle in git,
69    /// an artifact store, or anywhere else without changing what it identifies.
70    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    /// Replays a bundle into this store. Fail-fast rules: the store must be
81    /// empty (no merge semantics in v1 — ADR-011 rationale applies), the
82    /// header must match supported formats, and every event must reproduce
83    /// exactly the revision it was exported with.
84    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        // Projection derivation is pure. Prove every payload is rebuildable
101        // before the first append so a malformed later line cannot leave a
102        // half-restored store behind.
103        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
117/// Validates a bundle without opening or mutating a store. This is the
118/// recovery check: identity, range, about coverage and digest are all proved
119/// before an operator trusts a saved copy.
120pub fn verify_bundle(bundle: &str) -> Result<BundleHeader, PortError> {
121    parse_bundle(bundle).map(|verified| verified.header)
122}
123
124/// Merges only histories that have a deterministic answer: identical streams
125/// or one exact prefix of the other. Two branches that both appended at the
126/// same position are a semantic conflict, so KMP refuses to invent an order.
127pub 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    // Content-addressed head exports must be byte-identical across redb,
256    // SQLite, repeated exports and migrations. Their creation coordinate is
257    // therefore the newest event time. A named recovery point records when
258    // the operator created that point.
259    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    /// Replays a history into this store, in order, checking that every
341    /// event lands on the revision it was recorded with.
342    ///
343    /// That check is the whole point: a replay that silently renumbers
344    /// history would produce a store that reads plausibly and cites
345    /// revisions that never existed. Shared by import and migration, which
346    /// are the same operation seen from two different distances.
347    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}