Skip to main content

kmp_adapter_embedded/adapter/
snapshot_store.rs

1use 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        // The snapshot port is write-only; persist an auditable summary keyed
12        // by (root, role) so a stored decision context can be accounted for
13        // offline. Full bundle rendering stays above the ports.
14        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}