use crate::core::agent_context::{RUN_ID_KEY, STEP_CTX_KEY};
use crate::core::config::EngineCfg;
use crate::core::ctx::{Ctx, OperatorInfo, OperatorKind, SeniorBridge, SpawnHook};
use crate::core::errors::EngineError;
use crate::core::state::{
CapTokenRecord, DispatchOutcome, EngineState, Event, EventStream, OperatorSession, ResumeKey,
ResumePending, TaskSpec, TaskState, TaskStatus,
};
use crate::types::{
default_role_verb_table, now_unix, CapToken, Role, RoleVerbGate, RunId, SessionId, StepId,
TokenSigner, Verb,
};
use crate::worker::adapter::SpawnerAdapter;
use serde_json::Value;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::{broadcast, Mutex};
#[derive(Clone)]
pub struct Engine {
inner: Arc<EngineInner>,
}
struct EngineInner {
state: Mutex<EngineState>,
cfg: EngineCfg,
signer: TokenSigner,
gate: RoleVerbGate,
event_tx: broadcast::Sender<Event>,
senior_bridges: tokio::sync::RwLock<HashMap<String, Arc<dyn SeniorBridge>>>,
spawn_hooks: tokio::sync::RwLock<HashMap<String, Arc<dyn SpawnHook>>>,
operators: tokio::sync::RwLock<HashMap<String, Arc<dyn crate::operator::Operator>>>,
layer_registry: crate::middleware::LayerRegistry,
data_store: std::sync::RwLock<Option<Arc<dyn crate::store::output::OutputStore>>>,
verdict_contracts: std::sync::RwLock<HashMap<String, mlua_swarm_schema::VerdictContract>>,
}
pub(crate) fn render_directive_to_string(v: &Value) -> String {
match v {
Value::String(s) => s.clone(),
other => other.to_string(),
}
}
fn content_ref_to_value(content: crate::worker::output::ContentRef) -> Value {
match content {
crate::worker::output::ContentRef::Inline { value } => value,
crate::worker::output::ContentRef::FileRef {
path,
mime,
size_hint,
} => serde_json::json!({
"file_ref": path.to_string_lossy(),
"mime": mime,
"size_hint": size_hint,
}),
}
}
fn content_ref_to_comparable_string(content: crate::worker::output::ContentRef) -> String {
let value = content_ref_to_value(content);
match value {
Value::String(s) => s,
other => other.to_string(),
}
}
fn fold_final_and_parts(
tail: &[crate::worker::output::OutputEvent],
staged_names: &[String],
) -> Option<(Value, bool)> {
let (final_content, ok) = tail.iter().rev().find_map(|ev| match ev {
crate::worker::output::OutputEvent::Final { content, ok } => Some((content.clone(), *ok)),
_ => None,
})?;
let final_value = content_ref_to_value(final_content);
let mut parts = serde_json::Map::new();
for ev in tail {
if let crate::worker::output::OutputEvent::Artifact { name, content } = ev {
if staged_names.iter().any(|staged| staged == name) {
parts.insert(name.clone(), content_ref_to_value(content.clone()));
}
}
}
let value = if parts.is_empty() {
final_value
} else {
serde_json::json!({ "out": final_value, "parts": Value::Object(parts) })
};
Some((value, ok))
}
impl Engine {
pub fn new(cfg: EngineCfg) -> Self {
Self::new_with_layers(cfg, crate::middleware::LayerRegistry::new())
}
pub fn new_with_layers(
cfg: EngineCfg,
layer_registry: crate::middleware::LayerRegistry,
) -> Self {
let (event_tx, _) = broadcast::channel(256);
let signer = TokenSigner::new(&cfg.token_secret);
Self {
inner: Arc::new(EngineInner {
state: Mutex::new(EngineState::new()),
cfg,
signer,
gate: default_role_verb_table(),
event_tx,
senior_bridges: tokio::sync::RwLock::new(HashMap::new()),
spawn_hooks: tokio::sync::RwLock::new(HashMap::new()),
operators: tokio::sync::RwLock::new(HashMap::new()),
layer_registry,
data_store: std::sync::RwLock::new(None),
verdict_contracts: std::sync::RwLock::new(HashMap::new()),
}),
}
}
pub fn with_gate(self, gate: RoleVerbGate) -> Self {
let inner = Arc::new(EngineInner {
state: Mutex::new(EngineState::new()),
cfg: self.inner.cfg.clone(),
signer: self.inner.signer.clone(),
gate,
event_tx: self.inner.event_tx.clone(),
senior_bridges: tokio::sync::RwLock::new(HashMap::new()),
spawn_hooks: tokio::sync::RwLock::new(HashMap::new()),
operators: tokio::sync::RwLock::new(HashMap::new()),
layer_registry: self.inner.layer_registry.clone(),
data_store: std::sync::RwLock::new(None),
verdict_contracts: std::sync::RwLock::new(HashMap::new()),
});
Self { inner }
}
pub fn cfg(&self) -> &EngineCfg {
&self.inner.cfg
}
pub fn layer_registry(&self) -> &crate::middleware::LayerRegistry {
&self.inner.layer_registry
}
pub fn signer(&self) -> &TokenSigner {
&self.inner.signer
}
pub fn event_tx(&self) -> broadcast::Sender<Event> {
self.inner.event_tx.clone()
}
pub fn subscribe(&self) -> EventStream {
self.inner.event_tx.subscribe()
}
pub fn set_output_store(&self, store: Arc<dyn crate::store::output::OutputStore>) {
let mut guard = self
.inner
.data_store
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner());
*guard = Some(store);
}
fn output_store_backend(&self) -> Option<Arc<dyn crate::store::output::OutputStore>> {
self.inner
.data_store
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clone()
}
pub fn register_verdict_contracts(
&self,
contracts: HashMap<String, mlua_swarm_schema::VerdictContract>,
) {
let mut guard = self
.inner
.verdict_contracts
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner());
guard.extend(contracts);
}
pub async fn verdict_contract_for_task(
&self,
task_id: &StepId,
) -> Option<mlua_swarm_schema::VerdictContract> {
let tid = task_id.clone();
let agent = self
.with_state("verdict_contract_for_task", move |s| {
s.tasks.get(&tid).map(|t| t.spec.agent.clone())
})
.await
.ok()
.flatten()?;
self.inner
.verdict_contracts
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.get(&agent)
.cloned()
}
pub(crate) async fn staged_verdict_value_for(
&self,
task_id: &StepId,
attempt: u32,
) -> Option<String> {
let tail = self.output_tail(task_id, attempt).await;
tail.iter().rev().find_map(|ev| match ev {
crate::worker::output::OutputEvent::Artifact { name, content } if name == "verdict" => {
Some(content_ref_to_comparable_string(content.clone()))
}
_ => None,
})
}
async fn verdict_contract_completion_check(
&self,
task_id: &StepId,
attempt: u32,
ok: bool,
value: &str,
) -> Result<(), EngineError> {
if !ok {
return Ok(());
}
let Some(contract) = self.verdict_contract_for_task(task_id).await else {
return Ok(());
};
match contract.channel {
mlua_swarm_schema::VerdictChannel::Body => {
if contract.values.iter().any(|v| v == value) {
Ok(())
} else {
Err(EngineError::VerdictValueRejected {
value: value.to_string(),
allowed: contract.values.clone(),
})
}
}
mlua_swarm_schema::VerdictChannel::Part => {
match self.staged_verdict_value_for(task_id, attempt).await {
None => Err(EngineError::VerdictPartMissing {
allowed: contract.values.clone(),
}),
Some(staged) if contract.values.iter().any(|v| v == &staged) => Ok(()),
Some(staged) => Err(EngineError::VerdictValueRejected {
value: staged,
allowed: contract.values.clone(),
}),
}
}
}
}
pub async fn with_state<F, R>(&self, op: &'static str, f: F) -> Result<R, EngineError>
where
F: FnOnce(&mut EngineState) -> R,
{
let cfg = &self.inner.cfg;
let mut guard_opt = None;
for attempt in 0..=cfg.max_retry {
match self.inner.state.try_lock() {
Ok(g) => {
guard_opt = Some(g);
break;
}
Err(_) if cfg.try_only => return Err(EngineError::LockBusy(op)),
Err(_) => {
let backoff = cfg.backoff_ms_step * (attempt as u64 + 1);
tokio::time::sleep(Duration::from_millis(backoff)).await;
}
}
}
let mut guard = guard_opt.ok_or(EngineError::LockBusyAfterRetry(op))?;
let start = Instant::now();
let result = f(&mut guard);
let elapsed_ms = start.elapsed().as_millis();
drop(guard);
if elapsed_ms > cfg.max_hold_ms {
panic!(
"Engine.with_state('{op}') held {elapsed_ms}ms > max {}ms — suspected R3 violation (long op inside lock)",
cfg.max_hold_ms
);
}
Ok(result)
}
pub async fn verify_token(&self, token: &CapToken, verb: Verb) -> Result<(), EngineError> {
if !self.inner.signer.verify_sig(token) {
return Err(EngineError::BadSignature);
}
if token.is_expired(now_unix()) {
return Err(EngineError::TokenExpired);
}
if !self.inner.gate.is_allowed(token.role, verb) {
return Err(EngineError::RoleViolation {
role: token.role,
verb,
});
}
let fp = token.fingerprint();
self.with_state("token.consume", move |s| {
let rec = s
.tokens
.get_mut(&fp)
.ok_or_else(|| EngineError::TokenNotFound(fp.clone()))?;
rec.consume()
.map_err(|_: crate::core::state::CapTokenConsumeError| {
EngineError::TokenUsesExhausted
})?;
Ok::<(), EngineError>(())
})
.await??;
Ok(())
}
pub async fn verify_token_for_task(
&self,
token: &CapToken,
verb: Verb,
task_id: &StepId,
) -> Result<(), EngineError> {
self.verify_token(token, verb).await?;
if token.role != Role::Worker {
return Ok(());
}
let fp = token.fingerprint();
let arg_tid = task_id.clone();
self.with_state("token.ownership_gate", move |s| {
let bound = s.tokens.get(&fp).and_then(|r| r.task_id.as_ref()).cloned();
match bound {
Some(t) if t == arg_tid => Ok(()),
Some(t) => Err(EngineError::TokenTaskMismatch {
bound: t.into_string(),
arg: arg_tid.into_string(),
}),
None => Err(EngineError::TokenNotFound(fp.clone())),
}
})
.await??;
Ok(())
}
pub async fn task_id_from_token(&self, token: &CapToken) -> Result<StepId, EngineError> {
if token.role != Role::Worker {
return Err(EngineError::RoleViolation {
role: token.role,
verb: Verb::PostResult,
});
}
let fp = token.fingerprint();
self.with_state("task_id_from_token", move |s| {
s.tokens
.get(&fp)
.and_then(|r| r.task_id.as_ref())
.cloned()
.ok_or_else(|| EngineError::TokenNotFound(fp.clone()))
})
.await?
}
pub async fn task_id_from_handle(&self, handle: &str) -> Result<StepId, EngineError> {
let h = handle.to_string();
self.with_state("task_id_from_handle", move |s| {
let fp = s
.worker_handles
.get(&h)
.cloned()
.ok_or_else(|| EngineError::TokenNotFound(format!("handle={h}")))?;
s.tokens
.get(&fp)
.and_then(|r| r.task_id.as_ref())
.cloned()
.ok_or_else(|| EngineError::TokenNotFound(format!("fp={fp}")))
})
.await?
}
pub async fn submit_worker_result_trusted(
&self,
task_id: &StepId,
attempt: u32,
value: Value,
ok: bool,
) -> Result<(), EngineError> {
let comparable_value =
content_ref_to_comparable_string(crate::worker::output::ContentRef::Inline {
value: value.clone(),
});
self.verdict_contract_completion_check(task_id, attempt, ok, &comparable_value)
.await?;
let task_id_for_apply = task_id.clone();
let value_for_event = value.clone();
self.with_state("submit_worker_result_trusted.output", move |s| {
let ev = crate::worker::output::OutputEvent::Final {
content: crate::worker::output::ContentRef::Inline {
value: value_for_event,
},
ok,
};
s.output_store
.entry((task_id_for_apply.clone(), attempt))
.or_default()
.push(ev.clone());
s.push_event(crate::core::state::Event::WorkerOutput {
task_id: task_id_for_apply,
attempt,
event: ev,
});
})
.await?;
let task_id_for_result = task_id.clone();
let value_for_result = value.clone();
self.with_state("submit_worker_result_trusted.last_result", move |s| {
if let Some(t) = s.tasks.get_mut(&task_id_for_result) {
t.last_result = Some(value_for_result);
t.updated_at = now_unix();
}
})
.await?;
let content = crate::worker::output::ContentRef::Inline { value };
self.materialize_final_submission(task_id, attempt, &content, ok)
.await?;
Ok(())
}
pub async fn stage_worker_artifact_trusted(
&self,
task_id: &StepId,
attempt: u32,
name: String,
value: Value,
) -> Result<(), EngineError> {
let content = crate::worker::output::ContentRef::Inline { value };
let task_id_for_apply = task_id.clone();
let name_for_apply = name.clone();
let content_for_apply = content.clone();
self.with_state("stage_worker_artifact_trusted.output", move |s| {
let ev = crate::worker::output::OutputEvent::Artifact {
name: name_for_apply.clone(),
content: content_for_apply,
};
s.output_store
.entry((task_id_for_apply.clone(), attempt))
.or_default()
.push(ev.clone());
s.worker_artifact_names
.entry((task_id_for_apply.clone(), attempt))
.or_default()
.push(name_for_apply);
s.push_event(crate::core::state::Event::WorkerOutput {
task_id: task_id_for_apply,
attempt,
event: ev,
});
})
.await?;
self.materialize_artifact_submission(task_id, attempt, &name, &content)
.await?;
Ok(())
}
async fn worker_artifact_names_for(&self, task_id: &StepId, attempt: u32) -> Vec<String> {
let key = (task_id.clone(), attempt);
self.with_state("worker_artifact_names_for", move |s| {
s.worker_artifact_names
.get(&key)
.cloned()
.unwrap_or_default()
})
.await
.unwrap_or_default()
}
async fn mint_worker_handle(&self, worker_fp: String) -> Result<String, EngineError> {
let short = crate::types::secure_hex(4);
let handle = format!("wh-{short}");
let h = handle.clone();
self.with_state("mint_worker_handle", move |s| {
s.worker_handles.insert(h, worker_fp);
})
.await?;
Ok(handle)
}
pub async fn attach(
&self,
operator_id: impl Into<String>,
role: Role,
ttl: Duration,
) -> Result<CapToken, EngineError> {
self.attach_with(
operator_id,
role,
ttl,
crate::core::ctx::OperatorInfo::default(),
)
.await
}
pub async fn register_senior_bridge(
&self,
id: impl Into<String>,
bridge: Arc<dyn SeniorBridge>,
) {
self.inner
.senior_bridges
.write()
.await
.insert(id.into(), bridge);
}
pub async fn register_spawn_hook(&self, id: impl Into<String>, hook: Arc<dyn SpawnHook>) {
self.inner.spawn_hooks.write().await.insert(id.into(), hook);
}
pub async fn register_operator(
&self,
id: impl Into<String>,
operator: Arc<dyn crate::operator::Operator>,
) {
self.inner
.operators
.write()
.await
.insert(id.into(), operator);
}
pub async fn unregister_senior_bridge(&self, id: &str) {
self.inner.senior_bridges.write().await.remove(id);
}
pub async fn unregister_spawn_hook(&self, id: &str) {
self.inner.spawn_hooks.write().await.remove(id);
}
pub async fn unregister_operator(&self, id: &str) {
self.inner.operators.write().await.remove(id);
}
pub async fn list_spawn_hook_ids(&self) -> Vec<String> {
self.inner
.spawn_hooks
.read()
.await
.keys()
.cloned()
.collect()
}
pub async fn list_senior_bridge_ids(&self) -> Vec<String> {
self.inner
.senior_bridges
.read()
.await
.keys()
.cloned()
.collect()
}
pub async fn list_operator_ids(&self) -> Vec<String> {
self.inner.operators.read().await.keys().cloned().collect()
}
#[allow(clippy::too_many_arguments)]
pub async fn attach_with_ids(
&self,
operator_id: impl Into<String>,
role: Role,
ttl: Duration,
kind: Option<OperatorKind>,
bridge_id: Option<String>,
hook_id: Option<String>,
operator_backend_id: Option<String>,
operator_kind_overrides: HashMap<String, OperatorKind>,
bp_agent_kinds: HashMap<String, OperatorKind>,
bp_global_kind: Option<OperatorKind>,
) -> Result<CapToken, EngineError> {
let operator_id = operator_id.into();
let token = self
.inner
.signer
.session(operator_id.clone(), role, vec!["*".into()], ttl);
let session_id = SessionId::new();
let fp = token.fingerprint();
let now = now_unix();
let token_for_store = token.clone();
self.with_state("attach_with_ids", |s| {
s.tokens
.insert(fp.clone(), CapTokenRecord::from_token(token_for_store));
s.sessions.insert(
session_id.clone(),
OperatorSession {
id: session_id.clone(),
operator_id: operator_id.clone(),
role,
attached_at: now,
last_seen: now,
attached: true,
owned_task_ids: Vec::new(),
token_fp: fp.clone(),
operator_kind: kind,
runtime_agent_kinds: operator_kind_overrides,
bp_agent_kinds,
bp_global_kind,
bridge_id,
hook_id,
operator_backend_id,
},
);
s.push_event(Event::SessionAttached {
session_id: session_id.clone(),
role,
});
})
.await?;
let _ = self
.inner
.event_tx
.send(Event::SessionAttached { session_id, role });
Ok(token)
}
async fn resolve_operator_info(
&self,
session: &OperatorSession,
agent_name: &str,
) -> OperatorInfo {
let senior_bridge = if let Some(id) = &session.bridge_id {
self.inner.senior_bridges.read().await.get(id).cloned()
} else {
None
};
let spawn_hook = if let Some(id) = &session.hook_id {
self.inner.spawn_hooks.read().await.get(id).cloned()
} else {
None
};
let operator = if let Some(id) = &session.operator_backend_id {
self.inner.operators.read().await.get(id).cloned()
} else {
None
};
let runtime_agent = session.runtime_agent_kinds.get(agent_name).copied();
let runtime_global = session.operator_kind;
let bp_agent = session.bp_agent_kinds.get(agent_name).copied();
let bp_global = session.bp_global_kind;
let kind = crate::core::ctx::collapse_operator_kind(
runtime_agent,
runtime_global,
bp_agent,
bp_global,
);
OperatorInfo {
kind,
id: session.operator_id.clone(),
senior_bridge,
spawn_hook,
operator,
}
}
pub async fn attach_with(
&self,
operator_id: impl Into<String>,
role: Role,
ttl: Duration,
operator_info: crate::core::ctx::OperatorInfo,
) -> Result<CapToken, EngineError> {
let operator_id = operator_id.into();
let kind = operator_info.kind;
let bridge_id = if let Some(bridge) = operator_info.senior_bridge.clone() {
let id = format!("br-{}", crate::types::uid_hex(8));
self.inner
.senior_bridges
.write()
.await
.insert(id.clone(), bridge);
Some(id)
} else {
None
};
let hook_id = if let Some(hook) = operator_info.spawn_hook.clone() {
let id = format!("hk-{}", crate::types::uid_hex(8));
self.inner
.spawn_hooks
.write()
.await
.insert(id.clone(), hook);
Some(id)
} else {
None
};
let operator_backend_id = if let Some(operator) = operator_info.operator.clone() {
let id = format!("ob-{}", crate::types::uid_hex(8));
self.inner
.operators
.write()
.await
.insert(id.clone(), operator);
Some(id)
} else {
None
};
let token = self
.inner
.signer
.session(operator_id.clone(), role, vec!["*".into()], ttl);
let session_id = SessionId::new();
let fp = token.fingerprint();
let now = now_unix();
let token_for_store = token.clone();
self.with_state("attach_with", |s| {
s.tokens
.insert(fp.clone(), CapTokenRecord::from_token(token_for_store));
s.sessions.insert(
session_id.clone(),
OperatorSession {
id: session_id.clone(),
operator_id,
role,
attached_at: now,
last_seen: now,
attached: true,
owned_task_ids: Vec::new(),
token_fp: fp.clone(),
operator_kind: Some(kind),
runtime_agent_kinds: HashMap::new(),
bp_agent_kinds: HashMap::new(),
bp_global_kind: None,
bridge_id,
hook_id,
operator_backend_id,
},
);
s.push_event(Event::SessionAttached {
session_id: session_id.clone(),
role,
});
})
.await?;
let _ = self
.inner
.event_tx
.send(Event::SessionAttached { session_id, role });
Ok(token)
}
pub async fn detach(&self, token: &CapToken) -> Result<(), EngineError> {
self.verify_token(token, Verb::DetachSession).await?;
let fp = token.fingerprint();
self.with_state("detach", move |s| {
let sid = s
.sessions
.iter()
.find(|(_, sess)| sess.token_fp == fp)
.map(|(id, _)| id.clone());
if let Some(sid) = sid {
if let Some(sess) = s.sessions.get_mut(&sid) {
sess.attached = false;
}
s.push_event(Event::SessionDetached {
session_id: sid.clone(),
});
let _ = sid;
}
})
.await?;
Ok(())
}
pub async fn heartbeat(&self, token: &CapToken) -> Result<(), EngineError> {
self.verify_token(token, Verb::Heartbeat).await?;
let now = now_unix();
let fp = token.fingerprint();
self.with_state("heartbeat", move |s| {
if let Some(sess) = s.sessions.values_mut().find(|sess| sess.token_fp == fp) {
sess.last_seen = now;
sess.attached = true;
}
})
.await?;
Ok(())
}
pub async fn start_task(
&self,
token: &CapToken,
spec: TaskSpec,
) -> Result<StepId, EngineError> {
self.verify_token(token, Verb::StartTask).await?;
let task_id = StepId::new();
let initial_directive = spec.initial_directive.clone();
let task_id_clone = task_id.clone();
let fp = token.fingerprint();
let max_depth = self.inner.cfg.max_spawn_depth;
self.with_state("start_task", move |s| {
let parent_depth_opt = s
.tokens
.get(&fp)
.and_then(|rec| rec.task_id.as_ref())
.and_then(|tid| s.tasks.get(tid))
.map(|t| t.spawn_depth);
let depth = match parent_depth_opt {
Some(d) => {
if d + 1 >= max_depth {
return Err(EngineError::SpawnDepthExceeded {
current: d + 1,
max: max_depth,
});
}
d + 1
}
None => 0,
};
let mut task = TaskState::new(task_id_clone.clone(), spec);
task.spawn_depth = depth;
s.tasks.insert(task_id_clone.clone(), task);
s.prompts
.insert((task_id_clone.clone(), 1), initial_directive);
if let Some(sess) = s.sessions.values_mut().find(|sess| sess.token_fp == fp) {
sess.owned_task_ids.push(task_id_clone.clone());
}
s.push_event(Event::TaskCreated {
task_id: task_id_clone.clone(),
});
Ok::<(), EngineError>(())
})
.await??;
let _ = self.inner.event_tx.send(Event::TaskCreated {
task_id: task_id.clone(),
});
Ok(task_id)
}
pub async fn read_task_state(
&self,
token: &CapToken,
task_id: &StepId,
) -> Result<TaskState, EngineError> {
self.verify_token_for_task(token, Verb::ReadTaskState, task_id)
.await?;
let task_id = task_id.clone();
self.with_state("read_task_state", move |s| {
s.tasks
.get(&task_id)
.cloned()
.ok_or_else(|| EngineError::TaskNotFound(task_id.to_string()))
})
.await?
}
pub async fn cancel_task(&self, token: &CapToken, task_id: &StepId) -> Result<(), EngineError> {
self.verify_token_for_task(token, Verb::CancelTask, task_id)
.await?;
let tid = task_id.clone();
self.with_state("cancel_task", move |s| {
let task = s
.tasks
.get_mut(&tid)
.ok_or_else(|| EngineError::TaskNotFound(tid.to_string()))?;
task.status = TaskStatus::Cancelled;
task.updated_at = now_unix();
s.push_event(Event::TaskCancelled {
task_id: tid.clone(),
});
Ok::<(), EngineError>(())
})
.await??;
self.wake_task(task_id).await?;
Ok(())
}
pub async fn dispatch_attempt_with(
&self,
token: &CapToken,
task_id: &StepId,
spawner: &Arc<dyn SpawnerAdapter>,
run_id: Option<&RunId>,
) -> Result<DispatchOutcome, EngineError> {
self.verify_token(token, Verb::DispatchAttempt).await?;
let task_id = task_id.clone();
let fp = token.fingerprint();
let tid_for_prep = task_id.clone();
let (attempt, agent, session_snapshot, step_ctx) = self
.with_state("dispatch.prep", move |s| {
let task = s
.tasks
.get_mut(&tid_for_prep)
.ok_or_else(|| EngineError::TaskNotFound(tid_for_prep.to_string()))?;
task.attempt += 1;
task.status = TaskStatus::Running;
task.updated_at = now_unix();
let attempt = task.attempt;
let initial = task.spec.initial_directive.clone();
s.prompts
.entry((tid_for_prep.clone(), attempt))
.or_insert(initial);
let task = s
.tasks
.get(&tid_for_prep)
.ok_or_else(|| EngineError::TaskNotFound(tid_for_prep.to_string()))?;
let agent = task.spec.agent.clone();
let step_ctx = task.spec.step_ctx.clone();
let sess_clone = s
.sessions
.values()
.find(|sess| sess.token_fp == fp)
.cloned();
Ok::<_, EngineError>((attempt, agent, sess_clone, step_ctx))
})
.await??;
let operator_info = match session_snapshot {
Some(sess) => self.resolve_operator_info(&sess, &agent).await,
None => OperatorInfo::default(),
};
let worker_token = self.inner.signer.session(
format!("worker-of-{task_id}"),
Role::Worker,
vec!["*".into()],
Duration::from_secs(1800),
);
let worker_fp = worker_token.fingerprint();
let task_id_for_worker = task_id.clone();
let worker_token_for_store = worker_token.clone();
self.with_state("dispatch.mint_worker", move |s| {
s.tokens.insert(
worker_fp,
CapTokenRecord::from_worker_token(worker_token_for_store, task_id_for_worker),
);
})
.await?;
let worker_handle = self.mint_worker_handle(worker_token.fingerprint()).await?;
let mut ctx = Ctx::new(task_id.clone(), attempt, agent.clone());
ctx.operator = operator_info; ctx.meta
.runtime
.insert("worker_handle".to_string(), Value::String(worker_handle));
if let Some(rid) = run_id {
ctx.meta
.runtime
.insert(RUN_ID_KEY.to_string(), Value::String(rid.to_string()));
}
if let Some(step_ctx) = step_ctx {
ctx.meta.runtime.insert(STEP_CTX_KEY.to_string(), step_ctx);
}
let worker = spawner
.spawn(self, &ctx, task_id.clone(), attempt, worker_token)
.await
.map_err(|e| EngineError::DispatchFailed(e.to_string()))?;
let signal_result: Result<(), String> = worker.join().await.map_err(|e| e.to_string());
let value_ok: Result<(Value, bool), String> = match signal_result {
Ok(()) => {
let tail = self.output_tail(&task_id, attempt).await;
let staged_names = self.worker_artifact_names_for(&task_id, attempt).await;
fold_final_and_parts(&tail, &staged_names)
.ok_or_else(|| "no Final in output_tail".to_string())
}
Err(msg) => Err(msg),
};
let outcome = self
.with_state("dispatch.apply", |s| {
if !s.tasks.contains_key(&task_id) {
return Err(EngineError::TaskNotFound(task_id.to_string()));
}
match value_ok {
Ok((value, ok)) => {
let pass = ok;
{
let task = s.tasks.get_mut(&task_id).unwrap();
task.last_result = Some(value.clone());
task.updated_at = now_unix();
task.status = if pass {
TaskStatus::Pass
} else {
TaskStatus::Blocked
};
}
s.push_event(Event::TaskAttemptCompleted {
task_id: task_id.clone(),
attempt,
result: value.clone(),
});
if pass {
s.push_event(Event::TaskPass {
task_id: task_id.clone(),
result: value.clone(),
});
Ok::<_, EngineError>(DispatchOutcome::Pass(value))
} else {
s.push_event(Event::TaskBlocked {
task_id: task_id.clone(),
result: value.clone(),
});
Ok(DispatchOutcome::Blocked(value))
}
}
Err(msg) => {
let task = s.tasks.get_mut(&task_id).unwrap();
task.status = TaskStatus::Blocked;
task.updated_at = now_unix();
Err(EngineError::DispatchFailed(msg))
}
}
})
.await??;
let _ = self.inner.event_tx.send(Event::TaskAttemptCompleted {
task_id: task_id.clone(),
attempt,
result: match &outcome {
DispatchOutcome::Pass(v) | DispatchOutcome::Blocked(v) => v.clone(),
_ => Value::Null,
},
});
self.wake_task(&task_id).await?;
Ok(outcome)
}
pub async fn fetch_prompt(
&self,
token: &CapToken,
task_id: &StepId,
) -> Result<Value, EngineError> {
self.verify_token_for_task(token, Verb::FetchPrompt, task_id)
.await?;
let task_id = task_id.clone();
self.with_state("fetch_prompt", move |s| {
let task = s
.tasks
.get(&task_id)
.ok_or_else(|| EngineError::TaskNotFound(task_id.to_string()))?;
s.prompts
.get(&(task_id.clone(), task.attempt.max(1)))
.cloned()
.ok_or_else(|| {
EngineError::ResourceNotFound(format!(
"prompt({}, attempt={})",
task_id, task.attempt
))
})
})
.await?
}
pub async fn fetch_worker_payload(
&self,
token: &CapToken,
task_id: &StepId,
) -> Result<crate::types::WorkerPayload, EngineError> {
self.verify_token_for_task(token, Verb::FetchPrompt, task_id)
.await?;
let task_id_clone = task_id.clone();
let mut payload = self
.with_state("fetch_worker_payload", move |s| {
let task = s
.tasks
.get(&task_id_clone)
.ok_or_else(|| EngineError::TaskNotFound(task_id_clone.to_string()))?;
let attempt = task.attempt.max(1);
let prompt = s
.prompts
.get(&(task_id_clone.clone(), attempt))
.cloned()
.ok_or_else(|| {
EngineError::ResourceNotFound(format!(
"prompt({}, attempt={})",
task_id_clone, attempt
))
})?;
let system = s
.systems
.get(&(task_id_clone.clone(), attempt))
.cloned()
.unwrap_or(None);
let agent = task.spec.agent.clone();
let context = s
.agent_ctx
.get(&(task_id_clone.clone(), attempt))
.map(|e| e.view.clone());
Ok::<_, EngineError>(crate::types::WorkerPayload {
task_id: task_id_clone.clone(),
attempt,
agent,
prompt: render_directive_to_string(&prompt),
system,
context,
system_ref: None,
})
})
.await??;
self.apply_system_ref_threshold(&mut payload).await?;
Ok(payload)
}
pub async fn fetch_worker_payload_trusted(
&self,
task_id: &StepId,
) -> Result<crate::types::WorkerPayload, EngineError> {
let task_id_clone = task_id.clone();
let mut payload = self
.with_state("fetch_worker_payload_trusted", move |s| {
let task = s
.tasks
.get(&task_id_clone)
.ok_or_else(|| EngineError::TaskNotFound(task_id_clone.to_string()))?;
let attempt = task.attempt.max(1);
let prompt = s
.prompts
.get(&(task_id_clone.clone(), attempt))
.cloned()
.ok_or_else(|| {
EngineError::ResourceNotFound(format!(
"prompt({}, attempt={})",
task_id_clone, attempt
))
})?;
let system = s
.systems
.get(&(task_id_clone.clone(), attempt))
.cloned()
.unwrap_or(None);
let agent = task.spec.agent.clone();
let context = s
.agent_ctx
.get(&(task_id_clone.clone(), attempt))
.map(|e| e.view.clone());
Ok::<_, EngineError>(crate::types::WorkerPayload {
task_id: task_id_clone.clone(),
attempt,
agent,
prompt: render_directive_to_string(&prompt),
system,
context,
system_ref: None,
})
})
.await??;
self.apply_system_ref_threshold(&mut payload).await?;
Ok(payload)
}
async fn apply_system_ref_threshold(
&self,
payload: &mut crate::types::WorkerPayload,
) -> Result<(), EngineError> {
let Some(rendered) = payload.system.take() else {
return Ok(());
};
let cfg = self.cfg().system_ref.clone();
if rendered.len() <= cfg.threshold_bytes {
payload.system = Some(rendered);
return Ok(());
}
use sha2::Digest;
let size_bytes = rendered.len() as u64;
let sha256 = hex::encode(sha2::Sha256::digest(rendered.as_bytes()));
let task_id = &payload.task_id;
let attempt = payload.attempt;
let system_ref = match cfg.mode {
crate::types::SystemRefMode::Http => crate::types::SystemRef {
uri: format!("/v1/worker/prompt/system?task_id={task_id}&attempt={attempt}"),
sha256,
size_bytes,
mode: crate::types::SystemRefMode::Http,
},
crate::types::SystemRefMode::File => {
tokio::fs::create_dir_all(&cfg.store_dir).await?;
let path = cfg.store_dir.join(format!("{task_id}-{attempt}.md"));
tokio::fs::write(&path, rendered.as_bytes()).await?;
crate::types::SystemRef {
uri: format!("file://{}", path.display()),
sha256,
size_bytes,
mode: crate::types::SystemRefMode::File,
}
}
};
payload.system = None;
payload.system_ref = Some(system_ref);
Ok(())
}
pub async fn context_policy_for(
&self,
task_id: &StepId,
attempt: u32,
) -> mlua_swarm_schema::ContextPolicy {
let key = (task_id.clone(), attempt);
self.with_state("context_policy_for", move |s| {
s.agent_ctx
.get(&key)
.map(|e| e.policy.clone())
.unwrap_or_default()
})
.await
.unwrap_or_default()
}
pub async fn step_naming_for(
&self,
task_id: &StepId,
) -> Option<Arc<crate::core::step_naming::StepNaming>> {
let key = task_id.clone();
self.with_state("step_naming_for", move |s| {
s.step_namings.get(&key).cloned()
})
.await
.ok()
.flatten()
}
pub async fn projection_placement_for(
&self,
task_id: &StepId,
) -> Option<Arc<crate::core::projection_placement::ProjectionPlacement>> {
let key = task_id.clone();
self.with_state("projection_placement_for", move |s| {
s.projection_placements.get(&key).cloned()
})
.await
.ok()
.flatten()
}
pub async fn agent_context_for(
&self,
task_id: &StepId,
attempt: u32,
) -> Option<crate::core::agent_context::AgentContextView> {
let key = (task_id.clone(), attempt);
self.with_state("agent_context_for", move |s| {
s.agent_ctx.get(&key).map(|e| e.view.clone())
})
.await
.ok()
.flatten()
}
pub async fn task_attempt(&self, task_id: &StepId) -> Result<u32, EngineError> {
let task_id = task_id.clone();
self.with_state("task_attempt", move |s| {
s.tasks
.get(&task_id)
.map(|t| t.attempt)
.ok_or_else(|| EngineError::TaskNotFound(task_id.to_string()))
})
.await?
}
pub async fn bake_worker_system_prompt(
&self,
task_id: &StepId,
attempt: u32,
system: Option<String>,
) -> Result<(), EngineError> {
let task_id = task_id.clone();
self.with_state("bake_worker_system_prompt", move |s| {
if let Some(rendered) = system.as_ref() {
if let Some(agent) = s.tasks.get(&task_id).map(|t| t.spec.agent.clone()) {
s.agent_render_sizes.insert(agent, rendered.len());
}
}
s.systems.insert((task_id, attempt), system);
})
.await?;
Ok(())
}
pub async fn agent_last_rendered_size(&self, agent_name: &str) -> Option<usize> {
let agent_name = agent_name.to_string();
self.with_state("agent_last_rendered_size", move |s| {
s.agent_render_sizes.get(&agent_name).copied()
})
.await
.ok()
.flatten()
}
pub async fn raw_system_prompt(
&self,
task_id: &StepId,
attempt: u32,
) -> Result<Option<String>, EngineError> {
let task_id = task_id.clone();
self.with_state("raw_system_prompt", move |s| {
s.systems.get(&(task_id, attempt)).cloned().unwrap_or(None)
})
.await
}
pub async fn fetch_data(&self, token: &CapToken, key: &str) -> Result<Value, EngineError> {
self.verify_token(token, Verb::FetchData).await?;
let key = key.to_string();
self.with_state("fetch_data", move |s| {
s.resources
.get(&key)
.cloned()
.ok_or(EngineError::ResourceNotFound(key))
})
.await?
}
pub async fn submit_output(
&self,
token: &crate::types::CapToken,
task_id: &StepId,
attempt: u32,
event: crate::worker::output::OutputEvent,
) -> Result<(), EngineError> {
self.verify_token_for_task(token, crate::types::Verb::EmitOutput, task_id)
.await?;
if let crate::worker::output::OutputEvent::Final { content, ok } = &event {
let comparable_value = content_ref_to_comparable_string(content.clone());
self.verdict_contract_completion_check(task_id, attempt, *ok, &comparable_value)
.await?;
}
let task_id_for_apply = task_id.clone();
let event_clone = event.clone();
self.with_state("submit_output", move |s| {
s.output_store
.entry((task_id_for_apply.clone(), attempt))
.or_default()
.push(event_clone.clone());
s.push_event(crate::core::state::Event::WorkerOutput {
task_id: task_id_for_apply,
attempt,
event: event_clone,
});
})
.await?;
match &event {
crate::worker::output::OutputEvent::Final { content, ok } => {
self.materialize_final_submission(task_id, attempt, content, *ok)
.await?;
}
crate::worker::output::OutputEvent::Artifact { name, content } => {
self.materialize_artifact_submission(task_id, attempt, name, content)
.await?;
}
_ => {}
}
Ok(())
}
async fn materialize_final_submission(
&self,
task_id: &StepId,
attempt: u32,
content: &crate::worker::output::ContentRef,
ok: bool,
) -> Result<(), EngineError> {
let server_policy = self.cfg().check_policy;
let task_id_for_lookup = task_id.clone();
let lookup = self
.with_state("materialize_final_submission.lookup", move |s| {
let entry = s.tasks.get(&task_id_for_lookup);
let producer_agent = entry.map(|t| t.spec.agent.clone());
let task_policy = entry.and_then(|t| t.spec.check_policy);
let view = s
.agent_ctx
.get(&(task_id_for_lookup.clone(), attempt))
.map(|e| e.view.clone());
(producer_agent, task_policy, view)
})
.await;
let policy = lookup
.as_ref()
.ok()
.and_then(|(_, tp, _)| *tp)
.unwrap_or(server_policy);
let (producer_agent, view) = match lookup.map(|(pa, _, view)| (pa, view)) {
Ok(pair) => pair,
Err(err) => {
if !matches!(policy, crate::core::config::CheckPolicy::Silent) {
tracing::warn!(
%task_id,
error = %err,
"submit-time projection sink: state lookup failed; skipping (fail-open)"
);
}
apply_check_policy(
policy,
"submit-time projection sink: state lookup",
"state lookup failed; skipping (fail-open)",
)?;
return Ok(());
}
};
let Some(producer_agent) = producer_agent else {
return Ok(());
};
let placement = self
.projection_placement_for(task_id)
.await
.unwrap_or_default();
let root = view.and_then(|v| placement.resolve_root(&v));
let canonical_agent = self
.step_naming_for(task_id)
.await
.and_then(|naming| {
naming
.canonical_of_producer(&producer_agent)
.map(str::to_string)
})
.unwrap_or_else(|| producer_agent.clone());
if let Some(store) = self.output_store_backend() {
if let Err(err) = store
.append(
task_id.as_str(),
attempt,
&canonical_agent,
crate::worker::output::OutputEvent::Final {
content: content.clone(),
ok,
},
Vec::new(),
)
.await
{
if !matches!(policy, crate::core::config::CheckPolicy::Silent) {
tracing::warn!(
%task_id,
agent = %producer_agent,
canonical = %canonical_agent,
error = %err,
"submit-time projection sink: OutputStore dual-write failed (fail-open)"
);
}
apply_check_policy(
policy,
"submit-time projection sink: OutputStore dual-write",
"OutputStore dual-write failed (fail-open)",
)?;
}
}
let Some(root) = root else {
if !matches!(policy, crate::core::config::CheckPolicy::Silent) {
tracing::warn!(
%task_id,
agent = %producer_agent,
canonical = %canonical_agent,
"submit-time projection sink: no work_dir/project_root resolved; skipping file materialize (fail-open)"
);
}
apply_check_policy(
policy,
"submit-time projection sink: file materialize",
"no work_dir/project_root resolved; skipping file materialize (fail-open)",
)?;
return Ok(());
};
let value = match content {
crate::worker::output::ContentRef::Inline { value } => value.clone(),
crate::worker::output::ContentRef::FileRef {
path,
mime,
size_hint,
} => serde_json::json!({
"file_ref": path.to_string_lossy(),
"mime": mime,
"size_hint": size_hint,
}),
};
let key = crate::core::projection::ProjectionKey {
task_id: task_id.to_string(),
run_id: None,
step: Some(canonical_agent.clone()),
path: None,
};
let adapter = crate::core::projection::FileProjectionAdapter::with_placement(
root,
(*placement).clone(),
);
if let Err(err) = adapter.materialize_submission(&key, &value, attempt, ok) {
if !matches!(policy, crate::core::config::CheckPolicy::Silent) {
tracing::warn!(
%task_id,
agent = %producer_agent,
canonical = %canonical_agent,
error = %err,
"submit-time projection sink: file materialize failed (fail-open)"
);
}
apply_check_policy(
policy,
"submit-time projection sink: file materialize",
"file materialize failed (fail-open)",
)?;
}
Ok(())
}
async fn materialize_artifact_submission(
&self,
task_id: &StepId,
attempt: u32,
name: &str,
content: &crate::worker::output::ContentRef,
) -> Result<(), EngineError> {
let server_policy = self.cfg().check_policy;
let task_id_for_lookup = task_id.clone();
let lookup = self
.with_state("materialize_artifact_submission.lookup", move |s| {
let task_policy = s
.tasks
.get(&task_id_for_lookup)
.and_then(|t| t.spec.check_policy);
let view = s
.agent_ctx
.get(&(task_id_for_lookup.clone(), attempt))
.map(|e| e.view.clone());
(task_policy, view)
})
.await
.ok();
let policy = lookup
.as_ref()
.and_then(|(tp, _)| *tp)
.unwrap_or(server_policy);
let view = lookup.and_then(|(_, view)| view);
if let Some(store) = self.output_store_backend() {
if let Err(err) = store
.append(
task_id.as_str(),
attempt,
name,
crate::worker::output::OutputEvent::Artifact {
name: name.to_string(),
content: content.clone(),
},
Vec::new(),
)
.await
{
if !matches!(policy, crate::core::config::CheckPolicy::Silent) {
tracing::warn!(
%task_id,
artifact = %name,
error = %err,
"submit-time projection sink: OutputStore dual-write failed for Artifact (fail-open)"
);
}
apply_check_policy(
policy,
"submit-time projection sink: Artifact OutputStore dual-write",
"OutputStore dual-write failed for Artifact (fail-open)",
)?;
}
}
let placement = self
.projection_placement_for(task_id)
.await
.unwrap_or_default();
let Some(root) = view.and_then(|v| placement.resolve_root(&v)) else {
if !matches!(policy, crate::core::config::CheckPolicy::Silent) {
tracing::warn!(
%task_id,
artifact = %name,
"submit-time projection sink: no work_dir/project_root resolved; skipping part file materialize (fail-open)"
);
}
apply_check_policy(
policy,
"submit-time projection sink: part file materialize",
"no work_dir/project_root resolved; skipping part file materialize (fail-open)",
)?;
return Ok(());
};
let value = match content {
crate::worker::output::ContentRef::Inline { value } => value.clone(),
crate::worker::output::ContentRef::FileRef {
path,
mime,
size_hint,
} => serde_json::json!({
"file_ref": path.to_string_lossy(),
"mime": mime,
"size_hint": size_hint,
}),
};
let adapter = crate::core::projection::FileProjectionAdapter::with_placement(
root,
(*placement).clone(),
);
if let Err(err) = adapter.materialize_part(task_id.as_str(), name, &value) {
if !matches!(policy, crate::core::config::CheckPolicy::Silent) {
tracing::warn!(
%task_id,
artifact = %name,
error = %err,
"submit-time projection sink: part file materialize failed (fail-open)"
);
}
apply_check_policy(
policy,
"submit-time projection sink: part file materialize",
"part file materialize failed (fail-open)",
)?;
}
Ok(())
}
pub async fn output_tail(
&self,
task_id: &StepId,
attempt: u32,
) -> Vec<crate::worker::output::OutputEvent> {
let key = (task_id.clone(), attempt);
self.with_state("output_tail", move |s| {
s.output_store.get(&key).cloned().unwrap_or_default()
})
.await
.unwrap_or_default()
}
pub async fn post_result(
&self,
token: &CapToken,
task_id: &StepId,
result: Value,
) -> Result<(), EngineError> {
self.verify_token_for_task(token, Verb::PostResult, task_id)
.await?;
let task_id = task_id.clone();
let result_clone = result.clone();
self.with_state("post_result", move |s| {
let task = s
.tasks
.get_mut(&task_id)
.ok_or_else(|| EngineError::TaskNotFound(task_id.to_string()))?;
task.last_result = Some(result_clone);
task.updated_at = now_unix();
Ok::<(), EngineError>(())
})
.await??;
Ok(())
}
pub async fn set_resource(
&self,
key: impl Into<String>,
value: Value,
) -> Result<(), EngineError> {
let key = key.into();
self.with_state("set_resource", move |s| {
s.resources.insert(key, value);
})
.await?;
Ok(())
}
pub async fn query_senior(
&self,
token: &CapToken,
task_id: &StepId,
question: Value,
) -> Result<ResumeKey, EngineError> {
self.verify_token(token, Verb::QuerySenior).await?;
let task_id = task_id.clone();
let key = ResumeKey::for_senior(&task_id);
let task_notify = self
.with_state("query_senior.notify_ensure", |s| {
s.ensure_task_notify(&task_id)
})
.await?;
let key_clone = key.clone();
let task_id_inner = task_id.clone();
let question_clone = question.clone();
self.with_state("query_senior.suspend", move |s| {
let task = s
.tasks
.get_mut(&task_id_inner)
.ok_or_else(|| EngineError::TaskNotFound(task_id_inner.to_string()))?;
task.status = TaskStatus::Suspended;
task.suspended_on = Some(key_clone.clone());
task.updated_at = now_unix();
s.pending_resumes
.insert(key_clone.clone(), ResumePending::new());
s.push_event(Event::SeniorQueried {
task_id: task_id_inner.clone(),
question: question_clone.clone(),
});
s.push_event(Event::TaskSuspended {
task_id: task_id_inner.clone(),
key: key_clone.clone(),
});
Ok::<(), EngineError>(())
})
.await??;
task_notify.notify_waiters();
let _ = self
.inner
.event_tx
.send(Event::SeniorQueried { task_id, question });
Ok(key)
}
pub async fn resume(&self, key: ResumeKey, answer: Value) -> Result<(), EngineError> {
let answer_for_state = answer.clone();
let answer_for_event = answer.clone();
let key_clone = key.clone();
let (notify, task_notify, task_id_opt) = self
.with_state("resume.set", move |s| {
let pending = s
.pending_resumes
.get_mut(&key_clone)
.ok_or(EngineError::ResumeKeyNotFound)?;
pending.answer = Some(answer_for_state);
let notify = pending.notify.clone();
let task_id = s
.tasks
.iter()
.find(|(_, t)| t.suspended_on.as_ref() == Some(&key_clone))
.map(|(id, _)| id.clone());
let task_notify = task_id.as_ref().map(|tid| s.ensure_task_notify(tid));
if let Some(tid) = &task_id {
if let Some(task) = s.tasks.get_mut(tid) {
task.suspended_on = None;
task.status = TaskStatus::Running;
task.updated_at = now_unix();
}
s.push_event(Event::TaskResumed {
task_id: tid.clone(),
key: key_clone.clone(),
});
s.push_event(Event::SeniorAnswered {
task_id: tid.clone(),
answer: answer_for_event.clone(),
});
}
Ok::<_, EngineError>((notify, task_notify, task_id))
})
.await??;
notify.notify_waiters();
if let Some(n) = task_notify {
n.notify_waiters();
}
if let Some(tid) = task_id_opt {
let _ = self
.inner
.event_tx
.send(Event::TaskResumed { task_id: tid, key });
}
Ok(())
}
pub async fn await_resume(
&self,
key: ResumeKey,
timeout: Duration,
) -> Result<Value, EngineError> {
let key_clone = key.clone();
let (notify, existing) = self
.with_state("await_resume.snapshot", move |s| {
let pending = s
.pending_resumes
.get(&key_clone)
.ok_or(EngineError::ResumeKeyNotFound)?;
Ok::<_, EngineError>((pending.notify.clone(), pending.answer.clone()))
})
.await??;
if let Some(v) = existing {
return Ok(v);
}
if timeout.is_zero() {
return Err(EngineError::PollTimeout);
}
let waited = tokio::time::timeout(timeout, notify.notified()).await;
if waited.is_err() {
return Err(EngineError::PollTimeout);
}
let key_clone = key.clone();
self.with_state("await_resume.read", move |s| {
let pending = s
.pending_resumes
.get(&key_clone)
.ok_or(EngineError::ResumeKeyNotFound)?;
pending
.answer
.clone()
.ok_or_else(|| EngineError::Internal("notified but answer missing".into()))
})
.await?
}
pub async fn poll_task(
&self,
token: &CapToken,
task_id: &StepId,
hold: Duration,
) -> Result<TaskState, EngineError> {
self.verify_token_for_task(token, Verb::PollTask, task_id)
.await?;
let task_id_inner = task_id.clone();
let (state, notify) = self
.with_state("poll_task.snapshot", move |s| {
let task = s
.tasks
.get(&task_id_inner)
.cloned()
.ok_or_else(|| EngineError::TaskNotFound(task_id_inner.to_string()))?;
let notify = s.ensure_task_notify(&task_id_inner);
Ok::<_, EngineError>((task, notify))
})
.await??;
if matches!(
state.status,
TaskStatus::Pass | TaskStatus::Blocked | TaskStatus::Cancelled | TaskStatus::Suspended
) {
return Ok(state);
}
if hold.is_zero() {
return Ok(state);
}
let waited = tokio::time::timeout(hold, notify.notified()).await;
if waited.is_err() {
return Err(EngineError::PollTimeout);
}
let task_id_inner = task_id.clone();
self.with_state("poll_task.reread", move |s| {
s.tasks
.get(&task_id_inner)
.cloned()
.ok_or_else(|| EngineError::TaskNotFound(task_id_inner.to_string()))
})
.await?
}
pub fn start_detach_loop(&self) -> tokio::task::JoinHandle<()> {
let engine = self.clone();
let cfg = self.inner.cfg.long_hold.clone();
let interval = cfg.heartbeat_interval;
let miss_secs = cfg.heartbeat_interval.as_secs() * cfg.heartbeat_miss_threshold as u64;
tokio::spawn(async move {
let mut ticker = tokio::time::interval(interval);
ticker.tick().await; loop {
ticker.tick().await;
let now = now_unix();
let detached = engine
.with_state("detach_loop.scan", |s| {
let mut detached = Vec::new();
for (sid, sess) in s.sessions.iter_mut() {
if !sess.attached {
continue;
}
if now.saturating_sub(sess.last_seen) >= miss_secs {
sess.attached = false;
detached.push(sid.clone());
}
}
for sid in &detached {
s.push_event(Event::SessionDetached {
session_id: sid.clone(),
});
}
detached
})
.await
.unwrap_or_default();
for sid in detached {
let _ = engine
.inner
.event_tx
.send(Event::SessionDetached { session_id: sid });
}
}
})
}
async fn wake_task(&self, task_id: &StepId) -> Result<(), EngineError> {
let task_id = task_id.clone();
let notify_opt = self
.with_state("wake_task.get_notify", move |s| {
s.task_notifies.get(&task_id).cloned()
})
.await?;
if let Some(n) = notify_opt {
n.notify_waiters();
}
Ok(())
}
}
pub(crate) fn apply_check_policy(
policy: crate::core::config::CheckPolicy,
context: &str,
message: &str,
) -> Result<(), EngineError> {
match policy {
crate::core::config::CheckPolicy::Silent | crate::core::config::CheckPolicy::Warn => Ok(()),
crate::core::config::CheckPolicy::Strict => Err(EngineError::CheckPolicyStrict {
context: context.to_string(),
message: message.to_string(),
}),
}
}
#[cfg(test)]
mod check_policy_helper_tests {
use super::apply_check_policy;
use crate::core::config::CheckPolicy;
use crate::core::errors::EngineError;
#[test]
fn silent_returns_ok() {
let result = apply_check_policy(CheckPolicy::Silent, "call/site", "sink message");
assert!(matches!(result, Ok(())));
}
#[test]
fn warn_returns_ok() {
let result = apply_check_policy(CheckPolicy::Warn, "call/site", "sink message");
assert!(matches!(result, Ok(())));
}
#[test]
fn strict_returns_error_with_context_and_message() {
let result = apply_check_policy(
CheckPolicy::Strict,
"submit-time projection sink: file materialize",
"no work_dir/project_root resolved; skipping file materialize (fail-open)",
);
match result {
Err(EngineError::CheckPolicyStrict { context, message }) => {
assert_eq!(context, "submit-time projection sink: file materialize");
assert_eq!(
message,
"no work_dir/project_root resolved; skipping file materialize (fail-open)"
);
}
other => panic!("expected CheckPolicyStrict, got {:?}", other),
}
}
}
#[cfg(test)]
mod token_fingerprint_store_tests {
use super::*;
#[tokio::test]
async fn verify_unknown_token_reports_fingerprint_not_nonce() {
let engine = Engine::new(EngineCfg::default());
let token = engine.signer().session(
"ghost",
Role::Operator,
vec!["*".into()],
Duration::from_secs(60),
);
let err = engine
.verify_token(&token, Verb::ReadTaskState)
.await
.expect_err("token is not in the store");
let msg = err.to_string();
assert!(
msg.contains(&token.fingerprint()),
"error must carry the fingerprint: {msg}"
);
assert!(
!msg.contains(&token.nonce),
"error must not leak the nonce: {msg}"
);
}
#[tokio::test]
async fn attach_verify_heartbeat_detach_cycle_with_fp_keying() {
let engine = Engine::new(EngineCfg::default());
let token = engine
.attach("op-1", Role::Operator, Duration::from_secs(60))
.await
.expect("attach");
engine
.verify_token(&token, Verb::ReadTaskState)
.await
.expect("verify consumes via fp key");
engine
.heartbeat(&token)
.await
.expect("heartbeat finds the session by fp");
engine
.detach(&token)
.await
.expect("detach finds the session by fp");
}
}
#[cfg(test)]
mod resolve_operator_info_runtime_global_tests {
use super::*;
async fn attach_and_resolve(
runtime_global: Option<OperatorKind>,
bp_global: Option<OperatorKind>,
) -> OperatorInfo {
let engine = Engine::new(EngineCfg::default());
let token = engine
.attach_with_ids(
"ut-op",
Role::Operator,
Duration::from_secs(30),
runtime_global,
None,
None,
None,
HashMap::new(),
HashMap::new(),
bp_global,
)
.await
.expect("attach_with_ids ok");
let session = engine
.with_state("test.find_session", |s| {
s.sessions
.values()
.find(|sess| sess.token_fp == token.fingerprint())
.cloned()
})
.await
.expect("with_state ok")
.expect("session present after attach_with_ids");
engine.resolve_operator_info(&session, "agent-x").await
}
#[tokio::test]
async fn explicit_some_automate_outranks_bp_global_main_ai() {
let info =
attach_and_resolve(Some(OperatorKind::Automate), Some(OperatorKind::MainAi)).await;
assert_eq!(
info.kind,
OperatorKind::Automate,
"explicit Some(Automate) runtime_global must outrank bp_global MainAi"
);
}
#[tokio::test]
async fn none_lets_bp_global_main_ai_win() {
let info = attach_and_resolve(None, Some(OperatorKind::MainAi)).await;
assert_eq!(
info.kind,
OperatorKind::MainAi,
"None runtime_global must let bp_global MainAi win"
);
}
}
#[cfg(test)]
mod dispatch_attempt_with_run_id_tests {
use super::*;
use crate::worker::adapter::{SpawnError, SpawnerAdapter};
use crate::worker::Worker;
use std::sync::Mutex as StdMutex;
struct CtxProbe {
seen: Arc<StdMutex<Option<Ctx>>>,
}
#[async_trait::async_trait]
impl SpawnerAdapter for CtxProbe {
async fn spawn(
&self,
_engine: &Engine,
ctx: &Ctx,
_task_id: StepId,
_attempt: u32,
_token: CapToken,
) -> Result<Box<dyn Worker>, SpawnError> {
*self.seen.lock().unwrap() = Some(ctx.clone());
Err(SpawnError::Internal("probe stop".into()))
}
}
async fn dispatch_with_probe(run_id: Option<&RunId>) -> Ctx {
let engine = Engine::new(EngineCfg::default());
let token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let tid = engine
.start_task(
&token,
TaskSpec {
agent: "probe".into(),
initial_directive: "hi".into(),
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
let seen: Arc<StdMutex<Option<Ctx>>> = Arc::new(StdMutex::new(None));
let spawner: Arc<dyn SpawnerAdapter> = Arc::new(CtxProbe { seen: seen.clone() });
let _ = engine
.dispatch_attempt_with(&token, &tid, &spawner, run_id)
.await;
let captured = seen.lock().unwrap().clone();
captured.expect("inner ctx captured")
}
#[tokio::test]
async fn run_id_lands_in_ctx_meta_runtime_when_some() {
let run_id = RunId::new();
let observed = dispatch_with_probe(Some(&run_id)).await;
assert_eq!(
observed.meta.runtime.get("run_id").and_then(|v| v.as_str()),
Some(run_id.as_str()),
"ctx.meta.runtime[\"run_id\"] must carry the run_id passed to dispatch_attempt_with"
);
}
#[tokio::test]
async fn run_id_key_absent_when_none() {
let observed = dispatch_with_probe(None).await;
assert!(
!observed.meta.runtime.contains_key("run_id"),
"no run_id key must be injected when dispatch_attempt_with is called with None"
);
}
}
#[cfg(test)]
mod dispatch_attempt_with_step_ctx_tests {
use super::*;
use crate::worker::adapter::{SpawnError, SpawnerAdapter};
use crate::worker::Worker;
use std::sync::Mutex as StdMutex;
struct CtxProbe {
seen: Arc<StdMutex<Option<Ctx>>>,
}
#[async_trait::async_trait]
impl SpawnerAdapter for CtxProbe {
async fn spawn(
&self,
_engine: &Engine,
ctx: &Ctx,
_task_id: StepId,
_attempt: u32,
_token: CapToken,
) -> Result<Box<dyn Worker>, SpawnError> {
*self.seen.lock().unwrap() = Some(ctx.clone());
Err(SpawnError::Internal("probe stop".into()))
}
}
#[tokio::test]
async fn step_ctx_lands_in_ctx_meta_runtime_on_attempt_1_and_2() {
let engine = Engine::new(EngineCfg::default());
let token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let tid = engine
.start_task(
&token,
TaskSpec {
agent: "probe".into(),
initial_directive: "hi".into(),
step_ctx: Some(serde_json::json!({ "work_dir": "/step" })),
check_policy: None,
},
)
.await
.expect("start_task");
let seen: Arc<StdMutex<Option<Ctx>>> = Arc::new(StdMutex::new(None));
let spawner: Arc<dyn SpawnerAdapter> = Arc::new(CtxProbe { seen: seen.clone() });
let _ = engine
.dispatch_attempt_with(&token, &tid, &spawner, None)
.await;
let first = seen
.lock()
.unwrap()
.clone()
.expect("attempt 1 ctx captured");
assert_eq!(
first.meta.runtime.get(STEP_CTX_KEY),
Some(&serde_json::json!({ "work_dir": "/step" })),
"attempt 1 must carry TaskSpec.step_ctx in ctx.meta.runtime[STEP_CTX_KEY]"
);
let _ = engine
.dispatch_attempt_with(&token, &tid, &spawner, None)
.await;
let second = seen
.lock()
.unwrap()
.clone()
.expect("attempt 2 ctx captured");
assert_eq!(
second.meta.runtime.get(STEP_CTX_KEY),
Some(&serde_json::json!({ "work_dir": "/step" })),
"attempt 2 (retry) must ALSO carry TaskSpec.step_ctx — prep re-reads the spec every attempt"
);
}
#[tokio::test]
async fn step_ctx_key_absent_when_none() {
let engine = Engine::new(EngineCfg::default());
let token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let tid = engine
.start_task(
&token,
TaskSpec {
agent: "probe".into(),
initial_directive: "hi".into(),
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
let seen: Arc<StdMutex<Option<Ctx>>> = Arc::new(StdMutex::new(None));
let spawner: Arc<dyn SpawnerAdapter> = Arc::new(CtxProbe { seen: seen.clone() });
let _ = engine
.dispatch_attempt_with(&token, &tid, &spawner, None)
.await;
let observed = seen.lock().unwrap().clone().expect("ctx captured");
assert!(
!observed.meta.runtime.contains_key(STEP_CTX_KEY),
"no step_ctx key must be injected when TaskSpec.step_ctx is None"
);
}
}
#[cfg(test)]
mod initial_directive_value_passthrough_tests {
use super::*;
async fn seeded_engine(initial_directive: Value) -> (Engine, CapToken, StepId) {
let engine = Engine::new(EngineCfg::default());
let op_token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let task_id = engine
.start_task(
&op_token,
TaskSpec {
agent: "planner".to_string(),
initial_directive,
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
(engine, op_token, task_id)
}
async fn mint_worker_token(engine: &Engine, task_id: &StepId) -> CapToken {
let worker_token = engine.signer().session(
format!("worker-of-{task_id}"),
Role::Worker,
vec!["*".into()],
Duration::from_secs(600),
);
let fp = worker_token.fingerprint();
let record = CapTokenRecord::from_worker_token(worker_token.clone(), task_id.clone());
engine
.with_state("test.mint_worker", move |s| {
s.tokens.insert(fp, record);
})
.await
.expect("mint worker token");
worker_token
}
#[tokio::test]
async fn object_seed_passes_through_task_spec_unchanged() {
let seed = serde_json::json!({"key": "value"});
let (engine, token, task_id) = seeded_engine(seed.clone()).await;
let state = engine
.read_task_state(&token, &task_id)
.await
.expect("read_task_state");
assert_eq!(
state.spec.initial_directive, seed,
"TaskSpec.initial_directive must equal the raw Object seed, not a stringified copy"
);
}
#[tokio::test]
async fn object_seed_passes_through_fetch_prompt_as_value() {
let seed = serde_json::json!({"key": "value"});
let (engine, _token, task_id) = seeded_engine(seed.clone()).await;
let worker_token = mint_worker_token(&engine, &task_id).await;
let prompt = engine
.fetch_prompt(&worker_token, &task_id)
.await
.expect("fetch_prompt");
assert_eq!(
prompt, seed,
"fetch_prompt must return the raw Object Value, not a stringified copy"
);
}
#[tokio::test]
async fn object_seed_renders_as_json_literal_at_worker_payload_boundary() {
let seed = serde_json::json!({"key": "value"});
let (engine, _token, task_id) = seeded_engine(seed).await;
let worker_token = mint_worker_token(&engine, &task_id).await;
let payload = engine
.fetch_worker_payload(&worker_token, &task_id)
.await
.expect("fetch_worker_payload");
assert_eq!(
payload.prompt, r#"{"key":"value"}"#,
"WorkerPayload.prompt must be the JSON literal String render of the Value seed"
);
}
#[tokio::test]
async fn string_seed_passes_through_unchanged() {
let (engine, token, task_id) = seeded_engine(serde_json::json!("do the thing")).await;
let state = engine
.read_task_state(&token, &task_id)
.await
.expect("read_task_state");
assert_eq!(
state.spec.initial_directive,
serde_json::json!("do the thing")
);
let worker_token = mint_worker_token(&engine, &task_id).await;
let prompt = engine
.fetch_prompt(&worker_token, &task_id)
.await
.expect("fetch_prompt");
assert_eq!(prompt, serde_json::json!("do the thing"));
}
}
#[cfg(test)]
mod system_ref_threshold_tests {
use super::*;
async fn seeded_engine_with_cfg(cfg: EngineCfg) -> (Engine, CapToken, StepId) {
let engine = Engine::new(cfg);
let op_token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let task_id = engine
.start_task(
&op_token,
TaskSpec {
agent: "planner".to_string(),
initial_directive: serde_json::json!("do the thing"),
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
(engine, op_token, task_id)
}
async fn mint_worker_token(engine: &Engine, task_id: &StepId) -> CapToken {
let worker_token = engine.signer().session(
format!("worker-of-{task_id}"),
Role::Worker,
vec!["*".into()],
Duration::from_secs(600),
);
let fp = worker_token.fingerprint();
let record = CapTokenRecord::from_worker_token(worker_token.clone(), task_id.clone());
engine
.with_state("test.mint_worker", move |s| {
s.tokens.insert(fp, record);
})
.await
.expect("mint worker token");
worker_token
}
#[tokio::test]
async fn under_threshold_stays_inline() {
let (engine, _op_token, task_id) = seeded_engine_with_cfg(EngineCfg::default()).await;
let worker_token = mint_worker_token(&engine, &task_id).await;
let rendered = "a short system prompt".to_string();
engine
.bake_worker_system_prompt(&task_id, 1, Some(rendered.clone()))
.await
.expect("bake");
let payload = engine
.fetch_worker_payload(&worker_token, &task_id)
.await
.expect("fetch_worker_payload");
assert_eq!(payload.system, Some(rendered));
assert!(payload.system_ref.is_none());
}
#[tokio::test]
async fn over_threshold_switches_to_system_ref_with_matching_sha256() {
let mut cfg = EngineCfg::default();
cfg.system_ref.threshold_bytes = 16;
cfg.system_ref.mode = crate::types::SystemRefMode::File;
cfg.system_ref.store_dir =
std::env::temp_dir().join(format!("mse-system-ref-test-{}", crate::types::now_unix()));
let (engine, _op_token, task_id) = seeded_engine_with_cfg(cfg).await;
let rendered =
"this system prompt is deliberately longer than the 16 byte threshold".to_string();
engine
.bake_worker_system_prompt(&task_id, 1, Some(rendered.clone()))
.await
.expect("bake");
let payload = engine
.fetch_worker_payload_trusted(&task_id)
.await
.expect("fetch_worker_payload_trusted");
assert!(
payload.system.is_none(),
"over-threshold response must not also inline `system`"
);
let system_ref = payload
.system_ref
.expect("over-threshold response must populate system_ref");
assert_eq!(system_ref.size_bytes, rendered.len() as u64);
assert_eq!(system_ref.mode, crate::types::SystemRefMode::File);
use sha2::Digest;
let expected_sha256 = hex::encode(sha2::Sha256::digest(rendered.as_bytes()));
assert_eq!(system_ref.sha256, expected_sha256);
assert!(system_ref.uri.starts_with("file://"));
let written = tokio::fs::read_to_string(system_ref.uri.trim_start_matches("file://"))
.await
.expect("File mode must have written the referenced path");
assert_eq!(written, rendered);
}
#[tokio::test]
async fn over_threshold_http_mode_constructs_path_only_uri() {
let mut cfg = EngineCfg::default();
cfg.system_ref.threshold_bytes = 16;
cfg.system_ref.mode = crate::types::SystemRefMode::Http;
let (engine, _op_token, task_id) = seeded_engine_with_cfg(cfg).await;
let worker_token = mint_worker_token(&engine, &task_id).await;
let rendered =
"this system prompt is deliberately longer than the 16 byte threshold".to_string();
engine
.bake_worker_system_prompt(&task_id, 1, Some(rendered))
.await
.expect("bake");
let payload = engine
.fetch_worker_payload(&worker_token, &task_id)
.await
.expect("fetch_worker_payload");
let system_ref = payload.system_ref.expect("system_ref must be populated");
assert_eq!(system_ref.mode, crate::types::SystemRefMode::Http);
assert_eq!(
system_ref.uri,
format!("/v1/worker/prompt/system?task_id={task_id}&attempt=1")
);
}
#[tokio::test]
async fn bake_records_agent_render_size_last_write_wins() {
let (engine, _op_token, task_id) = seeded_engine_with_cfg(EngineCfg::default()).await;
assert_eq!(engine.agent_last_rendered_size("planner").await, None);
engine
.bake_worker_system_prompt(&task_id, 1, Some("a".repeat(10)))
.await
.expect("bake 1");
assert_eq!(engine.agent_last_rendered_size("planner").await, Some(10));
engine
.bake_worker_system_prompt(&task_id, 2, Some("b".repeat(20)))
.await
.expect("bake 2");
assert_eq!(
engine.agent_last_rendered_size("planner").await,
Some(20),
"most-recently-observed size wins, not the largest"
);
}
}
#[cfg(test)]
mod submit_time_projection_sink_tests {
use super::*;
use crate::core::agent_context::AgentContextView;
use crate::store::output::{ContentRef, InMemoryOutputStore, OutputEvent};
async fn seeded_task(agent: &str) -> (Engine, CapToken, StepId, CapToken) {
let engine = Engine::new(EngineCfg::default());
let op_token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let task_id = engine
.start_task(
&op_token,
TaskSpec {
agent: agent.to_string(),
initial_directive: Value::String("go".into()),
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
let worker_token = engine.signer().session(
format!("worker-of-{task_id}"),
Role::Worker,
vec!["*".into()],
Duration::from_secs(600),
);
let fp = worker_token.fingerprint();
let record = CapTokenRecord::from_worker_token(worker_token.clone(), task_id.clone());
engine
.with_state("test.mint_worker", move |s| {
s.tokens.insert(fp, record);
})
.await
.expect("mint worker token");
(engine, op_token, task_id, worker_token)
}
async fn seeded_task_with_policy(
agent: &str,
policy: crate::core::config::CheckPolicy,
) -> (Engine, CapToken, StepId, CapToken) {
let cfg = EngineCfg {
check_policy: policy,
..EngineCfg::default()
};
let engine = Engine::new(cfg);
let op_token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let task_id = engine
.start_task(
&op_token,
TaskSpec {
agent: agent.to_string(),
initial_directive: Value::String("go".into()),
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
let worker_token = engine.signer().session(
format!("worker-of-{task_id}"),
Role::Worker,
vec!["*".into()],
Duration::from_secs(600),
);
let fp = worker_token.fingerprint();
let record = CapTokenRecord::from_worker_token(worker_token.clone(), task_id.clone());
engine
.with_state("test.mint_worker", move |s| {
s.tokens.insert(fp, record);
})
.await
.expect("mint worker token");
(engine, op_token, task_id, worker_token)
}
async fn seed_agent_context(engine: &Engine, task_id: &StepId, attempt: u32, work_dir: &str) {
let task_id = task_id.clone();
let work_dir = work_dir.to_string();
engine
.with_state("test.seed_agent_context", move |s| {
s.agent_ctx.insert(
(task_id, attempt),
crate::core::state::AgentCtxEntry {
view: AgentContextView {
work_dir: Some(work_dir),
..Default::default()
},
policy: Default::default(),
},
);
})
.await
.expect("seed agent_ctx");
}
async fn seed_agent_context_roots(
engine: &Engine,
task_id: &StepId,
attempt: u32,
work_dir: Option<&str>,
project_root: Option<&str>,
) {
let task_id = task_id.clone();
let work_dir = work_dir.map(str::to_string);
let project_root = project_root.map(str::to_string);
engine
.with_state("test.seed_agent_context_roots", move |s| {
s.agent_ctx.insert(
(task_id, attempt),
crate::core::state::AgentCtxEntry {
view: AgentContextView {
work_dir,
project_root,
..Default::default()
},
policy: Default::default(),
},
);
})
.await
.expect("seed agent_ctx");
}
async fn seed_projection_placement(
engine: &Engine,
task_id: &StepId,
placement: crate::core::projection_placement::ProjectionPlacement,
) {
let task_id = task_id.clone();
let placement = Arc::new(placement);
engine
.with_state("test.seed_projection_placement", move |s| {
s.projection_placements.insert(task_id, placement);
})
.await
.expect("seed projection_placements");
}
async fn seed_step_naming(engine: &Engine, task_id: &StepId, producer: &str, canonical: &str) {
use crate::blueprint::{
current_schema_version, AgentDef, AgentKind, AgentMeta, Blueprint, BlueprintMetadata,
CompilerHints, CompilerStrategy,
};
use crate::core::step_naming::StepNaming;
use mlua_flow_ir::{Expr, Node};
let flow = Node::Step {
ref_: producer.to_string(),
in_: Expr::Path {
at: "$.in".parse().expect("literal test path: $.in"),
},
out: Expr::Path {
at: format!("$.{producer}_out")
.parse()
.expect("literal test path"),
},
};
let bp = Blueprint {
schema_version: current_schema_version(),
id: "sink-canonical-ut".into(),
flow,
agents: vec![AgentDef {
name: producer.to_string(),
kind: AgentKind::RustFn,
spec: serde_json::json!({ "fn_id": producer }),
profile: None,
meta: Some(AgentMeta {
projection_name: Some(canonical.to_string()),
..Default::default()
}),
runner: None,
runner_ref: None,
verdict: None,
}],
operators: vec![],
metas: vec![],
hints: CompilerHints::default(),
strategy: CompilerStrategy::default(),
metadata: BlueprintMetadata::default(),
spawner_hints: Default::default(),
default_agent_kind: AgentKind::Operator,
default_operator_kind: None,
default_init_ctx: None,
default_agent_ctx: None,
default_context_policy: None,
projection_placement: None,
audits: vec![],
degradation_policy: None,
runners: vec![],
default_runner: None,
check_policy: None,
};
let (naming, warnings) = StepNaming::from_blueprint(&bp).expect("no collision");
assert!(warnings.is_empty(), "single-step fixture has no collisions");
let naming = Arc::new(naming);
let task_id = task_id.clone();
engine
.with_state("test.seed_step_naming", move |s| {
s.step_namings.insert(task_id, naming);
})
.await
.expect("seed step_namings");
}
fn final_event(value: Value, ok: bool) -> crate::worker::output::OutputEvent {
crate::worker::output::OutputEvent::Final {
content: crate::worker::output::ContentRef::Inline { value },
ok,
}
}
#[tokio::test]
async fn submit_output_final_materializes_file_when_work_dir_resolved() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, worker_token) = seeded_task("planner").await;
seed_agent_context(&engine, &task_id, 1, &dir.path().to_string_lossy()).await;
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!({"plan": "do it"}), true),
)
.await
.expect("submit_output");
let expected_file = dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("ctx/planner.md");
assert!(
expected_file.exists(),
"materialized submission file missing at {expected_file:?}"
);
let body = std::fs::read_to_string(expected_file).unwrap();
assert!(body.contains(r#""plan": "do it""#), "body: {body}");
}
#[tokio::test]
async fn submit_output_final_skips_file_when_root_unresolved() {
let (engine, _op, task_id, worker_token) = seeded_task("planner").await;
let result = engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!("hi"), true),
)
.await;
assert!(
result.is_ok(),
"submit must succeed even with no resolvable root (fail-open, Invariant 1)"
);
}
#[tokio::test]
async fn submit_output_final_check_policy_warn_preserves_fail_open() {
let (engine, _op, task_id, worker_token) =
seeded_task_with_policy("planner", crate::core::config::CheckPolicy::Warn).await;
let result = engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!("hi"), true),
)
.await;
assert!(
result.is_ok(),
"Warn mode preserves fail-open: submit must succeed when root unresolved"
);
}
#[tokio::test]
async fn submit_output_final_check_policy_strict_surfaces_error_when_root_unresolved() {
let (engine, _op, task_id, worker_token) =
seeded_task_with_policy("planner", crate::core::config::CheckPolicy::Strict).await;
let err = engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!("hi"), true),
)
.await
.expect_err("Strict mode must return an error when root unresolved");
match err {
EngineError::CheckPolicyStrict { context, message } => {
assert!(
context.contains("file materialize"),
"context must identify the call site: {context}"
);
assert!(
message.contains("no work_dir/project_root resolved"),
"message must preserve the warn-log literal for log-parse compat: {message}"
);
}
other => panic!(
"expected EngineError::CheckPolicyStrict, got a different variant: {other:?}"
),
}
}
#[tokio::test]
async fn submit_output_final_check_policy_silent_returns_ok_when_root_unresolved() {
let (engine, _op, task_id, worker_token) =
seeded_task_with_policy("planner", crate::core::config::CheckPolicy::Silent).await;
let result = engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!("hi"), true),
)
.await;
assert!(
result.is_ok(),
"Silent mode returns Ok(()) at the error surface: submit must succeed"
);
}
#[tokio::test]
async fn resubmit_overwrites_materialized_file_with_latest() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, worker_token) = seeded_task("planner").await;
seed_agent_context(&engine, &task_id, 1, &dir.path().to_string_lossy()).await;
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!("first"), true),
)
.await
.expect("first submit");
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!("second"), true),
)
.await
.expect("second submit");
let expected_file = dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("ctx/planner.md");
let body = std::fs::read_to_string(expected_file).unwrap();
assert!(body.contains("second"), "body must reflect latest: {body}");
assert!(
!body.contains("first"),
"body must not carry the stale value: {body}"
);
}
#[tokio::test]
async fn submit_output_final_falls_back_to_project_root_when_work_dir_absent() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, worker_token) = seeded_task("planner").await;
seed_agent_context_roots(
&engine,
&task_id,
1,
None,
Some(&dir.path().to_string_lossy()),
)
.await;
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!({"plan": "via project_root"}), true),
)
.await
.expect("submit_output");
let expected_file = dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("ctx/planner.md");
assert!(
expected_file.exists(),
"materialized submission file missing at {expected_file:?} \
(work_dir absent must fall back to project_root)"
);
}
#[tokio::test]
async fn submit_output_final_uses_declared_projection_placement() {
let work_dir = tempfile::TempDir::new().unwrap();
let project_root = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, worker_token) = seeded_task("planner").await;
seed_agent_context_roots(
&engine,
&task_id,
1,
Some(&work_dir.path().to_string_lossy()),
Some(&project_root.path().to_string_lossy()),
)
.await;
seed_projection_placement(
&engine,
&task_id,
crate::core::projection_placement::ProjectionPlacement {
root_preference: crate::core::projection_placement::RootPreference::ProjectRoot,
dir_template: "custom/{task_id}/out".to_string(),
},
)
.await;
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!({"plan": "via custom placement"}), true),
)
.await
.expect("submit_output");
let expected_file = project_root
.path()
.join("custom")
.join(task_id.as_str())
.join("out/planner.md");
assert!(
expected_file.exists(),
"materialized submission file missing at custom placement target {expected_file:?}"
);
let unexpected_file = work_dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("ctx/planner.md");
assert!(
!unexpected_file.exists(),
"declared root_preference=ProjectRoot must not fall back to work_dir: {unexpected_file:?}"
);
}
#[tokio::test]
async fn submit_output_final_dual_writes_into_configured_output_store() {
let (engine, _op, task_id, worker_token) = seeded_task("reviewer").await;
let data_store: Arc<dyn crate::store::output::OutputStore> =
Arc::new(InMemoryOutputStore::new());
engine.set_output_store(data_store.clone());
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!({"verdict": "pass"}), true),
)
.await
.expect("submit_output");
let record = data_store
.get_latest_by_name("reviewer")
.await
.expect("dual-written record");
match record.event {
OutputEvent::Final { content, ok } => {
assert!(ok);
match content {
ContentRef::Inline { value } => {
assert_eq!(value, serde_json::json!({"verdict": "pass"}));
}
other => panic!("expected Inline content, got {other:?}"),
}
}
other => panic!("expected Final event, got {other:?}"),
}
}
#[tokio::test]
async fn submit_output_artifact_dual_writes_into_configured_output_store() {
let (engine, _op, task_id, worker_token) = seeded_task("echo").await;
let data_store: Arc<dyn crate::store::output::OutputStore> =
Arc::new(InMemoryOutputStore::new());
engine.set_output_store(data_store.clone());
engine
.submit_output(
&worker_token,
&task_id,
1,
OutputEvent::Artifact {
name: "audit:echo".to_string(),
content: ContentRef::Inline {
value: serde_json::json!({"finding": "clean"}),
},
},
)
.await
.expect("submit_output");
let record = data_store
.get_latest_by_name("audit:echo")
.await
.expect("dual-written artifact record");
match record.event {
OutputEvent::Artifact { name, content } => {
assert_eq!(name, "audit:echo");
match content {
ContentRef::Inline { value } => {
assert_eq!(value, serde_json::json!({"finding": "clean"}));
}
other => panic!("expected Inline content, got {other:?}"),
}
}
other => panic!("expected Artifact event, got {other:?}"),
}
assert!(
data_store.get_latest_by_name("echo").await.is_err(),
"artifact write must not fabricate a record under the raw producer_agent name"
);
}
#[tokio::test]
async fn submit_output_artifact_is_fail_open_when_no_output_store_configured() {
let (engine, _op, task_id, worker_token) = seeded_task("echo").await;
let result = engine
.submit_output(
&worker_token,
&task_id,
1,
OutputEvent::Artifact {
name: "audit:echo".to_string(),
content: ContentRef::Inline {
value: serde_json::json!("finding"),
},
},
)
.await;
assert!(
result.is_ok(),
"submit must succeed even with no OutputStore wired (fail-open, Invariant 1)"
);
}
#[tokio::test]
async fn submit_worker_result_trusted_also_triggers_projection_sink() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, _worker_token) = seeded_task("planner").await;
seed_agent_context(&engine, &task_id, 1, &dir.path().to_string_lossy()).await;
let data_store: Arc<dyn crate::store::output::OutputStore> =
Arc::new(InMemoryOutputStore::new());
engine.set_output_store(data_store.clone());
engine
.submit_worker_result_trusted(&task_id, 1, serde_json::json!("trusted-value"), true)
.await
.expect("submit_worker_result_trusted");
let expected_file = dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("ctx/planner.md");
assert!(expected_file.exists());
let record = data_store
.get_latest_by_name("planner")
.await
.expect("dual-written record");
assert!(matches!(record.event, OutputEvent::Final { ok: true, .. }));
}
#[tokio::test]
async fn submit_output_final_uses_canonical_name_when_step_naming_declares_one() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, worker_token) = seeded_task("reviewer").await;
seed_agent_context(&engine, &task_id, 1, &dir.path().to_string_lossy()).await;
seed_step_naming(&engine, &task_id, "reviewer", "verdict-final").await;
let data_store: Arc<dyn crate::store::output::OutputStore> =
Arc::new(InMemoryOutputStore::new());
engine.set_output_store(data_store.clone());
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!({"verdict": "pass"}), true),
)
.await
.expect("submit_output");
let record = data_store
.get_latest_by_name("verdict-final")
.await
.expect("dual-written record under canonical name");
assert!(matches!(record.event, OutputEvent::Final { ok: true, .. }));
assert!(
data_store.get_latest_by_name("reviewer").await.is_err(),
"raw producer_agent name must not be written once canonical resolves"
);
let expected_file = dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("ctx/verdict-final.md");
assert!(
expected_file.exists(),
"materialized file stem must be canonical at {expected_file:?}"
);
}
#[tokio::test]
async fn submit_output_final_falls_back_to_producer_agent_when_no_step_naming_table() {
let (engine, _op, task_id, worker_token) = seeded_task("reviewer").await;
let data_store: Arc<dyn crate::store::output::OutputStore> =
Arc::new(InMemoryOutputStore::new());
engine.set_output_store(data_store.clone());
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!({"verdict": "pass"}), true),
)
.await
.expect("submit_output");
let record = data_store
.get_latest_by_name("reviewer")
.await
.expect("fail-open dual-write under raw producer_agent name");
assert!(matches!(record.event, OutputEvent::Final { ok: true, .. }));
}
#[tokio::test]
async fn submit_output_final_is_resolvable_via_run_scoped_lookup() {
let (engine, _op, task_id, worker_token) = seeded_task("reviewer").await;
let data_store: Arc<dyn crate::store::output::OutputStore> =
Arc::new(InMemoryOutputStore::new());
engine.set_output_store(data_store.clone());
engine
.submit_output(
&worker_token,
&task_id,
1,
final_event(serde_json::json!({"verdict": "pass"}), true),
)
.await
.expect("submit_output");
let record = data_store
.get_latest_by_name_in_run(task_id.as_str(), 1, "reviewer")
.await
.expect("run-scoped lookup resolves the dual-written record");
assert!(matches!(record.event, OutputEvent::Final { ok: true, .. }));
assert!(
data_store
.get_latest_by_name_in_run(task_id.as_str(), 2, "reviewer")
.await
.is_err(),
"a different attempt must not resolve the same-named record"
);
}
#[tokio::test]
async fn stage_artifact_materializes_part_file_when_work_dir_resolved() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, _worker_token) = seeded_task("planner").await;
seed_agent_context(&engine, &task_id, 1, &dir.path().to_string_lossy()).await;
engine
.stage_worker_artifact_trusted(
&task_id,
1,
"plan.md".to_string(),
serde_json::json!("# Plan\n\nstep one\n"),
)
.await
.expect("stage artifact");
let expected_file = dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("ctx/plan.md");
assert!(
expected_file.exists(),
"materialized part file missing at {expected_file:?}"
);
let body = std::fs::read_to_string(expected_file).unwrap();
assert_eq!(body, "# Plan\n\nstep one\n");
}
#[tokio::test]
async fn stage_artifact_check_policy_warn_skips_part_file_when_root_unresolved() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, _worker_token) =
seeded_task_with_policy("planner", crate::core::config::CheckPolicy::Warn).await;
let result = engine
.stage_worker_artifact_trusted(
&task_id,
1,
"plan.md".to_string(),
serde_json::json!("x"),
)
.await;
assert!(
result.is_ok(),
"Warn mode preserves fail-open: stage must succeed when root unresolved"
);
assert!(
!dir.path().join("workspace").exists(),
"no part file may be materialized when root is unresolved"
);
}
#[tokio::test]
async fn stage_artifact_check_policy_strict_surfaces_error_when_root_unresolved() {
let (engine, _op, task_id, _worker_token) =
seeded_task_with_policy("planner", crate::core::config::CheckPolicy::Strict).await;
let err = engine
.stage_worker_artifact_trusted(
&task_id,
1,
"plan.md".to_string(),
serde_json::json!("x"),
)
.await
.expect_err("Strict mode must return an error when root unresolved");
match err {
EngineError::CheckPolicyStrict { context, message } => {
assert!(
context.contains("part file materialize"),
"context must identify the call site: {context}"
);
assert!(
message.contains("part file materialize"),
"message must identify the part-file sink: {message}"
);
assert!(
message.contains("no work_dir/project_root resolved"),
"message must preserve the warn-log literal: {message}"
);
}
other => panic!(
"expected EngineError::CheckPolicyStrict, got a different variant: {other:?}"
),
}
}
#[tokio::test]
async fn stage_artifact_traversal_name_is_fail_open_and_writes_nothing() {
let dir = tempfile::TempDir::new().unwrap();
let (engine, _op, task_id, _worker_token) = seeded_task("planner").await;
seed_agent_context(&engine, &task_id, 1, &dir.path().to_string_lossy()).await;
let result = engine
.stage_worker_artifact_trusted(
&task_id,
1,
"../evil.md".to_string(),
serde_json::json!("pwned"),
)
.await;
assert!(
result.is_ok(),
"default (Warn) policy is fail-open even on a rejected part name"
);
let escaped = dir
.path()
.join("workspace/tasks")
.join(task_id.as_str())
.join("evil.md");
assert!(
!escaped.exists(),
"a traversal name must never write outside the ctx dir: {escaped:?}"
);
}
}
#[cfg(test)]
mod named_multi_part_worker_output_tests {
use super::*;
use crate::worker::output::{ContentRef, OutputEvent};
fn artifact(name: &str, value: Value) -> OutputEvent {
OutputEvent::Artifact {
name: name.to_string(),
content: ContentRef::Inline { value },
}
}
fn final_ev(value: Value, ok: bool) -> OutputEvent {
OutputEvent::Final {
content: ContentRef::Inline { value },
ok,
}
}
fn names(list: &[&str]) -> Vec<String> {
list.iter().map(|s| s.to_string()).collect()
}
#[test]
fn fold_final_and_parts_assembles_out_and_parts_shape() {
let tail = vec![
artifact("summary", serde_json::json!("the summary")),
artifact("diff", serde_json::json!({"lines": 3})),
final_ev(serde_json::json!("final text"), true),
];
let staged = names(&["summary", "diff"]);
let (value, ok) = fold_final_and_parts(&tail, &staged).expect("Final present");
assert!(ok);
assert_eq!(
value,
serde_json::json!({
"out": "final text",
"parts": {
"summary": "the summary",
"diff": {"lines": 3},
}
})
);
}
#[test]
fn fold_final_and_parts_with_no_parts_returns_plain_final_value() {
let tail = vec![final_ev(serde_json::json!("plain value"), true)];
let (value, ok) = fold_final_and_parts(&tail, &[]).expect("Final present");
assert!(ok);
assert_eq!(value, serde_json::json!("plain value"));
}
#[test]
fn fold_final_and_parts_same_name_twice_last_write_wins() {
let tail = vec![
artifact("a", serde_json::json!("first")),
artifact("a", serde_json::json!("second")),
final_ev(serde_json::json!("f"), true),
];
let staged = names(&["a"]);
let (value, _ok) = fold_final_and_parts(&tail, &staged).expect("Final present");
assert_eq!(
value,
serde_json::json!({"out": "f", "parts": {"a": "second"}})
);
}
#[test]
fn fold_final_and_parts_returns_none_when_no_final_present() {
let tail = vec![artifact("a", serde_json::json!("v"))];
let staged = names(&["a"]);
assert!(fold_final_and_parts(&tail, &staged).is_none());
}
#[test]
fn fold_final_and_parts_ignores_artifacts_outside_the_staged_allowlist() {
let tail = vec![
final_ev(serde_json::json!({"echoed": "hi"}), true),
artifact("audit:echo", serde_json::json!({"finding": "clean"})),
];
let (value, ok) = fold_final_and_parts(&tail, &[]).expect("Final present");
assert!(ok);
assert_eq!(value, serde_json::json!({"echoed": "hi"}));
}
#[test]
fn fold_final_and_parts_folds_only_the_staged_subset_of_a_mixed_tail() {
let tail = vec![
artifact("summary", serde_json::json!("s")),
artifact("audit:echo", serde_json::json!({"finding": "clean"})),
final_ev(serde_json::json!("f"), true),
];
let staged = names(&["summary"]);
let (value, _ok) = fold_final_and_parts(&tail, &staged).expect("Final present");
assert_eq!(
value,
serde_json::json!({"out": "f", "parts": {"summary": "s"}})
);
}
#[tokio::test]
async fn stage_worker_artifact_trusted_is_isolated_per_attempt() {
let engine = Engine::new(EngineCfg::default());
let task_id = StepId::new();
engine
.stage_worker_artifact_trusted(&task_id, 1, "a".to_string(), serde_json::json!("v1"))
.await
.expect("stage attempt 1");
let attempt_1_tail = engine.output_tail(&task_id, 1).await;
assert_eq!(attempt_1_tail.len(), 1);
assert!(matches!(
&attempt_1_tail[0],
OutputEvent::Artifact { name, .. } if name == "a"
));
assert_eq!(
engine.worker_artifact_names_for(&task_id, 1).await,
vec!["a".to_string()]
);
let attempt_2_tail = engine.output_tail(&task_id, 2).await;
assert!(
attempt_2_tail.is_empty(),
"attempt 2 must not see attempt 1's staged part"
);
assert!(
engine
.worker_artifact_names_for(&task_id, 2)
.await
.is_empty(),
"attempt 2's allowlist must not see attempt 1's staged name"
);
}
}
#[cfg(test)]
mod verdict_contract_registry_tests {
use super::*;
async fn seeded_engine(agent: &str) -> (Engine, StepId) {
let engine = Engine::new(EngineCfg::default());
let op_token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let task_id = engine
.start_task(
&op_token,
TaskSpec {
agent: agent.to_string(),
initial_directive: serde_json::json!("x"),
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
(engine, task_id)
}
#[tokio::test]
async fn returns_none_when_no_contract_registered_for_the_agent() {
let (engine, task_id) = seeded_engine("gate").await;
assert_eq!(engine.verdict_contract_for_task(&task_id).await, None);
}
#[tokio::test]
async fn returns_the_registered_contract_for_the_running_agent() {
let (engine, task_id) = seeded_engine("gate").await;
let contract = mlua_swarm_schema::VerdictContract {
channel: mlua_swarm_schema::VerdictChannel::Body,
values: vec!["PASS".to_string(), "BLOCKED".to_string()],
};
engine.register_verdict_contracts(HashMap::from([("gate".to_string(), contract.clone())]));
assert_eq!(
engine.verdict_contract_for_task(&task_id).await,
Some(contract)
);
}
#[tokio::test]
async fn does_not_leak_a_contract_registered_for_a_different_agent() {
let (engine, task_id) = seeded_engine("gate").await;
engine.register_verdict_contracts(HashMap::from([(
"other-agent".to_string(),
mlua_swarm_schema::VerdictContract {
channel: mlua_swarm_schema::VerdictChannel::Body,
values: vec!["PASS".to_string()],
},
)]));
assert_eq!(engine.verdict_contract_for_task(&task_id).await, None);
}
#[tokio::test]
async fn returns_none_for_an_unknown_task_id() {
let engine = Engine::new(EngineCfg::default());
let unknown = StepId::new();
assert_eq!(engine.verdict_contract_for_task(&unknown).await, None);
}
#[tokio::test]
async fn register_verdict_contracts_is_additive_across_calls() {
let (engine, task_id) = seeded_engine("gate").await;
let contract = mlua_swarm_schema::VerdictContract {
channel: mlua_swarm_schema::VerdictChannel::Part,
values: vec!["ALLOW".to_string()],
};
engine.register_verdict_contracts(HashMap::from([("gate".to_string(), contract.clone())]));
engine.register_verdict_contracts(HashMap::from([(
"unrelated-agent".to_string(),
mlua_swarm_schema::VerdictContract {
channel: mlua_swarm_schema::VerdictChannel::Body,
values: vec!["X".to_string()],
},
)]));
assert_eq!(
engine.verdict_contract_for_task(&task_id).await,
Some(contract)
);
}
}
#[cfg(test)]
mod verdict_contract_completion_tests {
use super::*;
async fn seeded_task_with_worker_token(agent: &str) -> (Engine, CapToken, StepId) {
let engine = Engine::new(EngineCfg::default());
let op_token = engine
.attach("ut-op", Role::Operator, Duration::from_secs(30))
.await
.expect("attach");
let task_id = engine
.start_task(
&op_token,
TaskSpec {
agent: agent.to_string(),
initial_directive: serde_json::json!("x"),
step_ctx: None,
check_policy: None,
},
)
.await
.expect("start_task");
let worker_token = engine.signer().session(
format!("worker-of-{task_id}"),
Role::Worker,
vec!["*".into()],
Duration::from_secs(600),
);
let fp = worker_token.fingerprint();
let record = CapTokenRecord::from_worker_token(worker_token.clone(), task_id.clone());
engine
.with_state("test.mint_worker", move |s| {
s.tokens.insert(fp, record);
})
.await
.expect("mint worker token");
(engine, worker_token, task_id)
}
fn body_contract(values: &[&str]) -> mlua_swarm_schema::VerdictContract {
mlua_swarm_schema::VerdictContract {
channel: mlua_swarm_schema::VerdictChannel::Body,
values: values.iter().map(|v| v.to_string()).collect(),
}
}
fn part_contract(values: &[&str]) -> mlua_swarm_schema::VerdictContract {
mlua_swarm_schema::VerdictContract {
channel: mlua_swarm_schema::VerdictChannel::Part,
values: values.iter().map(|v| v.to_string()).collect(),
}
}
fn final_event(value: Value, ok: bool) -> crate::worker::output::OutputEvent {
crate::worker::output::OutputEvent::Final {
content: crate::worker::output::ContentRef::Inline { value },
ok,
}
}
#[tokio::test]
async fn submit_output_rejects_missing_verdict_part() {
let (engine, token, task_id) = seeded_task_with_worker_token("gate").await;
engine.register_verdict_contracts(HashMap::from([(
"gate".to_string(),
part_contract(&["PASS", "BLOCKED"]),
)]));
let err = engine
.submit_output(
&token,
&task_id,
1,
final_event(serde_json::json!("anything"), true),
)
.await
.expect_err("missing staged verdict part must be rejected");
assert!(
matches!(err, EngineError::VerdictPartMissing { .. }),
"unexpected error variant: {err:?}"
);
let tail = engine.output_tail(&task_id, 1).await;
assert!(
!tail
.iter()
.any(|ev| matches!(ev, crate::worker::output::OutputEvent::Final { .. })),
"a rejected completion must not write a Final onto output_tail"
);
}
#[tokio::test]
async fn submit_output_accepts_when_verdict_part_is_staged_and_a_member() {
let (engine, token, task_id) = seeded_task_with_worker_token("gate").await;
engine.register_verdict_contracts(HashMap::from([(
"gate".to_string(),
part_contract(&["PASS", "BLOCKED"]),
)]));
engine
.stage_worker_artifact_trusted(
&task_id,
1,
"verdict".to_string(),
serde_json::json!("PASS"),
)
.await
.expect("stage verdict part");
engine
.submit_output(
&token,
&task_id,
1,
final_event(serde_json::json!("full report"), true),
)
.await
.expect("staged + member verdict part must be accepted");
let tail = engine.output_tail(&task_id, 1).await;
assert!(
tail.iter()
.any(|ev| matches!(ev, crate::worker::output::OutputEvent::Final { .. })),
"an accepted completion must write its Final onto output_tail"
);
}
#[tokio::test]
async fn submit_output_rejects_body_value_outside_contract() {
let (engine, token, task_id) = seeded_task_with_worker_token("gate").await;
engine.register_verdict_contracts(HashMap::from([(
"gate".to_string(),
body_contract(&["PASS", "BLOCKED"]),
)]));
let err = engine
.submit_output(
&token,
&task_id,
1,
final_event(serde_json::json!("UNKNOWN"), true),
)
.await
.expect_err("out-of-contract body value must be rejected");
match err {
EngineError::VerdictValueRejected { value, allowed } => {
assert_eq!(value, "UNKNOWN");
assert_eq!(allowed, vec!["PASS".to_string(), "BLOCKED".to_string()]);
}
other => panic!("unexpected error variant: {other:?}"),
}
let tail = engine.output_tail(&task_id, 1).await;
assert!(
!tail
.iter()
.any(|ev| matches!(ev, crate::worker::output::OutputEvent::Final { .. })),
"a rejected completion must not write a Final onto output_tail"
);
}
#[tokio::test]
async fn submit_output_ok_false_bypasses_the_check() {
let (engine, token, task_id) = seeded_task_with_worker_token("gate").await;
engine.register_verdict_contracts(HashMap::from([(
"gate".to_string(),
body_contract(&["PASS", "BLOCKED"]),
)]));
engine
.submit_output(
&token,
&task_id,
1,
final_event(serde_json::json!("UNKNOWN"), false),
)
.await
.expect("ok=false must bypass the verdict contract check entirely");
let tail = engine.output_tail(&task_id, 1).await;
assert!(
tail.iter()
.any(|ev| matches!(ev, crate::worker::output::OutputEvent::Final { .. })),
"an ok=false completion is exempt, not rejected — its Final must still land"
);
}
#[tokio::test]
async fn staged_verdict_value_for_is_last_write_wins() {
let (engine, _token, task_id) = seeded_task_with_worker_token("gate").await;
engine
.stage_worker_artifact_trusted(
&task_id,
1,
"verdict".to_string(),
serde_json::json!("PASS"),
)
.await
.expect("stage first verdict part");
engine
.stage_worker_artifact_trusted(
&task_id,
1,
"verdict".to_string(),
serde_json::json!("BLOCKED"),
)
.await
.expect("stage second verdict part");
assert_eq!(
engine.staged_verdict_value_for(&task_id, 1).await,
Some("BLOCKED".to_string())
);
}
#[tokio::test]
async fn staged_verdict_value_for_ignores_other_artifact_names() {
let (engine, _token, task_id) = seeded_task_with_worker_token("gate").await;
engine
.stage_worker_artifact_trusted(
&task_id,
1,
"notes".to_string(),
serde_json::json!("irrelevant"),
)
.await
.expect("stage unrelated part");
assert_eq!(engine.staged_verdict_value_for(&task_id, 1).await, None);
}
#[tokio::test]
async fn staged_verdict_value_for_returns_none_when_nothing_staged() {
let (engine, _token, task_id) = seeded_task_with_worker_token("gate").await;
assert_eq!(engine.staged_verdict_value_for(&task_id, 1).await, None);
}
}