Skip to main content

kmp_adapter_embedded/adapter/
context_events.rs

1use kmp_domain::{ContextEventStore, ContextUpdatedEvent, IdempotentOutcome, PortError};
2
3use super::engine::{Key, Table};
4use super::serdes::{AggregateRecord, decode, encode};
5use super::store::{EmbeddedKernelStore, aggregate_key};
6
7impl ContextEventStore for EmbeddedKernelStore {
8    async fn append(
9        &self,
10        event: ContextUpdatedEvent,
11        expected_revision: u64,
12    ) -> Result<u64, PortError> {
13        self.run(move |store| {
14            let mut tx = store.begin_write()?;
15
16            let key = aggregate_key(&event.root_node_id, &event.role);
17            let current = match tx.get(Table::Aggregates, Key::Str(&key))? {
18                Some(raw) => decode::<AggregateRecord>("aggregate head", &raw)?.revision,
19                None => 0,
20            };
21            if current != expected_revision {
22                return Err(PortError::Conflict(format!(
23                    "expected revision {expected_revision}, current is {current}"
24                )));
25            }
26            let new_revision = current + 1;
27
28            // Stamp the assigned revision on the stored event so replay
29            // derives projections with the same revision the aggregate
30            // recorded.
31            let mut event = event;
32            event.revision = new_revision;
33
34            let aggregate_bytes = encode(
35                "aggregate head",
36                &AggregateRecord {
37                    revision: new_revision,
38                    content_hash: event.content_hash.clone(),
39                },
40            )?;
41            tx.insert(Table::Aggregates, Key::Str(&key), &aggregate_bytes)?;
42
43            let next_sequence = tx
44                .last_u64(Table::EventLog)?
45                .map_or(1, |(sequence, _)| sequence + 1);
46            let event_bytes = encode("context event", &event)?;
47            tx.insert(Table::EventLog, Key::U64(next_sequence), &event_bytes)?;
48
49            if let Some(idempotency_key) = event.idempotency_key.as_deref() {
50                let outcome_bytes = encode(
51                    "idempotency outcome",
52                    &IdempotentOutcome {
53                        revision: new_revision,
54                        content_hash: event.content_hash.clone(),
55                        logical_digest: event.logical_digest.clone(),
56                    },
57                )?;
58                tx.insert(
59                    Table::Idempotency,
60                    Key::Str(idempotency_key),
61                    &outcome_bytes,
62                )?;
63            }
64
65            tx.commit()?;
66            Ok(new_revision)
67        })
68        .await
69    }
70
71    async fn current_revision(&self, root_node_id: &str, role: &str) -> Result<u64, PortError> {
72        let key = aggregate_key(root_node_id, role);
73        self.run(move |store| {
74            let tx = store.begin_read()?;
75            match tx.get(Table::Aggregates, Key::Str(&key))? {
76                Some(raw) => Ok(decode::<AggregateRecord>("aggregate head", &raw)?.revision),
77                None => Ok(0),
78            }
79        })
80        .await
81    }
82
83    async fn current_content_hash(
84        &self,
85        root_node_id: &str,
86        role: &str,
87    ) -> Result<Option<String>, PortError> {
88        let key = aggregate_key(root_node_id, role);
89        self.run(move |store| {
90            let tx = store.begin_read()?;
91            match tx.get(Table::Aggregates, Key::Str(&key))? {
92                Some(raw) => Ok(Some(
93                    decode::<AggregateRecord>("aggregate head", &raw)?.content_hash,
94                )),
95                None => Ok(None),
96            }
97        })
98        .await
99    }
100
101    async fn find_by_idempotency_key(
102        &self,
103        key: &str,
104    ) -> Result<Option<IdempotentOutcome>, PortError> {
105        let key = key.to_string();
106        self.run(move |store| {
107            let tx = store.begin_read()?;
108            match tx.get(Table::Idempotency, Key::Str(&key))? {
109                Some(raw) => Ok(Some(decode("idempotency outcome", &raw)?)),
110                None => Ok(None),
111            }
112        })
113        .await
114    }
115}
116
117impl EmbeddedKernelStore {
118    /// The event log, read synchronously. The migration path uses this
119    /// before any runtime exists around the source copy.
120    pub(crate) fn read_event_log_blocking(&self) -> Result<Vec<ContextUpdatedEvent>, PortError> {
121        self.read_event_log()
122    }
123
124    /// Reads the full append-only event log in sequence order (audit and
125    /// replay surface).
126    pub(crate) fn read_event_log(&self) -> Result<Vec<ContextUpdatedEvent>, PortError> {
127        let tx = self.begin_read()?;
128        tx.scan_u64(Table::EventLog)?
129            .into_iter()
130            .map(|(_, raw)| decode("context event", &raw))
131            .collect()
132    }
133}