kmp_adapter_embedded/adapter/
memory_about_index.rs1use std::collections::BTreeSet;
2
3use kmp_domain::{MemoryAboutIndexReader, PortError};
4
5use super::dimension_lookup_header::DimensionLookupHeader;
6use super::engine::{LinkedJsonScan, Table};
7use super::projection_write::MEMORY_ANCHOR_KIND;
8use super::store::EmbeddedKernelStore;
9
10impl MemoryAboutIndexReader for EmbeddedKernelStore {
11 async fn list_memory_abouts(&self) -> Result<Vec<String>, PortError> {
12 self.run(|store| {
13 let tx = store.begin_read()?;
14 Ok(tx
15 .scan_str(Table::Anchors)?
16 .into_iter()
17 .map(|(anchor, _)| anchor)
18 .collect())
19 })
20 .await
21 }
22
23 async fn list_memory_abouts_by_dimensions(
24 &self,
25 dimension_ids: &[String],
26 ) -> Result<Vec<String>, PortError> {
27 let terms: BTreeSet<_> = dimension_ids.iter().cloned().collect();
28 self.run(move |store| {
29 let tx = store.begin_read()?;
30 let mut abouts = BTreeSet::new();
31 let mut after: Option<(String, String)> = None;
32 loop {
33 let rows = tx.scan_linked_json(&LinkedJsonScan {
34 roots: Table::Anchors,
35 links: Table::Relations,
36 records: Table::Nodes,
37 relation: "has_dimension",
38 fields: &["node_id", "node_kind", "properties.dimension_kind"],
39 after: after
40 .as_ref()
41 .map(|(source, target)| (source.as_str(), target.as_str())),
42 limit: 256,
43 })?;
44 let count = rows.len();
45 for row in rows {
46 after = Some((row.source.clone(), row.target));
47 if abouts.contains(&row.source) {
48 continue;
49 }
50 let Some(source) = DimensionLookupHeader::read(row.source_json.as_deref())?
51 else {
52 continue;
53 };
54 if source.node_kind != MEMORY_ANCHOR_KIND {
55 continue;
56 }
57 if DimensionLookupHeader::read(row.target_json.as_deref())?
58 .is_some_and(|target| target.matches(&terms))
59 {
60 abouts.insert(row.source);
61 }
62 }
63 if count < 256 {
64 break;
65 }
66 }
67 Ok(abouts.into_iter().collect())
68 })
69 .await
70 }
71}