use crate::core::agent_context::AgentContextView;
use crate::worker::adapter::{InProcSpawner, WorkerError, WorkerInvocation, WorkerResult};
use agent_block_core::bus::dispatcher::Handler;
use agent_block_core::host::{PromptSource, ScriptSource};
use agent_block_core::{run, BlockConfig};
use agent_block_types::error::BlockError;
use async_trait::async_trait;
use serde_json::Value;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::sync::oneshot;
struct WorkerResultCaptor {
tx: Mutex<Option<oneshot::Sender<WorkerResult>>>,
sink: Option<Arc<dyn crate::worker::output::OutputSink>>,
}
impl WorkerResultCaptor {
fn extract(payload: &Value) -> (Value, bool) {
let ok = payload.get("ok").and_then(|v| v.as_bool()).unwrap_or(true);
let value = payload
.get("content")
.cloned()
.or_else(|| payload.get("response").cloned())
.unwrap_or_else(|| payload.clone());
(value, ok)
}
fn extract_stats(payload: &Value) -> Option<crate::store::trace::WorkerStats> {
let usage_raw = payload.get("usage");
let usage = usage_raw.and_then(|u| {
let input = u.get("input_tokens").and_then(|v| v.as_u64());
let output = u.get("output_tokens").and_then(|v| v.as_u64());
match (input, output) {
(Some(i), Some(o)) => Some(crate::store::trace::TokenUsage {
input_tokens: i,
output_tokens: o,
total_tokens: u
.get("total_tokens")
.and_then(|v| v.as_u64())
.unwrap_or(i + o),
}),
_ => None,
}
});
let num_turns = payload
.get("num_turns")
.and_then(|v| v.as_u64())
.map(|n| n as u32);
if usage.is_none() && num_turns.is_none() {
return None;
}
Some(crate::store::trace::WorkerStats {
worker_kind: Some("agent_block".to_string()),
model: None,
usage,
num_turns,
adapter_data: usage_raw.cloned(),
})
}
async fn stage_artifact(&self, payload: &Value) -> Result<(), BlockError> {
let name = payload
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| {
BlockError::Runtime(format!(
"bus.emit(\"{ARTIFACT_EVENT_KIND}\", ...) requires a string `name` field \
naming the part (got: {payload})"
))
})?
.to_string();
let Some(sink) = self.sink.as_ref() else {
tracing::warn!(
artifact = %name,
"agent-block staged an artifact but no OutputSink is wired for this \
invocation; the part is dropped"
);
return Ok(());
};
let content = payload.get("content").cloned().unwrap_or(Value::Null);
sink.emit(crate::worker::output::OutputEvent::Artifact {
name: name.clone(),
content: crate::worker::output::ContentRef::Inline { value: content },
})
.await
.map_err(|e| BlockError::Runtime(format!("staging artifact '{name}': {e}")))?;
Ok(())
}
}
#[async_trait]
impl Handler for WorkerResultCaptor {
async fn call(
&self,
kind: String,
_id: String,
payload: Value,
_meta: Value,
) -> Result<Value, BlockError> {
if kind == ARTIFACT_EVENT_KIND {
self.stage_artifact(&payload).await?;
return Ok(Value::Null);
}
let (value, ok) = Self::extract(&payload);
let stats = Self::extract_stats(&payload);
let wr = WorkerResult { value, ok, stats }.ensure_worker_kind("agent_block");
if let Ok(mut guard) = self.tx.lock() {
if let Some(tx) = guard.take() {
let _ = tx.send(wr);
}
}
Ok(Value::Null)
}
}
pub const TASK_METADATA_GLOBAL: &str = "_TASK_METADATA";
pub const AGENT_CTX_GLOBAL: &str = "_AGENT_CTX";
pub fn context_globals(view: Option<&AgentContextView>) -> HashMap<String, Value> {
let mut globals = HashMap::new();
let Some(view) = view else {
return globals;
};
if let Some(meta) = view.task_metadata.clone() {
globals.insert(TASK_METADATA_GLOBAL.to_string(), meta);
}
if !view.extra.is_empty() {
globals.insert(
AGENT_CTX_GLOBAL.to_string(),
Value::Object(view.extra.clone()),
);
}
globals
}
pub const ARTIFACT_EVENT_KIND: &str = "artifact";
#[derive(Clone)]
struct AgentBlockSettings {
script: ScriptSource,
spec_project_root: PathBuf,
mcp_rpc_timeout: Duration,
profile_context: Option<String>,
}
async fn run_agent_block_worker(
settings: Arc<AgentBlockSettings>,
inv: WorkerInvocation,
) -> Result<WorkerResult, WorkerError> {
let (tx, rx) = oneshot::channel();
let captor: Arc<dyn Handler> = Arc::new(WorkerResultCaptor {
tx: Mutex::new(Some(tx)),
sink: inv.sink.clone(),
});
let project_root = resolve_project_root(inv.context.as_ref(), &settings.spec_project_root);
let shutdown_token = inv.cancel_token.clone().unwrap_or_default();
let mut builder = BlockConfig::builder(settings.script.clone(), project_root)
.mcp_rpc_timeout(settings.mcp_rpc_timeout)
.prompt(PromptSource::Inline(inv.prompt))
.host_handler(captor)
.auto_serve_bus(true)
.shutdown_token(shutdown_token.clone());
if let Some(system) = settings.profile_context.clone() {
builder = builder.context(PromptSource::Inline(system));
}
let globals = context_globals(inv.context.as_ref());
if !globals.is_empty() {
builder = builder.extra_globals(globals);
}
let config = builder.build();
let run_handle = tokio::spawn(run(config));
let run_result = run_handle
.await
.map_err(|e| WorkerError::Failed(format!("agent-block task join: {e}")))?;
run_result.map_err(|e| WorkerError::Failed(format!("agent-block run failed: {e}")))?;
rx.await.map_err(|_| {
WorkerError::Failed("agent-block script finished without emitting result via bus".into())
})
}
pub fn resolve_needed_mcp_servers(
declared_tools: &[String],
spec_mcp_servers: &[Value],
) -> Vec<Value> {
use std::collections::HashSet;
let needed: HashSet<&str> = declared_tools
.iter()
.filter_map(|t| {
let rest = t.strip_prefix("mcp__")?;
let idx = rest.find("__")?;
Some(&rest[..idx])
})
.collect();
spec_mcp_servers
.iter()
.filter(|cfg| {
cfg.get("name")
.and_then(|n| n.as_str())
.map(|name| needed.contains(name))
.unwrap_or(false)
})
.cloned()
.collect()
}
fn mcp_tools_of(tools: &[String]) -> Vec<&str> {
tools
.iter()
.filter(|t| t.starts_with("mcp__"))
.map(String::as_str)
.collect()
}
pub fn build_inline_agent_invoker(mcp_servers: &[Value]) -> ScriptSource {
let mcp_lua = json_array_to_lua_literal(mcp_servers);
let source = format!(
r##"local agent = require("agent")
local mcp_servers = {mcp_lua}
local r = agent.run({{
prompt = _PROMPT,
system = _CONTEXT,
mcp_servers = mcp_servers,
}})
bus.emit("agent_result", r)
"##
);
ScriptSource::Inline {
source,
name: "mlua_swarm_engine_default_agent_invoker.lua".into(),
}
}
fn json_to_lua_literal(v: &Value) -> String {
match v {
Value::Null => "nil".to_string(),
Value::Bool(b) => b.to_string(),
Value::Number(n) => n.to_string(),
Value::String(s) => format!("{s:?}"),
Value::Array(arr) => {
let items: Vec<String> = arr.iter().map(json_to_lua_literal).collect();
format!("{{{}}}", items.join(", "))
}
Value::Object(map) => {
let items: Vec<String> = map
.iter()
.map(|(k, v)| format!("[{k:?}]={}", json_to_lua_literal(v)))
.collect();
format!("{{{}}}", items.join(", "))
}
}
}
fn json_array_to_lua_literal(arr: &[Value]) -> String {
if arr.is_empty() {
return "{}".to_string();
}
let items: Vec<String> = arr.iter().map(json_to_lua_literal).collect();
format!("{{{}}}", items.join(", "))
}
fn resolve_spec_project_root(spec: &Value) -> PathBuf {
match spec.get("project_root").and_then(|v| v.as_str()) {
Some(s) => PathBuf::from(s),
None => std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")),
}
}
fn resolve_project_root(view: Option<&AgentContextView>, spec_fallback: &Path) -> PathBuf {
view.and_then(|v| v.work_dir.as_deref().or(v.project_root.as_deref()))
.map(PathBuf::from)
.unwrap_or_else(|| spec_fallback.to_path_buf())
}
pub struct AgentBlockInProcessSpawnerFactory;
impl Default for AgentBlockInProcessSpawnerFactory {
fn default() -> Self {
Self
}
}
impl AgentBlockInProcessSpawnerFactory {
pub fn new() -> Self {
Self
}
}
impl crate::blueprint::compiler::SpawnerFactoryKind for AgentBlockInProcessSpawnerFactory {
const KIND: crate::blueprint::AgentKind = crate::blueprint::AgentKind::AgentBlock;
type Worker = AgentBlockWorker;
}
impl crate::blueprint::compiler::SpawnerFactory for AgentBlockInProcessSpawnerFactory {
fn build(
&self,
agent_def: &crate::blueprint::AgentDef,
_hint: Option<&Value>,
) -> Result<
Arc<dyn crate::worker::adapter::SpawnerAdapter>,
crate::blueprint::compiler::CompileError,
> {
let agent_name = agent_def.name.clone();
let spec = &agent_def.spec;
let effective_tools: Vec<String> = agent_def
.profile
.as_ref()
.map(|p| p.tools.clone())
.unwrap_or_default();
let spec_mcp_servers: Vec<Value> = spec
.get("mcp_servers")
.and_then(|v| v.as_array())
.cloned()
.unwrap_or_default();
let needed_mcp_servers = resolve_needed_mcp_servers(&effective_tools, &spec_mcp_servers);
let script = match spec.get("script_path").and_then(|v| v.as_str()) {
Some(s) => {
let mcp_tools = mcp_tools_of(&effective_tools);
if !mcp_tools.is_empty() {
return Err(crate::blueprint::compiler::CompileError::InvalidSpec {
name: agent_name,
msg: format!(
"agent_block ScriptBasedAgent mode (spec.script_path = {s:?}) cannot \
enforce an MCP tool grant: the script opens its own connections via \
`mcp.connect`, so the declared tools ({}) would be unenforceable. \
Either drop spec.script_path to use PromptBasedAgent mode (where the \
declared servers ARE the only ones embedded into the invoker), or \
drop the mcp__ entries and let the script own its connections.",
mcp_tools.join(", ")
),
});
}
ScriptSource::Path(PathBuf::from(s))
}
None => build_inline_agent_invoker(&needed_mcp_servers),
};
let spec_project_root = resolve_spec_project_root(spec);
let mcp_rpc_timeout = match spec.get("mcp_rpc_timeout_ms").and_then(|v| v.as_u64()) {
Some(ms) => Duration::from_millis(ms),
None => Duration::from_secs(30),
};
let profile_context = agent_def.profile.as_ref().map(|p| p.system_prompt.clone());
let settings = Arc::new(AgentBlockSettings {
script,
spec_project_root,
mcp_rpc_timeout,
profile_context,
});
let worker_fn: crate::worker::adapter::WorkerFn = Arc::new(move |inv| {
let settings = settings.clone();
Box::pin(run_agent_block_worker(settings, inv))
});
let mut sp: InProcSpawner<AgentBlockWorker> = InProcSpawner::<AgentBlockWorker>::typed();
sp.registry.insert(agent_name, worker_fn);
Ok(Arc::new(sp))
}
}
pub struct AgentBlockWorker {
pub handler: crate::worker::WorkerJoinHandler,
}
impl From<crate::worker::WorkerJoinHandler> for AgentBlockWorker {
fn from(handler: crate::worker::WorkerJoinHandler) -> Self {
Self { handler }
}
}
#[async_trait]
impl crate::worker::Worker for AgentBlockWorker {
fn id(&self) -> &crate::types::WorkerId {
&self.handler.worker_id
}
fn cancel_token(&self) -> tokio_util::sync::CancellationToken {
self.handler.cancel.clone()
}
async fn join(self: Box<Self>) -> Result<(), WorkerError> {
self.handler.await_completion().await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::agent_context::{TASK_METADATA_KEY, TASK_PROJECT_ROOT_KEY, TASK_WORK_DIR_KEY};
#[test]
fn resolve_needed_mcp_servers_filters_by_tool_prefix() {
let tools = vec![
"mcp__semantic-scholar__search_papers".to_string(),
"mcp__semantic-scholar__get_paper".to_string(),
"Read".to_string(),
"mcp__outline__list_docs".to_string(),
"WebSearch".to_string(),
];
let spec_servers = vec![
serde_json::json!({"name": "semantic-scholar", "command": "ss-mcp", "args": []}),
serde_json::json!({"name": "outline", "command": "outline-mcp", "args": []}),
serde_json::json!({"name": "unused", "command": "nope", "args": []}),
];
let needed = resolve_needed_mcp_servers(&tools, &spec_servers);
assert_eq!(needed.len(), 2, "got: {needed:?}");
let names: Vec<&str> = needed
.iter()
.filter_map(|c| c.get("name").and_then(|n| n.as_str()))
.collect();
assert!(names.contains(&"semantic-scholar"));
assert!(names.contains(&"outline"));
assert!(!names.contains(&"unused"), "unused server is filtered out");
}
#[test]
fn resolve_needed_mcp_servers_returns_empty_when_no_mcp_tools() {
let tools = vec!["Read".to_string(), "WebSearch".to_string()];
let spec_servers =
vec![serde_json::json!({"name": "outline", "command": "outline-mcp", "args": []})];
let needed = resolve_needed_mcp_servers(&tools, &spec_servers);
assert!(
needed.is_empty(),
"no mcp__-prefixed tools → empty result, got: {needed:?}"
);
}
#[test]
fn build_inline_agent_invoker_embeds_mcp_servers_as_lua_literal() {
let servers =
vec![serde_json::json!({"name": "outline", "command": "outline-mcp", "args": []})];
let script = build_inline_agent_invoker(&servers);
match script {
ScriptSource::Inline { source, name } => {
assert!(name.ends_with(".lua"));
assert!(source.contains("require(\"agent\")"));
assert!(source.contains("mcp_servers = mcp_servers"));
assert!(source.contains("bus.emit(\"agent_result\""));
assert!(source.contains("[\"name\"]=\"outline\""));
assert!(source.contains("[\"command\"]=\"outline-mcp\""));
assert!(source.contains("[\"args\"]={}"), "args empty array literal");
}
other => panic!("expected Inline, got: {other:?}"),
}
}
#[test]
fn build_inline_agent_invoker_with_empty_servers_still_valid() {
let script = build_inline_agent_invoker(&[]);
match script {
ScriptSource::Inline { source, .. } => {
assert!(source.contains("local mcp_servers = {}"));
}
other => panic!("expected Inline, got: {other:?}"),
}
}
#[test]
fn json_to_lua_literal_handles_primitives_and_nested() {
assert_eq!(json_to_lua_literal(&serde_json::json!(null)), "nil");
assert_eq!(json_to_lua_literal(&serde_json::json!(true)), "true");
assert_eq!(json_to_lua_literal(&serde_json::json!(42)), "42");
assert_eq!(json_to_lua_literal(&serde_json::json!("hi")), "\"hi\"");
assert_eq!(
json_to_lua_literal(&serde_json::json!(["a", "b"])),
"{\"a\", \"b\"}"
);
assert_eq!(
json_to_lua_literal(&serde_json::json!({"k": 1})),
"{[\"k\"]=1}"
);
}
#[test]
fn extract_prefers_content_then_response_then_whole() {
let p = serde_json::json!({
"content": "Water boils at 100°C",
"messages": [{"role": "assistant"}],
"usage": {"input_tokens": 67, "output_tokens": 29},
"ok": true,
});
let (value, ok) = WorkerResultCaptor::extract(&p);
assert_eq!(value, serde_json::json!("Water boils at 100°C"));
assert!(ok);
let p = serde_json::json!({ "ok": false, "response": {"patch": "..."} });
let (value, ok) = WorkerResultCaptor::extract(&p);
assert_eq!(value, serde_json::json!({"patch": "..."}));
assert!(!ok);
let p = serde_json::json!({ "custom_field": 42 });
let (value, ok) = WorkerResultCaptor::extract(&p);
assert_eq!(value, serde_json::json!({"custom_field": 42}));
assert!(ok); }
#[tokio::test]
async fn captor_emits_worker_result_from_payload() {
let (tx, rx) = oneshot::channel();
let captor = WorkerResultCaptor {
tx: Mutex::new(Some(tx)),
sink: None,
};
let payload = serde_json::json!({ "ok": true, "response": "hello" });
let ack = captor
.call("worker_result".into(), "evt-1".into(), payload, Value::Null)
.await
.expect("handler ack");
assert_eq!(ack, Value::Null);
let wr = rx.await.expect("recv");
assert!(wr.ok);
assert_eq!(wr.value, serde_json::json!("hello"));
}
#[tokio::test]
async fn factory_builds_prompt_based_agent_when_script_path_absent() {
use crate::blueprint::compiler::SpawnerFactory;
use crate::blueprint::{AgentDef, AgentKind, AgentProfile};
let factory = AgentBlockInProcessSpawnerFactory::new();
let ad = AgentDef {
name: "writer".into(),
kind: AgentKind::AgentBlock,
spec: serde_json::json!({}),
profile: Some(AgentProfile {
system_prompt: "You are writer.".into(),
..Default::default()
}),
meta: None,
runner: None,
runner_ref: None,
verdict: None,
};
let _spawner = factory.build(&ad, None).expect("factory build");
}
fn agent_block_def(name: &str, spec: Value, tools: &[&str]) -> crate::blueprint::AgentDef {
use crate::blueprint::{AgentDef, AgentKind, AgentProfile};
AgentDef {
name: name.into(),
kind: AgentKind::AgentBlock,
spec,
profile: Some(AgentProfile {
system_prompt: "You are an auditor.".into(),
tools: tools.iter().map(|t| t.to_string()).collect(),
..Default::default()
}),
meta: None,
runner: None,
runner_ref: None,
verdict: None,
}
}
#[test]
fn mcp_tools_of_keeps_only_server_selecting_names() {
let tools = vec![
"Read".to_string(),
"mcp__outline__list_docs".to_string(),
"WebSearch".to_string(),
];
assert_eq!(mcp_tools_of(&tools), vec!["mcp__outline__list_docs"]);
assert!(mcp_tools_of(&["Read".to_string()]).is_empty());
}
#[tokio::test]
async fn effective_grant_narrows_the_embedded_mcp_servers() {
use crate::blueprint::compiler::SpawnerFactory;
let ad = agent_block_def(
"auditor",
serde_json::json!({
"mcp_servers": [
{"name": "outline", "command": "outline-mcp", "args": []},
{"name": "semantic-scholar", "command": "ss-mcp", "args": []},
]
}),
&["mcp__outline__list_docs"],
);
let effective = ad.profile.as_ref().unwrap().tools.clone();
let servers = resolve_needed_mcp_servers(
&effective,
ad.spec["mcp_servers"].as_array().expect("array"),
);
let names: Vec<&str> = servers
.iter()
.filter_map(|c| c.get("name").and_then(|n| n.as_str()))
.collect();
assert_eq!(
names,
vec!["outline"],
"semantic-scholar is declared in spec but not selected by the grant"
);
AgentBlockInProcessSpawnerFactory::new()
.build(&ad, None)
.expect("PromptBasedAgent mode accepts an MCP grant");
}
#[tokio::test]
async fn script_mode_rejects_a_declared_mcp_grant() {
use crate::blueprint::compiler::{CompileError, SpawnerFactory};
let ad = agent_block_def(
"gate-danger",
serde_json::json!({ "script_path": "gate.lua" }),
&["mcp__outline__list_docs"],
);
let err = AgentBlockInProcessSpawnerFactory::new()
.build(&ad, None)
.err()
.expect("must reject");
match err {
CompileError::InvalidSpec { name, msg } => {
assert_eq!(name, "gate-danger");
assert!(msg.contains("mcp.connect"), "explains why: {msg}");
assert!(
msg.contains("PromptBasedAgent"),
"names the actionable alternative: {msg}"
);
}
other => panic!("expected InvalidSpec, got: {other:?}"),
}
}
#[tokio::test]
async fn script_mode_accepts_empty_and_inert_grants() {
use crate::blueprint::compiler::SpawnerFactory;
let spec = serde_json::json!({ "script_path": "gate.lua" });
for tools in [&[][..], &["Read", "WebSearch"][..]] {
let ad = agent_block_def("gate-danger", spec.clone(), tools);
AgentBlockInProcessSpawnerFactory::new()
.build(&ad, None)
.unwrap_or_else(|e| panic!("script mode must accept tools {tools:?}: {e}"));
}
}
#[tokio::test]
async fn factory_builds_script_based_agent_when_script_path_present() {
use crate::blueprint::compiler::SpawnerFactory;
use crate::blueprint::{AgentDef, AgentKind, AgentProfile};
let factory = AgentBlockInProcessSpawnerFactory::new();
let ad = AgentDef {
name: "patch-spawner".into(),
kind: AgentKind::AgentBlock,
spec: serde_json::json!({
"script_path": "assets/operator_scripts/blueprint_patch_spawner.lua",
"project_root": ".",
}),
profile: Some(AgentProfile {
system_prompt: "Patch generator.".into(),
..Default::default()
}),
meta: None,
runner: None,
runner_ref: None,
verdict: None,
};
let _spawner = factory.build(&ad, None).expect("factory build");
}
#[test]
fn resolve_spec_project_root_uses_spec_value_when_present() {
let resolved =
resolve_spec_project_root(&serde_json::json!({ "project_root": "/spec-root" }));
assert_eq!(resolved, PathBuf::from("/spec-root"));
}
#[test]
fn resolve_spec_project_root_falls_back_to_env_current_dir_when_spec_absent() {
let resolved = resolve_spec_project_root(&serde_json::json!({}));
assert_eq!(
resolved,
std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))
);
}
fn view_with(pairs: &[(&str, Value)]) -> AgentContextView {
let mut ctx = crate::core::ctx::Ctx::new(
crate::types::StepId::parse("ST-project-root").unwrap(),
1,
"writer",
);
for (k, v) in pairs {
ctx.meta.runtime.insert((*k).to_string(), v.clone());
}
AgentContextView::from_ctx(&ctx)
}
#[test]
fn project_root_falls_back_to_spec_when_the_view_carries_neither() {
let view = view_with(&[]);
let resolved = resolve_project_root(Some(&view), Path::new("/spec-root"));
assert_eq!(resolved, PathBuf::from("/spec-root"));
}
#[test]
fn project_root_falls_back_to_spec_without_a_view() {
assert_eq!(
resolve_project_root(None, Path::new("/spec-root")),
PathBuf::from("/spec-root")
);
}
#[test]
fn project_root_prefers_the_view_over_spec() {
let view = view_with(&[(TASK_PROJECT_ROOT_KEY, serde_json::json!("/ctx-root"))]);
let resolved = resolve_project_root(Some(&view), Path::new("/spec-root"));
assert_eq!(resolved, PathBuf::from("/ctx-root"));
}
#[test]
fn project_root_prefers_work_dir_over_project_root() {
let view = view_with(&[
(TASK_PROJECT_ROOT_KEY, serde_json::json!("/ctx-root")),
(TASK_WORK_DIR_KEY, serde_json::json!("/ctx-work")),
]);
let resolved = resolve_project_root(Some(&view), Path::new("/spec-root"));
assert_eq!(resolved, PathBuf::from("/ctx-work"));
}
#[tokio::test]
async fn artifact_kind_stages_a_named_part_without_completing() {
use crate::worker::output::{ContentRef, OutputEvent, OutputSink};
#[derive(Default)]
struct RecordingSink(Mutex<Vec<OutputEvent>>);
#[async_trait]
impl OutputSink for RecordingSink {
async fn emit(&self, event: OutputEvent) -> Result<(), crate::EngineError> {
self.0.lock().unwrap().push(event);
Ok(())
}
}
let sink = Arc::new(RecordingSink::default());
let (tx, mut rx) = oneshot::channel();
let captor = WorkerResultCaptor {
tx: Mutex::new(Some(tx)),
sink: Some(sink.clone()),
};
captor
.call(
ARTIFACT_EVENT_KIND.into(),
"evt-1".into(),
serde_json::json!({ "name": "verdict", "content": "PASS" }),
Value::Null,
)
.await
.expect("staging must succeed");
let staged = sink.0.lock().unwrap().clone();
assert_eq!(staged.len(), 1, "exactly one artifact staged");
match &staged[0] {
OutputEvent::Artifact { name, content } => {
assert_eq!(name, "verdict");
match content {
ContentRef::Inline { value } => {
assert_eq!(value, &serde_json::json!("PASS"))
}
other => panic!("expected Inline content, got: {other:?}"),
}
}
other => panic!("expected Artifact, got: {other:?}"),
}
assert!(
rx.try_recv().is_err(),
"an artifact emit must NOT complete the invocation"
);
captor
.call(
"worker_result".into(),
"evt-2".into(),
serde_json::json!({ "ok": true, "response": "done" }),
Value::Null,
)
.await
.expect("terminal emit");
assert_eq!(
rx.await.expect("recv").value,
serde_json::json!("done"),
"the non-reserved kind still completes the invocation"
);
}
#[tokio::test]
async fn artifact_without_a_name_is_reported_to_the_script() {
let (tx, _rx) = oneshot::channel();
let captor = WorkerResultCaptor {
tx: Mutex::new(Some(tx)),
sink: None,
};
let err = captor
.call(
ARTIFACT_EVENT_KIND.into(),
"evt-1".into(),
serde_json::json!({ "content": "PASS" }),
Value::Null,
)
.await
.expect_err("a nameless artifact must fail loud");
assert!(
format!("{err}").contains("name"),
"names the missing field: {err}"
);
}
#[test]
fn context_globals_renders_task_metadata_and_agent_ctx() {
let mut view = view_with(&[(TASK_METADATA_KEY, serde_json::json!({"issue": 86}))]);
view.extra.insert(
"org_conventions".to_string(),
serde_json::json!("two-space indent"),
);
let globals = context_globals(Some(&view));
assert_eq!(
globals.get(TASK_METADATA_GLOBAL),
Some(&serde_json::json!({"issue": 86}))
);
assert_eq!(
globals.get(AGENT_CTX_GLOBAL),
Some(&serde_json::json!({"org_conventions": "two-space indent"})),
"Blueprint-declared agent ctx must reach the in-process lane too"
);
}
#[test]
fn context_globals_omits_absent_fields() {
assert!(context_globals(None).is_empty(), "no view → no globals");
assert!(
context_globals(Some(&view_with(&[]))).is_empty(),
"empty view → no globals"
);
let view = view_with(&[(TASK_METADATA_KEY, serde_json::json!({"issue": 86}))]);
let globals = context_globals(Some(&view));
assert!(globals.contains_key(TASK_METADATA_GLOBAL));
assert!(
!globals.contains_key(AGENT_CTX_GLOBAL),
"an empty `extra` must not render an empty _AGENT_CTX table"
);
}
#[test]
fn context_globals_use_names_the_sdk_does_not_reserve() {
for name in [TASK_METADATA_GLOBAL, AGENT_CTX_GLOBAL] {
for reserved in ["_PROMPT", "_CONTEXT", "_SCRIPT_NAME"] {
assert_ne!(name, reserved);
}
}
assert_ne!(TASK_METADATA_GLOBAL, AGENT_CTX_GLOBAL);
}
#[test]
fn task_metadata_global_does_not_shadow_an_sdk_reserved_name() {
for reserved in ["_PROMPT", "_CONTEXT", "_SCRIPT_NAME"] {
assert_ne!(TASK_METADATA_GLOBAL, reserved);
}
}
#[tokio::test]
async fn script_mode_never_rewrites_the_caller_chunk() {
use crate::blueprint::compiler::SpawnerFactory;
let ad = agent_block_def(
"gate-danger",
serde_json::json!({ "script_path": "/nonexistent/gate.lua" }),
&[],
);
AgentBlockInProcessSpawnerFactory::new()
.build(&ad, None)
.expect("script mode must not read the script at build time");
}
}