use mlua_swarm::core::config::EngineCfg;
use mlua_swarm::core::engine::Engine;
use mlua_swarm::core::state::{
is_skip_marker, unwrap_skip_marker, CapTokenRecord, TaskSpec, TaskState,
};
use mlua_swarm::{CapToken, Role, StepId};
async fn seed_task_with_handle(engine: &Engine, task_id: &StepId, agent: &str) -> String {
let handle = format!("wh-{}", mlua_swarm::types::secure_hex(4));
let task_id = task_id.clone();
let agent = agent.to_string();
let handle_clone = handle.clone();
engine
.with_state("test.seed_task_with_handle", move |s| {
let task = TaskState::new(
task_id.clone(),
TaskSpec {
agent: agent.clone(),
initial_directive: serde_json::json!("x"),
step_ctx: None,
check_policy: None,
},
);
s.tasks.insert(task_id.clone(), task);
let token = CapToken {
agent_id: agent,
role: Role::Worker,
scopes: vec!["*".to_string()],
issued_at: 0,
expire_at: u64::MAX,
max_uses: None,
nonce: format!("test-nonce-{task_id}"),
sig_hex: String::new(),
};
let fp = token.fingerprint();
s.tokens.insert(
fp.clone(),
CapTokenRecord {
token,
uses_left: None,
revoked: false,
task_id: Some(task_id),
},
);
s.worker_handles.insert(handle_clone, fp);
})
.await
.expect("seed_task_with_handle");
handle
}
async fn spawn_server(engine: Engine) -> String {
let router = mlua_swarm_server::build_router(engine);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind ephemeral port");
let addr = listener.local_addr().expect("local addr");
tokio::spawn(async move {
let _ = axum::serve(listener, router).await;
});
format!("http://{addr}")
}
#[tokio::test]
async fn post_worker_submit_verdict_skip_produces_skip_outcome_and_continues_flow() {
let engine = Engine::new(EngineCfg::default());
let task_id = StepId::new();
let handle = seed_task_with_handle(&engine, &task_id, "analyst").await;
let base_url = spawn_server(engine.clone()).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base_url}/v1/worker/submit?verdict=skip"))
.header("Authorization", format!("Bearer {handle}"))
.body("SKIP")
.send()
.await
.expect("request");
assert_eq!(
resp.status(),
reqwest::StatusCode::NO_CONTENT,
"verdict=skip is a successful completion (wire ack)"
);
let tail = engine.output_tail(&task_id, 0).await;
let (_, final_ok) = tail
.iter()
.rev()
.find_map(|ev| match ev {
mlua_swarm::OutputEvent::Final { content, ok } => Some((content.clone(), *ok)),
_ => None,
})
.expect("Final present after Skip submit");
assert!(
final_ok,
"Skip records Final.ok = true (flow continues, no error propagation)"
);
}
#[tokio::test]
async fn post_worker_submit_verdict_skip_downstream_binding_unresolved() {
let engine = Engine::new(EngineCfg::default());
let task_id = StepId::new();
let handle = seed_task_with_handle(&engine, &task_id, "analyst").await;
let base_url = spawn_server(engine.clone()).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base_url}/v1/worker/submit?verdict=skip"))
.header("Authorization", format!("Bearer {handle}"))
.body("SKIP")
.send()
.await
.expect("request");
assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT);
let last_result = engine
.with_state("test.read_last_result", {
let task_id = task_id.clone();
move |s| s.tasks.get(&task_id).and_then(|t| t.last_result.clone())
})
.await
.expect("read last_result");
let last_result = last_result.expect("last_result present after submit");
assert!(
is_skip_marker(&last_result),
"verdict=skip wraps last_result in the __mse_skip sentinel, got: {last_result}"
);
assert_eq!(
unwrap_skip_marker(&last_result),
Some(serde_json::json!("SKIP")),
"sentinel carries the raw body as its inner value"
);
}
#[tokio::test]
async fn post_worker_submit_verdict_invalid_returns_400() {
let engine = Engine::new(EngineCfg::default());
let task_id = StepId::new();
let handle = seed_task_with_handle(&engine, &task_id, "analyst").await;
let base_url = spawn_server(engine).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base_url}/v1/worker/submit?verdict=bogus"))
.header("Authorization", format!("Bearer {handle}"))
.body("payload")
.send()
.await
.expect("request");
assert_eq!(resp.status(), reqwest::StatusCode::BAD_REQUEST);
let body: serde_json::Value = resp.json().await.expect("json body");
let error = body["error"].as_str().expect("error string");
assert!(
error.contains("pass") && error.contains("blocked") && error.contains("skip"),
"error should enumerate the valid tier set: {error}"
);
}
#[tokio::test]
async fn post_worker_submit_verdict_skip_with_ok_false_returns_400() {
let engine = Engine::new(EngineCfg::default());
let task_id = StepId::new();
let handle = seed_task_with_handle(&engine, &task_id, "analyst").await;
let base_url = spawn_server(engine).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base_url}/v1/worker/submit?verdict=skip&ok=false"))
.header("Authorization", format!("Bearer {handle}"))
.body("SKIP")
.send()
.await
.expect("request");
assert_eq!(resp.status(), reqwest::StatusCode::BAD_REQUEST);
let body: serde_json::Value = resp.json().await.expect("json body");
let error = body["error"].as_str().expect("error string");
assert!(
error.contains("conflict") || error.contains("conflicting"),
"error should describe the ok/verdict conflict: {error}"
);
}
#[tokio::test]
async fn post_worker_submit_no_verdict_param_preserves_ok_true_pass_behavior() {
let engine = Engine::new(EngineCfg::default());
let task_id = StepId::new();
let handle = seed_task_with_handle(&engine, &task_id, "analyst").await;
let base_url = spawn_server(engine.clone()).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base_url}/v1/worker/submit"))
.header("Authorization", format!("Bearer {handle}"))
.body("PASS")
.send()
.await
.expect("request");
assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT);
let last_result = engine
.with_state("test.read_last_result", {
let task_id = task_id.clone();
move |s| s.tasks.get(&task_id).and_then(|t| t.last_result.clone())
})
.await
.expect("read last_result");
assert_eq!(
last_result,
Some(serde_json::json!("PASS")),
"no verdict param: raw body flows through unwrapped (pre-#76 shape)"
);
assert!(
!is_skip_marker(&last_result.expect("checked above")),
"no verdict param must NOT wrap in skip marker"
);
let tail = engine.output_tail(&task_id, 0).await;
let (_, final_ok) = tail
.iter()
.rev()
.find_map(|ev| match ev {
mlua_swarm::OutputEvent::Final { content, ok } => Some((content.clone(), *ok)),
_ => None,
})
.expect("Final present");
assert!(
final_ok,
"no verdict param, no ok override: Final.ok = true"
);
}