agentd-cli 2.6.0

agentd — a minimal, MCP-native, reactive agent runtime (the CLI binary over the agentd-core engine)
// SPDX-License-Identifier: Apache-2.0
//! **Steering over A2A** end to end (RFC 0029 §5/§7, now dispatched): a client
//! fires `workflow.signal` to resume a waiting run, pauses/resumes one run and
//! the whole instance (`a2a.pause`/`a2a.resume`), and reads a conversation's
//! plan — the control verbs a display client uses beyond cancel/drain.
#![cfg(all(unix, feature = "a2a"))]

mod common;

use std::io::{BufRead, BufReader, Read, Write};
use std::net::{TcpListener, TcpStream};
use std::process::{Child, Command, Stdio};
use std::time::{Duration, Instant};

use serde_json::{Value, json};

fn post_raw(addr: &str, body: &str) -> String {
    let mut s = TcpStream::connect(addr).expect("connect a2a http");
    s.set_read_timeout(Some(Duration::from_secs(130))).ok();
    let head = format!(
        "POST / HTTP/1.1\r\nHost: x\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
        body.len()
    );
    s.write_all(head.as_bytes()).unwrap();
    s.write_all(body.as_bytes()).unwrap();
    s.flush().unwrap();
    let mut reader = BufReader::new(s);
    let mut status = String::new();
    reader.read_line(&mut status).unwrap();
    loop {
        let mut l = String::new();
        reader.read_line(&mut l).unwrap();
        if l.trim().is_empty() {
            break;
        }
    }
    let mut b = String::new();
    reader.read_to_string(&mut b).unwrap();
    b
}

fn rpc(addr: &str, id: i64, method: &str, params: Value) -> Value {
    let body = json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params}).to_string();
    let resp = post_raw(addr, &body);
    let v: Value =
        serde_json::from_str(&resp).unwrap_or_else(|_| panic!("non-JSON A2A response: {resp:?}"));
    assert!(v.get("error").is_none(), "A2A rpc error for {method}: {v}");
    v["result"].clone()
}

fn command(addr: &str, id: i64, op: &str, extra: Value) -> Value {
    let mut data = json!({"op": op});
    if let (Value::Object(d), Value::Object(x)) = (&mut data, extra) {
        for (k, v) in x {
            d.insert(k, v);
        }
    }
    rpc(
        addr,
        id,
        "SendMessage",
        // Non-blocking: command tasks complete inline anyway, and workflow.run
        // must NOT be polled to terminal — the tests steer runs mid-flight.
        json!({"message": {"messageId": format!("m-{id}"), "parts": [{"data": {"agentd": data}}]},
               "configuration": {"blocking": false}}),
    )
}

/// Parse a command task's JSON artifact.
fn artifact_json(v: &Value) -> Value {
    v["task"]["artifacts"][0]["parts"][0]["text"]
        .as_str()
        .and_then(|t| serde_json::from_str(t).ok())
        .unwrap_or(Value::Null)
}

/// Poll `workflow.status` for one run until `pred` holds; returns the view.
/// The newest run's id, polled until the supervisor has registered it.
///
/// `workflow.run` returns as soon as the request is accepted, so a
/// `workflow.status` issued immediately after can still see an empty list —
/// rare locally, reliably reproducible under CI load.
fn wait_run_id(addr: &str, secs: u64) -> String {
    let deadline = Instant::now() + Duration::from_secs(secs);
    loop {
        let ws = command(addr, 902, "workflow.status", json!({}));
        if let Some(id) = artifact_json(&ws)["runs"][0]["run"].as_str() {
            return id.to_string();
        }
        assert!(
            Instant::now() < deadline,
            "timeout: no run appeared in workflow.status"
        );
        std::thread::sleep(Duration::from_millis(50));
    }
}

