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_evidence_paths(
224 &self,
225 request: &kmp_domain::EvidencePathRequest,
226 ) -> Result<kmp_domain::EvidencePathResult, PortError> {
227 let request = request.clone();
228 self.run(move |store| {
229 let tx = store.begin_read()?;
230 kmp_domain::search_evidence_paths(
231 &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
232 &request,
233 )
234 })
235 .await
236 }
237
238 async fn load_neighborhood(
239 &self,
240 root_node_id: &str,
241 depth: u32,
242 ) -> Result<Option<NodeNeighborhood>, PortError> {
243 self.load_scoped_neighborhood(&NeighborhoodRequest::new(root_node_id, depth))
244 .await
245 }
246
247 async fn load_scoped_neighborhood(
248 &self,
249 request: &NeighborhoodRequest,
250 ) -> Result<Option<NodeNeighborhood>, PortError> {
251 self.read_catalogue_neighborhood(request, false).await
252 }
253
254 async fn load_neighborhood_headers(
255 &self,
256 request: &NeighborhoodRequest,
257 ) -> Result<Option<NodeNeighborhood>, PortError> {
258 self.read_catalogue_neighborhood(request, true).await
259 }
260
261 async fn load_context_path(
262 &self,
263 root_node_id: &str,
264 target_node_id: &str,
265 subtree_depth: u32,
266 ) -> Result<Option<ContextPathNeighborhood>, PortError> {
267 let root_node_id = root_node_id.to_string();
268 let target_node_id = target_node_id.to_string();
269 self.run(move |store| {
270 let tx = store.begin_read()?;
271 let tx = tx.as_ref();
272
273 let Some(root) = load_node(tx, &root_node_id)? else {
274 return Ok(None);
275 };
276 if load_node(tx, &target_node_id)?.is_none() {
277 return Ok(None);
278 }
279 let Some(path_node_ids) = shortest_outward_path(tx, &root_node_id, &target_node_id)?
280 else {
281 return Ok(None);
282 };
283
284 let mut selected = path_node_ids.iter().cloned().collect::<BTreeSet<_>>();
285 selected.insert(target_node_id.clone());
286 selected.extend(reachable_outward(
287 tx,
288 &NeighborhoodRequest::new(&target_node_id, subtree_depth),
289 )?);
290
291 Ok(Some(ContextPathNeighborhood {
292 neighbors: selected_projections(tx, &selected, &root_node_id)?,
293 relations: relations_among(tx, &selected)?,
294 path_node_ids,
295 root,
296 }))
297 })
298 .await
299 }
300}
301
302impl NodeRelationshipReader for EmbeddedKernelStore {
303 async fn load_node_relationships(
304 &self,
305 node_id: &str,
306 ) -> Result<Option<NodeRelationships>, PortError> {
307 let node_id = node_id.to_string();
308 self.run(move |store| {
309 let tx = store.begin_read()?;
310 let tx = tx.as_ref();
311 if load_node(tx, &node_id)?.is_none() {
312 return Ok(None);
313 }
314
315 let mut incoming = Vec::new();
316 for ((target, source, relation_type), _) in
317 tx.scan_str3_by_first(Table::RelationsByTarget, &node_id)?
318 {
319 let Some(raw) = tx.get(
320 Table::Relations,
321 Key::Str3(&source, &target, &relation_type),
322 )?
323 else {
324 return Err(PortError::InvalidState(format!(
325 "embedded store adjacency index points at missing relation \
326 `{source}` -> `{target}` ({relation_type})"
327 )));
328 };
329 incoming.push(NodeRelationProjection {
330 explanation: decode_explanation(&raw)?,
331 source_node_id: source,
332 target_node_id: target,
333 relation_type,
334 });
335 }
336
337 Ok(Some(NodeRelationships {
338 incoming,
339 outgoing: outgoing_rows(tx, &node_id)?,
340 }))
341 })
342 .await
343 }
344}
345
346impl MemoryAboutIndexReader for EmbeddedKernelStore {
347 async fn list_memory_abouts(&self) -> Result<Vec<String>, PortError> {
348 self.run(|store| {
349 let tx = store.begin_read()?;
350 Ok(tx
351 .scan_str(Table::Anchors)?
352 .into_iter()
353 .map(|(anchor, _)| anchor)
354 .collect())
355 })
356 .await
357 }
358
359 async fn list_memory_abouts_by_dimensions(
360 &self,
361 dimension_ids: &[String],
362 ) -> Result<Vec<String>, PortError> {
363 let dimension_ids = dimension_ids.to_vec();
364 self.run(move |store| {
365 let tx = store.begin_read()?;
366 let tx = tx.as_ref();
367
368 let mut abouts = BTreeSet::new();
369 for (anchor, _) in tx.scan_str(Table::Anchors)? {
370 let is_anchor = load_node(tx, &anchor)?
371 .is_some_and(|node| node.node_kind == MEMORY_ANCHOR_KIND);
372 if !is_anchor {
373 continue;
374 }
375 for relation in outgoing_rows(tx, &anchor)? {
376 if relation.relation_type != "has_dimension" {
377 continue;
378 }
379 let matches =
380 load_node(tx, &relation.target_node_id)?.is_some_and(|dimension| {
381 dimension.node_kind == "memory_dimension"
382 && dimension_ids.iter().any(|dimension_id| {
383 dimension.node_id == *dimension_id
384 || kmp_domain::MemoryDimensionIdentity::parse(&dimension.node_id)
385 .is_some_and(|identity| identity.dimension_id() == dimension_id)
386 || dimension.properties.get("dimension_kind")
392 == Some(dimension_id)
393 })
394 });
395 if matches {
396 abouts.insert(anchor.clone());
397 break;
398 }
399 }
400 }
401 Ok(abouts.into_iter().collect())
402 })
403 .await
404 }
405}