kmp_adapter_embedded/adapter/
runtime_state.rs1use kmp_domain::{PortError, ProcessedEventStore, ProjectionCheckpoint, ProjectionCheckpointStore};
2
3use super::engine::{Key, Table};
4use super::serdes::{CheckpointRecord, decode, encode};
5use super::store::EmbeddedKernelStore;
6
7impl ProcessedEventStore for EmbeddedKernelStore {
8 async fn has_processed(&self, consumer_name: &str, event_id: &str) -> Result<bool, PortError> {
9 let consumer_name = consumer_name.to_string();
10 let event_id = event_id.to_string();
11 self.run(move |store| {
12 let tx = store.begin_read()?;
13 Ok(tx
14 .get(Table::Processed, Key::Str2(&consumer_name, &event_id))?
15 .is_some())
16 })
17 .await
18 }
19
20 async fn record_processed(&self, consumer_name: &str, event_id: &str) -> Result<(), PortError> {
21 let consumer_name = consumer_name.to_string();
22 let event_id = event_id.to_string();
23 self.run(move |store| {
24 let mut tx = store.begin_write()?;
25 tx.insert(Table::Processed, Key::Str2(&consumer_name, &event_id), &[])?;
26 tx.commit()
27 })
28 .await
29 }
30}
31
32impl ProjectionCheckpointStore for EmbeddedKernelStore {
33 async fn load_checkpoint(
34 &self,
35 consumer_name: &str,
36 stream_name: &str,
37 ) -> Result<Option<ProjectionCheckpoint>, PortError> {
38 let consumer_name = consumer_name.to_string();
39 let stream_name = stream_name.to_string();
40 self.run(move |store| {
41 let tx = store.begin_read()?;
42 match tx.get(Table::Checkpoints, Key::Str2(&consumer_name, &stream_name))? {
43 Some(raw) => Ok(Some(
44 decode::<CheckpointRecord>("projection checkpoint", &raw)?.into(),
45 )),
46 None => Ok(None),
47 }
48 })
49 .await
50 }
51
52 async fn save_checkpoint(&self, checkpoint: ProjectionCheckpoint) -> Result<(), PortError> {
53 self.run(move |store| {
54 let consumer_name = checkpoint.consumer_name.clone();
55 let stream_name = checkpoint.stream_name.clone();
56 let bytes = encode("projection checkpoint", &CheckpointRecord::from(checkpoint))?;
57 let mut tx = store.begin_write()?;
58 tx.insert(
59 Table::Checkpoints,
60 Key::Str2(&consumer_name, &stream_name),
61 &bytes,
62 )?;
63 tx.commit()
64 })
65 .await
66 }
67}