kmp_adapter_embedded/adapter/
snapshot_store.rs1use kmp_domain::{KmpBundle, PortError, SnapshotSaveOptions, SnapshotStore};
2
3use super::store::{EmbeddedKernelStore, SNAPSHOTS, commit_error, storage_error, table_error};
4
5impl SnapshotStore for EmbeddedKernelStore {
6 async fn save_bundle_with_options(
7 &self,
8 bundle: &KmpBundle,
9 options: SnapshotSaveOptions,
10 ) -> Result<(), PortError> {
11 let root_node_id = bundle.root_node().node_id().to_string();
15 let role = bundle.role().as_str().to_string();
16 let record = serde_json::json!({
17 "root_node_id": root_node_id,
18 "role": role,
19 "neighbor_node_ids": bundle
20 .neighbor_nodes()
21 .iter()
22 .map(|node| node.node_id())
23 .collect::<Vec<_>>(),
24 "relationship_count": bundle.relationships().len(),
25 "node_detail_count": bundle.node_details().len(),
26 "ttl_seconds": options.ttl_seconds(),
27 });
28 let bytes = serde_json::to_vec(&record).map_err(|error| {
29 PortError::InvalidState(format!(
30 "embedded store could not encode snapshot summary: {error}"
31 ))
32 })?;
33
34 self.run(move |store| {
35 let tx = store.begin_write()?;
36 {
37 let mut snapshots = tx.open_table(SNAPSHOTS).map_err(table_error)?;
38 snapshots
39 .insert((root_node_id.as_str(), role.as_str()), bytes.as_slice())
40 .map_err(storage_error)?;
41 }
42 tx.commit().map_err(commit_error)
43 })
44 .await
45 }
46}