Skip to main content

kmp_adapter_embedded/adapter/
graph_read.rs

1use std::collections::{BTreeMap, BTreeSet, VecDeque};
2
3use kmp_domain::{
4    ContextPathNeighborhood, GraphNeighborhoodReader, NeighborhoodRequest, NodeNeighborhood,
5    NodeProjection, NodeRelationProjection, NodeRelationshipReader, NodeRelationships, PortError,
6};
7
8use super::engine::{Key, ReadTx, Table};
9use super::serdes::{NodeRecord, decode, decode_explanation};
10use super::store::EmbeddedKernelStore;
11#[path = "node_admission_header.rs"]
12mod node_admission_header;
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 EmbeddedKernelStore {
142    async fn read_catalogue_neighborhood(
143        &self,
144        request: &NeighborhoodRequest,
145        headers: bool,
146    ) -> Result<Option<NodeNeighborhood>, PortError> {
147        let request = request.clone();
148        self.run(move |store| {
149            let tx = store.begin_read()?;
150            let tx = tx.as_ref();
151            let read = |id: &str| {
152                if headers {
153                    node_admission_header::read(tx, id)
154                } else {
155                    load_node(tx, id)
156                }
157            };
158            let Some(root) = read(request.root_node_id())? else {
159                return Ok(None);
160            };
161            let reachable = reachable_outward(tx, &request)?;
162            let relations = if reachable.is_empty() {
163                Vec::new()
164            } else {
165                let mut selected = reachable.clone();
166                selected.insert(request.root_node_id().to_string());
167                relations_among(tx, &selected)?
168            };
169            let mut neighbors = Vec::with_capacity(reachable.len());
170            for id in reachable {
171                if let Some(node) = read(&id)? {
172                    neighbors.push(node);
173                }
174            }
175            Ok(Some(NodeNeighborhood {
176                root,
177                neighbors,
178                relations,
179            }))
180        })
181        .await
182    }
183}
184
185impl GraphNeighborhoodReader for EmbeddedKernelStore {
186    async fn graph_read_revision(
187        &self,
188    ) -> Result<Option<kmp_domain::GraphReadRevision>, PortError> {
189        Ok(self.read_revision())
190    }
191
192    async fn load_nodes_batch(
193        &self,
194        node_ids: Vec<String>,
195    ) -> Result<Vec<Option<NodeProjection>>, PortError> {
196        self.run(move |store| {
197            let tx = store.begin_read()?;
198            node_ids
199                .iter()
200                .map(|id| load_node(tx.as_ref(), id))
201                .collect()
202        })
203        .await
204    }
205
206    async fn load_bounded_trace(
207        &self,
208        request: &kmp_domain::TraceSearchRequest,
209    ) -> Result<kmp_domain::TraceSearchResult, PortError> {
210        let request = request.clone();
211        self.run(move |store| {
212            let tx = store.begin_read()?;
213            kmp_domain::bounded_trace_search(
214                &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
215                &request,
216            )
217        })
218        .await
219    }
220
221    async fn load_memory_nodes(
222        &self,
223        request: &kmp_domain::MemoryNodesRequest,
224    ) -> Result<kmp_domain::MemoryNodesResult, PortError> {
225        let request = request.clone();
226        self.run(move |store| {
227            let pinned = store.pin_snapshot()?;
228            let revision = pinned.read_revision().ok_or_else(|| {
229                PortError::Conflict(
230                    "store changed while pinning the node batch; retry the original read".into(),
231                )
232            })?;
233            if request
234                .expect_snapshot
235                .as_ref()
236                .is_some_and(|expected| expected != &revision)
237            {
238                return Err(PortError::Conflict(
239                    "node batch snapshot changed; discard previous batches and restart the focus"
240                        .into(),
241                ));
242            }
243            let tx = pinned.begin_read()?;
244            let mut result = kmp_domain::read_memory_nodes(
245                &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
246                &request,
247            )?;
248            result.snapshot = Some(revision);
249            Ok(result)
250        })
251        .await
252    }
253
254    async fn load_evidence_paths(
255        &self,
256        request: &kmp_domain::EvidencePathRequest,
257    ) -> Result<kmp_domain::EvidencePathResult, PortError> {
258        let request = request.clone();
259        self.run(move |store| {
260            let tx = store.begin_read()?;
261            kmp_domain::search_evidence_paths(
262                &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
263                &request,
264            )
265        })
266        .await
267    }
268
269    async fn load_neighborhood(
270        &self,
271        root_node_id: &str,
272        depth: u32,
273    ) -> Result<Option<NodeNeighborhood>, PortError> {
274        self.load_scoped_neighborhood(&NeighborhoodRequest::new(root_node_id, depth))
275            .await
276    }
277
278    async fn load_scoped_neighborhood(
279        &self,
280        request: &NeighborhoodRequest,
281    ) -> Result<Option<NodeNeighborhood>, PortError> {
282        self.read_catalogue_neighborhood(request, false).await
283    }
284
285    async fn load_neighborhood_headers(
286        &self,
287        request: &NeighborhoodRequest,
288    ) -> Result<Option<NodeNeighborhood>, PortError> {
289        self.read_catalogue_neighborhood(request, true).await
290    }
291
292    async fn load_context_path(
293        &self,
294        root_node_id: &str,
295        target_node_id: &str,
296        subtree_depth: u32,
297    ) -> Result<Option<ContextPathNeighborhood>, PortError> {
298        let root_node_id = root_node_id.to_string();
299        let target_node_id = target_node_id.to_string();
300        self.run(move |store| {
301            let tx = store.begin_read()?;
302            let tx = tx.as_ref();
303
304            let Some(root) = load_node(tx, &root_node_id)? else {
305                return Ok(None);
306            };
307            if load_node(tx, &target_node_id)?.is_none() {
308                return Ok(None);
309            }
310            let Some(path_node_ids) = shortest_outward_path(tx, &root_node_id, &target_node_id)?
311            else {
312                return Ok(None);
313            };
314
315            let mut selected = path_node_ids.iter().cloned().collect::<BTreeSet<_>>();
316            selected.insert(target_node_id.clone());
317            selected.extend(reachable_outward(
318                tx,
319                &NeighborhoodRequest::new(&target_node_id, subtree_depth),
320            )?);
321
322            Ok(Some(ContextPathNeighborhood {
323                neighbors: selected_projections(tx, &selected, &root_node_id)?,
324                relations: relations_among(tx, &selected)?,
325                path_node_ids,
326                root,
327            }))
328        })
329        .await
330    }
331}
332
333impl NodeRelationshipReader for EmbeddedKernelStore {
334    async fn load_node_relationships(
335        &self,
336        node_id: &str,
337    ) -> Result<Option<NodeRelationships>, PortError> {
338        let node_id = node_id.to_string();
339        self.run(move |store| {
340            let tx = store.begin_read()?;
341            let tx = tx.as_ref();
342            if load_node(tx, &node_id)?.is_none() {
343                return Ok(None);
344            }
345
346            let mut incoming = Vec::new();
347            for ((target, source, relation_type), _) in
348                tx.scan_str3_by_first(Table::RelationsByTarget, &node_id)?
349            {
350                let Some(raw) = tx.get(
351                    Table::Relations,
352                    Key::Str3(&source, &target, &relation_type),
353                )?
354                else {
355                    return Err(PortError::InvalidState(format!(
356                        "embedded store adjacency index points at missing relation \
357                         `{source}` -> `{target}` ({relation_type})"
358                    )));
359                };
360                incoming.push(NodeRelationProjection {
361                    explanation: decode_explanation(&raw)?,
362                    source_node_id: source,
363                    target_node_id: target,
364                    relation_type,
365                });
366            }
367
368            Ok(Some(NodeRelationships {
369                incoming,
370                outgoing: outgoing_rows(tx, &node_id)?,
371            }))
372        })
373        .await
374    }
375}