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