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