Skip to main content

kmp_adapter_embedded/adapter/
runtime_state.rs

1use kmp_domain::{PortError, ProcessedEventStore, ProjectionCheckpoint, ProjectionCheckpointStore};
2
3use super::serdes::{CheckpointRecord, decode, encode};
4use super::store::{
5    CHECKPOINTS, EmbeddedKernelStore, PROCESSED, commit_error, storage_error, table_error,
6};
7
8impl ProcessedEventStore for EmbeddedKernelStore {
9    async fn has_processed(&self, consumer_name: &str, event_id: &str) -> Result<bool, PortError> {
10        let consumer_name = consumer_name.to_string();
11        let event_id = event_id.to_string();
12        self.run(move |store| {
13            let tx = store.begin_read()?;
14            let processed = tx.open_table(PROCESSED).map_err(table_error)?;
15            Ok(processed
16                .get((consumer_name.as_str(), event_id.as_str()))
17                .map_err(storage_error)?
18                .is_some())
19        })
20        .await
21    }
22
23    async fn record_processed(&self, consumer_name: &str, event_id: &str) -> Result<(), PortError> {
24        let consumer_name = consumer_name.to_string();
25        let event_id = event_id.to_string();
26        self.run(move |store| {
27            let tx = store.begin_write()?;
28            {
29                let mut processed = tx.open_table(PROCESSED).map_err(table_error)?;
30                processed
31                    .insert((consumer_name.as_str(), event_id.as_str()), ())
32                    .map_err(storage_error)?;
33            }
34            tx.commit().map_err(commit_error)
35        })
36        .await
37    }
38}
39
40impl ProjectionCheckpointStore for EmbeddedKernelStore {
41    async fn load_checkpoint(
42        &self,
43        consumer_name: &str,
44        stream_name: &str,
45    ) -> Result<Option<ProjectionCheckpoint>, PortError> {
46        let consumer_name = consumer_name.to_string();
47        let stream_name = stream_name.to_string();
48        self.run(move |store| {
49            let tx = store.begin_read()?;
50            let checkpoints = tx.open_table(CHECKPOINTS).map_err(table_error)?;
51            match checkpoints
52                .get((consumer_name.as_str(), stream_name.as_str()))
53                .map_err(storage_error)?
54            {
55                Some(guard) => Ok(Some(
56                    decode::<CheckpointRecord>("projection checkpoint", guard.value())?.into(),
57                )),
58                None => Ok(None),
59            }
60        })
61        .await
62    }
63
64    async fn save_checkpoint(&self, checkpoint: ProjectionCheckpoint) -> Result<(), PortError> {
65        self.run(move |store| {
66            let tx = store.begin_write()?;
67            {
68                let key = (
69                    checkpoint.consumer_name.clone(),
70                    checkpoint.stream_name.clone(),
71                );
72                let bytes = encode("projection checkpoint", &CheckpointRecord::from(checkpoint))?;
73                let mut checkpoints = tx.open_table(CHECKPOINTS).map_err(table_error)?;
74                checkpoints
75                    .insert((key.0.as_str(), key.1.as_str()), bytes.as_slice())
76                    .map_err(storage_error)?;
77            }
78            tx.commit().map_err(commit_error)
79        })
80        .await
81    }
82}