af-workflow 0.5.0

Spec-driven workflow chassis: typed node expressions composed into a branched DAG. Port of agent_core/workflow.
Documentation
//! Unit coverage for the pure-logic runtime surfaces: every builtin node,
//! Event, WorkflowContext, template resolution, and the registry.

use std::sync::Arc;

use af_workflow::{
    resolve_config, Event, MemoryState, NodeRegistry, Spec, State, StepResult, Terminal,
    WorkflowContext, WorkflowHost,
};
use serde_json::json;

// ── Event ──────────────────────────────────────────────────────────────────

#[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)));

    // non-object payload is wrapped under "value"
    let wrapped = Event::from_json(json!(42));
    assert_eq!(wrapped.payload_path("value"), Some(&json!(42)));

    // with_payload / with_metadata keep the id, add fields
    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")));

    // nested + missing paths
    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);
}

// ── WorkflowContext ──────────────────────────────────────────────────────────

#[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");
}

// ── template resolution ──────────────────────────────────────────────────────

#[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));

    // Owner namespaces are strict; unresolved + event-out-of-scope are errors.
    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());
}

// ── registry ─────────────────────────────────────────────────────────────────

#[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")); // prefix rule

    // unknown node type → build error
    assert!(r.build_step("nope.node", &json!({})).is_err());

    // fan-out tracking
    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);
}

// ── builtin nodes (via compile + run) ────────────────────────────────────────

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, &reg).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() {
    // state_set writes; a second branch's state_read reads it back into payload.
    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() {
    // missing key + default → uses default
    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);

    // missing key + no default → drop
    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 { .. }));
}

// ── StepResult helper ────────────────────────────────────────────────────────

#[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!(),
    }
}