Skip to main content

af_workflow/
executor.rs

1//! Compile a spec branch into a runnable step chain and drive events through it.
2//!
3//! Port of the execution half of `platform/compiler.py`. The spec DAG is
4//! linear day-1 (per R9): one ingress head, then a topologically-ordered chain
5//! of step nodes. The ingress node is the event *source* (owned by the runner,
6//! next phase); compilation separates it from the steps it feeds.
7
8use 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/// Why a branch could not be compiled.
20#[derive(Debug, thiserror::Error)]
21pub enum CompileError {
22    /// A node could not be built.
23    #[error(transparent)]
24    Node(#[from] NodeError),
25    /// Branch '`branch_id`' has no ingress node (chain needs a source).
26    #[error("branch '{branch_id}' has no ingress node (chain needs a source)")]
27    /// The branch has no ingress node.
28    NoIngress {
29        /// Branch id.
30        branch_id: String,
31    },
32    /// Branch '`branch_id`' has `count` ingress nodes; day-1 supports one.
33    #[error("branch '{branch_id}' has {count} ingress nodes; day-1 supports one")]
34    /// The branch has more than one ingress node.
35    MultipleIngress {
36        /// Branch id.
37        branch_id: String,
38        /// Ingress nodes found.
39        count: usize,
40    },
41    /// Branch '`branch_id`' edges form a cycle (DAG required).
42    #[error("branch '{branch_id}' edges form a cycle (DAG required)")]
43    /// The branch is not a DAG.
44    Cycle {
45        /// Branch id.
46        branch_id: String,
47    },
48}
49
50/// A compiled, runnable branch: an ingress source + an ordered step chain.
51pub(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/// Where a run ended.
87#[derive(Debug, Clone, PartialEq, Eq)]
88pub enum Terminal {
89    /// All events fell out (every branch of the chain Dropped).
90    Dropped {
91        /// Node that dropped the last event.
92        node_id: String,
93        /// Why.
94        reason: String,
95    },
96    /// The chain ran to the end with at least one surviving event.
97    Completed,
98}
99
100/// Summary of one event driven through a compiled branch.
101#[derive(Debug, Clone)]
102pub struct RunOutcome {
103    /// Steps the event passed through.
104    pub steps_run: usize,
105    /// Where the run ended.
106    pub terminal: Terminal,
107    /// Events still alive at the end (empty if everything dropped).
108    pub survivors: Vec<Event>,
109    /// External effects the steps prepared, in emission order.
110    pub actions: Vec<PreparedAction>,
111    /// At least one event passed a material step or reached the branch end.
112    pub matched: bool,
113    /// The branch reached a sink or retained an event through its final step.
114    pub succeeded: bool,
115}
116
117impl CompiledBranch {
118    /// Compile a spec branch against a node registry.
119    pub(crate) fn compile(branch: &Branch, registry: &NodeRegistry) -> Result<Self, CompileError> {
120        let order = topo_order(branch)?;
121
122        // Split ingress head from step chain.
123        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; // ingress is the source, not a step
150            }
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    /// Drive one event through the step chain.
168    ///
169    /// `Pass` forwards, `Drop` removes that event, `FanOut` multiplies it. The
170    /// run completes when the chain is exhausted or every event has dropped.
171    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 &current {
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
378/// Kahn topological sort over a branch's edges. Nodes with no edges keep spec
379/// order. Errors on a cycle.
380fn 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        // Dangling edges are caught by Spec::validate_structure; ignore here.
388        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    // Seed queue in spec order to keep deterministic output.
398    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, &registry).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, &registry)
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}