1use std::collections::{HashMap, VecDeque};
9
10use serde_json::{json, Value};
11
12use crate::event::Event;
13use crate::node::{StepNode, WorkflowContext};
14use crate::recorder::{RunStatus, StepStatus};
15use crate::registry::{NodeError, NodeRegistry};
16use crate::result::{PreparedAction, StepResult};
17use crate::spec::Branch;
18
19#[derive(Debug, thiserror::Error)]
21pub enum CompileError {
22 #[error(transparent)]
24 Node(#[from] NodeError),
25 #[error("branch '{branch_id}' has no ingress node (chain needs a source)")]
27 NoIngress {
29 branch_id: String,
31 },
32 #[error("branch '{branch_id}' has {count} ingress nodes; day-1 supports one")]
34 MultipleIngress {
36 branch_id: String,
38 count: usize,
40 },
41 #[error("branch '{branch_id}' edges form a cycle (DAG required)")]
43 Cycle {
45 branch_id: String,
47 },
48}
49
50pub(crate) struct CompiledBranch {
52 steps: Vec<CompiledStep>,
53}
54
55struct CompiledStep {
56 node_id: String,
57 node_type: String,
58 node: Box<dyn StepNode>,
59 fan_out_limit: usize,
60 action_capable: bool,
61}
62
63const MAX_RUN_FAN_OUT: usize = 1_000;
64
65fn fan_out_limit(config: &Value) -> usize {
66 ["count", "levels", "fanout"]
67 .into_iter()
68 .find_map(|key| config.get(key)?.as_u64())
69 .and_then(|value| usize::try_from(value).ok())
70 .unwrap_or(MAX_RUN_FAN_OUT)
71 .min(MAX_RUN_FAN_OUT)
72}
73
74fn is_material(node_type: &str) -> bool {
75 node_type.starts_with("execute.") || node_type.starts_with("notify.")
76}
77
78fn event_detail(event: &Event) -> Value {
79 json!({
80 "event_id": event.id,
81 "payload": event.payload,
82 "metadata": event.metadata,
83 })
84}
85
86#[derive(Debug, Clone, PartialEq, Eq)]
88pub enum Terminal {
89 Dropped {
91 node_id: String,
93 reason: String,
95 },
96 Completed,
98}
99
100#[derive(Debug, Clone)]
102pub struct RunOutcome {
103 pub steps_run: usize,
105 pub terminal: Terminal,
107 pub survivors: Vec<Event>,
109 pub actions: Vec<PreparedAction>,
111 pub matched: bool,
113 pub succeeded: bool,
115}
116
117impl CompiledBranch {
118 pub(crate) fn compile(branch: &Branch, registry: &NodeRegistry) -> Result<Self, CompileError> {
120 let order = topo_order(branch)?;
121
122 let ingress_nodes: Vec<&crate::spec::Node> = branch
124 .nodes
125 .iter()
126 .filter(|n| registry.is_ingress(&n.node_type))
127 .collect();
128 match ingress_nodes.len() {
129 0 => {
130 return Err(CompileError::NoIngress {
131 branch_id: branch.branch_id.clone(),
132 })
133 }
134 1 => {}
135 n => {
136 return Err(CompileError::MultipleIngress {
137 branch_id: branch.branch_id.clone(),
138 count: n,
139 })
140 }
141 }
142 let by_id: HashMap<&str, &crate::spec::Node> =
143 branch.nodes.iter().map(|n| (n.id.as_str(), n)).collect();
144
145 let mut steps = Vec::new();
146 for node_id in order {
147 let node = by_id[node_id.as_str()];
148 if registry.is_ingress(&node.node_type) {
149 continue; }
151 let built = registry.build_step(&node.node_type, &node.config)?;
152 steps.push(CompiledStep {
153 node_id: node.id.clone(),
154 node_type: node.node_type.clone(),
155 node: built,
156 fan_out_limit: fan_out_limit(&node.config),
157 action_capable: node.node_type.starts_with("execute.")
158 || registry
159 .capability(&node.node_type)
160 .is_some_and(|manifest| manifest.kind == crate::CapabilityKind::Action),
161 });
162 }
163
164 Ok(Self { steps })
165 }
166
167 pub(crate) async fn run_event(&self, ctx: &WorkflowContext, event: Event) -> RunOutcome {
172 self.run_event_with_continuation_from(ctx, event, 0).await.0
173 }
174
175 pub(crate) async fn run_event_from(
176 &self,
177 ctx: &WorkflowContext,
178 event: Event,
179 start_step: usize,
180 ) -> RunOutcome {
181 self.run_event_with_continuation_from(ctx, event, start_step)
182 .await
183 .0
184 }
185
186 pub(crate) async fn run_event_with_continuation_from(
187 &self,
188 ctx: &WorkflowContext,
189 event: Event,
190 start_step: usize,
191 ) -> (RunOutcome, Option<usize>) {
192 let trigger = Value::Object(event.payload.clone());
193 let mut run = ctx.recorder.start(&ctx.trigger_kind, trigger).await;
194 let mut current = vec![event];
195 let mut steps_run = 0;
196 let mut material_steps = 0u32;
197 let mut actions = Vec::new();
198
199 for (step_index, step) in self.steps.iter().enumerate().skip(start_step) {
200 let mut next = Vec::new();
201 let mut last_drop: Option<(String, Option<String>)> = None;
202 let mut fan_out_error = None;
203 for ev in ¤t {
204 match step.node.process(ev, ctx).await {
205 StepResult::Pass(event) => {
206 if is_material(&step.node_type) {
207 run.record_step(
208 &step.node_id,
209 &step.node_type,
210 StepStatus::Ok,
211 None,
212 event_detail(&event),
213 )
214 .await;
215 material_steps += 1;
216 }
217 next.push(event);
218 }
219 StepResult::Drop {
220 reason,
221 exit_reason,
222 } => last_drop = Some((reason, exit_reason)),
223 StepResult::FanOut(evs) => {
224 if evs.len() > step.fan_out_limit
225 || next.len().saturating_add(evs.len()) > MAX_RUN_FAN_OUT
226 {
227 fan_out_error = Some(format!(
228 "fan-out exceeded node limit {} or run limit {MAX_RUN_FAN_OUT}",
229 step.fan_out_limit
230 ));
231 break;
232 }
233 next.extend(evs);
234 }
235 StepResult::Action { event, action } => {
236 if !step.action_capable {
237 fan_out_error = Some(format!(
238 "node '{}' emitted an action without an action capability",
239 step.node_id
240 ));
241 break;
242 }
243 run.record_step(
244 &step.node_id,
245 &step.node_type,
246 StepStatus::Ok,
247 None,
248 event_detail(&event),
249 )
250 .await;
251 material_steps += 1;
252 actions.push(*action);
253 next.push(event);
254 }
255 }
256 }
257 steps_run += 1;
258 if let Some(reason) = fan_out_error {
259 run.record_step(
260 &step.node_id,
261 &step.node_type,
262 StepStatus::Error,
263 Some("fanout_limit_exceeded"),
264 json!({ "reason": reason }),
265 )
266 .await;
267 run.end(RunStatus::Error, Some("fanout_limit_exceeded"))
268 .await;
269 return (
270 RunOutcome {
271 steps_run,
272 terminal: Terminal::Dropped {
273 node_id: step.node_id.clone(),
274 reason,
275 },
276 survivors: Vec::new(),
277 actions: Vec::new(),
278 matched: false,
279 succeeded: false,
280 },
281 None,
282 );
283 }
284 if next.is_empty() {
285 let (reason, exit_reason) = last_drop.unwrap_or_else(|| ("dropped".into(), None));
286
287 if step.node_type.starts_with("sink.") || material_steps > 0 {
288 run.end(RunStatus::Ok, Some("natural")).await;
289 } else if let Some(code) = exit_reason {
290 let step_status = if code.starts_with("invalid_") {
291 StepStatus::Error
292 } else {
293 StepStatus::Skipped
294 };
295 run.record_step(
296 &step.node_id,
297 &step.node_type,
298 step_status,
299 Some(&code),
300 json!({ "reason": reason }),
301 )
302 .await;
303 run.end(
304 if step_status == StepStatus::Error {
305 RunStatus::Error
306 } else {
307 RunStatus::Skipped
308 },
309 Some(&code),
310 )
311 .await;
312 } else {
313 run.mark_filtered(&step.node_id, &step.node_type, &reason)
314 .await;
315 run.end(RunStatus::Skipped, None).await;
316 }
317
318 let sink_completed = step.node_type.starts_with("sink.");
319 return (
320 RunOutcome {
321 steps_run,
322 terminal: if sink_completed {
323 Terminal::Completed
324 } else {
325 Terminal::Dropped {
326 node_id: step.node_id.clone(),
327 reason,
328 }
329 },
330 survivors: Vec::new(),
331 actions,
332 matched: sink_completed || material_steps > 0,
333 succeeded: sink_completed,
334 },
335 None,
336 );
337 }
338 current = next;
339 if !actions.is_empty() {
340 run.end(RunStatus::Ok, Some("action_pending")).await;
341 return (
342 RunOutcome {
343 steps_run,
344 terminal: Terminal::Completed,
345 survivors: current,
346 actions,
347 matched: true,
348 succeeded: false,
349 },
350 (step_index + 1 < self.steps.len()).then_some(step_index + 1),
351 );
352 }
353 }
354
355 run.end(
356 if material_steps > 0 {
357 RunStatus::Ok
358 } else {
359 RunStatus::Skipped
360 },
361 Some("natural"),
362 )
363 .await;
364 (
365 RunOutcome {
366 steps_run,
367 terminal: Terminal::Completed,
368 survivors: current,
369 actions,
370 matched: true,
371 succeeded: true,
372 },
373 None,
374 )
375 }
376}
377
378fn topo_order(branch: &Branch) -> Result<Vec<String>, CompileError> {
381 let ids: Vec<&str> = branch.nodes.iter().map(|n| n.id.as_str()).collect();
382
383 let mut indegree: HashMap<&str, usize> = ids.iter().map(|id| (*id, 0)).collect();
384 let mut adj: HashMap<&str, Vec<&str>> = ids.iter().map(|id| (*id, Vec::new())).collect();
385
386 for edge in &branch.edges {
387 if let (Some(successors), Some(indegree)) = (
389 adj.get_mut(edge.source.as_str()),
390 indegree.get_mut(edge.target.as_str()),
391 ) {
392 successors.push(&edge.target);
393 *indegree += 1;
394 }
395 }
396
397 let mut queue: VecDeque<&str> = ids.iter().copied().filter(|id| indegree[id] == 0).collect();
399
400 let mut order = Vec::with_capacity(ids.len());
401 while let Some(id) = queue.pop_front() {
402 order.push(id.to_string());
403 for &next in &adj[id] {
404 let Some(d) = indegree.get_mut(next) else {
405 continue;
406 };
407 *d -= 1;
408 if *d == 0 {
409 queue.push_back(next);
410 }
411 }
412 }
413
414 if order.len() != ids.len() {
415 return Err(CompileError::Cycle {
416 branch_id: branch.branch_id.clone(),
417 });
418 }
419 Ok(order)
420}
421
422#[cfg(test)]
423mod tests {
424 use std::sync::Arc;
425
426 use async_trait::async_trait;
427 use serde_json::json;
428
429 use super::*;
430 use crate::node::StepNode;
431 use crate::spec::{Edge, Node};
432 use crate::state::MemoryState;
433
434 struct OverProducingMap;
435
436 struct Pass;
437
438 #[async_trait]
439 impl StepNode for Pass {
440 async fn process(&self, event: &Event, _: &WorkflowContext) -> StepResult {
441 StepResult::Pass(event.clone())
442 }
443 }
444
445 struct Drop;
446
447 #[async_trait]
448 impl StepNode for Drop {
449 async fn process(&self, _: &Event, _: &WorkflowContext) -> StepResult {
450 StepResult::drop("filtered")
451 }
452 }
453
454 #[async_trait]
455 impl StepNode for OverProducingMap {
456 fn produces_fan_out(&self) -> bool {
457 true
458 }
459
460 async fn process(&self, event: &Event, _: &WorkflowContext) -> StepResult {
461 StepResult::FanOut(vec![event.clone(), event.clone(), event.clone()])
462 }
463 }
464
465 fn build_over_producing(_: &Value) -> Result<Box<dyn StepNode>, NodeError> {
466 Ok(Box::new(OverProducingMap))
467 }
468
469 #[tokio::test]
470 async fn runtime_rejects_more_fanout_than_the_static_declaration() {
471 let mut registry = NodeRegistry::empty();
472 registry.register_ingress("ingress.event");
473 registry.register_step("map.test", build_over_producing);
474 registry.register_fan_out("map.test");
475 let branch = Branch {
476 branch_id: "root".into(),
477 nodes: vec![
478 Node {
479 id: "in".into(),
480 node_type: "ingress.event".into(),
481 config: json!({}),
482 },
483 Node {
484 id: "map".into(),
485 node_type: "map.test".into(),
486 config: json!({"count": 2}),
487 },
488 ],
489 edges: vec![Edge {
490 source: "in".into(),
491 target: "map".into(),
492 }],
493 };
494 let compiled = CompiledBranch::compile(&branch, ®istry).unwrap();
495 let context = WorkflowContext::new("root", Arc::new(MemoryState::new()));
496 let outcome = compiled
497 .run_event(&context, Event::from_json(json!({})))
498 .await;
499 assert!(matches!(
500 outcome.terminal,
501 Terminal::Dropped { ref reason, .. } if reason.contains("fan-out exceeded")
502 ));
503 assert!(outcome.survivors.is_empty());
504 }
505
506 #[tokio::test]
507 async fn a_late_filter_matches_without_claiming_success() {
508 let mut registry = NodeRegistry::empty();
509 registry.register_ingress("ingress.event");
510 registry.register_step("execute.pass", |_| Ok(Box::new(Pass)));
511 registry.register_step("filter.drop", |_| Ok(Box::new(Drop)));
512 let branch = Branch {
513 branch_id: "root".into(),
514 nodes: vec![
515 Node {
516 id: "in".into(),
517 node_type: "ingress.event".into(),
518 config: json!({}),
519 },
520 Node {
521 id: "material".into(),
522 node_type: "execute.pass".into(),
523 config: json!({}),
524 },
525 Node {
526 id: "drop".into(),
527 node_type: "filter.drop".into(),
528 config: json!({}),
529 },
530 ],
531 edges: vec![
532 Edge {
533 source: "in".into(),
534 target: "material".into(),
535 },
536 Edge {
537 source: "material".into(),
538 target: "drop".into(),
539 },
540 ],
541 };
542 let outcome = CompiledBranch::compile(&branch, ®istry)
543 .unwrap()
544 .run_event(
545 &WorkflowContext::new("root", Arc::new(MemoryState::new())),
546 Event::from_json(json!({})),
547 )
548 .await;
549 assert!(outcome.matched);
550 assert!(!outcome.succeeded);
551 assert!(matches!(outcome.terminal, Terminal::Dropped { .. }));
552 }
553}