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.event");
27 r.register_ingress("ingress.manual");
28
29 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
130fn 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
166struct 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
205struct 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
239struct 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
278struct 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
402struct 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
433struct 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
520struct 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
555struct 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}