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    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    /// Reads the full append-only event log in sequence order (audit and
105    /// replay surface).
106    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    // Stamp the assigned revision on the stored event so replay
133    // derives projections with the same revision the aggregate
134    // recorded.
135    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}