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