kmp_adapter_embedded/adapter/
graph_read.rs1use 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 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}