From 540f79e7e98668ad6951ba8aa4b3f4fd3a59ffb7 Mon Sep 17 00:00:00 2001 From: Nathan Herald Date: Mon, 10 Aug 2026 18:06:13 +0200 Subject: [PATCH] Add declared native delivery selector --- crates/agent-spec/src/kdl_format.rs | 17 ++++++ crates/agent-spec/src/lib.rs | 4 +- crates/agent-spec/src/spec.rs | 58 ++++++++++++++++++- crates/agent-spec/tests/discovery.rs | 87 +++++++++++++++++++++++++++- src/catalog_transaction.rs | 6 ++ src/eval_run.rs | 1 + src/lib.rs | 4 +- src/main.rs | 10 ++++ src/run.rs | 3 + tests/catalog_diff.rs | 25 ++++++++ tests/doctor.rs | 54 +++++++++++++++++ tests/reconcile.rs | 1 + tests/run.rs | 1 + 13 files changed, 262 insertions(+), 9 deletions(-) diff --git a/crates/agent-spec/src/kdl_format.rs b/crates/agent-spec/src/kdl_format.rs index 6540668e..d4541af5 100644 --- a/crates/agent-spec/src/kdl_format.rs +++ b/crates/agent-spec/src/kdl_format.rs @@ -124,6 +124,23 @@ fn agent_node_to_raw(node: &DeclaredNode) -> anyhow::Result { "command" => raw.command = arg_string(child), "argv" => raw.argv = Some(argv(child)?), "ding" => raw.ding = true, + "deliver" => { + anyhow::ensure!( + raw.deliver.is_none(), + "agent declares `deliver` more than once" + ); + anyhow::ensure!( + child.type_name.is_none() + && child.children.is_empty() + && child.entries.len() == 1 + && child.entries[0].name.is_none(), + "agent `deliver` must contain exactly one positional string" + ); + raw.deliver = Some(Some( + arg_string(child) + .ok_or_else(|| anyhow::anyhow!("agent `deliver` value must be a string"))?, + )); + } "env" => {} "pty" => { if let Some(name) = arg_string(child) { diff --git a/crates/agent-spec/src/lib.rs b/crates/agent-spec/src/lib.rs index e9824715..cfe26b4e 100644 --- a/crates/agent-spec/src/lib.rs +++ b/crates/agent-spec/src/lib.rs @@ -43,6 +43,6 @@ pub use discovery::{ path_defaults, }; pub use spec::{ - AgentDesiredState, AgentSpec, JobType, Resource, Restart, RestartMode, Task, TaskKind, - TaskLifecycle, parse_duration, validate_desired_state_reason, + AgentDesiredState, AgentSpec, DeliveryTransport, JobType, Resource, Restart, RestartMode, Task, + TaskKind, TaskLifecycle, parse_duration, validate_desired_state_reason, }; diff --git a/crates/agent-spec/src/spec.rs b/crates/agent-spec/src/spec.rs index a0cf9a29..d4b0afc8 100644 --- a/crates/agent-spec/src/spec.rs +++ b/crates/agent-spec/src/spec.rs @@ -5,9 +5,10 @@ //! stage's script; must NOT allocate a terminal, R09). st2 reads only the runner-normative subset: //! `identity`, presentation (`name`, `description`), `host`, `role` (metadata only), `type`, //! `workspace`, whole-agent desired state (plus legacy `retired`), `keep`, `supervisor`, -//! `restart{}`, task lifecycle, Resource bindings (declaration metadata), and the tasks. Everything render-only -//! (`harness`, `model`, `persona`, `permissions`, `transport`, `strategy`, `meta{}`) is baked into -//! the tasks/commands by the render layer and ignored here. +//! `restart{}`, `deliver`, task lifecycle, Resource bindings (declaration metadata), and the tasks. +//! Everything render-only (`harness`, `model`, `persona`, `permissions`, legacy `transport` +//! metadata, `strategy`, `meta{}`) is baked into the tasks/commands by the render layer and ignored +//! here. //! //! Three on-disk formats lower to this model: KDL (canonical, parsed by hand in `kdl_format`), and //! TOML/JSON (serde). Every spec is a `service` — `type = batch` is retired; evals run through the @@ -36,6 +37,32 @@ pub enum AgentDesiredState { Retired { reason: Option }, } +/// One provider-native message delivery transport declared by an agent. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DeliveryTransport { + Mcp, + AppServer, +} + +impl DeliveryTransport { + pub fn as_str(self) -> &'static str { + match self { + Self::Mcp => "mcp", + Self::AppServer => "app-server", + } + } + + fn parse(value: &str) -> anyhow::Result { + match value { + "mcp" => Ok(Self::Mcp), + "app-server" => Ok(Self::AppServer), + _ => anyhow::bail!( + "unsupported `deliver` value '{value}' (expected `mcp` or `app-server`)" + ), + } + } +} + impl AgentDesiredState { pub fn as_str(&self) -> &'static str { match self { @@ -91,6 +118,8 @@ pub struct AgentSpec { pub keep: bool, /// Crash/restart policy (§4). `None` → the runner's default policy. pub restart: Option, + /// Provider-native delivery selected by `deliver`; `None` means legacy `ding` or no delivery. + pub delivery: Option, /// Named typed references used by the agent. st2 preserves these for readers but does not /// resolve them or assign launch, readiness, access, or lifecycle semantics. pub resources: Vec, @@ -327,6 +356,15 @@ impl AgentSpec { .any(|task| !task.derived && (task.command.is_some() || task.argv.is_some())) } + /// True when the declaration selected legacy screen delivery or one native transport. + pub fn has_delivery_transport(&self) -> bool { + self.delivery.is_some() + || self + .tasks + .iter() + .any(|task| task.derived && task.kind == TaskKind::Exec && task.name == "ding") + } + /// The restart policy in effect (declared, else the runner default). pub fn restart_policy(&self) -> Restart { self.restart.clone().unwrap_or_default() @@ -399,6 +437,9 @@ pub(crate) struct RawSpec { /// Compact catalog form: include the built-in `st2 ding` sidecar. #[serde(default)] pub ding: bool, + /// Compact catalog form: select one provider-native delivery transport. + #[serde(default, deserialize_with = "deserialize_explicit_optional")] + pub deliver: Option>, /// Compact catalog form: reconciliation policy for the generated agent PTY. pub lifecycle: Option, /// `pty "" {}` / `[pty.]` — interactive tasks. @@ -750,6 +791,7 @@ impl RawSpec { || self.command.is_some() || self.argv.is_some() || self.ding + || self.deliver.is_some() || !self.resource.0.is_empty() || !self.pty.is_empty() || !self.exec.is_empty() @@ -778,6 +820,15 @@ impl RawSpec { desired_state_value.as_deref(), desired_state_reason, )?; + let deliver = reject_explicit_null("deliver", self.deliver)?; + let delivery = deliver + .as_deref() + .map(DeliveryTransport::parse) + .transpose()?; + anyhow::ensure!( + !(self.ding && delivery.is_some()), + "agent '{identity}' declares both `ding` and `deliver`; choose one transport" + ); validate_launch( &identity, self.command.as_ref(), @@ -850,6 +901,7 @@ impl RawSpec { desired_state, keep: self.keep, restart: self.restart.map(RawRestart::lower), + delivery, resources, tasks, path, diff --git a/crates/agent-spec/tests/discovery.rs b/crates/agent-spec/tests/discovery.rs index af08b9a3..baa5e654 100644 --- a/crates/agent-spec/tests/discovery.rs +++ b/crates/agent-spec/tests/discovery.rs @@ -8,7 +8,7 @@ use std::fs; use std::path::Path; use std::time::Duration; -use agent_spec::spec::{TaskKind, TaskLifecycle}; +use agent_spec::spec::{DeliveryTransport, TaskKind, TaskLifecycle}; use agent_spec::{ AgentDesiredState, AgentSpec, JobType, Resource, Task, discover, discover_strict, }; @@ -139,10 +139,11 @@ fn lifecycle_fields_make_path_placed_files_agent_candidates() { } #[test] -fn explicit_json_null_lifecycle_fields_are_rejected_instead_of_granting_running_intent() { +fn explicit_json_null_fields_are_rejected_instead_of_granting_default_behavior() { for (name, lifecycle) in [ ("null-retired", r#""retired":null"#), ("null-state", r#""desired_state":null"#), + ("null-deliver", r#""deliver":null"#), ( "null-reason", r#""desired_state":"suspended","desired_state_reason":null"#, @@ -328,6 +329,88 @@ agent "cos" { ding.env.get("ST_AGENT").map(String::as_str), Some("Silber.cos") ); + assert!(spec.delivery.is_none()); + assert!(spec.has_delivery_transport()); +} + +#[test] +fn deliver_is_typed_without_lowering_to_the_legacy_ding_task() { + let tmp = tempfile::tempdir().unwrap(); + write( + tmp.path(), + "agents/h/claude/agent.kdl", + r#"agent "claude" { host "h"; command "claude"; deliver "mcp" }"#, + ); + write( + tmp.path(), + "agents/h/codex/agent.kdl", + r#"agent "codex" { host "h"; command "codex"; deliver "app-server" }"#, + ); + + let found = discover(tmp.path()); + assert!(found.errors.is_empty(), "{:?}", found.errors); + let claude = find(&found.specs, "claude"); + let codex = find(&found.specs, "codex"); + assert_eq!(claude.delivery, Some(DeliveryTransport::Mcp)); + assert_eq!(codex.delivery, Some(DeliveryTransport::AppServer)); + assert_eq!(claude.delivery.unwrap().as_str(), "mcp"); + assert_eq!(codex.delivery.unwrap().as_str(), "app-server"); + for spec in [claude, codex] { + assert!(spec.has_delivery_transport()); + assert_eq!(spec.tasks.len(), 1); + assert!(spec.tasks.iter().all(|task| !task.derived)); + } +} + +#[test] +fn deliver_rejects_unknown_duplicate_mixed_and_malformed_declarations() { + for (name, declaration, expected) in [ + ( + "unknown", + r#"agent "worker" { command "true"; deliver "socket" }"#, + "unsupported `deliver` value 'socket'", + ), + ( + "duplicate", + r#"agent "worker" { command "true"; deliver "mcp"; deliver "app-server" }"#, + "declares `deliver` more than once", + ), + ( + "mixed", + r#"agent "worker" { command "true"; ding; deliver "mcp" }"#, + "declares both `ding` and `deliver`", + ), + ( + "missing", + r#"agent "worker" { command "true"; deliver }"#, + "must contain exactly one positional string", + ), + ( + "non-string", + r#"agent "worker" { command "true"; deliver #true }"#, + "value must be a string", + ), + ( + "property", + r#"agent "worker" { command "true"; deliver "mcp" mode="extra" }"#, + "must contain exactly one positional string", + ), + ] { + let tmp = tempfile::tempdir().unwrap(); + write( + tmp.path(), + &format!("agents/h/{name}/agent.kdl"), + declaration, + ); + let found = discover(tmp.path()); + assert!(found.specs.is_empty(), "{name}: {:?}", found.specs); + assert_eq!(found.errors.len(), 1, "{name}: {:?}", found.errors); + assert!( + found.errors[0].message.contains(expected), + "{name}: expected {expected:?}, got {:?}", + found.errors[0] + ); + } } #[test] diff --git a/src/catalog_transaction.rs b/src/catalog_transaction.rs index cadc52b6..ff153459 100644 --- a/src/catalog_transaction.rs +++ b/src/catalog_transaction.rs @@ -629,6 +629,12 @@ fn normalize_agent(spec: &agent_spec::AgentSpec) -> Result Vec desired_state: AgentDesiredState::Running, keep: false, restart: None, + delivery: None, resources: Vec::new(), tasks, path: path.clone(), diff --git a/src/lib.rs b/src/lib.rs index 386cf746..0dd18146 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -41,8 +41,8 @@ pub use agent_spec::{discovery, spec}; pub use agent_spec::discovery::{Discovered, SpecError, discover, discover_strict}; pub use agent_spec::spec::{ - AgentDesiredState, AgentSpec, JobType, Resource, Restart, RestartMode, Task, TaskKind, - TaskLifecycle, parse_duration, + AgentDesiredState, AgentSpec, DeliveryTransport, JobType, Resource, Restart, RestartMode, Task, + TaskKind, TaskLifecycle, parse_duration, }; pub use catalog_lock::CatalogLock; pub use exec_backend::ExecBackend; diff --git a/src/main.rs b/src/main.rs index c68324e2..32d1c77f 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1433,6 +1433,12 @@ fn doctor_cmd(root: &Path, host: Option, require_supervisor: bool) -> Re ); continue; } + if !spec.has_delivery_transport() { + report_advisory( + &format!("{bus_id} delivery transport missing"), + "declare `ding` or `deliver`; agent receives no DING", + ); + } for task in &spec.tasks { let id = task .id @@ -1546,6 +1552,10 @@ fn report_check(problems: &mut usize, ok: bool, label: &str, detail: &str) { } } +fn report_advisory(label: &str, detail: &str) { + println!(" ⚠ {label} — {detail}"); +} + fn presentation_cmd( field: st2::agent_author::PresentationField, args: PresentationArgs, diff --git a/src/run.rs b/src/run.rs index 33d3d90c..3ec307c3 100644 --- a/src/run.rs +++ b/src/run.rs @@ -2164,6 +2164,7 @@ mod tests { desired_state: crate::AgentDesiredState::Running, keep: false, restart: None, + delivery: None, resources: vec![], tasks: vec![Task { kind: TaskKind::Pty, @@ -2213,6 +2214,7 @@ mod tests { desired_state: crate::AgentDesiredState::Running, keep: false, restart: None, + delivery: None, resources: vec![], tasks: vec![Task { kind: TaskKind::Pty, @@ -2469,6 +2471,7 @@ mod tests { desired_state: crate::AgentDesiredState::Running, keep: false, restart: None, + delivery: None, resources: vec![], tasks: vec![], path: std::path::PathBuf::from("/x"), diff --git a/tests/catalog_diff.rs b/tests/catalog_diff.rs index d9de8858..5a3b4da8 100644 --- a/tests/catalog_diff.rs +++ b/tests/catalog_diff.rs @@ -158,6 +158,31 @@ fn desired_state_and_reason_have_distinct_secret_safe_semantic_addresses() { assert!(!rendered.contains(reason)); } +#[test] +fn declared_delivery_has_one_exact_semantic_address() { + let (_temp, catalog, prepared, root) = fixture(); + fs::write( + prepared.join("agents/host/worker/agent.kdl"), + r#"agent "worker" { + host "host" + deliver "app-server" + argv "tool" "arg" +} +"#, + ) + .unwrap(); + + let receipt = parsed(&diff(&catalog, &prepared, &root)); + let fields = agent_fields(&receipt); + assert_eq!( + fields + .iter() + .filter(|field| field.as_str() == "/agents/host/worker/delivery") + .count(), + 1 + ); +} + #[test] fn effective_task_id_and_cwd_defaults_normalize_to_explicit_values() { let (_temp, catalog, prepared, _root) = fixture(); diff --git a/tests/doctor.rs b/tests/doctor.rs index 6ef50b28..edb46f83 100644 --- a/tests/doctor.rs +++ b/tests/doctor.rs @@ -361,3 +361,57 @@ fn suspended_declaration_distinguishes_live_dead_keep_and_dead_nonkeep() { assert!(!stdout.contains("h.idle presence"), "{stdout}"); } } + +#[test] +fn missing_delivery_is_advisory_while_an_invalid_delivery_is_a_catalog_problem() { + let tmp = tempfile::tempdir().unwrap(); + let catalog = tmp.path().join("catalog"); + let declaration = catalog.join("agents/h/worker/agent.kdl"); + let bin = tmp.path().join("bin"); + fs::create_dir_all(declaration.parent().unwrap()).unwrap(); + fs::create_dir_all(&bin).unwrap(); + fs::write( + &declaration, + r#"agent "worker" { host "h"; command "true" }"#, + ) + .unwrap(); + fs::write(declaration.parent().unwrap().join("status"), "available\n").unwrap(); + executable( + &bin.join("pty"), + "#!/bin/sh\nif [ \"$1\" = list ]; then printf '[{\"name\":\"h.worker\",\"status\":\"running\"}]\\n'; fi\n", + ); + + let missing = doctor(&catalog, &bin, &tmp.path().join("state")); + let stdout = String::from_utf8_lossy(&missing.stdout); + assert!( + missing.status.success(), + "stdout:\n{stdout}\nstderr:\n{}", + String::from_utf8_lossy(&missing.stderr) + ); + assert!( + stdout.contains( + "⚠ h.worker delivery transport missing — declare `ding` or `deliver`; agent receives no DING" + ), + "{stdout}" + ); + + fs::write( + &declaration, + r#"agent "worker" { host "h"; command "true"; deliver "mcp" }"#, + ) + .unwrap(); + let declared = doctor(&catalog, &bin, &tmp.path().join("state")); + let stdout = String::from_utf8_lossy(&declared.stdout); + assert!(declared.status.success(), "{stdout}"); + assert!(!stdout.contains("delivery transport missing"), "{stdout}"); + + fs::write( + &declaration, + r#"agent "worker" { host "h"; command "true"; deliver "mpc" }"#, + ) + .unwrap(); + let invalid = doctor(&catalog, &bin, &tmp.path().join("state")); + let stdout = String::from_utf8_lossy(&invalid.stdout); + assert!(!invalid.status.success(), "{stdout}"); + assert!(stdout.contains("unsupported `deliver` value 'mpc'"), "{stdout}"); +} diff --git a/tests/reconcile.rs b/tests/reconcile.rs index c0850f4c..a9f670ac 100644 --- a/tests/reconcile.rs +++ b/tests/reconcile.rs @@ -398,6 +398,7 @@ fn spec( }, keep: false, restart: None, + delivery: None, resources: Vec::new(), tasks, path: PathBuf::from(format!( diff --git a/tests/run.rs b/tests/run.rs index 67630946..c35f97f1 100644 --- a/tests/run.rs +++ b/tests/run.rs @@ -236,6 +236,7 @@ fn task_spec(identity: &str, host: Option<&str>, id: &str) -> AgentSpec { desired_state: AgentDesiredState::Running, keep: false, restart: None, + delivery: None, resources: vec![], tasks: vec![Task { kind: TaskKind::Exec,