use std::sync::Arc;
use af_workflow::{
resolve_config, Event, MemoryState, NodeRegistry, Spec, State, StepResult, Terminal,
WorkflowContext, WorkflowHost,
};
use serde_json::json;
#[test]
fn event_construction_and_derivation() {
let e = Event::from_json(json!({"a": 1}));
assert!(e.id.starts_with("evt_"));
assert_eq!(e.payload_path("a"), Some(&json!(1)));
let wrapped = Event::from_json(json!(42));
assert_eq!(wrapped.payload_path("value"), Some(&json!(42)));
let mut patch = serde_json::Map::new();
patch.insert("b".into(), json!(2));
let e2 = e.with_payload(patch);
assert_eq!(e2.id, e.id);
assert_eq!(e2.payload_path("b"), Some(&json!(2)));
let mut meta = serde_json::Map::new();
meta.insert("trace".into(), json!("x"));
let e3 = e2.with_metadata(meta);
assert_eq!(e3.metadata.get("trace"), Some(&json!("x")));
let nested = Event::from_json(json!({"outer": {"inner": 7}}));
assert_eq!(nested.payload_path("outer.inner"), Some(&json!(7)));
assert_eq!(nested.payload_path("outer.nope"), None);
assert_eq!(nested.payload_path("missing"), None);
}
#[tokio::test]
async fn context_scoping_and_config() {
let state = Arc::new(MemoryState::new());
let ctx = WorkflowContext::new("branchX", state).with_config(json!({"k": "v"}));
assert_eq!(ctx.scoped("key"), "branchX.key");
assert_eq!(ctx.config, json!({"k": "v"}));
assert_eq!(ctx.branch_id, "branchX");
}
#[test]
fn template_all_paths() {
let cfg = json!({
"fromcfg": "$config.x",
"frominst": "$instance.run",
"fromevent": "$event.p",
"literal": "plain",
"list": ["$config.x", "lit"],
});
let out = resolve_config(
&cfg,
&json!({"x": 1}),
&json!({"run": "r1"}),
Some(&json!({"p": 9})),
)
.unwrap();
assert_eq!(out["fromcfg"], json!(1));
assert_eq!(out["frominst"], json!("r1"));
assert_eq!(out["fromevent"], json!(9));
assert_eq!(out["literal"], json!("plain"));
assert_eq!(out["list"][0], json!(1));
assert!(resolve_config(&json!("$instance.y"), &json!({"y": 5}), &json!({}), None).is_err());
assert!(resolve_config(&json!("$config.nope"), &json!({}), &json!({}), None).is_err());
assert!(resolve_config(&json!("$event.x"), &json!({}), &json!({}), None).is_err());
}
#[test]
fn registry_surfaces() {
let mut r = NodeRegistry::empty();
assert!(!r.is_step("sink.log"));
r = NodeRegistry::with_builtins();
assert!(r.is_step("sink.log"));
assert!(r.is_step("transform.state_append"));
assert!(r.is_ingress("ingress.cron"));
assert!(r.is_ingress("ingress.anything"));
assert!(r.build_step("nope.node", &json!({})).is_err());
r.register_fan_out("transform.fork");
assert!(r.is_fan_out_capable("transform.fork"));
assert!(!r.is_fan_out_capable("sink.log"));
assert!(r.known_step_types().count() >= 6);
}
async fn run_single(spec_json: &str, event: Event) -> (Terminal, Arc<MemoryState>) {
let spec = Spec::from_json(spec_json).unwrap();
let reg = NodeRegistry::with_builtins();
let host = WorkflowHost::from_spec(&spec, ®).unwrap();
let state = Arc::new(MemoryState::new());
let ctx = host.context(state.clone());
let outcome = host.run_event(&ctx, event).await.unwrap();
(outcome.terminal, state)
}
#[tokio::test]
async fn builtin_state_set_and_read() {
let spec = r#"{"spec_id":"s","version":"1.0","branches":[{"branch_id":"__root__",
"nodes":[
{"id":"t","type":"ingress.cron","config":{}},
{"id":"set","type":"transform.state_set","config":{"key":"cfg","value":{"n":5}}},
{"id":"read","type":"transform.state_read","config":{"key":"cfg","into":"loaded"}},
{"id":"need","type":"filter.required_fields","config":{"fields":["loaded"]}},
{"id":"log","type":"sink.log","config":{"level":"warning","message":"done"}}],
"edges":[{"source":"t","target":"set"},{"source":"set","target":"read"},
{"source":"read","target":"need"},{"source":"need","target":"log"}]}]}"#;
let (term, state) = run_single(spec, Event::from_json(json!({}))).await;
assert_eq!(term, Terminal::Completed);
assert_eq!(state.get("__root__.cfg").await.unwrap(), json!({"n": 5}));
}
#[tokio::test]
async fn builtin_state_read_default_and_missing() {
let with_default = r#"{"spec_id":"d","version":"1.0","branches":[{"branch_id":"__root__",
"nodes":[
{"id":"t","type":"ingress.cron","config":{}},
{"id":"r","type":"transform.state_read","config":{"key":"absent","into":"x","default":7}},
{"id":"need","type":"filter.required_fields","config":{"fields":["x"]}}],
"edges":[{"source":"t","target":"r"},{"source":"r","target":"need"}]}]}"#;
let (term, _) = run_single(with_default, Event::from_json(json!({}))).await;
assert_eq!(term, Terminal::Completed);
let no_default = r#"{"spec_id":"d2","version":"1.0","branches":[{"branch_id":"__root__",
"nodes":[
{"id":"t","type":"ingress.cron","config":{}},
{"id":"r","type":"transform.state_read","config":{"key":"absent","into":"x"}}],
"edges":[{"source":"t","target":"r"}]}]}"#;
let (term, _) = run_single(no_default, Event::from_json(json!({}))).await;
assert!(matches!(term, Terminal::Dropped { .. }));
}
#[tokio::test]
async fn builtin_set_fields_merges_payload() {
let spec = r#"{"spec_id":"sf","version":"1.0","branches":[{"branch_id":"__root__",
"nodes":[
{"id":"t","type":"ingress.cron","config":{}},
{"id":"sf","type":"transform.set_fields","config":{"fields":{"added":"yes"}}},
{"id":"need","type":"filter.required_fields","config":{"fields":["added"]}}],
"edges":[{"source":"t","target":"sf"},{"source":"sf","target":"need"}]}]}"#;
let (term, _) = run_single(spec, Event::from_json(json!({}))).await;
assert_eq!(term, Terminal::Completed);
}
#[tokio::test]
async fn builtin_filter_drops_when_field_absent() {
let spec = r#"{"spec_id":"f","version":"1.0","branches":[{"branch_id":"__root__",
"nodes":[
{"id":"t","type":"ingress.cron","config":{}},
{"id":"need","type":"filter.required_fields","config":{"fields":["nope"]}},
{"id":"log","type":"sink.log","config":{"message":"x"}}],
"edges":[{"source":"t","target":"need"},{"source":"need","target":"log"}]}]}"#;
let (term, _) = run_single(spec, Event::from_json(json!({"present": 1}))).await;
match term {
Terminal::Dropped { node_id, .. } => assert_eq!(node_id, "need"),
_ => panic!("expected drop at filter"),
}
}
#[tokio::test]
async fn builtin_state_append_drops_on_missing_path() {
let spec = r#"{"spec_id":"a","version":"1.0","branches":[{"branch_id":"__root__",
"nodes":[
{"id":"t","type":"ingress.cron","config":{}},
{"id":"app","type":"transform.state_append","config":{"path":"missing","key":"k"}}],
"edges":[{"source":"t","target":"app"}]}]}"#;
let (term, _) = run_single(spec, Event::from_json(json!({"other": 1}))).await;
assert!(matches!(term, Terminal::Dropped { .. }));
}
#[test]
fn step_result_drop_helper() {
let d = StepResult::drop("because");
match d {
StepResult::Drop {
reason,
exit_reason,
} => {
assert_eq!(reason, "because");
assert!(exit_reason.is_none());
}
_ => panic!(),
}
}