fn wait_run<F: Fn(&Value) -> bool>(addr: &str, run: &str, secs: u64, what: &str, pred: F) -> Value {
    let deadline = Instant::now() + Duration::from_secs(secs);
    loop {
        let ws = command(addr, 901, "workflow.status", json!({"run": run}));
        let view = artifact_json(&ws)["runs"][0].clone();
        if pred(&view) {
            return view;
        }
        assert!(Instant::now() < deadline, "timeout: {what}; last: {view}");
        std::thread::sleep(Duration::from_millis(80));
    }
}

struct MockLlm {
    child: Child,
    addr_file: String,
    uri: String,
}
impl Drop for MockLlm {
    fn drop(&mut self) {
        let _ = self.child.kill();
        let _ = self.child.wait();
        let _ = std::fs::remove_file(&self.addr_file);
    }
}
fn spawn_mock_llm(playbook: &Value) -> MockLlm {
    let pb = common::unique_path("steer-playbook", "json");
    std::fs::write(&pb, playbook.to_string()).unwrap();
    let addr_file = common::unique_path("steer-mock-llm", "addr");
    let _ = std::fs::remove_file(&addr_file);
    let child = Command::new(env!("CARGO_BIN_EXE_agentd"))
        .args(["--internal-mock-llm", &addr_file, &format!("file:{pb}")])
        .stdout(Stdio::null())
        .stderr(Stdio::null())
        .spawn()
        .expect("spawn mock llm");
    let addr = common::read_addr_file(&addr_file);
    MockLlm {
        child,
        addr_file,
        uri: format!("http://{addr}"),
    }
}

struct Daemon {
    child: Child,
    stderr_path: String,
}
impl Drop for Daemon {
    fn drop(&mut self) {
        unsafe { libc::kill(self.child.id() as i32, libc::SIGTERM) };
        let deadline = Instant::now() + Duration::from_secs(3);
        while Instant::now() < deadline {
            if matches!(self.child.try_wait(), Ok(Some(_))) {
                break;
            }
            std::thread::sleep(Duration::from_millis(20));
        }
        let _ = self.child.kill();
        let _ = self.child.wait();
        let _ = std::fs::remove_file(&self.stderr_path);
    }
}
fn spawn_daemon(config: &str) -> Daemon {
    let stderr_path = common::unique_path("steer-daemon", "log");
    let errf = std::fs::File::create(&stderr_path).unwrap();
    let child = Command::new(env!("CARGO_BIN_EXE_agentd"))
        .args(["--config", config])
        .stdin(Stdio::null())
        .stdout(Stdio::null())
        .stderr(Stdio::from(errf))
        .spawn()
        .expect("spawn agentd daemon");
    Daemon { child, stderr_path }
}

fn free_port() -> u16 {
    TcpListener::bind("127.0.0.1:0")
        .unwrap()
        .local_addr()
        .unwrap()
        .port()
}

fn write_config(yaml: &str) -> String {
    let path = common::unique_path("agentd-steer", "yaml");
    std::fs::write(&path, yaml).unwrap();
    path
}

/// Spawn the daemon on a probed free port and return the authority IT actually
/// bound. The probe→bind gap is a real race under parallel CI (another process
/// can take the port), so a daemon whose bind lost is retried on a fresh port
/// rather than leaving the test talking to a stranger's listener.
fn spawn_bound(cfg_for: impl Fn(u16) -> String) -> (Daemon, String, String) {
    spawn_bound_with(cfg_for, spawn_daemon)
}

fn spawn_bound_with(
    cfg_for: impl Fn(u16) -> String,
    spawn: impl Fn(&str) -> Daemon,
) -> (Daemon, String, String) {
    for _ in 0..5 {
        let cfg = write_config(&cfg_for(free_port()));
        let daemon = spawn(&cfg);
        if let Some(addr) = common::try_a2a_bound(&daemon.stderr_path, Duration::from_secs(15)) {
            return (daemon, addr, cfg);
        }
        std::fs::remove_file(&cfg).ok();
    }
    panic!("the daemon never bound an A2A listener (5 attempts)");
}

