Skip to main content

kmp_adapter_embedded/adapter/
runtime_state.rs

1use 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}