use async_trait::async_trait;
use mlua_swarm::core::agent_context::{AgentContextView, PROJECTION_PLACEMENT_KEY};
use mlua_swarm::core::projection::{
FileProjectionAdapter, ProjectionAdapter, ProjectionKey, ProjectionRef,
};
use mlua_swarm::core::projection_placement::ProjectionPlacement;
use mlua_swarm::{
CapToken, Ctx, Operator, SeniorBridge, SessionId, SpawnHook, StepId, WorkerBinding,
WorkerError, WorkerResult,
};
use serde_json::Value;
use std::collections::HashMap;
use tokio::sync::{mpsc, oneshot, Mutex};
use super::protocol::{current_parent_req_id, PendingReply, ServerMsg};
pub struct WSOperatorSession {
sid: SessionId,
tx: Mutex<Option<mpsc::UnboundedSender<ServerMsg>>>,
pending: Mutex<HashMap<String, oneshot::Sender<PendingReply>>>,
base_url: Option<std::sync::Arc<str>>,
}
impl WSOperatorSession {
pub(super) fn new_with_base_url(
sid: SessionId,
tx: mpsc::UnboundedSender<ServerMsg>,
base_url: Option<std::sync::Arc<str>>,
) -> Self {
Self {
sid,
tx: Mutex::new(Some(tx)),
pending: Mutex::new(HashMap::new()),
base_url,
}
}
pub(super) async fn replace_tx(&self, new_tx: mpsc::UnboundedSender<ServerMsg>) {
*self.tx.lock().await = Some(new_tx);
}
pub(super) async fn is_connected(&self) -> bool {
self.tx.lock().await.is_some()
}
pub(super) async fn clear_tx_if(&self, expected: &mpsc::UnboundedSender<ServerMsg>) {
let mut current = self.tx.lock().await;
if current
.as_ref()
.is_some_and(|sender| sender.same_channel(expected))
{
*current = None;
}
}
pub(crate) async fn clear_tx(&self) {
*self.tx.lock().await = None;
}
pub(crate) async fn fail_pending(&self, reason: &str) {
let drained: Vec<(String, oneshot::Sender<PendingReply>)> =
self.pending.lock().await.drain().collect();
if !drained.is_empty() {
tracing::warn!(
sid = %self.sid,
count = drained.len(),
reason,
"ws operator session: failing in-flight pending replies"
);
}
}
pub(super) async fn resolve_pending(&self, req_id: &str, reply: PendingReply) {
if let Some(otx) = self.pending.lock().await.remove(req_id) {
let _ = otx.send(reply);
}
}
async fn send_and_await(&self, req_id: String, msg: ServerMsg) -> Result<PendingReply, String> {
let (otx, orx) = oneshot::channel::<PendingReply>();
self.pending.lock().await.insert(req_id.clone(), otx);
let send_result = {
let guard = self.tx.lock().await;
match guard.as_ref() {
Some(tx) => tx
.send(msg)
.map_err(|_| "ws send channel closed".to_string()),
None => Err("ws operator disconnected".to_string()),
}
};
if let Err(e) = send_result {
self.pending.lock().await.remove(&req_id);
return Err(e);
}
orx.await
.map_err(|_| "ws operator: oneshot cancelled (= reply path closed)".to_string())
}
async fn send_oneway(&self, msg: ServerMsg) -> Result<(), String> {
let guard = self.tx.lock().await;
match guard.as_ref() {
Some(tx) => tx
.send(msg)
.map_err(|_| "ws send channel closed".to_string()),
None => Err("ws operator disconnected".to_string()),
}
}
}
#[async_trait]
impl SeniorBridge for WSOperatorSession {
async fn ask(&self, task_id: &StepId, question: Value) -> Result<Value, String> {
let req_id = format!("{}-ask-{}", self.sid, uuid::Uuid::new_v4());
let msg = ServerMsg::Ask {
req_id: req_id.clone(),
parent_req_id: current_parent_req_id(),
task_id: task_id.clone(),
question,
};
match self.send_and_await(req_id, msg).await? {
PendingReply::Answer(v) => Ok(v),
PendingReply::HookAck { .. } => {
Err("ws operator: unexpected hook_ack reply to ask".into())
}
PendingReply::SpawnAck { .. } => {
Err("ws operator: unexpected spawn_ack reply to ask".into())
}
PendingReply::SpawnHalt { .. } => {
Err("ws operator: unexpected spawn_halt reply to ask".into())
}
}
}
}
#[async_trait]
impl SpawnHook for WSOperatorSession {
async fn before(&self, ctx: &Ctx) -> Result<(), String> {
let req_id = format!("{}-hb-{}", self.sid, uuid::Uuid::new_v4());
let msg = ServerMsg::HookBefore {
req_id: req_id.clone(),
parent_req_id: current_parent_req_id(),
task_id: ctx.task_id.clone(),
agent: ctx.agent.clone(),
attempt: ctx.attempt,
};
match self.send_and_await(req_id, msg).await? {
PendingReply::HookAck { ok: true, .. } => Ok(()),
PendingReply::HookAck { ok: false, reason } => {
Err(reason.unwrap_or_else(|| "ws operator: spawn rejected".into()))
}
PendingReply::Answer(_) => {
Err("ws operator: unexpected answer reply to hook_before".into())
}
PendingReply::SpawnAck { .. } => {
Err("ws operator: unexpected spawn_ack reply to hook_before".into())
}
PendingReply::SpawnHalt { .. } => {
Err("ws operator: unexpected spawn_halt reply to hook_before".into())
}
}
}
async fn after(&self, ctx: &Ctx, result: &Value) -> Result<(), String> {
let req_id = format!("{}-ha-{}", self.sid, uuid::Uuid::new_v4());
let msg = ServerMsg::HookAfter {
req_id,
parent_req_id: current_parent_req_id(),
task_id: ctx.task_id.clone(),
agent: ctx.agent.clone(),
attempt: ctx.attempt,
result: result.clone(),
};
let _ = self.send_oneway(msg).await;
Ok(())
}
}
#[async_trait]
impl Operator for WSOperatorSession {
async fn execute(
&self,
ctx: &Ctx,
_system: Option<String>,
prompt: Value,
worker: Option<WorkerBinding>,
worker_token: CapToken,
) -> Result<WorkerResult, WorkerError> {
let Some(worker) = worker else {
return Err(WorkerError::Failed(format!(
"agent '{}' has no worker_binding; WS thin-path requires one \
(Blueprint AgentDef.profile.worker_binding)",
ctx.agent
)));
};
let req_id = format!("{}-spawn-{}", self.sid, uuid::Uuid::new_v4());
let worker_handle = ctx
.meta
.runtime
.get("worker_handle")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
let data_sink_endpoint = ctx
.meta
.runtime
.get("data_sink_endpoint")
.and_then(|v| v.as_str());
let run_id = ctx.meta.runtime.get("run_id").and_then(|v| v.as_str());
let view = AgentContextView::materialized_or_from_ctx(ctx);
let directive = default_spawn_directive_with_task_directive(
&ctx.agent,
ctx.task_id.as_str(),
&worker.variant,
&view,
data_sink_endpoint,
self.base_url.as_deref(),
run_id,
&prompt,
);
let projection_placement = ctx
.meta
.runtime
.get(PROJECTION_PLACEMENT_KEY)
.and_then(|v| serde_json::from_value::<ProjectionPlacement>(v.clone()).ok())
.unwrap_or_default();
let directive = append_projection_pointer(
directive,
&ctx.task_id,
&view,
run_id,
&projection_placement,
);
let msg = ServerMsg::Spawn {
req_id: req_id.clone(),
parent_req_id: current_parent_req_id(),
task_id: ctx.task_id.clone(),
agent: ctx.agent.clone(),
attempt: ctx.attempt,
capability_token: worker_token.encode(),
worker_handle,
worker: Some(worker),
directive,
};
match self.send_and_await(req_id, msg).await {
Ok(PendingReply::SpawnAck {
value,
ok,
error: None,
stats,
}) => Ok(WorkerResult {
value,
ok,
stats: stats.and_then(decode_ack_stats),
}),
Ok(PendingReply::SpawnAck {
error: Some(msg), ..
}) => Err(WorkerError::Failed(msg)),
Ok(PendingReply::SpawnHalt { value, reason }) => {
let marker = serde_json::json!({
"halted": true,
"reason": reason,
"value": value,
});
Ok(WorkerResult {
value: marker,
ok: true,
stats: None,
})
}
Ok(_) => Err(WorkerError::Failed(
"ws operator: unexpected non-spawn reply".into(),
)),
Err(e) => Err(WorkerError::Failed(format!("ws operator spawn: {e}"))),
}
}
fn requires_worker_binding(&self) -> bool {
true
}
}
fn decode_ack_stats(v: serde_json::Value) -> Option<mlua_swarm::store::trace::WorkerStats> {
let mut stats: mlua_swarm::store::trace::WorkerStats = serde_json::from_value(v).ok()?;
if stats.worker_kind.is_none() {
stats.worker_kind = Some("operator".to_string());
}
if stats.is_empty() {
None
} else {
Some(stats)
}
}
#[allow(clippy::too_many_arguments)]
pub(super) fn default_spawn_directive(
agent: &str,
task_id: &str,
subagent_type: &str,
view: &AgentContextView,
data_sink_endpoint: Option<&str>,
base_url: Option<&str>,
run_id: Option<&str>,
) -> String {
let context_header = view.to_directive_header();
let data_endpoint_block = match data_sink_endpoint {
Some(base) => format!(
"\n\
[Data path endpoint — MainAgent reminder]\n\
When you kick a SubAgent, inject the following two lines into\n\
its prompt / environment so Big Response payloads (4k+ tokens,\n\
files, intermediate artifacts) flow directly to the Store owner,\n\
bypassing the MainAgent (context stays small; only the out_id\n\
ref is passed around).\n \
DATA_EMIT: {base}/v1/data/emit (POST, auth = Bearer worker_handle or ?token=)\n \
DATA_GET: {base}/v1/data/<out_id|out_name> (the next SubAgent fetches from $IN_REFS)\n\
When a SubAgent produces a Big Response, POST it to DATA_EMIT\n\
and return only the one-line out_id ref (do not mix the body\n\
in; the MainAgent must not answer directly).\n\
\n"
),
None => String::new(),
};
let main_ai_reminder = match view.project_name_alias.as_deref() {
Some(a) => format!(
"\n\
[LDS Session Alias Reminder — MainAI mandatory]\n\
Before kicking the SubAgent below, call:\n \
mcp__lds__session_create(root=<working_dir>, alias=\"{a}\")\n\
(= establish a single task-level lds session; reuse on repeated dispatch).\n\
Then add this literal line to the SubAgent prompt body below:\n \
LDS Session Alias: {a}\n\
The SubAgent will call mcp__lds__session_start(alias=\"{a}\") on init,\n\
keeping worktree ownership unified across dispatches.\n\
(Full discipline rationale is inlined above; reach is via this directive itself,\n\
not via any external doc path. The 2 steps above are the complete contract.)\n\
\n"
),
None => String::new(),
};
let base_url_line = match base_url {
Some(u) => u.to_string(),
None => "<your server's actual bind — check with mse_doctor>".to_string(),
};
let run_route_line = match run_id {
Some(rid) => format!("GET <base_url>/v1/runs/{rid}"),
None => "GET <base_url>/v1/runs/<run_id>".to_string(),
};
format!(
"[agent_primitive dispatch=@{agent}]\n\
worker endpoint:\n \
GET <base_url>/v1/worker/prompt?task_id={task_id}\n \
POST <base_url>/v1/worker/submit\n\
auth: Bearer <worker_handle from THIS Spawn payload (= short `wh-XXXXXXXX` form)>\n\
task_id: {task_id}\n\
agent_id: {agent}\n\
{context_header}\
{data_endpoint_block}\
{main_ai_reminder}\
Kick a SubAgent via Agent tool with subagent_type=\"{subagent_type}\" (= project-local \
`.claude/agents/{subagent_type}.md`, this agent's Blueprint-declared worker binding). \
The prompt you pass to it MUST be EXACTLY these 4 lines (no preamble, no extra text):\n\
\n \
agent_id: {agent}\n \
worker_handle: <THIS Spawn payload's `worker_handle` field (short string `wh-XXXXXXXX`)>\n \
base_url: {base_url_line}\n \
task_id: {task_id}\n\
\n\
The SubAgent self-fetches system + prompt via GET (Bearer = handle), \
executes as agent @{agent}, POSTs raw body to /v1/worker/submit (Bearer = handle, \
server resolves task_id from handle), and replies `OUTPUT` 1 word. You then forward \
SpawnAck {{req_id, value:{{}}, ok:true}} through your operator client — MCP path: \
mse_ack(sid, req_id, kind=\"spawn_ack\", ok=true) (= empty value because canonical \
body lives in output_tail via the POST). \
Do NOT fetch /v1/worker/prompt yourself. Do NOT wrap, summarize, or field-select \
the SubAgent reply. Observation / debug is a separate channel (= agent-inspect MCP / \
{run_route_line}), do NOT mix it into the forward path. \
If the SubAgent type is not registered, FAIL LOUD: reply SpawnAck ok=false with an \
error explaining the missing `.claude/agents/{subagent_type}.md` — do NOT fall back \
to another subagent_type."
)
}
#[allow(clippy::too_many_arguments)]
pub(super) fn default_spawn_directive_with_task_directive(
agent: &str,
task_id: &str,
subagent_type: &str,
view: &AgentContextView,
data_sink_endpoint: Option<&str>,
base_url: Option<&str>,
run_id: Option<&str>,
task_directive: &Value,
) -> String {
let base = default_spawn_directive(
agent,
task_id,
subagent_type,
view,
data_sink_endpoint,
base_url,
run_id,
);
let task_directive_line = match task_directive {
Value::Null => String::new(),
Value::String(s) => format!("task_directive: {s}\n"),
other => format!("task_directive: {other}\n"),
};
format!("{base}{task_directive_line}")
}
fn append_projection_pointer(
directive: String,
task_id: &StepId,
view: &AgentContextView,
run_id: Option<&str>,
placement: &ProjectionPlacement,
) -> String {
let Some(root) = placement.resolve_root(view) else {
return directive;
};
match serde_json::to_value(view) {
Ok(ctx_data) => {
let key = ProjectionKey {
task_id: task_id.to_string(),
run_id: run_id.map(str::to_string),
step: None,
path: None,
};
let adapter = FileProjectionAdapter::with_placement(root, placement.clone());
match adapter.project(&key, &ctx_data) {
Ok(reference) => {
let pointer_value = match &reference {
ProjectionRef::File { path } => serde_json::json!({ "file": path }),
ProjectionRef::Query { endpoint, key } => {
serde_json::json!({ "endpoint": endpoint, "key": key })
}
};
format!("{directive}ctx_projection: {pointer_value}\n")
}
Err(err) => {
tracing::warn!(
%task_id,
error = %err,
"projection hook: materialize failed, spawning without a pointer"
);
directive
}
}
}
Err(err) => {
tracing::warn!(
%task_id,
error = %err,
"projection hook: AgentContextView serialize failed, spawning without a pointer"
);
directive
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use mlua_swarm::core::agent_context::{
TASK_METADATA_KEY, TASK_PROJECT_ROOT_KEY, TASK_WORK_DIR_KEY,
};
fn view_with(
alias: Option<&str>,
project_root: Option<&str>,
work_dir: Option<&str>,
) -> AgentContextView {
AgentContextView {
project_name_alias: alias.map(String::from),
project_root: project_root.map(String::from),
work_dir: work_dir.map(String::from),
..AgentContextView::default()
}
}
#[tokio::test]
async fn connection_state_tracks_the_current_sender() {
let (tx, _rx) = mpsc::unbounded_channel();
let session = WSOperatorSession::new_with_base_url(
SessionId::parse("S-connection-state").unwrap(),
tx.clone(),
None,
);
assert!(session.is_connected().await);
session.clear_tx_if(&tx).await;
assert!(!session.is_connected().await);
}
#[tokio::test]
async fn stale_disconnect_does_not_clear_a_reconnected_sender() {
let (old_tx, _old_rx) = mpsc::unbounded_channel();
let (new_tx, _new_rx) = mpsc::unbounded_channel();
let session = WSOperatorSession::new_with_base_url(
SessionId::parse("S-reconnect-state").unwrap(),
old_tx.clone(),
None,
);
session.replace_tx(new_tx).await;
session.clear_tx_if(&old_tx).await;
assert!(session.is_connected().await);
}
#[test]
fn directive_omits_project_name_alias_when_none() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert!(!d.contains("project_name_alias:"));
assert!(!d.contains("LDS Session Alias"));
assert!(!d.contains("session_create"));
}
#[test]
fn directive_emits_project_name_alias_when_some() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(Some("mse-task-7785"), None, None),
None,
None,
None,
);
assert!(
d.contains("project_name_alias: mse-task-7785"),
"directive missing project_name_alias header: {d}"
);
assert!(
d.contains("mcp__lds__session_create(root=<working_dir>, alias=\"mse-task-7785\")"),
"directive missing session_create reminder: {d}"
);
assert!(
d.contains("LDS Session Alias: mse-task-7785"),
"directive missing SubAgent prompt inject line: {d}"
);
assert!(
d.contains("inlined above") || d.contains("complete contract"),
"directive should inline rationale rather than point at external doc: {d}"
);
let forbidden_doc_ref = format!(".{}/CLAUDE.md", "claude");
assert!(
!d.contains(&forbidden_doc_ref),
"directive must not reference {forbidden_doc_ref} (out of MainAI scope): {d}"
);
}
#[test]
fn directive_omits_data_endpoint_when_none() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert!(!d.contains("[Data path endpoint"));
assert!(!d.contains("DATA_EMIT"));
assert!(!d.contains("DATA_GET"));
}
#[test]
fn directive_emits_data_endpoint_when_some() {
let base = "http://127.0.0.1:7785";
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
Some(base),
None,
None,
);
assert!(
d.contains("[Data path endpoint"),
"directive missing data endpoint block header: {d}"
);
assert!(
d.contains(&format!("DATA_EMIT: {base}/v1/data/emit")),
"directive missing single-mouth emit line: {d}"
);
assert!(
d.contains("Bearer worker_handle or ?token="),
"directive missing auth transport hint: {d}"
);
assert!(
d.contains(&format!("DATA_GET: {base}/v1/data/<out_id|out_name>")),
"directive missing GET line: {d}"
);
assert!(
!d.contains("emit-auth"),
"old split endpoint must not leak into directive: {d}"
);
assert!(
d.contains("bypassing the MainAgent") && d.contains("out_id ref"),
"directive should carry the ownership + bypass reasoning: {d}"
);
}
#[test]
fn directive_carries_declared_subagent_type_and_has_no_fallback() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert!(
d.contains("subagent_type=\"code-worker\""),
"directive must carry the Blueprint-declared subagent_type literally: {d}"
);
assert!(
d.contains(".claude/agents/code-worker.md"),
"directive must reference the declared subagent's own .md path: {d}"
);
assert!(
!d.contains("general-purpose"),
"directive must not fall back to subagent_type=\"general-purpose\": {d}"
);
assert!(
!d.contains("mse-worker\""),
"directive must not carry the old hardcoded \"mse-worker\" literal: {d}"
);
assert!(
d.contains("FAIL LOUD"),
"directive must instruct the MainAI to fail loud instead of falling back: {d}"
);
}
#[test]
fn directive_renders_actual_base_url_when_some() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
Some("http://127.0.0.1:8888"),
None,
);
assert!(
d.contains("base_url: http://127.0.0.1:8888"),
"directive must render the actual bind literally: {d}"
);
assert!(
!d.contains("mse_doctor"),
"no mse_doctor detour when bind is known: {d}"
);
}
#[test]
fn directive_falls_back_to_mse_doctor_pointer_when_none() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert!(
d.contains("check with mse_doctor"),
"fallback must point at mse_doctor: {d}"
);
}
#[test]
fn directive_never_contains_stale_example_port_7786() {
for base in [
None,
Some("http://127.0.0.1:7777"),
Some("http://192.0.2.1:9000"),
] {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(Some("mse-task-alias"), None, None),
Some("http://127.0.0.1:7785"),
base,
None,
);
assert!(
!d.contains("7786"),
"stale example port 7786 leaked: base={base:?}, d={d}"
);
}
}
#[test]
fn directive_never_contains_stale_tasks_id_route() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
Some("R-abc123"),
);
assert!(
!d.contains("/v1/tasks/{id}") && !d.contains("/v1/tasks/{{id}}"),
"stale /v1/tasks/{{id}} observation hint leaked: {d}"
);
}
#[test]
fn directive_renders_actual_run_id_when_some() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
Some("R-abc123"),
);
assert!(
d.contains("GET <base_url>/v1/runs/R-abc123"),
"directive missing real run_id in observation route: {d}"
);
}
#[test]
fn directive_falls_back_to_run_id_placeholder_when_none() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert!(
d.contains("GET <base_url>/v1/runs/<run_id>"),
"directive missing placeholder observation route: {d}"
);
}
#[test]
fn directive_omits_project_root_and_work_dir_when_both_none() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert!(!d.contains("project_root:"));
assert!(!d.contains("work_dir:"));
}
#[test]
fn directive_splices_project_root_and_work_dir_when_both_present() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, Some("/repo"), Some("/repo/work")),
None,
None,
None,
);
assert!(
d.contains("project_root: /repo"),
"directive missing project_root header: {d}"
);
assert!(
d.contains("work_dir: /repo/work"),
"directive missing work_dir header: {d}"
);
}
#[test]
fn directive_splices_project_root_only_when_work_dir_absent() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, Some("/repo"), None),
None,
None,
None,
);
assert!(
d.contains("project_root: /repo"),
"directive missing project_root header: {d}"
);
assert!(!d.contains("work_dir:"));
}
#[test]
fn directive_splices_task_metadata_when_some() {
let view = AgentContextView {
task_metadata: Some(serde_json::json!({"issue": 20})),
..view_with(None, Some("/repo"), None)
};
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view,
None,
None,
None,
);
assert!(
d.contains(r#"task_metadata: {"issue":20}"#),
"directive missing task_metadata header: {d}"
);
assert!(d.contains("project_root: /repo"));
}
#[test]
fn directive_omits_task_metadata_when_none() {
let d = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert!(!d.contains("task_metadata:"));
}
fn test_ctx(task_id: &str) -> mlua_swarm::Ctx {
mlua_swarm::Ctx::new(mlua_swarm::StepId::parse(task_id).unwrap(), 1, "a")
}
fn test_worker_binding() -> mlua_swarm::WorkerBinding {
mlua_swarm::WorkerBinding {
variant: "test-variant".into(),
tools: vec![],
request_digest: None,
requested_model: None,
}
}
fn test_cap_token() -> mlua_swarm::CapToken {
mlua_swarm::CapToken {
agent_id: "a".into(),
role: mlua_swarm::Role::Worker,
scopes: vec!["*".into()],
issued_at: 0,
expire_at: u64::MAX / 2,
max_uses: None,
nonce: "test-nonce".into(),
sig_hex: "".into(),
}
}
#[tokio::test]
async fn spawn_halt_reply_lands_as_ok_worker_result_with_marker() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-halt").unwrap(),
tx,
None,
));
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&test_ctx("ST-halt"),
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn { req_id, .. } => req_id,
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnHalt {
value: serde_json::json!({"partial": "abc"}),
reason: Some("shape verified".into()),
},
)
.await;
let result = handle.await.expect("join").expect("execute Ok");
assert!(
result.ok,
"spawn_halt must land as ok=true (normal termination), got: {result:?}"
);
assert_eq!(result.value["halted"], true);
assert_eq!(result.value["reason"], "shape verified");
assert_eq!(result.value["value"], serde_json::json!({"partial": "abc"}));
}
#[tokio::test]
async fn spawn_ack_with_error_still_lands_as_worker_error() {
use mlua_swarm::{Operator, WorkerError};
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-err").unwrap(),
tx,
None,
));
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&test_ctx("ST-err"),
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn { req_id, .. } => req_id,
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: false,
error: Some("real crash".into()),
stats: None,
},
)
.await;
let err = handle.await.expect("join").expect_err("must be error");
assert!(matches!(err, WorkerError::Failed(msg) if msg.contains("real crash")));
}
#[tokio::test]
async fn fail_pending_unblocks_a_parked_spawn_with_worker_error() {
use mlua_swarm::{Operator, WorkerError};
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-teardown").unwrap(),
tx,
None,
));
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&test_ctx("ST-teardown"),
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let _sent = rx.recv().await.expect("Spawn sent");
session.fail_pending("operator session torn down").await;
let err = handle
.await
.expect("join")
.expect_err("a parked spawn must fail once pending is drained");
assert!(
matches!(err, WorkerError::Failed(_)),
"fail_pending must surface a WorkerError::Failed, got: {err:?}"
);
}
#[tokio::test]
async fn execute_splices_project_root_and_work_dir_from_ctx_meta_runtime() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-ctxroot").unwrap(),
tx,
None,
));
let mut ctx = test_ctx("ST-ctxroot");
ctx.meta.runtime.insert(
TASK_PROJECT_ROOT_KEY.to_string(),
serde_json::json!("/repo"),
);
ctx.meta.runtime.insert(
TASK_WORK_DIR_KEY.to_string(),
serde_json::json!("/repo/work"),
);
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
let directive = directive.as_str();
assert!(
directive.contains("project_root: /repo"),
"directive missing project_root splice: {directive}"
);
assert!(
directive.contains("work_dir: /repo/work"),
"directive missing work_dir splice: {directive}"
);
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
}
#[tokio::test]
async fn execute_splices_project_root_only_when_ctx_meta_runtime_partial() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-ctxpartial").unwrap(),
tx,
None,
));
let mut ctx = test_ctx("ST-ctxpartial");
ctx.meta.runtime.insert(
TASK_PROJECT_ROOT_KEY.to_string(),
serde_json::json!("/repo"),
);
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
let directive = directive.as_str();
assert!(
directive.contains("project_root: /repo"),
"directive missing project_root splice: {directive}"
);
assert!(!directive.contains("work_dir:"));
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
}
#[tokio::test]
async fn execute_omits_project_root_and_work_dir_when_ctx_meta_runtime_absent() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-ctxabsent").unwrap(),
tx,
None,
));
let ctx = test_ctx("ST-ctxabsent");
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
let directive = directive.as_str();
assert!(!directive.contains("project_root:"));
assert!(!directive.contains("work_dir:"));
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
}
#[tokio::test]
async fn execute_splices_task_metadata_from_ctx_meta_runtime() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-ctxmeta").unwrap(),
tx,
None,
));
let mut ctx = test_ctx("ST-ctxmeta");
ctx.meta.runtime.insert(
TASK_METADATA_KEY.to_string(),
serde_json::json!({"issue": 20}),
);
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
let directive = directive.as_str();
assert!(
directive.contains(r#"task_metadata: {"issue":20}"#),
"directive missing task_metadata splice: {directive}"
);
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
}
#[test]
fn with_task_directive_splices_string_seed_verbatim() {
let directive = default_spawn_directive_with_task_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
&serde_json::json!("do the thing"),
);
let text = directive.as_str();
assert!(
text.contains("task_directive: do the thing"),
"missing task_directive line for a String seed: {text}"
);
}
#[test]
fn with_task_directive_renders_object_seed_as_json_literal() {
let directive = default_spawn_directive_with_task_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
&serde_json::json!({"key": "value"}),
);
let text = directive.as_str();
assert!(
text.contains(r#"task_directive: {"key":"value"}"#),
"missing JSON-literal task_directive line for an Object seed: {text}"
);
}
#[test]
fn with_task_directive_omits_line_when_null() {
let wrapped = default_spawn_directive_with_task_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
&serde_json::Value::Null,
);
let plain = default_spawn_directive(
"implementer",
"task-x",
"code-worker",
&view_with(None, None, None),
None,
None,
None,
);
assert_eq!(
wrapped,
serde_json::Value::String(plain),
"Value::Null seed must not add a task_directive line"
);
}
#[tokio::test]
async fn execute_splices_json_literal_task_directive_for_object_seed() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-objseed").unwrap(),
tx,
None,
));
let ctx = test_ctx("ST-objseed");
let rendered_prompt = serde_json::json!({"key": "value"});
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
rendered_prompt,
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
let directive = directive.as_str();
assert!(
directive.contains(r#"task_directive: {"key":"value"}"#),
"directive missing JSON-literal task_directive splice: {directive}"
);
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
}
#[tokio::test]
async fn execute_with_work_dir_appends_ctx_projection_pointer_and_materializes_file() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let dir = tempfile::TempDir::new().unwrap();
let mut ctx = test_ctx("ST-proj-1");
ctx.meta.runtime.insert(
TASK_WORK_DIR_KEY.to_string(),
Value::String(dir.path().to_string_lossy().into_owned()),
);
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-proj-1").unwrap(),
tx,
None,
));
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
assert!(
directive.contains("ctx_projection:"),
"directive missing ctx_projection pointer line: {directive}"
);
assert!(
!directive.contains("ctx_step_dir:"),
"directive must not carry the retired ctx_step_dir line: {directive}"
);
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
let expected_file = dir.path().join("workspace/tasks/ST-proj-1/ctx/_ctx.md");
assert!(
expected_file.exists(),
"materialized projection file missing at {expected_file:?}"
);
}
#[tokio::test]
async fn execute_without_work_dir_spawns_without_ctx_projection_pointer() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-proj-2").unwrap(),
tx,
None,
));
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&test_ctx("ST-proj-2"),
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
assert!(
!directive.contains("ctx_projection:"),
"directive must not carry a pointer line when work_dir is absent \
(fallback): {directive}"
);
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle
.await
.expect("join")
.expect("execute Ok — a materialize skip must not fail the spawn");
}
#[tokio::test]
async fn execute_with_project_root_only_appends_ctx_projection_pointer_default_placement() {
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let dir = tempfile::TempDir::new().unwrap();
let mut ctx = test_ctx("ST-proj-3");
ctx.meta.runtime.insert(
TASK_PROJECT_ROOT_KEY.to_string(),
Value::String(dir.path().to_string_lossy().into_owned()),
);
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-proj-3").unwrap(),
tx,
None,
));
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
assert!(
directive.contains("ctx_projection:"),
"work_dir absent must still fall back to project_root: {directive}"
);
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
let expected_file = dir.path().join("workspace/tasks/ST-proj-3/ctx/_ctx.md");
assert!(
expected_file.exists(),
"materialized projection file missing at {expected_file:?}"
);
}
#[tokio::test]
async fn execute_with_custom_projection_placement_uses_declared_root_and_template() {
use mlua_swarm::core::projection_placement::{ProjectionPlacement, RootPreference};
use mlua_swarm::Operator;
use tokio::sync::mpsc;
let work_dir = tempfile::TempDir::new().unwrap();
let project_root = tempfile::TempDir::new().unwrap();
let mut ctx = test_ctx("ST-proj-4");
ctx.meta.runtime.insert(
TASK_WORK_DIR_KEY.to_string(),
Value::String(work_dir.path().to_string_lossy().into_owned()),
);
ctx.meta.runtime.insert(
TASK_PROJECT_ROOT_KEY.to_string(),
Value::String(project_root.path().to_string_lossy().into_owned()),
);
let placement = ProjectionPlacement {
root_preference: RootPreference::ProjectRoot,
dir_template: "custom/{task_id}/out".to_string(),
};
ctx.meta.runtime.insert(
PROJECTION_PLACEMENT_KEY.to_string(),
serde_json::to_value(&placement).expect("placement serializes"),
);
let (tx, mut rx) = mpsc::unbounded_channel();
let session = std::sync::Arc::new(WSOperatorSession::new_with_base_url(
SessionId::parse("S-proj-4").unwrap(),
tx,
None,
));
let session_bg = session.clone();
let handle = tokio::spawn(async move {
session_bg
.execute(
&ctx,
None,
"".into(),
Some(test_worker_binding()),
test_cap_token(),
)
.await
});
let sent = rx.recv().await.expect("Spawn sent");
let req_id = match sent {
ServerMsg::Spawn {
req_id, directive, ..
} => {
assert!(
directive.contains("ctx_projection:"),
"directive missing ctx_projection pointer line: {directive}"
);
req_id
}
other => panic!("expected Spawn, got {other:?}"),
};
session
.resolve_pending(
&req_id,
PendingReply::SpawnAck {
value: serde_json::json!({}),
ok: true,
error: None,
stats: None,
},
)
.await;
handle.await.expect("join").expect("execute Ok");
let expected_file = project_root.path().join("custom/ST-proj-4/out/_ctx.md");
assert!(
expected_file.exists(),
"materialized projection file missing at custom placement target {expected_file:?}"
);
let unexpected_file = work_dir
.path()
.join("workspace/tasks/ST-proj-4/ctx/_ctx.md");
assert!(
!unexpected_file.exists(),
"declared root_preference=ProjectRoot must not fall back to work_dir: {unexpected_file:?}"
);
}
}