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
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 GraphNeighborhoodReader for EmbeddedKernelStore {
142 async fn load_nodes_batch(
143 &self,
144 node_ids: Vec<String>,
145 ) -> Result<Vec<Option<NodeProjection>>, PortError> {
146 self.run(move |store| {
147 let tx = store.begin_read()?;
148 node_ids
149 .iter()
150 .map(|id| load_node(tx.as_ref(), id))
151 .collect()
152 })
153 .await
154 }
155
156 async fn load_bounded_trace(
157 &self,
158 request: &kmp_domain::TraceSearchRequest,
159 ) -> Result<kmp_domain::TraceSearchResult, PortError> {
160 let request = request.clone();
161 self.run(move |store| {
162 let tx = store.begin_read()?;
163 kmp_domain::bounded_trace_search(
164 &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
165 &request,
166 )
167 })
168 .await
169 }
170
171 async fn load_evidence_paths(
172 &self,
173 request: &kmp_domain::EvidencePathRequest,
174 ) -> Result<kmp_domain::EvidencePathResult, PortError> {
175 let request = request.clone();
176 self.run(move |store| {
177 let tx = store.begin_read()?;
178 kmp_domain::search_evidence_paths(
179 &super::trace_snapshot::TraceSnapshot(tx.as_ref()),
180 &request,
181 )
182 })
183 .await
184 }
185
186 async fn load_neighborhood(
187 &self,
188 root_node_id: &str,
189 depth: u32,
190 ) -> Result<Option<NodeNeighborhood>, PortError> {
191 self.load_scoped_neighborhood(&NeighborhoodRequest::new(root_node_id, depth))
192 .await
193 }
194
195 async fn load_scoped_neighborhood(
196 &self,
197 request: &NeighborhoodRequest,
198 ) -> Result<Option<NodeNeighborhood>, PortError> {
199 let request = request.clone();
200 let root_node_id = request.root_node_id().to_string();
201 self.run(move |store| {
202 let tx = store.begin_read()?;
203 let tx = tx.as_ref();
204
205 let Some(root) = load_node(tx, &root_node_id)? else {
206 return Ok(None);
207 };
208
209 let reachable = reachable_outward(tx, &request)?;
210 let relation_rows = if reachable.is_empty() {
213 Vec::new()
214 } else {
215 let mut selected = reachable.clone();
216 selected.insert(root_node_id.clone());
217 relations_among(tx, &selected)?
218 };
219
220 Ok(Some(NodeNeighborhood {
221 neighbors: selected_projections(tx, &reachable, &root_node_id)?,
222 relations: relation_rows,
223 root,
224 }))
225 })
226 .await
227 }
228
229 async fn load_context_path(
230 &self,
231 root_node_id: &str,
232 target_node_id: &str,
233 subtree_depth: u32,
234 ) -> Result<Option<ContextPathNeighborhood>, PortError> {
235 let root_node_id = root_node_id.to_string();
236 let target_node_id = target_node_id.to_string();
237 self.run(move |store| {
238 let tx = store.begin_read()?;
239 let tx = tx.as_ref();
240
241 let Some(root) = load_node(tx, &root_node_id)? else {
242 return Ok(None);
243 };
244 if load_node(tx, &target_node_id)?.is_none() {
245 return Ok(None);
246 }
247 let Some(path_node_ids) = shortest_outward_path(tx, &root_node_id, &target_node_id)?
248 else {
249 return Ok(None);
250 };
251
252 let mut selected = path_node_ids.iter().cloned().collect::<BTreeSet<_>>();
253 selected.insert(target_node_id.clone());
254 selected.extend(reachable_outward(
255 tx,
256 &NeighborhoodRequest::new(&target_node_id, subtree_depth),
257 )?);
258
259 Ok(Some(ContextPathNeighborhood {
260 neighbors: selected_projections(tx, &selected, &root_node_id)?,
261 relations: relations_among(tx, &selected)?,
262 path_node_ids,
263 root,
264 }))
265 })
266 .await
267 }
268}
269
270impl NodeRelationshipReader for EmbeddedKernelStore {
271 async fn load_node_relationships(
272 &self,
273 node_id: &str,
274 ) -> Result<Option<NodeRelationships>, PortError> {
275 let node_id = node_id.to_string();
276 self.run(move |store| {
277 let tx = store.begin_read()?;
278 let tx = tx.as_ref();
279 if load_node(tx, &node_id)?.is_none() {
280 return Ok(None);
281 }
282
283 let mut incoming = Vec::new();
284 for ((target, source, relation_type), _) in
285 tx.scan_str3_by_first(Table::RelationsByTarget, &node_id)?
286 {
287 let Some(raw) = tx.get(
288 Table::Relations,
289 Key::Str3(&source, &target, &relation_type),
290 )?
291 else {
292 return Err(PortError::InvalidState(format!(
293 "embedded store adjacency index points at missing relation \
294 `{source}` -> `{target}` ({relation_type})"
295 )));
296 };
297 incoming.push(NodeRelationProjection {
298 explanation: decode_explanation(&raw)?,
299 source_node_id: source,
300 target_node_id: target,
301 relation_type,
302 });
303 }
304
305 Ok(Some(NodeRelationships {
306 incoming,
307 outgoing: outgoing_rows(tx, &node_id)?,
308 }))
309 })
310 .await
311 }
312}
313
314impl MemoryAboutIndexReader for EmbeddedKernelStore {
315 async fn list_memory_abouts(&self) -> Result<Vec<String>, PortError> {
316 self.run(|store| {
317 let tx = store.begin_read()?;
318 Ok(tx
319 .scan_str(Table::Anchors)?
320 .into_iter()
321 .map(|(anchor, _)| anchor)
322 .collect())
323 })
324 .await
325 }
326
327 async fn list_memory_abouts_by_dimensions(
328 &self,
329 dimension_ids: &[String],
330 ) -> Result<Vec<String>, PortError> {
331 let dimension_ids = dimension_ids.to_vec();
332 self.run(move |store| {
333 let tx = store.begin_read()?;
334 let tx = tx.as_ref();
335
336 let mut abouts = BTreeSet::new();
337 for (anchor, _) in tx.scan_str(Table::Anchors)? {
338 let is_anchor = load_node(tx, &anchor)?
339 .is_some_and(|node| node.node_kind == MEMORY_ANCHOR_KIND);
340 if !is_anchor {
341 continue;
342 }
343 for relation in outgoing_rows(tx, &anchor)? {
344 if relation.relation_type != "has_dimension" {
345 continue;
346 }
347 let matches =
348 load_node(tx, &relation.target_node_id)?.is_some_and(|dimension| {
349 dimension.node_kind == "memory_dimension"
350 && dimension_ids.iter().any(|dimension_id| {
351 dimension.node_id == *dimension_id
352 || kmp_domain::MemoryDimensionIdentity::parse(&dimension.node_id)
353 .is_some_and(|identity| identity.dimension_id() == dimension_id)
354 || dimension.properties.get("dimension_kind")
360 == Some(dimension_id)
361 })
362 });
363 if matches {
364 abouts.insert(anchor.clone());
365 break;
366 }
367 }
368 }
369 Ok(abouts.into_iter().collect())
370 })
371 .await
372 }
373}