fn steer_config(llm: &str, port: u16, extra: &str) -> String {
    format!(
        "config_version: \"2\"\n\
         agent:\n  name: steer-e2e\n  instruction: Test agent.\n  preflight: never\n\
         intelligence:\n  endpoints: {llm}\n  model: mock\n\
         store:\n  kind: memory\n\
         a2a:\n  listen: http://127.0.0.1:{port}\n\
         interface:\n  enabled: true\n  debug: true\n\
         lifecycle:\n  run_until: drained\n{extra}"
    )
}

#[test]
fn a_signal_resumes_a_waiting_run() {
    let llm = spawn_mock_llm(&json!({"turns": [{"content": "unused"}]}));
    // A run that WAITS for a named signal, then finishes with its payload.
    let extra = "workflows:\n  - name: waiter\n    steps:\n      s: {kind: manual}\n      w: {kind: wait, on: signal, signal: go, depends_on: [s]}\n      f: {kind: finish, depends_on: [w], output: \"released\"}\n";
    let (_daemon, addr, cfg) = spawn_bound(|port| steer_config(&llm.uri, port, extra));

    let started = command(&addr, 1, "workflow.run", json!({"name": "waiter"}));
    let task_id = started["task"]["id"].as_str().unwrap().to_string();
    // Find the run id + confirm it parks on the wait.
    let run_id = wait_run_id(&addr, 10);
    wait_run(&addr, &run_id, 10, "suspended on the signal", |v| {
        v["status"] == "suspended" || v["status"] == "running"
    });

    // The steering verb: `workflow.signal` fires the named signal.
    let sig = command(
        &addr,
        3,
        "workflow.signal",
        json!({"name": "go", "payload": {"by": "e2e"}}),
    );
    assert_eq!(artifact_json(&sig)["delivered"], 1, "{sig}");

    wait_run(&addr, &run_id, 10, "run completed after the signal", |v| {
        v["status"] == "completed"
    });
    // The tracking task went terminal too.
    let t = rpc(&addr, 4, "GetTask", json!({"id": task_id}));
    assert_eq!(t["status"]["state"], "TASK_STATE_COMPLETED", "{t}");
    std::fs::remove_file(&cfg).ok();
}

#[test]
fn a_single_run_pauses_and_resumes() {
    let llm = spawn_mock_llm(&json!({"turns": [{"content": "unused"}]}));
    // A run with a 1s sleep between steps — enough of a window to pause it.
    let extra = "workflows:\n  - name: slow\n    steps:\n      s: {kind: manual}\n      z: {kind: sleep, duration: 1s, depends_on: [s]}\n      f: {kind: finish, depends_on: [z], output: \"done\"}\n";
    let (_daemon, addr, cfg) = spawn_bound(|port| steer_config(&llm.uri, port, extra));

    command(&addr, 1, "workflow.run", json!({"name": "slow"}));
    let run_id = wait_run_id(&addr, 10);

    // Pause the run mid-flight; it must NOT complete while paused.
    let paused = rpc(&addr, 3, "a2a.pause", json!({"run": run_id}));
    assert_eq!(paused["paused"], run_id, "{paused}");
    std::thread::sleep(Duration::from_millis(1600)); // past the sleep deadline
    let view = artifact_json(&command(
        &addr,
        4,
        "workflow.status",
        json!({"run": run_id}),
    ))["runs"][0]
        .clone();
    assert_ne!(
        view["status"], "completed",
        "paused runs don't advance: {view}"
    );

    // Resume → completes.
    let resumed = rpc(&addr, 5, "a2a.resume", json!({"run": run_id}));
    assert_eq!(resumed["resumed"], run_id, "{resumed}");
    wait_run(&addr, &run_id, 10, "completion after resume", |v| {
        v["status"] == "completed"
    });
    std::fs::remove_file(&cfg).ok();
}

