kmp_adapter_embedded/adapter/
context_events.rs1use 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 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 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}