Skip to main content

kmp_adapter_embedded/adapter/
context_events.rs

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