1use crate::components::component_for_path;
9use crate::lifecycle::component_candidates;
10use crate::{RealityGraph, Result};
11use scc_core::kinds;
12use scc_core::{entity_id, Entity, Flow, FlowKind, FlowStep, Provenance};
13use scc_store::Store;
14use serde_json::json;
15use std::collections::{BTreeMap, HashSet};
16
17pub fn compile_workflows(graph: &RealityGraph, store: &Store) -> Result<Vec<Flow>> {
20 let mut out: Vec<Flow> = Vec::new();
21 let sequences: Vec<Flow> = store
22 .flows()?
23 .into_iter()
24 .filter(|f| f.kind == FlowKind::Sequence)
25 .collect();
26
27 let declared: HashSet<String> = store
30 .intent_claims()?
31 .into_iter()
32 .filter(|(source, _)| source == "flow")
33 .filter_map(|(_, claim)| {
34 if claim.get("kind").and_then(|v| v.as_str()) == Some("workflow") {
35 claim.get("name").and_then(|v| v.as_str()).map(|s| s.to_string())
36 } else {
37 None
38 }
39 })
40 .collect();
41 for seq in &sequences {
42 if declared.contains(&seq.name) {
43 let mut clone = seq.clone();
44 clone.kind = FlowKind::Workflow;
45 clone.id = entity_id(&store.repo_id, kinds::WORKFLOW, &seq.name);
46 out.push(clone);
47 }
48 }
49
50 for seq in &sequences {
56 if !seq.attributes.contains_key("branches") {
57 continue;
58 }
59 let name = format!("{}-workflow", seq.name);
60 out.push(Flow {
61 id: entity_id(&store.repo_id, kinds::FLOW, &name),
62 kind: FlowKind::Workflow,
63 name,
64 trigger: seq.trigger.clone(),
65 steps: seq.steps.clone(),
66 attributes: seq.attributes.clone(),
67 });
68 }
69
70 let candidates = component_candidates(graph);
73 let comp_name_to_id: BTreeMap<String, String> = graph
74 .components
75 .iter()
76 .map(|c| (c.name.clone(), c.id.clone()))
77 .collect();
78 let mut signals_by_comp: BTreeMap<String, Vec<&Entity>> = BTreeMap::new();
79 for e in graph.entities_of_kind(kinds::SYMBOL) {
80 let name_l = e.name.to_ascii_lowercase();
81 if !e.attributes.contains_key("retry_policy")
82 && !name_l.contains("retry")
83 && !name_l.contains("fallback")
84 && !name_l.contains("backoff")
85 {
86 continue;
87 }
88 let comp = e
89 .attributes
90 .get("file")
91 .and_then(|v| v.as_str())
92 .map(|f| component_for_path(f, &candidates))
93 .unwrap_or_else(|| "root".to_string());
94 signals_by_comp.entry(comp).or_default().push(e);
95 }
96 for (comp, mut signals) in signals_by_comp {
97 if signals.len() < 2 {
98 continue;
99 }
100 signals.sort_by(|a, b| a.name.cmp(&b.name));
101 let comp_id = comp_name_to_id
102 .get(&comp)
103 .cloned()
104 .unwrap_or_else(|| format!("component:{comp}"));
105 let name = format!("{comp}-workflow");
106
107 let mut steps: Vec<FlowStep> = Vec::new();
108 steps.push(FlowStep {
109 id: "step:1".to_string(),
110 order: 1,
111 actor: comp_id.clone(),
112 operation: comp.clone(),
113 condition: None,
114 r#async: None,
115 timeout_ms: None,
116 retry_policy: None,
117 failure_outcome: None,
118 provenance: Some(Provenance::Inferred),
119 evidence: Vec::new(),
120 });
121 let mut retries = 0usize;
122 let mut fallbacks = 0usize;
123 for (i, s) in signals.iter().enumerate() {
124 let name_l = s.name.to_ascii_lowercase();
125 let is_retry = s.attributes.contains_key("retry_policy")
126 || name_l.contains("retry")
127 || name_l.contains("backoff");
128 let is_fallback = name_l.contains("fallback");
129 if is_retry {
130 retries += 1;
131 }
132 if is_fallback {
133 fallbacks += 1;
134 }
135 steps.push(FlowStep {
136 id: format!("step:{}", i + 2),
137 order: (i + 2) as u32,
138 actor: comp_id.clone(),
139 operation: s.name.clone(),
140 condition: None,
141 r#async: None,
142 timeout_ms: None,
143 retry_policy: s
144 .attributes
145 .get("retry_policy")
146 .and_then(|v| v.as_str())
147 .map(|p| p.to_string()),
148 failure_outcome: if is_fallback {
149 Some("fallback".to_string())
150 } else {
151 None
152 },
153 provenance: Some(Provenance::Extracted),
154 evidence: Vec::new(),
155 });
156 }
157
158 let mut attributes = BTreeMap::new();
159 attributes.insert("retries".to_string(), json!(retries));
160 attributes.insert("fallbacks".to_string(), json!(fallbacks));
161 out.push(Flow {
162 id: entity_id(&store.repo_id, kinds::FLOW, &name),
163 kind: FlowKind::Workflow,
164 name,
165 trigger: None,
166 steps,
167 attributes,
168 });
169 }
170
171 let mut seen: HashSet<String> = HashSet::new();
174 out.retain(|f| seen.insert(f.id.clone()));
175 out.sort_by(|a, b| a.name.cmp(&b.name));
176 Ok(out)
177}
178
179#[cfg(test)]
180mod tests {
181 use super::*;
182 use scc_core::symbol_id;
183
184 fn setup() -> (tempfile::TempDir, Store) {
185 let dir = tempfile::TempDir::new().unwrap();
186 let root = dir.path().join("repo");
187 std::fs::create_dir_all(&root).unwrap();
188 let store = Store::open(&dir.path().join("scc.db"), &root).unwrap();
189 (dir, store)
190 }
191
192 fn put_components(store: &Store, comps: &[(&str, &[&str])]) {
193 let entities: Vec<Entity> = comps
194 .iter()
195 .map(|(name, paths)| {
196 let mut c = Entity::new(
197 entity_id(&store.repo_id, kinds::COMPONENT, name),
198 kinds::COMPONENT,
199 *name,
200 );
201 c.attr("implementation", json!({ "paths": paths, "symbols": [] }));
202 c
203 })
204 .collect();
205 store.replace_components(&entities).unwrap();
206 }
207
208 fn put_symbol(
209 store: &Store,
210 name: &str,
211 file: &str,
212 attrs: &[(&str, serde_json::Value)],
213 ) {
214 let mut e = Entity::new(symbol_id(&store.repo_id, file, name), kinds::SYMBOL, name);
215 e.attr("kind", serde_json::json!("function"));
216 e.attr("file", serde_json::json!(file));
217 for (k, v) in attrs {
218 e.attr(k, v.clone());
219 }
220 store.insert_entity(&e, &[file.to_string()]).unwrap();
221 }
222
223 fn seq_flow(store: &Store, name: &str, ops: &[&str], attrs: serde_json::Value) -> Flow {
224 let steps: Vec<FlowStep> = ops
225 .iter()
226 .enumerate()
227 .map(|(i, op)| FlowStep {
228 id: format!("step:{}", i + 1),
229 order: (i + 1) as u32,
230 actor: "actor".to_string(),
231 operation: op.to_string(),
232 condition: None,
233 r#async: None,
234 timeout_ms: None,
235 retry_policy: None,
236 failure_outcome: None,
237 provenance: Some(Provenance::Resolved),
238 evidence: Vec::new(),
239 })
240 .collect();
241 Flow {
242 id: entity_id(&store.repo_id, kinds::FLOW, name),
243 kind: FlowKind::Sequence,
244 name: name.to_string(),
245 trigger: Some("t".to_string()),
246 steps,
247 attributes: serde_json::from_value(attrs).unwrap(),
248 }
249 }
250
251 #[test]
252 fn intent_declared_workflow_clone() {
253 let (_dir, store) = setup();
254 let seq = seq_flow(&store, "onboard", &["validate", "create"], json!({}));
255 store.replace_flows(&[seq]).unwrap();
256 store
257 .replace_intent_claims(&[(
258 "flow".to_string(),
259 json!({ "name": "onboard", "kind": "workflow", "entrypoint": "onboard_user" }),
260 )])
261 .unwrap();
262 let graph = RealityGraph::load(&store).unwrap();
263
264 let flows = compile_workflows(&graph, &store).unwrap();
265 assert_eq!(flows.len(), 1);
266 let f = &flows[0];
267 assert_eq!(f.kind, FlowKind::Workflow);
268 assert_eq!(f.name, "onboard");
269 assert_eq!(f.id, entity_id(&store.repo_id, kinds::WORKFLOW, "onboard"));
270 assert_ne!(f.id, entity_id(&store.repo_id, kinds::FLOW, "onboard"));
271 assert_eq!(f.steps.len(), 2);
272 assert_eq!(f.steps[0].operation, "validate");
273 assert_eq!(f.steps[1].operation, "create");
274 }
275
276 #[test]
277 fn non_workflow_intent_is_ignored() {
278 let (_dir, store) = setup();
279 let seq = seq_flow(&store, "checkout", &["validate"], json!({}));
280 store.replace_flows(&[seq]).unwrap();
281 store
282 .replace_intent_claims(&[(
283 "flow".to_string(),
284 json!({ "name": "checkout", "kind": "sequence", "entrypoint": "checkout_fn" }),
285 )])
286 .unwrap();
287 let graph = RealityGraph::load(&store).unwrap();
288 let flows = compile_workflows(&graph, &store).unwrap();
289 assert!(flows.is_empty());
290 }
291
292 #[test]
293 fn collapsed_step_does_not_yield_branch_workflow() {
294 let (_dir, store) = setup();
297 let seq = seq_flow(
298 &store,
299 "checkout",
300 &["validate", "charge, refund"],
301 json!({}),
302 );
303 store.replace_flows(&[seq]).unwrap();
304 let graph = RealityGraph::load(&store).unwrap();
305
306 let flows = compile_workflows(&graph, &store).unwrap();
307 assert!(flows.is_empty(), "no evidence -> no branch workflow: {flows:?}");
308 }
309
310 #[test]
311 fn branches_attribute_yields_workflow_without_branch_step() {
312 let (_dir, store) = setup();
313 let seq = seq_flow(
314 &store,
315 "pipeline",
316 &["validate", "execute"],
317 json!({ "branches": ["fast", "full"] }),
318 );
319 store.replace_flows(&[seq]).unwrap();
320 let graph = RealityGraph::load(&store).unwrap();
321
322 let flows = compile_workflows(&graph, &store).unwrap();
323 assert_eq!(flows.len(), 1);
324 let f = &flows[0];
325 assert_eq!(f.name, "pipeline-workflow");
326 assert_eq!(f.steps.len(), 2, "no collapsed step -> no branch step appended");
327 assert_eq!(f.attributes["branches"], json!(["fast", "full"]));
328 }
329
330 #[test]
331 fn retry_fallback_component_workflow() {
332 let (_dir, store) = setup();
333 put_components(&store, &[("ingest", &["ingest"])]);
334 put_symbol(&store, "retry_upload", "ingest/upload.py", &[]);
335 put_symbol(
336 &store,
337 "retry_download",
338 "ingest/download.py",
339 &[("retry_policy", serde_json::json!("exponential"))],
340 );
341 put_symbol(&store, "fallback_queue", "ingest/queue.py", &[]);
342 put_components(&store, &[("ingest", &["ingest"]), ("jobs", &["jobs"])]);
344 put_symbol(&store, "retry_job", "jobs/run.py", &[]);
345 let graph = RealityGraph::load(&store).unwrap();
346
347 let flows = compile_workflows(&graph, &store).unwrap();
348 assert_eq!(flows.len(), 1);
349 let f = &flows[0];
350 assert_eq!(f.kind, FlowKind::Workflow);
351 assert_eq!(f.name, "ingest-workflow");
352 assert_eq!(f.id, entity_id(&store.repo_id, kinds::FLOW, "ingest-workflow"));
353 assert_eq!(f.steps.len(), 4);
354 assert_eq!(f.steps[0].operation, "ingest");
355 assert_eq!(f.steps[0].provenance, Some(Provenance::Inferred));
356 assert_eq!(f.steps[0].actor, entity_id(&store.repo_id, kinds::COMPONENT, "ingest"));
357 let step = |op: &str| f.steps.iter().find(|s| s.operation == op).unwrap();
358 assert_eq!(step("retry_upload").retry_policy, None);
359 assert_eq!(step("retry_upload").failure_outcome, None);
360 assert_eq!(step("retry_download").retry_policy.as_deref(), Some("exponential"));
361 assert_eq!(step("retry_download").failure_outcome, None);
362 assert_eq!(step("fallback_queue").retry_policy, None);
363 assert_eq!(step("fallback_queue").failure_outcome.as_deref(), Some("fallback"));
364 for s in &f.steps[1..] {
365 assert_eq!(s.provenance, Some(Provenance::Extracted));
366 }
367 assert_eq!(f.attributes["retries"], json!(2));
368 assert_eq!(f.attributes["fallbacks"], json!(1));
369 }
370
371 #[test]
372 fn workflow_views_sorted_by_name() {
373 let (_dir, store) = setup();
374 put_components(&store, &[("b", &["b"])]);
376 put_symbol(&store, "retry_one", "b/x.py", &[]);
377 put_symbol(&store, "backoff_two", "b/y.py", &[]);
378 let seq = seq_flow(
379 &store,
380 "a",
381 &["validate", "execute"],
382 json!({ "branches": ["fast", "full"] }),
383 );
384 store.replace_flows(&[seq]).unwrap();
385 let graph = RealityGraph::load(&store).unwrap();
386
387 let flows = compile_workflows(&graph, &store).unwrap();
388 assert_eq!(flows.len(), 2);
389 let names: Vec<&str> = flows.iter().map(|f| f.name.as_str()).collect();
390 assert_eq!(names, vec!["a-workflow", "b-workflow"]);
391 assert_eq!(flows[0].attributes["branches"], json!(["fast", "full"]));
392 }
393}