1use crate::flows::{collect_entrypoints, walk_calls};
20use crate::components::prov_rank;
21use crate::RealityGraph;
22use scc_core::{
23 entity_id, FlowEdge, FlowEdgeKind, FlowGraph, FlowKind, FlowNode, Provenance,
24};
25use scc_store::Store;
26use std::collections::{BTreeMap, BTreeSet, HashMap};
27
28type NodeKey = (String, String);
30
31fn op_of(graph: &RealityGraph, sym: &str) -> String {
32 graph
33 .entities
34 .get(sym)
35 .map(|e| e.name.clone())
36 .unwrap_or_else(|| sym.to_string())
37}
38
39struct EdgeTable {
42 edges: Vec<FlowEdge>,
43 seen: BTreeSet<(u32, u32, String)>,
44 in_degree: HashMap<u32, u32>,
45}
46
47impl EdgeTable {
48 fn new() -> EdgeTable {
49 EdgeTable {
50 edges: Vec::new(),
51 seen: BTreeSet::new(),
52 in_degree: HashMap::new(),
53 }
54 }
55
56 fn push(
57 &mut self,
58 from: u32,
59 to: u32,
60 kind: FlowEdgeKind,
61 condition: Option<String>,
62 prov: Provenance,
63 evidence: Vec<String>,
64 ) {
65 let k = (from, to, format!("{kind:?}"));
66 if let Some(existing) = self
72 .edges
73 .iter_mut()
74 .find(|e| e.from == from && e.to == to && format!("{:?}", e.kind) == format!("{kind:?}"))
75 {
76 let have = existing.provenance.unwrap_or(Provenance::Inferred);
77 if prov_rank(prov) > prov_rank(have) {
78 existing.provenance = Some(prov);
79 existing.confidence = prov.default_confidence();
80 }
81 if existing.condition.is_none() {
82 existing.condition = condition;
83 }
84 for e in evidence {
85 if !existing.evidence.contains(&e) {
86 existing.evidence.push(e);
87 }
88 }
89 return;
90 }
91 self.seen.insert(k);
92 *self.in_degree.entry(to).or_insert(0) += 1;
93 self.edges.push(FlowEdge {
94 from,
95 to,
96 kind,
97 condition,
98 provenance: Some(prov),
99 confidence: prov.default_confidence(),
100 evidence,
101 });
102 }
103}
104
105struct NodeTable {
108 nodes: Vec<FlowNode>,
109 by_key: HashMap<NodeKey, u32>,
110}
111
112impl NodeTable {
113 fn new() -> NodeTable {
114 NodeTable {
115 nodes: Vec::new(),
116 by_key: HashMap::new(),
117 }
118 }
119
120 fn get(&mut self, key: &NodeKey) -> u32 {
121 if let Some(id) = self.by_key.get(key) {
122 return *id;
123 }
124 let id = self.nodes.len() as u32;
125 self.by_key.insert(key.clone(), id);
126 self.nodes.push(FlowNode {
127 id,
128 actor: key.0.clone(),
129 operation: key.1.clone(),
130 evidence: Vec::new(),
131 });
132 id
133 }
134
135 fn lookup(&self, key: &NodeKey) -> Option<u32> {
136 self.by_key.get(key).copied()
137 }
138}
139
140pub fn compile_flow_graphs(
147 graph: &RealityGraph,
148 store: &Store,
149 intent: &[(String, serde_json::Value)],
150 symbol_comp: &HashMap<String, String>,
151) -> crate::Result<Vec<FlowGraph>> {
152 let entrypoints = collect_entrypoints(graph, store, intent);
153 let mut out: Vec<FlowGraph> = Vec::new();
154
155 for ep in entrypoints {
156 if ep.symbol_id.is_empty() {
157 continue;
158 }
159 let paths = walk_calls(graph, &ep.symbol_id);
160 let mut table = NodeTable::new();
161 let mut node_evidence: BTreeMap<NodeKey, (Provenance, BTreeSet<String>)> = BTreeMap::new();
162
163 let entry_actor = symbol_comp
165 .get(&ep.symbol_id)
166 .cloned()
167 .unwrap_or_else(|| "component:root".into());
168 let entry_key = (entry_actor.clone(), op_of(graph, &ep.symbol_id));
169 let entry_id = table.get(&entry_key);
170
171 let mut successors: BTreeMap<String, BTreeSet<String>> = BTreeMap::new();
174 let mut call_rel: HashMap<(String, String), (Provenance, Vec<String>)> = HashMap::new();
175 let mut syms: BTreeSet<String> = paths
176 .iter()
177 .flatten()
178 .cloned()
179 .collect();
180 syms.insert(ep.symbol_id.clone());
181 for sym in &syms {
182 for r in graph.out_pred(sym, scc_core::predicates::CALLS) {
183 if matches!(r.provenance, Provenance::Extracted | Provenance::Resolved) {
184 successors
185 .entry(sym.clone())
186 .or_default()
187 .insert(r.object.clone());
188 let key = (sym.clone(), r.object.clone());
189 let entry = call_rel.entry(key).or_insert((r.provenance, Vec::new()));
190 if prov_rank(r.provenance) > prov_rank(entry.0) {
191 entry.0 = r.provenance;
192 }
193 for e in &r.evidence {
194 if !entry.1.contains(e) {
195 entry.1.push(e.clone());
196 }
197 }
198 }
199 }
200 }
201
202 let mut edges = EdgeTable::new();
203
204 for sym in &syms {
208 let key = (
209 symbol_comp
210 .get(sym)
211 .cloned()
212 .unwrap_or_else(|| "component:root".into()),
213 op_of(graph, sym),
214 );
215 table.get(&key);
216 }
217
218 for (sym, targets) in &successors {
219 let key = (
220 symbol_comp.get(sym).cloned().unwrap_or_else(|| "component:root".into()),
221 op_of(graph, sym),
222 );
223 let from = table.lookup(&key).unwrap_or(entry_id);
224 if let Some(e) = graph.entities.get(sym) {
226 if let Some(nk) = table
227 .by_key
228 .iter()
229 .find(|(_, v)| **v == from)
230 .map(|(k, _)| k.clone())
231 {
232 node_evidence
233 .entry(nk)
234 .or_insert((Provenance::Resolved, BTreeSet::new()))
235 .1
236 .extend(e.evidence.clone());
237 }
238 }
239 let cond_attr = |name: &str| -> Option<serde_json::Value> {
246 graph.entities.get(sym).and_then(|e| e.attributes.get(name).cloned())
247 };
248 let conditional_calls: std::collections::HashSet<String> = cond_attr("conditional_calls")
249 .and_then(|v| serde_json::from_value(v).ok())
250 .unwrap_or_default();
251 let call_blocks: BTreeMap<String, String> = cond_attr("call_blocks")
252 .and_then(|v| serde_json::from_value(v).ok())
253 .unwrap_or_default();
254 let call_order: BTreeMap<String, u32> = cond_attr("call_order")
255 .and_then(|v| serde_json::from_value(v).ok())
256 .unwrap_or_default();
257 let awaited_calls: std::collections::BTreeSet<String> = cond_attr("awaited_calls")
258 .and_then(|v| serde_json::from_value(v).ok())
259 .unwrap_or_default();
260 let match_callee = |attr: &str, t_op: &str| -> bool {
261 attr == t_op || attr.ends_with(&format!(".{t_op}"))
262 };
263 let mut ordered: Vec<(u32, &String)> = targets
268 .iter()
269 .map(|t| {
270 let t_op = op_of(graph, t);
271 let order = call_order
272 .iter()
273 .filter(|(c, _)| match_callee(c, &t_op))
274 .map(|(_, o)| *o)
275 .min()
276 .unwrap_or(u32::MAX);
277 (order, t)
278 })
279 .collect();
280 ordered.sort_by(|a, b| (a.0, a.1).cmp(&(b.0, b.1)));
281 for (_, t) in ordered {
282 let actor = symbol_comp
283 .get(t)
284 .cloned()
285 .unwrap_or_else(|| "component:root".into());
286 let tkey = (actor, op_of(graph, t));
287 let to = table.get(&tkey);
288 let (prov, evidence) = call_rel
289 .get(&(sym.clone(), t.clone()))
290 .cloned()
291 .unwrap_or((Provenance::Resolved, Vec::new()));
292 let t_op = op_of(graph, t);
293 let block = call_blocks
294 .iter()
295 .find(|(c, _)| match_callee(c, &t_op))
296 .map(|(_, b)| b.clone());
297 let fallback = conditional_calls
298 .iter()
299 .any(|c| match_callee(c, &t_op));
300 let is_branch = block.is_some() || fallback;
301 edges.push(
302 from,
303 to,
304 if is_branch {
305 FlowEdgeKind::Branch
306 } else {
307 FlowEdgeKind::Next
308 },
309 if let Some(b) = &block {
310 Some(b.clone())
311 } else if is_branch {
312 Some(format!("conditional: {t_op}"))
313 } else {
314 None
315 },
316 prov,
317 evidence,
318 );
319 if awaited_calls.iter().any(|c| match_callee(c, &t_op)) {
322 edges.push(
323 from,
324 to,
325 FlowEdgeKind::Async,
326 Some(format!("awaited: {t_op}")),
327 prov,
328 Vec::new(),
329 );
330 }
331 }
332 }
333
334 let mut retry_syms: Vec<(u32, String, String)> = Vec::new();
340 for (sym, comp) in symbol_comp {
341 if let Some(e) = graph.entities.get(sym) {
342 if let Some(rp) = e.attributes.get("retry_policy").and_then(|v| v.as_str()) {
343 let key = (comp.clone(), op_of(graph, sym));
344 if let Some(id) = table.lookup(&key) {
345 retry_syms.push((id, sym.clone(), rp.to_string()));
346 }
347 }
348 }
349 }
350 retry_syms.sort();
351 for (id, sym, rp) in retry_syms {
352 let evidence: Vec<String> = graph
353 .entities
354 .get(&sym)
355 .map(|e| e.evidence.clone())
356 .unwrap_or_default();
357 edges.push(
358 id,
359 id,
360 FlowEdgeKind::Retry,
361 Some(format!("attempt ({rp})")),
362 Provenance::Extracted,
363 evidence.clone(),
364 );
365 let succ: Vec<u32> = edges
367 .edges
368 .iter()
369 .filter(|e| e.from == id && e.kind == FlowEdgeKind::Next)
370 .map(|e| e.to)
371 .collect();
372 if let Some(s) = succ.first() {
373 edges.push(
374 id,
375 *s,
376 FlowEdgeKind::Error,
377 Some("exhausted".into()),
378 Provenance::Extracted,
379 evidence,
380 );
381 }
382 }
383
384 for (sym, comp) in symbol_comp {
386 for r in graph.out_pred(sym, scc_core::predicates::PUBLISHES) {
387 let from_key = (comp.clone(), op_of(graph, sym));
388 if let Some(from) = table.lookup(&from_key) {
389 let q_key = (r.object.clone(), format!("queue {}", op_of(graph, &r.object)));
390 let to = table.get(&q_key);
391 edges.push(
392 from,
393 to,
394 FlowEdgeKind::Publish,
395 None,
396 r.provenance,
397 r.evidence.clone(),
398 );
399 }
400 }
401 for r in graph.out_pred(sym, scc_core::predicates::CONSUMES) {
402 let from_key = (comp.clone(), op_of(graph, sym));
403 if let Some(from) = table.lookup(&from_key) {
404 let q_key = (r.object.clone(), format!("queue {}", op_of(graph, &r.object)));
405 let to = table.get(&q_key);
406 edges.push(
407 from,
408 to,
409 FlowEdgeKind::Consume,
410 None,
411 r.provenance,
412 r.evidence.clone(),
413 );
414 }
415 }
416 }
417
418 for (sym, comp) in symbol_comp {
422 for (pred, kind) in [
423 (
424 scc_core::predicates::READS,
425 FlowEdgeKind::Read,
426 ),
427 (
428 scc_core::predicates::WRITES,
429 FlowEdgeKind::Write,
430 ),
431 ] {
432 for r in graph.out_pred(sym, pred) {
433 let from_key = (comp.clone(), op_of(graph, sym));
434 if let Some(from) = table.lookup(&from_key) {
435 let to_key =
436 (r.object.clone(), format!("state {}", op_of(graph, &r.object)));
437 let to = table.get(&to_key);
438 edges.push(from, to, kind, None, r.provenance, r.evidence.clone());
439 }
440 }
441 }
442 }
443
444 for (sym, comp) in symbol_comp {
448 for r in graph.out_pred(sym, scc_core::predicates::CALLS) {
449 if r.provenance != Provenance::Resolved {
450 continue;
451 }
452 let is_async = graph
453 .entities
454 .get(&r.object)
455 .and_then(|e| e.attributes.get("async"))
456 .and_then(|v| v.as_bool())
457 .unwrap_or(false);
458 if !is_async {
459 continue;
460 }
461 let from_key = (comp.clone(), op_of(graph, sym));
462 let to_key = (
463 symbol_comp
464 .get(&r.object)
465 .cloned()
466 .unwrap_or_else(|| "component:root".into()),
467 op_of(graph, &r.object),
468 );
469 if let (Some(from), Some(to)) =
470 (table.lookup(&from_key), table.lookup(&to_key))
471 {
472 edges.push(
473 from,
474 to,
475 FlowEdgeKind::Async,
476 None,
477 r.provenance,
478 r.evidence.clone(),
479 );
480 }
481 }
482 }
483
484 let convergents: Vec<u32> = {
486 let mut v: Vec<u32> = edges
487 .in_degree
488 .iter()
489 .filter(|(_, n)| **n > 1)
490 .map(|(id, _)| *id)
491 .collect();
492 v.sort_unstable();
493 v
494 };
495 for to in convergents {
496 let mut preds: Vec<u32> = edges
497 .edges
498 .iter()
499 .filter(|e| e.to == to && e.kind == FlowEdgeKind::Next)
500 .map(|e| e.from)
501 .collect();
502 preds.sort_unstable();
503 preds.dedup();
504 for from in preds {
505 edges.push(
506 from,
507 to,
508 FlowEdgeKind::Join,
509 Some("converge".into()),
510 Provenance::Resolved,
511 Vec::new(),
512 );
513 }
514 }
515
516 for (key, (prov, ev)) in &node_evidence {
518 if let Some(id) = table.lookup(key) {
519 let n = table.nodes.get_mut(id as usize).expect("node exists");
520 n.evidence = ev.iter().cloned().collect();
521 n.evidence.sort();
522 }
523 let _ = prov;
524 }
525
526 let has_out: BTreeSet<u32> = edges.edges.iter().map(|e| e.from).collect();
528 let mut exits: Vec<u32> = table
529 .nodes
530 .iter()
531 .map(|n| n.id)
532 .filter(|id| !has_out.contains(id))
533 .collect();
534 exits.sort_unstable();
535
536 let mut provenance_summary: BTreeMap<String, usize> = BTreeMap::new();
538 for e in &edges.edges {
539 if let Some(p) = e.provenance {
540 *provenance_summary.entry(p.as_str().to_string()).or_insert(0) += 1;
541 }
542 }
543
544 let graph_id = entity_id(&store.repo_id, scc_core::kinds::FLOW, &ep.name);
545 out.push(FlowGraph {
546 id: graph_id,
547 kind: FlowKind::Sequence,
548 name: ep.name.clone(),
549 trigger: Some(ep.trigger.clone()),
550 nodes: table.nodes,
551 edges: edges.edges,
552 entrypoints: vec![entry_id],
553 exits,
554 provenance_summary,
555 });
556 }
557
558 out.sort_by(|a, b| a.name.cmp(&b.name));
559 Ok(out)
560}
561
562#[cfg(test)]
563mod tests {
564 use super::*;
565
566 fn key(comp: &str, op: &str) -> NodeKey {
567 (comp.to_string(), op.to_string())
568 }
569
570 #[test]
571 fn node_key_and_op() {
572 assert_eq!(key("services", "save"), ("services".to_string(), "save".to_string()));
573 assert_eq!(op_of(&RealityGraph::empty(), "repo://r/symbol/x.py/f"), "repo://r/symbol/x.py/f");
574 }
575
576 #[test]
577 fn edge_dedupe_is_provenance_rank_wins() {
578 let mut t = EdgeTable::new();
584 t.push(0, 1, FlowEdgeKind::Next, None, Provenance::Extracted, vec!["ev:native".into()]);
585 t.push(0, 1, FlowEdgeKind::Next, None, Provenance::Resolved, vec!["ev:lsp".into()]);
586 assert_eq!(t.edges.len(), 1, "duplicate triple dedupes to one edge");
587 assert_eq!(
588 t.edges[0].provenance,
589 Some(Provenance::Resolved),
590 "RESOLVED wins over an earlier EXTRACTED candidate"
591 );
592 assert!(t.edges[0].evidence.contains(&"ev:native".to_string()));
593 assert!(t.edges[0].evidence.contains(&"ev:lsp".to_string()));
594 assert_eq!(t.in_degree.get(&1), Some(&1), "in-degree counted once");
595
596 let mut native = EdgeTable::new();
597 native.push(0, 1, FlowEdgeKind::Next, None, Provenance::Extracted, Vec::new());
598 assert_eq!(
599 native.edges[0].provenance,
600 Some(Provenance::Extracted),
601 "native-only edge keeps EXTRACTED provenance"
602 );
603 }
604
605 #[test]
606 fn branch_detection_is_structural() {
607 assert_ne!(FlowEdgeKind::Branch, FlowEdgeKind::Next);
610 }
611}