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            let outcome = IdempotentOutcome::for_event(&event, new_revision)?;
34
35            let aggregate_bytes = encode(
36                "aggregate head",
37                &AggregateRecord {
38                    revision: new_revision,
39                    content_hash: event.content_hash.clone(),
40                },
41            )?;
42            tx.insert(Table::Aggregates, Key::Str(&key), &aggregate_bytes)?;
43
44            let next_sequence = tx
45                .last_u64(Table::EventLog)?
46                .map_or(1, |(sequence, _)| sequence + 1);
47            let event_bytes = encode("context event", &event)?;
48            tx.insert(Table::EventLog, Key::U64(next_sequence), &event_bytes)?;
49
50            if let Some(idempotency_key) = event.idempotency_key.as_deref() {
51                let outcome_bytes = encode("idempotency outcome", &outcome)?;
52                tx.insert(
53                    Table::Idempotency,
54                    Key::Str(idempotency_key),
55                    &outcome_bytes,
56                )?;
57            }
58
59            tx.commit()?;
60            Ok(new_revision)
61        })
62        .await
63    }
64
65    async fn current_revision(&self, root_node_id: &str, role: &str) -> Result<u64, PortError> {
66        let key = aggregate_key(root_node_id, role);
67        self.run(move |store| {
68            let tx = store.begin_read()?;
69            match tx.get(Table::Aggregates, Key::Str(&key))? {
70                Some(raw) => Ok(decode::<AggregateRecord>("aggregate head", &raw)?.revision),
71                None => Ok(0),
72            }
73        })
74        .await
75    }
76
77    async fn current_content_hash(
78        &self,
79        root_node_id: &str,
80        role: &str,
81    ) -> Result<Option<String>, PortError> {
82        let key = aggregate_key(root_node_id, role);
83        self.run(move |store| {
84            let tx = store.begin_read()?;
85            match tx.get(Table::Aggregates, Key::Str(&key))? {
86                Some(raw) => Ok(Some(
87                    decode::<AggregateRecord>("aggregate head", &raw)?.content_hash,
88                )),
89                None => Ok(None),
90            }
91        })
92        .await
93    }
94
95    async fn find_by_idempotency_key(
96        &self,
97        key: &str,
98    ) -> Result<Option<IdempotentOutcome>, PortError> {
99        let key = key.to_string();
100        self.run(move |store| {
101            let tx = store.begin_read()?;
102            match tx.get(Table::Idempotency, Key::Str(&key))? {
103                Some(raw) => Ok(Some(decode("idempotency outcome", &raw)?)),
104                None => Ok(None),
105            }
106        })
107        .await
108    }
109}
110
111impl EmbeddedKernelStore {
112    /// Reads the full append-only event log in sequence order (audit and
113    /// replay surface).
114    pub(crate) fn read_event_log(&self) -> Result<Vec<ContextUpdatedEvent>, PortError> {
115        let tx = self.begin_read()?;
116        tx.scan_u64(Table::EventLog)?
117            .into_iter()
118            .map(|(_, raw)| decode("context event", &raw))
119            .collect()
120    }
121}