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 fn commits_projections_atomically(&self) -> bool {
9 true
10 }
11
12 async fn append_projected(
13 &self,
14 event: ContextUpdatedEvent,
15 expected_revision: u64,
16 read_revisions: Vec<kmp_domain::ContextRevision>,
17 mutations: Vec<kmp_domain::ProjectionMutation>,
18 ) -> Result<u64, PortError> {
19 self.run(move |store| {
20 let mut tx = store.begin_write()?;
21 for revision in read_revisions {
22 let key = aggregate_key(&revision.root_node_id, &revision.role);
23 let current = match tx.get(Table::Aggregates, Key::Str(&key))? {
24 Some(raw) => decode::<AggregateRecord>("aggregate head", &raw)?.revision,
25 None => 0,
26 };
27 if current != revision.revision {
28 return Err(PortError::Conflict(format!(
29 "reviewed context {} changed: expected {}, current {current}",
30 revision.root_node_id, revision.revision
31 )));
32 }
33 }
34 let new_revision = append_in_transaction(tx.as_mut(), event, expected_revision)?;
35 super::projection_write::apply_mutations_in_transaction(tx.as_mut(), mutations)?;
36 tx.commit()?;
37 Ok(new_revision)
38 })
39 .await
40 }
41
42 async fn append(
43 &self,
44 event: ContextUpdatedEvent,
45 expected_revision: u64,
46 ) -> Result<u64, PortError> {
47 self.run(move |store| {
48 let mut tx = store.begin_write()?;
49
50 let new_revision = append_in_transaction(tx.as_mut(), event, expected_revision)?;
51 tx.commit()?;
52 Ok(new_revision)
53 })
54 .await
55 }
56
57 async fn current_revision(&self, root_node_id: &str, role: &str) -> Result<u64, PortError> {
58 let key = aggregate_key(root_node_id, role);
59 self.run(move |store| {
60 let tx = store.begin_read()?;
61 match tx.get(Table::Aggregates, Key::Str(&key))? {
62 Some(raw) => Ok(decode::<AggregateRecord>("aggregate head", &raw)?.revision),
63 None => Ok(0),
64 }
65 })
66 .await
67 }
68
69 async fn current_content_hash(
70 &self,
71 root_node_id: &str,
72 role: &str,
73 ) -> Result<Option<String>, PortError> {
74 let key = aggregate_key(root_node_id, role);
75 self.run(move |store| {
76 let tx = store.begin_read()?;
77 match tx.get(Table::Aggregates, Key::Str(&key))? {
78 Some(raw) => Ok(Some(
79 decode::<AggregateRecord>("aggregate head", &raw)?.content_hash,
80 )),
81 None => Ok(None),
82 }
83 })
84 .await
85 }
86
87 async fn find_by_idempotency_key(
88 &self,
89 key: &str,
90 ) -> Result<Option<IdempotentOutcome>, PortError> {
91 let key = key.to_string();
92 self.run(move |store| {
93 let tx = store.begin_read()?;
94 match tx.get(Table::Idempotency, Key::Str(&key))? {
95 Some(raw) => Ok(Some(decode("idempotency outcome", &raw)?)),
96 None => Ok(None),
97 }
98 })
99 .await
100 }
101}
102
103impl EmbeddedKernelStore {
104 pub(crate) fn read_event_log(&self) -> Result<Vec<ContextUpdatedEvent>, PortError> {
107 let tx = self.begin_read()?;
108 tx.scan_u64(Table::EventLog)?
109 .into_iter()
110 .map(|(_, raw)| decode("context event", &raw))
111 .collect()
112 }
113}
114
115fn append_in_transaction(
116 tx: &mut dyn super::engine::WriteTx,
117 event: ContextUpdatedEvent,
118 expected_revision: u64,
119) -> Result<u64, PortError> {
120 let key = aggregate_key(&event.root_node_id, &event.role);
121 let current = match tx.get(Table::Aggregates, Key::Str(&key))? {
122 Some(raw) => decode::<AggregateRecord>("aggregate head", &raw)?.revision,
123 None => 0,
124 };
125 if current != expected_revision {
126 return Err(PortError::Conflict(format!(
127 "expected revision {expected_revision}, current is {current}"
128 )));
129 }
130 let new_revision = current + 1;
131
132 let mut event = event;
136 event.revision = new_revision;
137 let outcome = IdempotentOutcome::for_event(&event, new_revision)?;
138
139 let aggregate_bytes = encode(
140 "aggregate head",
141 &AggregateRecord {
142 revision: new_revision,
143 content_hash: event.content_hash.clone(),
144 },
145 )?;
146 tx.insert(Table::Aggregates, Key::Str(&key), &aggregate_bytes)?;
147
148 let next_sequence = tx
149 .last_u64(Table::EventLog)?
150 .map_or(1, |(sequence, _)| sequence + 1);
151 let event_bytes = encode("context event", &event)?;
152 tx.insert(Table::EventLog, Key::U64(next_sequence), &event_bytes)?;
153
154 if let Some(idempotency_key) = event.idempotency_key.as_deref() {
155 let outcome_bytes = encode("idempotency outcome", &outcome)?;
156 tx.insert(
157 Table::Idempotency,
158 Key::Str(idempotency_key),
159 &outcome_bytes,
160 )?;
161 }
162
163 Ok(new_revision)
164}