Skip to main content

af_workflow/
builtins.rs

1//! Generic (non-business) node implementations.
2//!
3//! These are foundation nodes every product can share. External execution stays
4//! behind the governed Agent platform; workflow nodes remain deterministic.
5//!
6//! Each node holds its raw config `Value` and resolves templates per-event in
7//! `process` (matching the Python `resolved_config` discipline), then reads
8//! fields. Required-key checks happen in the factory at compile time.
9
10use std::collections::HashMap;
11
12use async_trait::async_trait;
13use serde_json::{Map, Value};
14
15use crate::event::Event;
16use crate::node::{StepNode, WorkflowContext};
17use crate::registry::{NodeError, NodeRegistry};
18use crate::result::StepResult;
19use crate::template::resolve_config;
20
21/// Register every builtin step + ingress type.
22pub fn register_builtins(r: &mut NodeRegistry) {
23    // Ingress sources — execution owned by the runner; registered by name so
24    // the compiler recognizes them as legal chain heads.
25    r.register_ingress("ingress.cron");
26    r.register_ingress("ingress.event");
27    r.register_ingress("ingress.manual");
28
29    // Step nodes.
30    r.register_step("transform.state_append", TransformStateAppend::factory);
31    r.register_step("transform.state_set", TransformStateSet::factory);
32    r.register_step("transform.state_read", TransformStateRead::factory);
33    r.register_step("transform.set_fields", TransformSetFields::factory);
34    r.register_step(
35        "transform.state_publish_cross_branch",
36        TransformStatePublishCrossBranch::factory,
37    );
38    r.register_step(
39        "transform.state_read_cross_branch",
40        TransformStateReadCrossBranch::factory,
41    );
42    r.register_step(
43        "transform.state_append_cross_branch",
44        TransformStateAppendCrossBranch::factory,
45    );
46    r.register_step("decision.consensus_voting", ConsensusVoting::factory);
47    r.register_step("filter.required_fields", FilterRequiredFields::factory);
48    r.register_step("sink.log", SinkLog::factory);
49
50    register_builtin_schemas(r);
51}
52
53fn register_builtin_schemas(registry: &mut NodeRegistry) {
54    use crate::registry::{FieldSpec as Field, FieldType, NodeSchema};
55
56    let schema = |fields| NodeSchema { fields };
57    registry.register_schema(
58        "transform.state_append",
59        schema(vec![
60            Field::required("path", FieldType::String),
61            Field::required("key", FieldType::String),
62            Field::optional("max_len", FieldType::Number),
63            Field::optional("ttl_seconds", FieldType::Number),
64        ]),
65    );
66    registry.register_schema(
67        "transform.state_set",
68        schema(vec![
69            Field::required("key", FieldType::String),
70            Field::required("value", FieldType::Any),
71            Field::optional("ttl_seconds", FieldType::Number),
72        ]),
73    );
74    registry.register_schema(
75        "transform.state_read",
76        schema(vec![
77            Field::required("key", FieldType::String),
78            Field::required("into", FieldType::String),
79            Field::optional("default", FieldType::Any),
80        ]),
81    );
82    registry.register_schema(
83        "transform.set_fields",
84        schema(vec![Field::required("fields", FieldType::Object)]),
85    );
86    for node_type in [
87        "transform.state_publish_cross_branch",
88        "transform.state_append_cross_branch",
89    ] {
90        registry.register_schema(
91            node_type,
92            schema(vec![
93                Field::required("path", FieldType::String),
94                Field::required("key", FieldType::String),
95                Field::optional("max_len", FieldType::Number),
96                Field::optional("ttl_seconds", FieldType::Number),
97            ]),
98        );
99    }
100    registry.register_schema(
101        "transform.state_read_cross_branch",
102        schema(vec![
103            Field::required("key", FieldType::String),
104            Field::required("into", FieldType::String),
105            Field::optional("default", FieldType::Any),
106        ]),
107    );
108    registry.register_schema(
109        "decision.consensus_voting",
110        schema(vec![
111            Field::required("signals_path", FieldType::String),
112            Field::optional("mode", FieldType::String),
113            Field::optional("quorum", FieldType::Number),
114            Field::optional("threshold", FieldType::Number),
115        ]),
116    );
117    registry.register_schema(
118        "filter.required_fields",
119        schema(vec![Field::required("fields", FieldType::Array)]),
120    );
121    registry.register_schema(
122        "sink.log",
123        schema(vec![
124            Field::required("message", FieldType::String),
125            Field::optional("level", FieldType::String),
126        ]),
127    );
128}
129
130// ── helpers ──────────────────────────────────────────────────────────────
131
132fn require_keys(node_type: &str, config: &Value, keys: &[&str]) -> Result<(), NodeError> {
133    let obj = config.as_object().ok_or_else(|| NodeError::InvalidConfig {
134        node_type: node_type.to_string(),
135        reason: "config must be an object".to_string(),
136    })?;
137    for k in keys {
138        if !obj.contains_key(*k) {
139            return Err(NodeError::InvalidConfig {
140                node_type: node_type.to_string(),
141                reason: format!("missing required key '{k}'"),
142            });
143        }
144    }
145    Ok(())
146}
147
148fn resolve(config: &Value, ctx: &WorkflowContext, event: &Event) -> Result<Value, StepResult> {
149    let payload = Value::Object(event.payload.clone());
150    resolve_config(config, &ctx.config, &ctx.instance_metadata, Some(&payload))
151        .map_err(|e| StepResult::drop(e.to_string()))
152}
153
154fn cfg_str<'a>(cfg: &'a Value, key: &str) -> Option<&'a str> {
155    cfg.get(key).and_then(|v| v.as_str())
156}
157
158fn cfg_usize(cfg: &Value, key: &str) -> Option<usize> {
159    cfg.get(key).and_then(|v| v.as_u64()).map(|n| n as usize)
160}
161
162fn cross_branch_key(key: &str) -> String {
163    format!("__cross_branch__.{key}")
164}
165
166// ── transform.state_append ─────────────────────────────────────────────────
167
168/// Append `payload[path]` to a branch-scoped sliding-window list `key`.
169struct TransformStateAppend {
170    config: Value,
171}
172
173impl TransformStateAppend {
174    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
175        require_keys("transform.state_append", config, &["path", "key"])?;
176        Ok(Box::new(Self {
177            config: config.clone(),
178        }))
179    }
180}
181
182#[async_trait]
183impl StepNode for TransformStateAppend {
184    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
185        let cfg = match resolve(&self.config, ctx, event) {
186            Ok(c) => c,
187            Err(drop) => return drop,
188        };
189        let path = cfg_str(&cfg, "path").unwrap_or_default();
190        let key = cfg_str(&cfg, "key").unwrap_or_default();
191        let item = match event.payload_path(path) {
192            Some(v) => v.clone(),
193            None => return StepResult::drop(format!("state_append: path '{path}' missing")),
194        };
195        let max_len = cfg_usize(&cfg, "max_len");
196        let ttl = cfg
197            .get("ttl_seconds")
198            .and_then(|v| v.as_u64())
199            .map(std::time::Duration::from_secs);
200        ctx.state.append(&ctx.scoped(key), item, max_len, ttl).await;
201        StepResult::Pass(event.clone())
202    }
203}
204
205// ── transform.state_set ────────────────────────────────────────────────────
206
207/// Set a branch-scoped key to a (possibly templated) value.
208struct TransformStateSet {
209    config: Value,
210}
211
212impl TransformStateSet {
213    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
214        require_keys("transform.state_set", config, &["key", "value"])?;
215        Ok(Box::new(Self {
216            config: config.clone(),
217        }))
218    }
219}
220
221#[async_trait]
222impl StepNode for TransformStateSet {
223    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
224        let cfg = match resolve(&self.config, ctx, event) {
225            Ok(c) => c,
226            Err(drop) => return drop,
227        };
228        let key = cfg_str(&cfg, "key").unwrap_or_default().to_string();
229        let value = cfg.get("value").cloned().unwrap_or(Value::Null);
230        let ttl = cfg
231            .get("ttl_seconds")
232            .and_then(|v| v.as_u64())
233            .map(std::time::Duration::from_secs);
234        ctx.state.set(&ctx.scoped(&key), value, ttl).await;
235        StepResult::Pass(event.clone())
236    }
237}
238
239// ── transform.state_read ───────────────────────────────────────────────────
240
241/// Read a branch-scoped key into `payload[into]`. Drops if the key is absent
242/// and no `default` is configured.
243struct TransformStateRead {
244    config: Value,
245}
246
247impl TransformStateRead {
248    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
249        require_keys("transform.state_read", config, &["key", "into"])?;
250        Ok(Box::new(Self {
251            config: config.clone(),
252        }))
253    }
254}
255
256#[async_trait]
257impl StepNode for TransformStateRead {
258    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
259        let cfg = match resolve(&self.config, ctx, event) {
260            Ok(c) => c,
261            Err(drop) => return drop,
262        };
263        let key = cfg_str(&cfg, "key").unwrap_or_default();
264        let into = cfg_str(&cfg, "into").unwrap_or_default().to_string();
265        let value = match ctx.state.get(&ctx.scoped(key)).await {
266            Some(v) => v,
267            None => match cfg.get("default") {
268                Some(d) => d.clone(),
269                None => return StepResult::drop(format!("state_read: key '{key}' absent")),
270            },
271        };
272        let mut patch = Map::new();
273        patch.insert(into, value);
274        StepResult::Pass(event.with_payload(patch))
275    }
276}
277
278// ── explicit cross-branch state ────────────────────────────────────────────
279
280struct TransformStatePublishCrossBranch {
281    config: Value,
282}
283
284impl TransformStatePublishCrossBranch {
285    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
286        require_keys(
287            "transform.state_publish_cross_branch",
288            config,
289            &["path", "key"],
290        )?;
291        Ok(Box::new(Self {
292            config: config.clone(),
293        }))
294    }
295}
296
297#[async_trait]
298impl StepNode for TransformStatePublishCrossBranch {
299    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
300        let cfg = match resolve(&self.config, ctx, event) {
301            Ok(cfg) => cfg,
302            Err(drop) => return drop,
303        };
304        let path = cfg_str(&cfg, "path").unwrap_or_default();
305        let key = cfg_str(&cfg, "key").unwrap_or_default();
306        let Some(value) = event.payload_path(path).cloned() else {
307            return StepResult::drop(format!("cross-branch publish: path '{path}' missing"));
308        };
309        let ttl = cfg
310            .get("ttl_seconds")
311            .and_then(Value::as_u64)
312            .map(std::time::Duration::from_secs);
313        ctx.state.set(&cross_branch_key(key), value, ttl).await;
314        StepResult::Pass(event.clone())
315    }
316}
317
318struct TransformStateReadCrossBranch {
319    config: Value,
320}
321
322impl TransformStateReadCrossBranch {
323    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
324        require_keys(
325            "transform.state_read_cross_branch",
326            config,
327            &["key", "into"],
328        )?;
329        Ok(Box::new(Self {
330            config: config.clone(),
331        }))
332    }
333}
334
335#[async_trait]
336impl StepNode for TransformStateReadCrossBranch {
337    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
338        let cfg = match resolve(&self.config, ctx, event) {
339            Ok(cfg) => cfg,
340            Err(drop) => return drop,
341        };
342        let key = cfg_str(&cfg, "key").unwrap_or_default();
343        let target = cfg_str(&cfg, "into").unwrap_or_default().to_string();
344        let value = match ctx.state.get(&cross_branch_key(key)).await {
345            Some(value) => value,
346            None => match cfg.get("default") {
347                Some(value) => value.clone(),
348                None => return StepResult::drop(format!("cross-branch key '{key}' absent")),
349            },
350        };
351        let mut patch = Map::new();
352        patch.insert(target, value);
353        StepResult::Pass(event.with_payload(patch))
354    }
355}
356
357struct TransformStateAppendCrossBranch {
358    config: Value,
359}
360
361impl TransformStateAppendCrossBranch {
362    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
363        require_keys(
364            "transform.state_append_cross_branch",
365            config,
366            &["path", "key"],
367        )?;
368        Ok(Box::new(Self {
369            config: config.clone(),
370        }))
371    }
372}
373
374#[async_trait]
375impl StepNode for TransformStateAppendCrossBranch {
376    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
377        let cfg = match resolve(&self.config, ctx, event) {
378            Ok(cfg) => cfg,
379            Err(drop) => return drop,
380        };
381        let path = cfg_str(&cfg, "path").unwrap_or_default();
382        let key = cfg_str(&cfg, "key").unwrap_or_default();
383        let Some(value) = event.payload_path(path).cloned() else {
384            return StepResult::drop(format!("cross-branch append: path '{path}' missing"));
385        };
386        let ttl = cfg
387            .get("ttl_seconds")
388            .and_then(Value::as_u64)
389            .map(std::time::Duration::from_secs);
390        ctx.state
391            .append(
392                &cross_branch_key(key),
393                value,
394                cfg_usize(&cfg, "max_len"),
395                ttl,
396            )
397            .await;
398        StepResult::Pass(event.clone())
399    }
400}
401
402// ── transform.set_fields ───────────────────────────────────────────────────
403
404/// Merge a static/templated `fields` object into the event payload.
405struct TransformSetFields {
406    config: Value,
407}
408
409impl TransformSetFields {
410    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
411        require_keys("transform.set_fields", config, &["fields"])?;
412        Ok(Box::new(Self {
413            config: config.clone(),
414        }))
415    }
416}
417
418#[async_trait]
419impl StepNode for TransformSetFields {
420    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
421        let cfg = match resolve(&self.config, ctx, event) {
422            Ok(c) => c,
423            Err(drop) => return drop,
424        };
425        let fields = match cfg.get("fields").and_then(|v| v.as_object()) {
426            Some(f) => f.clone(),
427            None => return StepResult::drop("set_fields: 'fields' must be an object"),
428        };
429        StepResult::Pass(event.with_payload(fields))
430    }
431}
432
433// ── generic consensus ──────────────────────────────────────────────────────
434
435/// Vote over `{choice, weight?}` objects. Quorum counts votes; weighted mode
436/// selects the highest total only when its lead exceeds `threshold`.
437struct ConsensusVoting {
438    config: Value,
439}
440
441impl ConsensusVoting {
442    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
443        require_keys("decision.consensus_voting", config, &["signals_path"])?;
444        Ok(Box::new(Self {
445            config: config.clone(),
446        }))
447    }
448}
449
450#[async_trait]
451impl StepNode for ConsensusVoting {
452    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
453        let cfg = match resolve(&self.config, ctx, event) {
454            Ok(cfg) => cfg,
455            Err(drop) => return drop,
456        };
457        let path = cfg_str(&cfg, "signals_path").unwrap_or("signals");
458        let Some(signals) = event.payload_path(path).and_then(Value::as_array) else {
459            return StepResult::drop("consensus: no signals");
460        };
461        if signals.is_empty() {
462            return StepResult::drop("consensus: no signals");
463        }
464
465        let mut totals: HashMap<&str, (u64, f64)> = HashMap::new();
466        for signal in signals {
467            let Some(choice) = signal.get("choice").and_then(Value::as_str) else {
468                continue;
469            };
470            let entry = totals.entry(choice).or_default();
471            entry.0 += 1;
472            entry.1 += signal.get("weight").and_then(Value::as_f64).unwrap_or(1.0);
473        }
474        if totals.is_empty() {
475            return StepResult::drop("consensus: no valid choices");
476        }
477
478        let weighted = cfg_str(&cfg, "mode") == Some("weighted");
479        let mut ranked = totals.into_iter().collect::<Vec<_>>();
480        ranked.sort_by(|left, right| {
481            let left_score = if weighted {
482                left.1 .1
483            } else {
484                left.1 .0 as f64
485            };
486            let right_score = if weighted {
487                right.1 .1
488            } else {
489                right.1 .0 as f64
490            };
491            right_score.total_cmp(&left_score)
492        });
493        let (choice, (votes, weight)) = ranked[0];
494        let score = if weighted { weight } else { votes as f64 };
495        let runner_up = ranked
496            .get(1)
497            .map(|(_, (votes, weight))| if weighted { *weight } else { *votes as f64 })
498            .unwrap_or(0.0);
499        let required = if weighted {
500            cfg.get("threshold").and_then(Value::as_f64).unwrap_or(0.0)
501        } else {
502            cfg.get("quorum").and_then(Value::as_u64).unwrap_or(1) as f64
503        };
504        if (!weighted && score < required) || (weighted && score - runner_up <= required) {
505            return StepResult::drop("consensus: threshold not met");
506        }
507        if score == runner_up {
508            return StepResult::drop("consensus: tie");
509        }
510
511        let mut patch = Map::new();
512        patch.insert(
513            "consensus".into(),
514            serde_json::json!({ "choice": choice, "score": score }),
515        );
516        StepResult::Pass(event.with_payload(patch))
517    }
518}
519
520// ── filter.required_fields ─────────────────────────────────────────────────
521
522/// Pass only if every listed payload path is present; otherwise Drop.
523struct FilterRequiredFields {
524    fields: Vec<String>,
525}
526
527impl FilterRequiredFields {
528    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
529        let fields = config
530            .get("fields")
531            .and_then(|v| v.as_array())
532            .ok_or_else(|| NodeError::InvalidConfig {
533                node_type: "filter.required_fields".to_string(),
534                reason: "'fields' must be an array of paths".to_string(),
535            })?
536            .iter()
537            .filter_map(|v| v.as_str().map(str::to_string))
538            .collect();
539        Ok(Box::new(Self { fields }))
540    }
541}
542
543#[async_trait]
544impl StepNode for FilterRequiredFields {
545    async fn process(&self, event: &Event, _ctx: &WorkflowContext) -> StepResult {
546        for path in &self.fields {
547            if event.payload_path(path).is_none() {
548                return StepResult::drop(format!("required field '{path}' missing"));
549            }
550        }
551        StepResult::Pass(event.clone())
552    }
553}
554
555// ── sink.log ───────────────────────────────────────────────────────────────
556
557/// Emit one structured log line and stop propagation (terminal).
558struct SinkLog {
559    config: Value,
560}
561
562impl SinkLog {
563    fn factory(config: &Value) -> Result<Box<dyn StepNode>, NodeError> {
564        require_keys("sink.log", config, &["message"])?;
565        Ok(Box::new(Self {
566            config: config.clone(),
567        }))
568    }
569}
570
571#[async_trait]
572impl StepNode for SinkLog {
573    async fn process(&self, event: &Event, ctx: &WorkflowContext) -> StepResult {
574        let cfg = match resolve(&self.config, ctx, event) {
575            Ok(c) => c,
576            Err(drop) => return drop,
577        };
578        let message = cfg_str(&cfg, "message").unwrap_or("sink.log");
579        let level = cfg_str(&cfg, "level").unwrap_or("info");
580        match level {
581            "debug" => tracing::debug!(target: "wf.sink", event_id = %event.id, "{message}"),
582            "warning" | "warn" => {
583                tracing::warn!(target: "wf.sink", event_id = %event.id, "{message}")
584            }
585            "error" => tracing::error!(target: "wf.sink", event_id = %event.id, "{message}"),
586            _ => tracing::info!(target: "wf.sink", event_id = %event.id, "{message}"),
587        }
588        StepResult::drop("sink.log emitted")
589    }
590}