kmp_adapter_embedded/adapter/
runtime_state.rs1use 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}