Skip to main content

kmp_adapter_embedded/adapter/
memory_about_index.rs

1use 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}