Skip to main content

scc_graph/
flowgraph.rs

1//! Canonical causal FlowGraph compiler (Wave 3, docs/SYSTEM_DESIGN.md §9).
2//!
3//! One graph per entrypoint, built ONLY from evidence:
4//!
5//! - Next edges from resolved CALLS paths
6//! - Branch edges from call fanout (one operation -> 2+ distinct callees)
7//! - Retry edges from RETRIES predicates (extracted failure behavior)
8//! - Error edges to failure outcomes (retry exhaustion)
9//! - Async edges from async call attributes
10//! - Publish/Consume edges from queue predicates
11//! - Join edges at convergence; Return/exit detection
12//!
13//! Branches NEVER come from generated-text heuristics (e.g. splitting on
14//! ", "). The canonical graph preserves topology exactly: alternate
15//! execution paths can never be flattened into false sequential causality.
16//! Individual operations are retained; component grouping is a display-time
17//! concern only (ComponentSpan).
18
19use 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
28/// Node key: (actor component id, operation).
29type 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
39/// Mutable edge registry for one graph build (dedup + join detection).
40// trace:exempt reason=internal-detail
41struct 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        // Dedupe is provenance-rank-wins, never first-wins: a target the
67        // language server RESOLVED must not be demoted by an earlier
68        // native EXTRACTED candidate for the same (from, to, kind) triple
69        // (and vice versa, a native candidate keeps the edge when no LSP
70        // evidence exists — holdout repos). Evidence merges.
71        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
105/// Mutable node registry for one graph build.
106// trace:exempt reason=internal-detail
107struct 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
140/// Compile the canonical flow graphs for every entrypoint.
141///
142/// `symbol_comp` maps symbol id -> component entity id (same mapping the
143/// projection compilers use, so displays agree).
144// trace:v1 id=impl.scc.flowgraph work=WORK-SCC-005 satisfies=REQ-SCC-FLOW
145// trace:v1 id=impl.scc.graph.state-flow-read-write work=WORK-SI-MMMJA4G6 satisfies=REQ-SI-503JSBGP
146pub 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        // entry node
164        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        // ---- edges from call paths ----
172        // successors[sym] = distinct callees resolved from CALLS
173        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        // Materialize every reachable symbol as a node up front so node ids
205        // are deterministic (sorted symbol order) and every caller's
206        // from-lookup succeeds regardless of edge-creation order.
207        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            // node evidence from the entity itself
225            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            // P1 §19: call fanout is NOT control-flow branching. A call
240            // graph says "A may call B and C", not "A chooses B or C".
241            // Branch edges exist ONLY where the extractor recorded the call
242            // inside a conditional/loop/try body — `call_blocks` (CFG
243            // evidence, condition = block kind) with `conditional_calls`
244            // as fallback for older indexes.
245            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            // Deterministic order: CFG lexical order within the caller
264            // (min over matching call sites), then callee name. Straight-
265            // line calls keep their source sequence; unknown callees sort
266            // after known ones.
267            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                // Async edge for awaited/spawned call sites (CFG evidence):
320                // `await`, `.await`, `go` statement, `Promise.all`.
321                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        // ---- retry / failure edges from extracted behavior ----
335        // Evidence: the `retry_policy` attribute stamped on decorated
336        // symbols (@retry / @tenacity.retry / @backoff decorators). The
337        // retrying op gets a Retry self-loop (attempt) and an Error edge to
338        // its successor (exhausted -> failure path).
339        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            // exhausted -> failure: Error edge to the first successor
366            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        // ---- publish/consume edges (queue semantics) ----
385        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        // ---- read/write edges (state authority: handler READS/WRITES
419        // stores). Same mechanical pattern as publish/consume: the store
420        // relationship is the evidence, provenance rides along.
421        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        // ---- async attributes on call edges ----
445        // (extractors set `async: true` on CALLS relationships when the
446        // call is fire-and-forget; honor them if present)
447        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        // ---- join edges at convergence ----
485        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        // ---- attach per-node evidence ----
517        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        // ---- exits (no outgoing edges) ----
527        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        // provenance summary
537        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        // The same (from, to, kind) triple pushed twice — once with a
579        // native EXTRACTED candidate, once with LSP RESOLVED proof — keeps
580        // the RESOLVED provenance (dedupe is rank-wins, never first-wins),
581        // and a native-only triple keeps its EXTRACTED provenance (flows
582        // exist without LSP).
583        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        // Branch edges come from call FANOUT — never from text heuristics.
608        // (topology is exercised end-to-end in the fixture tests)
609        assert_ne!(FlowEdgeKind::Branch, FlowEdgeKind::Next);
610    }
611}