#[test]
fn a_global_pause_holds_new_work_and_resume_releases_it() {
    let llm = spawn_mock_llm(&json!({"turns": [{"content": "Answered after the hold."}]}));
    let (_daemon, addr, cfg) = spawn_bound(|port| steer_config(&llm.uri, port, ""));

    // Pause the instance; intake continues but nothing dispatches.
    let paused = rpc(&addr, 1, "a2a.pause", json!({}));
    assert_eq!(paused["state"], "paused");
    let st = command(&addr, 2, "status", json!({}));
    assert_eq!(artifact_json(&st)["paused"], true, "{st}");

    let sent = rpc(
        &addr,
        3,
        "SendMessage",
        json!({"message": {"messageId": "m1", "parts": [{"text": "Hello during the pause"}]},
               "configuration": {"blocking": false}}),
    );
    let task_id = sent["task"]["id"].as_str().unwrap().to_string();
    std::thread::sleep(Duration::from_millis(900));
    let held = rpc(&addr, 4, "GetTask", json!({"id": task_id}));
    assert_eq!(
        held["status"]["state"], "TASK_STATE_WORKING",
        "the turn is queued, not run, while paused: {held}"
    );

    // Resume → the queued turn dispatches and completes.
    rpc(&addr, 5, "a2a.resume", json!({}));
    let deadline = Instant::now() + Duration::from_secs(15);
    loop {
        let t = rpc(&addr, 6, "GetTask", json!({"id": task_id}));
        if t["status"]["state"] == "TASK_STATE_COMPLETED" {
            assert!(
                t["artifacts"][0]["parts"][0]["text"]
                    .as_str()
                    .unwrap_or("")
                    .contains("after the hold"),
                "{t}"
            );
            break;
        }
        assert!(Instant::now() < deadline, "resume released the turn: {t}");
        std::thread::sleep(Duration::from_millis(100));
    }
    std::fs::remove_file(&cfg).ok();
}

#[test]
fn subagent_send_injects_into_a_warm_subagent_and_plan_get_reads_the_plan() {
    // The root spawns a WARM subagent, then the e2e steers it over A2A.
    let llm = spawn_mock_llm(&json!({
        "turns": [
            {"tool_calls": [{"name": "subagent.run", "arguments": {"instruction": "stand by for instructions", "mode": "warm"}}]},
            {"content": "Warm helper started."}
        ],
        "match": [
            {"when_contains": "You are agentd, an autonomous agent.", "content": "standing by"}
        ]
    }));
    let (_daemon, addr, cfg) = spawn_bound(|port| steer_config(&llm.uri, port, ""));

    let sent = rpc(
        &addr,
        1,
        "SendMessage",
        json!({"message": {"messageId": "m1", "parts": [{"text": "Start a warm helper"}]}}),
    );
    assert_eq!(
        sent["task"]["status"]["state"], "TASK_STATE_COMPLETED",
        "{sent}"
    );
    // Find the warm handle.
    let st = command(&addr, 2, "status", json!({}));
    let subs = artifact_json(&st)["subagents"].clone();
    let handle = subs[0]["handle"]
        .as_str()
        .expect("warm subagent")
        .to_string();

    // Steer it: inject a message over A2A.
    let injected = command(
        &addr,
        3,
        "subagent.send",
        json!({"handle": handle, "message": "focus on the staging cluster"}),
    );
    assert_eq!(artifact_json(&injected)["ok"], true, "{injected}");

    // Unknown handle → clean error.
    let body = json!({"jsonrpc": "2.0", "id": 4, "method": "SendMessage", "params": {"message": {"messageId": "m4", "parts": [{"data": {"agentd": {"op": "subagent.send", "handle": "nope", "message": "x"}}}]}}}).to_string();
    let resp: Value = serde_json::from_str(&post_raw(&addr, &body)).unwrap();
    assert_eq!(resp["error"]["code"], -32602, "{resp}");

    // plan.get on the root conversation (operator).
    let plan = command(&addr, 5, "plan.get", json!({}));
    assert!(
        artifact_json(&plan).get("plan").is_some(),
        "plan.get answers (plan may be null): {plan}"
    );
    std::fs::remove_file(&cfg).ok();
}