Skip to main content

kmp_adapter_embedded/adapter/
graph_read.rs

1use std::collections::{BTreeMap, BTreeSet, VecDeque};
2
3use kmp_domain::{
4    ContextPathNeighborhood, GraphNeighborhoodReader, MemoryAboutIndexReader, NeighborhoodRequest,
5    NodeNeighborhood, NodeProjection, NodeRelationProjection, NodeRelationshipReader,
6    NodeRelationships, PortError,
7};
8
9use super::engine::{Key, ReadTx, Table};
10use super::projection_write::MEMORY_ANCHOR_KIND;
11use super::serdes::{NodeRecord, decode, decode_explanation};
12use super::store::EmbeddedKernelStore;
13
14fn load_node(tx: &dyn ReadTx, node_id: &str) -> Result<Option<NodeProjection>, PortError> {
15    match tx.get(Table::Nodes, Key::Str(node_id))? {
16        Some(raw) => Ok(Some(
17            decode::<NodeRecord>("graph node", &raw)?.into_projection()?,
18        )),
19        None => Ok(None),
20    }
21}
22
23fn outgoing_rows(tx: &dyn ReadTx, source: &str) -> Result<Vec<NodeRelationProjection>, PortError> {
24    tx.scan_str3_by_first(Table::Relations, source)?
25        .into_iter()
26        .map(|((source_node_id, target_node_id, relation_type), raw)| {
27            Ok(NodeRelationProjection {
28                source_node_id,
29                target_node_id,
30                relation_type,
31                explanation: decode_explanation(&raw)?,
32            })
33        })
34        .collect()
35}
36
37fn outgoing_targets(tx: &dyn ReadTx, source: &str) -> Result<Vec<String>, PortError> {
38    Ok(tx
39        .scan_str3_by_first(Table::Relations, source)?
40        .into_iter()
41        .map(|((_, target, _), _)| target)
42        .collect())
43}
44
45fn reachable_outward(
46    tx: &dyn ReadTx,
47    request: &NeighborhoodRequest,
48) -> Result<BTreeSet<String>, PortError> {
49    let root_node_id = request.root_node_id();
50    let mut visited = BTreeSet::from([root_node_id.to_string()]);
51    let mut reachable = BTreeSet::new();
52    let mut frontier = VecDeque::from([(root_node_id.to_string(), 0u32)]);
53
54    while let Some((node_id, hops)) = frontier.pop_front() {
55        if hops == request.depth() {
56            continue;
57        }
58        for target in outgoing_targets(tx, &node_id)? {
59            // A dimension the caller did not ask for is not descended into,
60            // and everything hanging from it is thereby never loaded. Nothing
61            // else is refused: the narrowing is on the axis, not on the
62            // contents.
63            if !request.admits(&target) {
64                continue;
65            }
66            if visited.insert(target.clone()) {
67                reachable.insert(target.clone());
68                frontier.push_back((target, hops + 1));
69            }
70        }
71    }
72
73    reachable.remove(root_node_id);
74    Ok(reachable)
75}
76
77fn relations_among(
78    tx: &dyn ReadTx,
79    selected: &BTreeSet<String>,
80) -> Result<Vec<NodeRelationProjection>, PortError> {
81    let mut rows = Vec::new();
82    for source in selected {
83        for relation in outgoing_rows(tx, source)? {
84            if selected.contains(&relation.target_node_id) {
85                rows.push(relation);
86            }
87        }
88    }
89    Ok(rows)
90}
91
92fn selected_projections(
93    tx: &dyn ReadTx,
94    selected: &BTreeSet<String>,
95    root_node_id: &str,
96) -> Result<Vec<NodeProjection>, PortError> {
97    let mut projections = Vec::new();
98    for node_id in selected {
99        if node_id == root_node_id {
100            continue;
101        }
102        if let Some(projection) = load_node(tx, node_id)? {
103            projections.push(projection);
104        }
105    }
106    Ok(projections)
107}
108
109fn shortest_outward_path(
110    tx: &dyn ReadTx,
111    root_node_id: &str,
112    target_node_id: &str,
113) -> Result<Option<Vec<String>>, PortError> {
114    let mut predecessors = BTreeMap::<String, String>::new();
115    let mut visited = BTreeSet::from([root_node_id.to_string()]);
116    let mut frontier = VecDeque::from([root_node_id.to_string()]);
117
118    while let Some(node_id) = frontier.pop_front() {
119        for target in outgoing_targets(tx, &node_id)? {
120            if !visited.insert(target.clone()) {
121                continue;
122            }
123            predecessors.insert(target.clone(), node_id.clone());
124            if target == target_node_id {
125                let mut path = vec![target.clone()];
126                let mut current = target.as_str();
127                while let Some(previous) = predecessors.get(current) {
128                    path.push(previous.clone());
129                    current = previous;
130                }
131                path.reverse();
132                return Ok(Some(path));
133            }
134            frontier.push_back(target);
135        }
136    }
137
138    Ok(None)
139}
140
141impl GraphNeighborhoodReader for EmbeddedKernelStore {
142    async fn load_nodes_batch(
143        &self,
144        node_ids: Vec<String>,
145    ) -> Result<Vec<Option<NodeProjection>>, PortError> {
146        self.run(move |store| {
147            let tx = store.begin_read()?;
148            node_ids
149                .iter()
150                .map(|id| load_node(tx.as_ref(), id))
151                .collect()
152        })
153        .await
154    }
155
156    async fn load_bounded_trace(
157        &self,
158        request: &kmp_domain::TraceSearchRequest,
159    ) -> Result<kmp_domain::TraceSearchResult, PortError> {
160        let request = request.clone();
161        self.run(move |store| {
162            let tx = store.begin_read()?;
163            kmp_domain::bounded_trace_search(
164                &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
165                &request,
166            )
167        })
168        .await
169    }
170
171    async fn load_evidence_paths(
172        &self,
173        request: &kmp_domain::EvidencePathRequest,
174    ) -> Result<kmp_domain::EvidencePathResult, PortError> {
175        let request = request.clone();
176        self.run(move |store| {
177            let tx = store.begin_read()?;
178            kmp_domain::search_evidence_paths(
179                &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
180                &request,
181            )
182        })
183        .await
184    }
185
186    async fn load_neighborhood(
187        &self,
188        root_node_id: &str,
189        depth: u32,
190    ) -> Result<Option<NodeNeighborhood>, PortError> {
191        self.load_scoped_neighborhood(&NeighborhoodRequest::new(root_node_id, depth))
192            .await
193    }
194
195    async fn load_scoped_neighborhood(
196        &self,
197        request: &NeighborhoodRequest,
198    ) -> Result<Option<NodeNeighborhood>, PortError> {
199        let request = request.clone();
200        let root_node_id = request.root_node_id().to_string();
201        self.run(move |store| {
202            let tx = store.begin_read()?;
203            let tx = tx.as_ref();
204
205            let Some(root) = load_node(tx, &root_node_id)? else {
206                return Ok(None);
207            };
208
209            let reachable = reachable_outward(tx, &request)?;
210            // Mirrors the Neo4j neighborhood query: an empty neighborhood
211            // reports no relations, even for self-referential root edges.
212            let relation_rows = if reachable.is_empty() {
213                Vec::new()
214            } else {
215                let mut selected = reachable.clone();
216                selected.insert(root_node_id.clone());
217                relations_among(tx, &selected)?
218            };
219
220            Ok(Some(NodeNeighborhood {
221                neighbors: selected_projections(tx, &reachable, &root_node_id)?,
222                relations: relation_rows,
223                root,
224            }))
225        })
226        .await
227    }
228
229    async fn load_context_path(
230        &self,
231        root_node_id: &str,
232        target_node_id: &str,
233        subtree_depth: u32,
234    ) -> Result<Option<ContextPathNeighborhood>, PortError> {
235        let root_node_id = root_node_id.to_string();
236        let target_node_id = target_node_id.to_string();
237        self.run(move |store| {
238            let tx = store.begin_read()?;
239            let tx = tx.as_ref();
240
241            let Some(root) = load_node(tx, &root_node_id)? else {
242                return Ok(None);
243            };
244            if load_node(tx, &target_node_id)?.is_none() {
245                return Ok(None);
246            }
247            let Some(path_node_ids) = shortest_outward_path(tx, &root_node_id, &target_node_id)?
248            else {
249                return Ok(None);
250            };
251
252            let mut selected = path_node_ids.iter().cloned().collect::<BTreeSet<_>>();
253            selected.insert(target_node_id.clone());
254            selected.extend(reachable_outward(
255                tx,
256                &NeighborhoodRequest::new(&target_node_id, subtree_depth),
257            )?);
258
259            Ok(Some(ContextPathNeighborhood {
260                neighbors: selected_projections(tx, &selected, &root_node_id)?,
261                relations: relations_among(tx, &selected)?,
262                path_node_ids,
263                root,
264            }))
265        })
266        .await
267    }
268}
269
270impl NodeRelationshipReader for EmbeddedKernelStore {
271    async fn load_node_relationships(
272        &self,
273        node_id: &str,
274    ) -> Result<Option<NodeRelationships>, PortError> {
275        let node_id = node_id.to_string();
276        self.run(move |store| {
277            let tx = store.begin_read()?;
278            let tx = tx.as_ref();
279            if load_node(tx, &node_id)?.is_none() {
280                return Ok(None);
281            }
282
283            let mut incoming = Vec::new();
284            for ((target, source, relation_type), _) in
285                tx.scan_str3_by_first(Table::RelationsByTarget, &node_id)?
286            {
287                let Some(raw) = tx.get(
288                    Table::Relations,
289                    Key::Str3(&source, &target, &relation_type),
290                )?
291                else {
292                    return Err(PortError::InvalidState(format!(
293                        "embedded store adjacency index points at missing relation \
294                         `{source}` -> `{target}` ({relation_type})"
295                    )));
296                };
297                incoming.push(NodeRelationProjection {
298                    explanation: decode_explanation(&raw)?,
299                    source_node_id: source,
300                    target_node_id: target,
301                    relation_type,
302                });
303            }
304
305            Ok(Some(NodeRelationships {
306                incoming,
307                outgoing: outgoing_rows(tx, &node_id)?,
308            }))
309        })
310        .await
311    }
312}
313
314impl MemoryAboutIndexReader for EmbeddedKernelStore {
315    async fn list_memory_abouts(&self) -> Result<Vec<String>, PortError> {
316        self.run(|store| {
317            let tx = store.begin_read()?;
318            Ok(tx
319                .scan_str(Table::Anchors)?
320                .into_iter()
321                .map(|(anchor, _)| anchor)
322                .collect())
323        })
324        .await
325    }
326
327    async fn list_memory_abouts_by_dimensions(
328        &self,
329        dimension_ids: &[String],
330    ) -> Result<Vec<String>, PortError> {
331        let dimension_ids = dimension_ids.to_vec();
332        self.run(move |store| {
333            let tx = store.begin_read()?;
334            let tx = tx.as_ref();
335
336            let mut abouts = BTreeSet::new();
337            for (anchor, _) in tx.scan_str(Table::Anchors)? {
338                let is_anchor = load_node(tx, &anchor)?
339                    .is_some_and(|node| node.node_kind == MEMORY_ANCHOR_KIND);
340                if !is_anchor {
341                    continue;
342                }
343                for relation in outgoing_rows(tx, &anchor)? {
344                    if relation.relation_type != "has_dimension" {
345                        continue;
346                    }
347                    let matches =
348                        load_node(tx, &relation.target_node_id)?.is_some_and(|dimension| {
349                            dimension.node_kind == "memory_dimension"
350                                && dimension_ids.iter().any(|dimension_id| {
351                                    dimension.node_id == *dimension_id
352                                        || kmp_domain::MemoryDimensionIdentity::parse(&dimension.node_id)
353                                            .is_some_and(|identity| identity.dimension_id() == dimension_id)
354                                        // A selection names dimensions by kind
355                                        // (`incident`) as readily as by id
356                                        // (`incident:north-outage`); the filter
357                                        // that follows reads kinds, so the
358                                        // index that picks the abouts must too.
359                                        || dimension.properties.get("dimension_kind")
360                                            == Some(dimension_id)
361                                })
362                        });
363                    if matches {
364                        abouts.insert(anchor.clone());
365                        break;
366                    }
367                }
368            }
369            Ok(abouts.into_iter().collect())
370        })
371        .await
372    }
373}