kmp_adapter_embedded/adapter/
context_events.rs1use 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 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 pub(crate) fn read_event_log_blocking(&self) -> Result<Vec<ContextUpdatedEvent>, PortError> {
121 self.read_event_log()
122 }
123
124 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}