1use axum::{
71 extract::{Query, State},
72 http::{header, header::AUTHORIZATION, HeaderMap, StatusCode},
73 Json,
74};
75use mlua_swarm::core::agent_context::StepPointer;
76use mlua_swarm::core::state::SubmitOutcome;
77use mlua_swarm::core::step_naming::StepNaming;
78use mlua_swarm::store::run::{DegradationEntry, RunStatus, RunStoreError};
79use mlua_swarm::{CapToken, ContentRef, EngineError, OutputEvent, RunId, StepId, WorkerPayload};
80use mlua_swarm_schema::{ContextPolicy, VerdictChannel};
81use serde::Deserialize;
82use serde_json::Value;
83
84use crate::projection::McpQueryAdapter;
85use crate::{ApiError, AppState};
86
87#[derive(Debug, Deserialize)]
89pub struct PromptQuery {
90 pub task_id: StepId,
94}
95
96pub async fn worker_prompt(
102 State(state): State<AppState>,
103 headers: HeaderMap,
104 Query(q): Query<PromptQuery>,
105) -> Result<Json<WorkerPayload>, ApiError> {
106 let task_id = q.task_id;
107 let bearer = extract_bearer_raw(&headers)?;
108 let mut payload = if let Some(handle) = parse_worker_handle(&bearer) {
109 let resolved = state
111 .engine
112 .task_id_from_handle(handle)
113 .await
114 .map_err(map_handle_lookup_err)?;
115 if resolved != task_id {
116 return Err(ApiError::bad_request(format!(
117 "handle {handle} is bound to task {resolved}, not {task_id}"
118 )));
119 }
120 state
121 .engine
122 .fetch_worker_payload_trusted(&task_id)
123 .await
124 .map_err(|e| ApiError::engine(format!("fetch_worker_payload_trusted: {e}")))?
125 } else {
126 let token = CapToken::decode(bearer.trim())
128 .map_err(|e| ApiError::bad_request(format!("invalid token: {e}")))?;
129 state
130 .engine
131 .fetch_worker_payload(&token, &task_id)
132 .await
133 .map_err(|e| ApiError::engine(format!("fetch_worker_payload: {e}")))?
134 };
135 assemble_step_pointers(&state, &mut payload).await;
136 Ok(Json(payload))
137}
138
139async fn assemble_step_pointers(state: &AppState, payload: &mut WorkerPayload) {
167 let Some(context) = payload.context.as_mut() else {
168 return;
169 };
170 let Some(run_id_str) = context.run_id.clone() else {
171 return;
172 };
173 let Ok(run_id) = RunId::parse(run_id_str) else {
174 return;
175 };
176
177 let adapter = McpQueryAdapter::new(
178 state.data_store.clone(),
179 state.run_store.clone(),
180 state.engine.clone(),
181 );
182 let Ok((run, resolved_steps)) = adapter.list_steps_by_run_id(&run_id).await else {
183 return;
184 };
185
186 let naming = state.engine.step_naming_for(&payload.task_id).await;
187 let policy = state
188 .engine
189 .context_policy_for(&payload.task_id, payload.attempt)
190 .await;
191 let self_canonical = naming
192 .as_deref()
193 .and_then(|n| n.canonical_of_producer(&payload.agent))
194 .map(str::to_string)
195 .unwrap_or_else(|| payload.agent.clone());
196
197 let mut pointers = Vec::new();
198 for step in &resolved_steps {
199 if step.name == self_canonical
200 || !allows_step_canonical(&policy, naming.as_deref(), &step.name)
201 {
202 continue;
203 }
204 if let Some((size_bytes, file_path, content_url, sha256)) =
205 crate::projection::resolve_step_pointer_fields(state, &run, step).await
206 {
207 pointers.push(StepPointer {
208 name: step.name.clone(),
209 size_bytes,
210 file_path,
211 content_url,
212 sha256,
213 });
214 }
215 }
216 context.steps = pointers;
217}
218
219fn allows_step_canonical(
234 policy: &ContextPolicy,
235 naming: Option<&StepNaming>,
236 canonical_name: &str,
237) -> bool {
238 let resolves_to = |raw: &str| -> bool {
239 match naming {
240 Some(n) => n
241 .resolve(raw)
242 .map(|c| c == canonical_name)
243 .unwrap_or(raw == canonical_name),
244 None => raw == canonical_name,
245 }
246 };
247 if policy
248 .steps_exclude
249 .iter()
250 .any(|excluded| resolves_to(excluded))
251 {
252 return false;
253 }
254 match &policy.steps {
255 None => true,
256 Some(list) => list.iter().any(|included| resolves_to(included)),
257 }
258}
259
260#[derive(Debug, Deserialize)]
262pub struct WorkerResultReq {
263 pub task_id: StepId,
266 pub value: Value,
268 #[serde(default = "default_ok_true")]
272 pub ok: bool,
273 #[serde(default)]
276 pub attempt: Option<u32>,
277}
278
279fn default_ok_true() -> bool {
280 true
281}
282
283pub async fn worker_result(
286 State(state): State<AppState>,
287 headers: HeaderMap,
288 Json(req): Json<WorkerResultReq>,
289) -> Result<StatusCode, ApiError> {
290 let token = decode_worker_bearer(&headers)?;
291 let task_id = req.task_id.clone();
292
293 let attempt = match req.attempt {
295 Some(n) => n,
296 None => state
297 .engine
298 .task_attempt(&task_id)
299 .await
300 .map_err(|e| ApiError::engine(format!("task_attempt: {e}")))?,
301 };
302
303 let event = OutputEvent::Final {
304 content: ContentRef::Inline {
305 value: req.value.clone(),
306 },
307 ok: req.ok,
308 };
309 map_completion_result(
314 state
315 .engine
316 .submit_output(&token, &task_id, attempt, event)
317 .await,
318 "submit_output",
319 )?;
320 state
321 .engine
322 .post_result(&token, &task_id, req.value)
323 .await
324 .map_err(|e| ApiError::engine(format!("post_result: {e}")))?;
325 Ok(StatusCode::NO_CONTENT)
326}
327
328const FILE_SENTINEL_PREFIX: &str = "@file:";
336
337const FILE_SENTINEL_MAX_BYTES: u64 = 2 * 1024 * 1024;
343
344const FILE_SENTINEL_ALLOW_KEY: &str = "allow_file_submit";
356
357async fn resolve_file_sentinel(
390 state: &AppState,
391 task_id: &StepId,
392 attempt: u32,
393 body_str: String,
394) -> Result<String, ApiError> {
395 let Some(rest) = body_str.strip_prefix(FILE_SENTINEL_PREFIX) else {
396 return Ok(body_str);
397 };
398 let path_str = rest.trim();
399 if path_str.is_empty() {
400 return Err(ApiError::bad_request(
401 "@file: sentinel: empty path".to_string(),
402 ));
403 }
404 if path_str.contains('\n') || path_str.contains('\r') {
405 return Err(ApiError::bad_request(
406 "@file: sentinel: path must be a single line".to_string(),
407 ));
408 }
409 let path = std::path::Path::new(path_str);
410 if !path.is_absolute() {
411 return Err(ApiError::bad_request(format!(
412 "@file: sentinel: path must be absolute (got {path_str:?})"
413 )));
414 }
415 let view = state
416 .engine
417 .agent_context_for(task_id, attempt)
418 .await
419 .ok_or_else(|| {
420 ApiError::bad_request(
421 "@file: sentinel: no AgentContextView for this task/attempt \
422 (spawn must run through AgentContextMiddleware to enable \
423 sentinel resolution)"
424 .to_string(),
425 )
426 })?;
427 if view.extra.get(FILE_SENTINEL_ALLOW_KEY) != Some(&Value::Bool(true)) {
431 return Err(ApiError::bad_request(format!(
432 "@file: sentinel: file submission is not allowed for this step \
433 (declare `{FILE_SENTINEL_ALLOW_KEY}: true` via `$step_meta` / \
434 `AgentMeta.ctx` / `Blueprint.metas`; strict boolean `true` \
435 required)"
436 )));
437 }
438 let work_dir = view.work_dir.ok_or_else(|| {
439 ApiError::bad_request("@file: sentinel: task has no resolved work_dir".to_string())
440 })?;
441 let work_dir_canon = tokio::fs::canonicalize(&work_dir).await.map_err(|e| {
442 ApiError::engine(format!(
443 "@file: sentinel: canonicalize work_dir {work_dir:?}: {e}"
444 ))
445 })?;
446 let path_canon = match tokio::fs::canonicalize(path).await {
447 Ok(p) => p,
448 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
449 return Err(ApiError::not_found(format!(
450 "@file: sentinel: file not found: {path_str}"
451 )));
452 }
453 Err(e) => {
454 return Err(ApiError::engine(format!(
455 "@file: sentinel: canonicalize {path_str:?}: {e}"
456 )));
457 }
458 };
459 if !path_canon.starts_with(&work_dir_canon) {
460 return Err(ApiError::bad_request(format!(
461 "@file: sentinel: path {} is not under work_dir {} (canonicalized: {} vs {})",
462 path_str,
463 work_dir,
464 path_canon.display(),
465 work_dir_canon.display(),
466 )));
467 }
468 let meta = tokio::fs::metadata(&path_canon)
469 .await
470 .map_err(|e| ApiError::engine(format!("@file: sentinel: metadata {path_str:?}: {e}")))?;
471 if meta.len() > FILE_SENTINEL_MAX_BYTES {
472 return Err(ApiError::payload_too_large(format!(
473 "@file: sentinel: file size {} exceeds limit {}",
474 meta.len(),
475 FILE_SENTINEL_MAX_BYTES
476 )));
477 }
478 let bytes = tokio::fs::read(&path_canon)
479 .await
480 .map_err(|e| ApiError::engine(format!("@file: sentinel: read {path_str:?}: {e}")))?;
481 Ok(String::from_utf8_lossy(&bytes).trim_end().to_string())
483}
484
485async fn check_verdict_contract(
506 state: &AppState,
507 task_id: &StepId,
508 channel: VerdictChannel,
509 value: &str,
510) -> Result<(), ApiError> {
511 let Some(contract) = state.engine.verdict_contract_for_task(task_id).await else {
512 return Ok(());
513 };
514 if contract.channel != channel {
515 return Ok(());
516 }
517 if contract.values.iter().any(|v| v == value) {
518 return Ok(());
519 }
520 Err(ApiError::unprocessable(format!(
521 "verdict contract violation: {value:?} is not a member of the declared values {:?}",
522 contract.values
523 )))
524}
525
526fn map_completion_result<T>(result: Result<T, EngineError>, context: &str) -> Result<T, ApiError> {
543 result.map_err(|e| match e {
544 EngineError::VerdictValueRejected { value, allowed } => ApiError::unprocessable(format!(
545 "verdict contract violation: {value:?} is not a member of the declared values {allowed:?}"
546 )),
547 EngineError::VerdictPartMissing { allowed } => ApiError::unprocessable(format!(
548 "verdict contract violation: no staged \"verdict\" part found for this attempt; declared values {allowed:?}"
549 )),
550 other => ApiError::engine(format!("{context}: {other}")),
551 })
552}
553
554#[derive(Debug, Deserialize, Default)]
576pub struct SubmitQuery {
577 #[serde(default)]
581 pub ok: Option<bool>,
582 #[serde(default)]
592 pub verdict: Option<String>,
593}
594
595fn resolve_submit_outcome(
607 verdict: Option<&str>,
608 ok: Option<bool>,
609) -> Result<SubmitOutcome, String> {
610 match verdict {
611 None => Ok(if ok.unwrap_or(true) {
612 SubmitOutcome::Pass
613 } else {
614 SubmitOutcome::Blocked
615 }),
616 Some(v) => match v {
617 "pass" => {
618 if ok == Some(false) {
619 Err(
620 "conflicting signal: verdict=pass with ok=false; drop one of them"
621 .to_string(),
622 )
623 } else {
624 Ok(SubmitOutcome::Pass)
625 }
626 }
627 "blocked" => {
628 if ok == Some(true) {
629 Err(
630 "conflicting signal: verdict=blocked with ok=true; drop one of them"
631 .to_string(),
632 )
633 } else {
634 Ok(SubmitOutcome::Blocked)
635 }
636 }
637 "skip" => {
638 if ok == Some(false) {
639 Err(
640 "conflicting signal: verdict=skip with ok=false; drop one of them"
641 .to_string(),
642 )
643 } else {
644 Ok(SubmitOutcome::Skip)
645 }
646 }
647 other => Err(format!(
648 "verdict must be one of: pass, blocked, skip (got {other:?})"
649 )),
650 },
651 }
652}
653
654pub async fn worker_submit(
660 State(state): State<AppState>,
661 headers: HeaderMap,
662 Query(q): Query<SubmitQuery>,
663 body: axum::body::Bytes,
664) -> Result<StatusCode, ApiError> {
665 let bearer = extract_bearer_raw(&headers)?;
668 let task_id = if let Some(handle) = parse_worker_handle(&bearer) {
669 state
670 .engine
671 .task_id_from_handle(handle)
672 .await
673 .map_err(map_handle_lookup_err)?
674 } else {
675 let token = CapToken::decode(bearer.trim())
676 .map_err(|e| ApiError::bad_request(format!("invalid token: {e}")))?;
677 state
678 .engine
679 .task_id_from_token(&token)
680 .await
681 .map_err(|e| ApiError::engine(format!("task_id_from_token: {e}")))?
682 };
683 let attempt = state
684 .engine
685 .task_attempt(&task_id)
686 .await
687 .map_err(|e| ApiError::engine(format!("task_attempt: {e}")))?;
688 reject_if_run_terminal(&state, &task_id, attempt).await?;
691 let body_str = String::from_utf8_lossy(&body).trim_end().to_string();
696 let body_str = resolve_file_sentinel(&state, &task_id, attempt, body_str).await?;
699 let value = Value::String(body_str);
708
709 let outcome =
718 resolve_submit_outcome(q.verdict.as_deref(), q.ok).map_err(ApiError::bad_request)?;
719 let submit_result = state
720 .engine
721 .submit_worker_result_trusted(&task_id, attempt, value, outcome)
722 .await;
723 map_completion_result(submit_result, "submit_worker_result_trusted")?;
724 Ok(StatusCode::NO_CONTENT)
725}
726
727#[derive(Debug, Deserialize)]
729pub struct ArtifactQuery {
730 pub name: String,
737}
738
739pub async fn worker_artifact(
761 State(state): State<AppState>,
762 headers: HeaderMap,
763 Query(q): Query<ArtifactQuery>,
764 body: axum::body::Bytes,
765) -> Result<StatusCode, ApiError> {
766 let name = q.name.trim();
767 if name.is_empty() {
768 return Err(ApiError::bad_request("name must not be empty".into()));
769 }
770 let name = name.to_string();
771
772 let bearer = extract_bearer_raw(&headers)?;
773 let task_id = if let Some(handle) = parse_worker_handle(&bearer) {
774 state
775 .engine
776 .task_id_from_handle(handle)
777 .await
778 .map_err(map_handle_lookup_err)?
779 } else {
780 let token = CapToken::decode(bearer.trim())
781 .map_err(|e| ApiError::bad_request(format!("invalid token: {e}")))?;
782 state
783 .engine
784 .task_id_from_token(&token)
785 .await
786 .map_err(|e| ApiError::engine(format!("task_id_from_token: {e}")))?
787 };
788 let attempt = state
789 .engine
790 .task_attempt(&task_id)
791 .await
792 .map_err(|e| ApiError::engine(format!("task_attempt: {e}")))?;
793 reject_if_run_terminal(&state, &task_id, attempt).await?;
796 let body_str = String::from_utf8_lossy(&body).trim_end().to_string();
797 let body_str = resolve_file_sentinel(&state, &task_id, attempt, body_str).await?;
799 if name == "verdict" {
805 check_verdict_contract(&state, &task_id, VerdictChannel::Part, &body_str).await?;
806 }
807 let value = Value::String(body_str);
808
809 state
810 .engine
811 .stage_worker_artifact_trusted(&task_id, attempt, name, value)
812 .await
813 .map_err(|e| ApiError::engine(format!("stage_worker_artifact_trusted: {e}")))?;
814 Ok(StatusCode::NO_CONTENT)
815}
816
817#[derive(Debug, Deserialize)]
819pub struct DegradationBody {
820 pub tool: String,
822 pub error: String,
824 pub fallback: String,
826 #[serde(default)]
828 pub note: Option<String>,
829}
830
831pub async fn worker_degradation(
859 State(state): State<AppState>,
860 headers: HeaderMap,
861 Json(body): Json<DegradationBody>,
862) -> Result<StatusCode, ApiError> {
863 let bearer = extract_bearer_raw(&headers)?;
864 let task_id = if let Some(handle) = parse_worker_handle(&bearer) {
865 state
866 .engine
867 .task_id_from_handle(handle)
868 .await
869 .map_err(map_handle_lookup_err)?
870 } else {
871 let token = CapToken::decode(bearer.trim())
872 .map_err(|e| ApiError::bad_request(format!("invalid token: {e}")))?;
873 state
874 .engine
875 .task_id_from_token(&token)
876 .await
877 .map_err(|e| ApiError::engine(format!("task_id_from_token: {e}")))?
878 };
879 let attempt = state
880 .engine
881 .task_attempt(&task_id)
882 .await
883 .map_err(|e| ApiError::engine(format!("task_attempt: {e}")))?;
884 reject_if_run_terminal(&state, &task_id, attempt).await?;
887
888 let tid = task_id.clone();
893 let (run_id_str, agent) = match state
894 .engine
895 .with_state("worker_degradation_run_lookup", move |s| {
896 s.agent_ctx.get(&(tid, attempt)).and_then(|e| {
897 e.view
898 .run_id
899 .clone()
900 .map(|run_id| (run_id, e.view.agent.clone()))
901 })
902 })
903 .await
904 {
905 Ok(Some(pair)) => pair,
906 _ => {
907 tracing::warn!(%task_id, "worker_degradation: no run linkage for this task; entry dropped");
908 return Ok(StatusCode::NO_CONTENT);
909 }
910 };
911 let Ok(run_id) = RunId::parse(run_id_str) else {
912 tracing::warn!(%task_id, "worker_degradation: run_id failed to parse; entry dropped");
913 return Ok(StatusCode::NO_CONTENT);
914 };
915
916 mlua_swarm::store::trace::TraceHandle::new(run_id.clone(), state.run_trace_store.clone())
920 .append(
921 mlua_swarm::store::trace::kind::WORKER_DEGRADATION,
922 Some(agent.as_str()),
923 Some(attempt),
924 serde_json::json!({
925 "tool": body.tool.as_str(),
926 "error": body.error.as_str(),
927 "fallback": body.fallback.as_str(),
928 }),
929 )
930 .await;
931
932 let entry = DegradationEntry {
933 tool: body.tool,
934 error: body.error,
935 fallback: body.fallback,
936 note: body.note,
937 step_ref: Some(agent),
938 attempt: Some(attempt),
939 at: crate::tasks::now_secs(),
940 };
941 match state.run_store.append_degradation(&run_id, entry).await {
942 Ok(()) => Ok(StatusCode::NO_CONTENT),
943 Err(RunStoreError::NotFound(_)) => {
944 tracing::warn!(%task_id, %run_id, "worker_degradation: run not found in run_store; entry dropped");
945 Ok(StatusCode::NO_CONTENT)
946 }
947 Err(e) => Err(ApiError::engine(format!("append_degradation: {e}"))),
948 }
949}
950
951#[derive(Debug, Deserialize)]
955pub struct StatsBody {
956 #[serde(default)]
960 pub worker_kind: Option<String>,
961 #[serde(default)]
963 pub model: Option<String>,
964 #[serde(default)]
966 pub usage: Option<mlua_swarm::store::trace::TokenUsage>,
967 #[serde(default)]
969 pub num_turns: Option<u32>,
970 #[serde(default)]
972 pub adapter_data: Option<Value>,
973}
974
975pub async fn worker_stats(
986 State(state): State<AppState>,
987 headers: HeaderMap,
988 Json(body): Json<StatsBody>,
989) -> Result<StatusCode, ApiError> {
990 let bearer = extract_bearer_raw(&headers)?;
991 let task_id = if let Some(handle) = parse_worker_handle(&bearer) {
992 state
993 .engine
994 .task_id_from_handle(handle)
995 .await
996 .map_err(map_handle_lookup_err)?
997 } else {
998 let token = CapToken::decode(bearer.trim())
999 .map_err(|e| ApiError::bad_request(format!("invalid token: {e}")))?;
1000 state
1001 .engine
1002 .task_id_from_token(&token)
1003 .await
1004 .map_err(|e| ApiError::engine(format!("task_id_from_token: {e}")))?
1005 };
1006 let attempt = state
1007 .engine
1008 .task_attempt(&task_id)
1009 .await
1010 .map_err(|e| ApiError::engine(format!("task_attempt: {e}")))?;
1011 reject_if_run_terminal(&state, &task_id, attempt).await?;
1014
1015 let stats = mlua_swarm::store::trace::WorkerStats {
1016 worker_kind: Some(body.worker_kind.unwrap_or_else(|| "operator".to_string())),
1017 model: body.model,
1018 usage: body.usage,
1019 num_turns: body.num_turns,
1020 adapter_data: body.adapter_data,
1021 };
1022 state
1023 .engine
1024 .record_worker_stats(&task_id, attempt, stats)
1025 .await;
1026 Ok(StatusCode::NO_CONTENT)
1027}
1028
1029async fn reject_if_run_terminal(
1046 state: &AppState,
1047 task_id: &StepId,
1048 attempt: u32,
1049) -> Result<(), ApiError> {
1050 let tid = task_id.clone();
1051 let run_id_str = match state
1052 .engine
1053 .with_state("worker_terminal_run_guard", move |s| {
1054 s.agent_ctx
1055 .get(&(tid, attempt))
1056 .and_then(|e| e.view.run_id.clone())
1057 })
1058 .await
1059 {
1060 Ok(Some(rid)) => rid,
1061 _ => return Ok(()),
1062 };
1063 let Ok(run_id) = RunId::parse(run_id_str) else {
1064 return Ok(());
1065 };
1066 let Ok(rec) = state.run_store.get(&run_id).await else {
1067 return Ok(());
1068 };
1069 match rec.status {
1070 RunStatus::Done | RunStatus::Failed | RunStatus::Interrupted | RunStatus::Cancelled => {
1071 Err(ApiError::gone(format!(
1072 "run {run_id} is already terminal ({:?}): this attempt's output cannot be \
1073 delivered to a flow context; re-kick the task (POST /v1/tasks/:id/runs) and \
1074 fetch a fresh prompt",
1075 rec.status
1076 )))
1077 }
1078 RunStatus::Pending | RunStatus::Running => Ok(()),
1079 }
1080}
1081
1082#[derive(Debug, Deserialize)]
1087pub struct PromptSystemQuery {
1088 pub task_id: StepId,
1091 pub attempt: u32,
1093}
1094
1095pub async fn worker_prompt_system(
1105 State(state): State<AppState>,
1106 headers: HeaderMap,
1107 Query(q): Query<PromptSystemQuery>,
1108) -> Result<impl axum::response::IntoResponse, ApiError> {
1109 let task_id = q.task_id;
1110 let attempt = q.attempt;
1111 let bearer = extract_bearer_raw(&headers)?;
1112 if let Some(handle) = parse_worker_handle(&bearer) {
1113 let resolved = state
1114 .engine
1115 .task_id_from_handle(handle)
1116 .await
1117 .map_err(map_handle_lookup_err)?;
1118 if resolved != task_id {
1119 return Err(ApiError::bad_request(format!(
1120 "handle {handle} is bound to task {resolved}, not {task_id}"
1121 )));
1122 }
1123 } else {
1124 let token = CapToken::decode(bearer.trim())
1125 .map_err(|e| ApiError::bad_request(format!("invalid token: {e}")))?;
1126 state
1127 .engine
1128 .verify_token_for_task(&token, mlua_swarm::Verb::FetchPrompt, &task_id)
1129 .await
1130 .map_err(|e| ApiError::engine(format!("verify_token_for_task: {e}")))?;
1131 }
1132 let system = state
1133 .engine
1134 .raw_system_prompt(&task_id, attempt)
1135 .await
1136 .map_err(|e| ApiError::engine(format!("raw_system_prompt: {e}")))?
1137 .ok_or_else(|| {
1138 ApiError::not_found(format!(
1139 "no baked system prompt for task {task_id} attempt {attempt}"
1140 ))
1141 })?;
1142 Ok((
1143 [(header::CONTENT_TYPE, "text/plain; charset=utf-8")],
1144 system,
1145 ))
1146}
1147
1148#[derive(Debug, serde::Serialize)]
1150pub struct AgentRenderSizeResponse {
1151 pub agent: String,
1153 pub last_rendered_bytes: Option<usize>,
1157}
1158
1159pub async fn agent_render_size(
1169 State(state): State<AppState>,
1170 axum::extract::Path(name): axum::extract::Path<String>,
1171) -> Json<AgentRenderSizeResponse> {
1172 let last_rendered_bytes = state.engine.agent_last_rendered_size(&name).await;
1173 Json(AgentRenderSizeResponse {
1174 agent: name,
1175 last_rendered_bytes,
1176 })
1177}
1178
1179fn extract_bearer_raw(headers: &HeaderMap) -> Result<String, ApiError> {
1183 let v = headers
1184 .get(AUTHORIZATION)
1185 .ok_or_else(|| ApiError::bad_request("missing Authorization header".into()))?
1186 .to_str()
1187 .map_err(|_| ApiError::bad_request("invalid Authorization header encoding".into()))?;
1188 let s = v
1189 .strip_prefix("Bearer ")
1190 .ok_or_else(|| ApiError::bad_request("Authorization must be 'Bearer <token>'".into()))?
1191 .trim();
1192 if s.is_empty() {
1193 return Err(ApiError::bad_request("Bearer is empty".into()));
1194 }
1195 Ok(s.to_string())
1196}
1197
1198fn map_handle_lookup_err(e: EngineError) -> ApiError {
1211 match e {
1212 EngineError::TokenNotFound(_) => ApiError::gone(
1213 "worker handle is no longer valid (the engine's in-flight state was reset, \
1214 e.g. by a server restart): re-kick the task (POST /v1/tasks/:id/runs) and \
1215 fetch a fresh prompt/handle"
1216 .to_string(),
1217 ),
1218 other => ApiError::engine(format!("task_id_from_handle: {other}")),
1219 }
1220}
1221
1222fn parse_worker_handle(s: &str) -> Option<&str> {
1226 let s = s.trim();
1227 if s.starts_with("wh-")
1228 && s.len() >= 5
1229 && s.len() <= 64
1230 && s[3..].chars().all(|c| c.is_ascii_alphanumeric())
1231 {
1232 Some(s)
1233 } else {
1234 None
1235 }
1236}
1237
1238fn decode_worker_bearer(headers: &HeaderMap) -> Result<CapToken, ApiError> {
1242 let v = headers
1243 .get(AUTHORIZATION)
1244 .ok_or_else(|| ApiError::bad_request("missing Authorization header".into()))?
1245 .to_str()
1246 .map_err(|_| ApiError::bad_request("invalid Authorization header encoding".into()))?;
1247 let encoded = v
1248 .strip_prefix("Bearer ")
1249 .ok_or_else(|| ApiError::bad_request("Authorization must be 'Bearer <token>'".into()))?
1250 .trim();
1251 if encoded.is_empty() {
1252 return Err(ApiError::bad_request("Bearer token is empty".into()));
1253 }
1254 CapToken::decode(encoded).map_err(|e| ApiError::bad_request(format!("invalid token: {e}")))
1255}
1256
1257#[cfg(test)]
1262mod tests {
1263 use super::*;
1264 use axum::response::IntoResponse;
1265 use mlua_swarm::core::agent_context::AgentContextView;
1266 use mlua_swarm::core::config::EngineCfg;
1267 use mlua_swarm::core::engine::Engine;
1268 use mlua_swarm::store::output::{InMemoryOutputStore, OutputStore};
1269 use mlua_swarm::store::run::{InMemoryRunStore, RunRecord, RunStatus, RunStore, StepEntry};
1270 use mlua_swarm::store::task::InMemoryTaskStore;
1271 use mlua_swarm::{RunId, StepId, TaskId};
1272 use serde_json::json;
1273 use std::collections::HashMap;
1274 use std::sync::Arc;
1275 use tokio::sync::Mutex;
1276
1277 fn test_state(data_store: Arc<dyn OutputStore>, run_store: Arc<dyn RunStore>) -> AppState {
1283 let engine = Engine::new(EngineCfg::default());
1284 let compiler = mlua_swarm::Compiler::new(crate::default_registry());
1285 let launch = Arc::new(mlua_swarm::TaskLaunchService::new(engine.clone(), compiler));
1286 AppState {
1287 engine,
1288 sessions: Arc::new(Mutex::new(crate::SessionStore::default())),
1289 task_app: Arc::new(mlua_swarm::TaskApplication::new_inline_only(launch)),
1290 ws_operator_factory: None,
1291 data_store,
1292 operator_sessions: Arc::new(Mutex::new(HashMap::new())),
1293 roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
1294 task_store: Arc::new(InMemoryTaskStore::new()),
1295 run_store,
1296 replay_store: Arc::new(mlua_swarm::store::replay::InMemoryReplayStore::new()),
1297 run_trace_store: Arc::new(mlua_swarm::store::trace::InMemoryRunTraceStore::new()),
1298 base_url: None,
1299 sync_timeout_secs: 300,
1300 }
1301 }
1302
1303 #[test]
1307 fn resolve_submit_outcome_absent_verdict_preserves_pre_gh76_wire() {
1308 assert!(matches!(
1310 resolve_submit_outcome(None, None),
1311 Ok(SubmitOutcome::Pass)
1312 ));
1313 assert!(matches!(
1314 resolve_submit_outcome(None, Some(true)),
1315 Ok(SubmitOutcome::Pass)
1316 ));
1317 assert!(matches!(
1318 resolve_submit_outcome(None, Some(false)),
1319 Ok(SubmitOutcome::Blocked)
1320 ));
1321 }
1322
1323 #[test]
1324 fn resolve_submit_outcome_verdict_pass_and_blocked_match_ok_bool_or_default() {
1325 assert!(matches!(
1326 resolve_submit_outcome(Some("pass"), None),
1327 Ok(SubmitOutcome::Pass)
1328 ));
1329 assert!(matches!(
1330 resolve_submit_outcome(Some("pass"), Some(true)),
1331 Ok(SubmitOutcome::Pass)
1332 ));
1333 assert!(resolve_submit_outcome(Some("pass"), Some(false)).is_err());
1334
1335 assert!(matches!(
1336 resolve_submit_outcome(Some("blocked"), None),
1337 Ok(SubmitOutcome::Blocked)
1338 ));
1339 assert!(matches!(
1340 resolve_submit_outcome(Some("blocked"), Some(false)),
1341 Ok(SubmitOutcome::Blocked)
1342 ));
1343 assert!(resolve_submit_outcome(Some("blocked"), Some(true)).is_err());
1344 }
1345
1346 #[test]
1347 fn resolve_submit_outcome_verdict_skip_is_ok_true_only() {
1348 assert!(matches!(
1349 resolve_submit_outcome(Some("skip"), None),
1350 Ok(SubmitOutcome::Skip)
1351 ));
1352 assert!(matches!(
1353 resolve_submit_outcome(Some("skip"), Some(true)),
1354 Ok(SubmitOutcome::Skip)
1355 ));
1356 let err = resolve_submit_outcome(Some("skip"), Some(false))
1358 .expect_err("skip + ok=false must be a conflict");
1359 assert!(
1360 err.contains("conflict") || err.contains("conflicting"),
1361 "err should name the conflict: {err}"
1362 );
1363 }
1364
1365 #[test]
1366 fn resolve_submit_outcome_invalid_verdict_names_valid_set() {
1367 let err = resolve_submit_outcome(Some("bogus"), None)
1368 .expect_err("unknown verdict must be an error");
1369 assert!(
1370 err.contains("pass") && err.contains("blocked") && err.contains("skip"),
1371 "err should enumerate the valid tier set: {err}"
1372 );
1373 }
1374
1375 async fn append_final(
1376 data_store: &Arc<dyn OutputStore>,
1377 task_id: &str,
1378 producer: &str,
1379 value: Value,
1380 ) {
1381 data_store
1382 .append(
1383 task_id,
1384 1,
1385 producer,
1386 OutputEvent::Final {
1387 content: ContentRef::Inline { value },
1388 ok: true,
1389 },
1390 vec![],
1391 )
1392 .await
1393 .expect("append final");
1394 }
1395
1396 fn step_entry(step_id: &StepId, step_ref: &str) -> StepEntry {
1397 StepEntry::basic(
1398 step_id.clone(),
1399 Some(step_ref.to_string()),
1400 Some("passed".to_string()),
1401 None,
1402 0,
1403 )
1404 }
1405
1406 fn run_record(task_id: &TaskId, run_id: &RunId, step_entries: Vec<StepEntry>) -> RunRecord {
1407 RunRecord {
1408 id: run_id.clone(),
1409 task_id: task_id.clone(),
1410 status: RunStatus::Running,
1411 step_entries,
1412 degradations: Vec::new(),
1413 operator_sid: None,
1414 result_ref: None,
1415 input_json: None,
1416 created_at: 0,
1417 updated_at: 0,
1418 }
1419 }
1420
1421 fn consumer_payload(consumer_step_id: &StepId, run_id: &RunId) -> WorkerPayload {
1422 WorkerPayload {
1423 task_id: consumer_step_id.clone(),
1424 attempt: 1,
1425 agent: "consumer".to_string(),
1426 system: None,
1427 prompt: String::new(),
1428 context: Some(AgentContextView {
1429 task_id: consumer_step_id.to_string(),
1430 agent: "consumer".to_string(),
1431 attempt: 1,
1432 run_id: Some(run_id.to_string()),
1433 ..Default::default()
1434 }),
1435 system_ref: None,
1436 }
1437 }
1438
1439 #[tokio::test]
1444 async fn context_policy_unspecified_yields_every_submitted_step() {
1445 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1446 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1447 let task_id = TaskId::new();
1448 let run_id = RunId::new();
1449 let planner_id = StepId::new();
1450 let coder_id = StepId::new();
1451
1452 append_final(
1453 &data_store,
1454 planner_id.as_str(),
1455 "planner",
1456 json!({"plan": "x"}),
1457 )
1458 .await;
1459 append_final(
1460 &data_store,
1461 coder_id.as_str(),
1462 "coder",
1463 json!({"code": "y"}),
1464 )
1465 .await;
1466 run_store
1467 .create(run_record(
1468 &task_id,
1469 &run_id,
1470 vec![
1471 step_entry(&planner_id, "planner"),
1472 step_entry(&coder_id, "coder"),
1473 ],
1474 ))
1475 .await
1476 .expect("create run");
1477
1478 let state = test_state(data_store, run_store);
1479 let consumer_id = StepId::new();
1480 let mut payload = consumer_payload(&consumer_id, &run_id);
1481 assemble_step_pointers(&state, &mut payload).await;
1482
1483 let names: Vec<&str> = payload
1484 .context
1485 .as_ref()
1486 .expect("context")
1487 .steps
1488 .iter()
1489 .map(|p| p.name.as_str())
1490 .collect();
1491 assert!(names.contains(&"planner"), "names: {names:?}");
1492 assert!(names.contains(&"coder"), "names: {names:?}");
1493 }
1494
1495 #[tokio::test]
1497 async fn context_policy_steps_include_list_filters_to_named_steps() {
1498 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1499 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1500 let task_id = TaskId::new();
1501 let run_id = RunId::new();
1502 let planner_id = StepId::new();
1503 let coder_id = StepId::new();
1504 append_final(&data_store, planner_id.as_str(), "planner", json!("x")).await;
1505 append_final(&data_store, coder_id.as_str(), "coder", json!("y")).await;
1506 run_store
1507 .create(run_record(
1508 &task_id,
1509 &run_id,
1510 vec![
1511 step_entry(&planner_id, "planner"),
1512 step_entry(&coder_id, "coder"),
1513 ],
1514 ))
1515 .await
1516 .expect("create run");
1517
1518 let state = test_state(data_store, run_store);
1519 let consumer_id = StepId::new();
1520 state
1521 .engine
1522 .with_state("test.seed_policy", {
1523 let consumer_id = consumer_id.clone();
1524 move |s| {
1525 s.agent_ctx.insert(
1526 (consumer_id, 1),
1527 mlua_swarm::core::state::AgentCtxEntry {
1528 policy: mlua_swarm_schema::ContextPolicy {
1529 steps: Some(vec!["planner".to_string()]),
1530 ..Default::default()
1531 },
1532 ..Default::default()
1533 },
1534 );
1535 }
1536 })
1537 .await
1538 .expect("seed policy");
1539
1540 let mut payload = consumer_payload(&consumer_id, &run_id);
1541 assemble_step_pointers(&state, &mut payload).await;
1542
1543 let names: Vec<&str> = payload
1544 .context
1545 .as_ref()
1546 .expect("context")
1547 .steps
1548 .iter()
1549 .map(|p| p.name.as_str())
1550 .collect();
1551 assert_eq!(names, vec!["planner"], "names: {names:?}");
1552 }
1553
1554 #[tokio::test]
1556 async fn context_policy_steps_empty_list_yields_no_pointers() {
1557 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1558 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1559 let task_id = TaskId::new();
1560 let run_id = RunId::new();
1561 let planner_id = StepId::new();
1562 append_final(&data_store, planner_id.as_str(), "planner", json!("x")).await;
1563 run_store
1564 .create(run_record(
1565 &task_id,
1566 &run_id,
1567 vec![step_entry(&planner_id, "planner")],
1568 ))
1569 .await
1570 .expect("create run");
1571
1572 let state = test_state(data_store, run_store);
1573 let consumer_id = StepId::new();
1574 state
1575 .engine
1576 .with_state("test.seed_policy", {
1577 let consumer_id = consumer_id.clone();
1578 move |s| {
1579 s.agent_ctx.insert(
1580 (consumer_id, 1),
1581 mlua_swarm::core::state::AgentCtxEntry {
1582 policy: mlua_swarm_schema::ContextPolicy {
1583 steps: Some(vec![]),
1584 ..Default::default()
1585 },
1586 ..Default::default()
1587 },
1588 );
1589 }
1590 })
1591 .await
1592 .expect("seed policy");
1593
1594 let mut payload = consumer_payload(&consumer_id, &run_id);
1595 assemble_step_pointers(&state, &mut payload).await;
1596
1597 assert!(payload.context.expect("context").steps.is_empty());
1598 }
1599
1600 #[tokio::test]
1602 async fn context_policy_steps_exclude_wins_over_steps() {
1603 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1604 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1605 let task_id = TaskId::new();
1606 let run_id = RunId::new();
1607 let planner_id = StepId::new();
1608 let coder_id = StepId::new();
1609 append_final(&data_store, planner_id.as_str(), "planner", json!("x")).await;
1610 append_final(&data_store, coder_id.as_str(), "coder", json!("y")).await;
1611 run_store
1612 .create(run_record(
1613 &task_id,
1614 &run_id,
1615 vec![
1616 step_entry(&planner_id, "planner"),
1617 step_entry(&coder_id, "coder"),
1618 ],
1619 ))
1620 .await
1621 .expect("create run");
1622
1623 let state = test_state(data_store, run_store);
1624 let consumer_id = StepId::new();
1625 state
1626 .engine
1627 .with_state("test.seed_policy", {
1628 let consumer_id = consumer_id.clone();
1629 move |s| {
1630 s.agent_ctx.insert(
1631 (consumer_id, 1),
1632 mlua_swarm::core::state::AgentCtxEntry {
1633 policy: mlua_swarm_schema::ContextPolicy {
1634 steps: Some(vec!["planner".to_string(), "coder".to_string()]),
1635 steps_exclude: vec!["planner".to_string()],
1636 ..Default::default()
1637 },
1638 ..Default::default()
1639 },
1640 );
1641 }
1642 })
1643 .await
1644 .expect("seed policy");
1645
1646 let mut payload = consumer_payload(&consumer_id, &run_id);
1647 assemble_step_pointers(&state, &mut payload).await;
1648
1649 let names: Vec<&str> = payload
1650 .context
1651 .as_ref()
1652 .expect("context")
1653 .steps
1654 .iter()
1655 .map(|p| p.name.as_str())
1656 .collect();
1657 assert_eq!(names, vec!["coder"], "names: {names:?}");
1658 }
1659
1660 #[tokio::test]
1669 async fn in_flight_step_output_is_visible_before_run_finalizes() {
1670 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1671 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1672 let task_id = TaskId::new();
1673 let run_id = RunId::new();
1674 let step1_id = StepId::new();
1675 append_final(
1676 &data_store,
1677 step1_id.as_str(),
1678 "step1",
1679 json!({"step1_out": "hi"}),
1680 )
1681 .await;
1682 let mut run = run_record(&task_id, &run_id, vec![step_entry(&step1_id, "step1")]);
1683 run.status = RunStatus::Running;
1684 run.result_ref = None; run_store.create(run).await.expect("create run");
1686
1687 let state = test_state(data_store, run_store);
1688 let consumer_id = StepId::new();
1689 let mut payload = consumer_payload(&consumer_id, &run_id);
1690 assemble_step_pointers(&state, &mut payload).await;
1691
1692 let steps = &payload.context.expect("context").steps;
1693 assert_eq!(steps.len(), 1);
1694 assert_eq!(steps[0].name, "step1");
1695 }
1696
1697 #[tokio::test]
1701 async fn self_agent_name_is_always_excluded() {
1702 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1703 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1704 let task_id = TaskId::new();
1705 let run_id = RunId::new();
1706 let planner_id = StepId::new();
1707 let consumer_prior_id = StepId::new();
1708 append_final(&data_store, planner_id.as_str(), "planner", json!("x")).await;
1709 append_final(
1710 &data_store,
1711 consumer_prior_id.as_str(),
1712 "consumer",
1713 json!("self"),
1714 )
1715 .await;
1716 run_store
1717 .create(run_record(
1718 &task_id,
1719 &run_id,
1720 vec![
1721 step_entry(&planner_id, "planner"),
1722 step_entry(&consumer_prior_id, "consumer"),
1723 ],
1724 ))
1725 .await
1726 .expect("create run");
1727
1728 let state = test_state(data_store, run_store);
1729 let consumer_id = StepId::new();
1730 let mut payload = consumer_payload(&consumer_id, &run_id);
1731 assemble_step_pointers(&state, &mut payload).await;
1732
1733 let names: Vec<&str> = payload
1734 .context
1735 .as_ref()
1736 .expect("context")
1737 .steps
1738 .iter()
1739 .map(|p| p.name.as_str())
1740 .collect();
1741 assert!(!names.contains(&"consumer"), "names: {names:?}");
1742 assert!(names.contains(&"planner"), "names: {names:?}");
1743 }
1744
1745 #[tokio::test]
1749 async fn step_pointer_serializes_with_no_preview_or_content_bytes() {
1750 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1751 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1752 let task_id = TaskId::new();
1753 let run_id = RunId::new();
1754 let planner_id = StepId::new();
1755 append_final(
1756 &data_store,
1757 planner_id.as_str(),
1758 "planner",
1759 json!({"plan": "do the thing, at length".repeat(50)}),
1760 )
1761 .await;
1762 run_store
1763 .create(run_record(
1764 &task_id,
1765 &run_id,
1766 vec![step_entry(&planner_id, "planner")],
1767 ))
1768 .await
1769 .expect("create run");
1770
1771 let state = test_state(data_store, run_store);
1772 let consumer_id = StepId::new();
1773 let mut payload = consumer_payload(&consumer_id, &run_id);
1774 assemble_step_pointers(&state, &mut payload).await;
1775
1776 let steps = &payload.context.expect("context").steps;
1777 assert_eq!(steps.len(), 1);
1778 let json_value = serde_json::to_value(&steps[0]).expect("serialize StepPointer");
1779 let obj = json_value.as_object().expect("object");
1780 for forbidden in ["preview", "content", "value", "bytes"] {
1781 assert!(
1782 !obj.contains_key(forbidden),
1783 "StepPointer must not carry a {forbidden:?} field: {obj:?}"
1784 );
1785 }
1786 assert!(obj.contains_key("name"));
1787 assert!(obj.contains_key("size_bytes"));
1788 assert!(obj.contains_key("content_url"));
1789 assert!(obj.contains_key("sha256"));
1790 }
1791
1792 fn declared_name_bp() -> mlua_swarm::blueprint::Blueprint {
1800 use mlua_flow_ir::{Expr, Node};
1801 use mlua_swarm::blueprint::{
1802 current_schema_version, AgentDef, AgentKind, AgentMeta, Blueprint, BlueprintMetadata,
1803 CompilerHints, CompilerStrategy,
1804 };
1805 Blueprint {
1806 schema_version: current_schema_version(),
1807 id: "worker-test-declared-name-bp".into(),
1808 flow: Node::Step {
1809 ref_: "planner".to_string(),
1810 in_: Expr::Path {
1811 at: "$.in".parse().expect("literal test path: $.in"),
1812 },
1813 out: Expr::Path {
1814 at: "$.plan".parse().expect("literal test path: $.plan"),
1815 },
1816 },
1817 agents: vec![AgentDef {
1818 name: "planner".to_string(),
1819 kind: AgentKind::RustFn,
1820 spec: json!({"fn_id": "planner"}),
1821 profile: None,
1822 meta: Some(AgentMeta {
1823 projection_name: Some("plan-out".to_string()),
1824 ..Default::default()
1825 }),
1826 runner: None,
1827 runner_ref: None,
1828 verdict: None,
1829 }],
1830 operators: vec![],
1831 metas: vec![],
1832 hints: CompilerHints::default(),
1833 strategy: CompilerStrategy::default(),
1834 metadata: BlueprintMetadata::default(),
1835 spawner_hints: Default::default(),
1836 default_agent_kind: AgentKind::Operator,
1837 default_operator_kind: None,
1838 default_init_ctx: None,
1839 default_agent_ctx: None,
1840 default_context_policy: None,
1841 projection_placement: None,
1842 audits: vec![],
1843 degradation_policy: None,
1844 runners: vec![],
1845 default_runner: None,
1846 subprocesses: vec![],
1847 check_policy: None,
1848 blueprint_ref_includes: Vec::new(),
1849 }
1850 }
1851
1852 #[tokio::test]
1858 async fn declared_projection_name_pointer_name_is_canonical_and_policy_matches_it() {
1859 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
1860 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1861 let task_id = TaskId::new();
1862 let run_id = RunId::new();
1863 let planner_id = StepId::new();
1864
1865 append_final(
1868 &data_store,
1869 planner_id.as_str(),
1870 "plan-out",
1871 json!({"plan": "x"}),
1872 )
1873 .await;
1874 run_store
1875 .create(run_record(
1876 &task_id,
1877 &run_id,
1878 vec![step_entry(&planner_id, "planner")],
1879 ))
1880 .await
1881 .expect("create run");
1882
1883 let state = test_state(data_store, run_store);
1884
1885 let (naming, _warnings) =
1891 mlua_swarm::core::step_naming::StepNaming::from_blueprint(&declared_name_bp())
1892 .expect("no collision");
1893 let naming = Arc::new(naming);
1894 let consumer_id = StepId::new();
1895 state
1896 .engine
1897 .with_state("test.seed_step_naming", {
1898 let naming = naming.clone();
1899 let planner_id = planner_id.clone();
1900 let consumer_id = consumer_id.clone();
1901 move |s| {
1902 s.step_namings.insert(planner_id, naming.clone());
1903 s.step_namings.insert(consumer_id, naming);
1904 }
1905 })
1906 .await
1907 .expect("seed step naming");
1908 state
1909 .engine
1910 .with_state("test.seed_policy", {
1911 let consumer_id = consumer_id.clone();
1912 move |s| {
1913 s.agent_ctx.insert(
1914 (consumer_id, 1),
1915 mlua_swarm::core::state::AgentCtxEntry {
1916 policy: mlua_swarm_schema::ContextPolicy {
1917 steps: Some(vec!["plan-out".to_string()]),
1918 ..Default::default()
1919 },
1920 ..Default::default()
1921 },
1922 );
1923 }
1924 })
1925 .await
1926 .expect("seed policy");
1927
1928 let mut payload = consumer_payload(&consumer_id, &run_id);
1929 assemble_step_pointers(&state, &mut payload).await;
1930
1931 let steps = &payload.context.expect("context").steps;
1932 assert_eq!(steps.len(), 1, "steps: {steps:?}");
1933 assert_eq!(
1934 steps[0].name, "plan-out",
1935 "StepPointer.name must be the canonical name"
1936 );
1937 }
1938
1939 async fn seed_task_with_handle(
1949 state: &AppState,
1950 task_id: &StepId,
1951 agent: &str,
1952 attempt: u32,
1953 system: Option<String>,
1954 ) -> String {
1955 let handle = format!("wh-{}", mlua_swarm::types::secure_hex(4));
1956 let task_id = task_id.clone();
1957 let agent = agent.to_string();
1958 let handle_clone = handle.clone();
1959 state
1960 .engine
1961 .with_state("test.seed_task_with_handle", move |s| {
1962 let mut task = mlua_swarm::core::state::TaskState::new(
1963 task_id.clone(),
1964 mlua_swarm::core::state::TaskSpec {
1965 agent: agent.clone(),
1966 initial_directive: json!("x"),
1967 step_ctx: None,
1968 check_policy: None,
1969 },
1970 );
1971 task.attempt = attempt;
1972 s.tasks.insert(task_id.clone(), task);
1973 s.systems.insert((task_id.clone(), attempt), system);
1974 let token = CapToken {
1975 agent_id: agent,
1976 role: mlua_swarm::Role::Worker,
1977 scopes: vec!["*".to_string()],
1978 issued_at: 0,
1979 expire_at: u64::MAX,
1980 max_uses: None,
1981 nonce: format!("test-nonce-{task_id}"),
1982 sig_hex: String::new(),
1983 };
1984 let fp = token.fingerprint();
1985 s.tokens.insert(
1986 fp.clone(),
1987 mlua_swarm::core::state::CapTokenRecord {
1988 token,
1989 uses_left: None,
1990 revoked: false,
1991 task_id: Some(task_id),
1992 },
1993 );
1994 s.worker_handles.insert(handle_clone, fp);
1995 })
1996 .await
1997 .expect("seed_task_with_handle");
1998 handle
1999 }
2000
2001 fn bearer_headers(handle: &str) -> HeaderMap {
2002 let mut headers = HeaderMap::new();
2003 headers.insert(
2004 AUTHORIZATION,
2005 format!("Bearer {handle}").parse().expect("header value"),
2006 );
2007 headers
2008 }
2009
2010 #[tokio::test]
2014 async fn worker_prompt_system_returns_raw_bytes_for_baked_system() {
2015 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2016 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2017 let state = test_state(data_store, run_store);
2018 let task_id = StepId::new();
2019 let rendered = "# Hello\n\nThis is the baked system prompt.".to_string();
2020 let handle =
2021 seed_task_with_handle(&state, &task_id, "planner", 1, Some(rendered.clone())).await;
2022
2023 let resp = worker_prompt_system(
2024 State(state.clone()),
2025 bearer_headers(&handle),
2026 Query(PromptSystemQuery {
2027 task_id: task_id.clone(),
2028 attempt: 1,
2029 }),
2030 )
2031 .await
2032 .expect("worker_prompt_system")
2033 .into_response();
2034
2035 assert_eq!(resp.status(), StatusCode::OK);
2036 let content_type = resp
2037 .headers()
2038 .get(header::CONTENT_TYPE)
2039 .expect("content-type header")
2040 .to_str()
2041 .expect("ascii");
2042 assert_eq!(content_type, "text/plain; charset=utf-8");
2043 let body_bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
2044 .await
2045 .expect("body bytes");
2046 assert_eq!(body_bytes.as_ref(), rendered.as_bytes());
2047 }
2048
2049 #[tokio::test]
2052 async fn worker_prompt_system_404s_when_no_baked_system() {
2053 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2054 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2055 let state = test_state(data_store, run_store);
2056 let task_id = StepId::new();
2057 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2058
2059 let result = worker_prompt_system(
2060 State(state.clone()),
2061 bearer_headers(&handle),
2062 Query(PromptSystemQuery {
2063 task_id: task_id.clone(),
2064 attempt: 1,
2065 }),
2066 )
2067 .await;
2068 let err = match result {
2069 Ok(_) => panic!("expected 404 ApiError, got Ok"),
2070 Err(e) => e,
2071 };
2072 assert_eq!(err.into_response().status(), StatusCode::NOT_FOUND);
2073 }
2074
2075 #[tokio::test]
2078 async fn worker_prompt_system_rejects_handle_task_mismatch() {
2079 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2080 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2081 let state = test_state(data_store, run_store);
2082 let task_id = StepId::new();
2083 let other_task_id = StepId::new();
2084 let handle =
2085 seed_task_with_handle(&state, &task_id, "planner", 1, Some("x".to_string())).await;
2086
2087 let result = worker_prompt_system(
2088 State(state.clone()),
2089 bearer_headers(&handle),
2090 Query(PromptSystemQuery {
2091 task_id: other_task_id,
2092 attempt: 1,
2093 }),
2094 )
2095 .await;
2096 let err = match result {
2097 Ok(_) => panic!("expected 400 ApiError for task mismatch, got Ok"),
2098 Err(e) => e,
2099 };
2100 assert_eq!(err.into_response().status(), StatusCode::BAD_REQUEST);
2101 }
2102
2103 #[tokio::test]
2107 async fn agent_render_size_returns_null_for_unknown_agent() {
2108 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2109 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2110 let state = test_state(data_store, run_store);
2111
2112 let Json(body) = agent_render_size(
2113 State(state.clone()),
2114 axum::extract::Path("never-dispatched".to_string()),
2115 )
2116 .await;
2117 assert_eq!(body.agent, "never-dispatched");
2118 assert_eq!(body.last_rendered_bytes, None);
2119 }
2120
2121 #[tokio::test]
2124 async fn agent_render_size_reports_last_rendered_bytes() {
2125 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2126 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2127 let state = test_state(data_store, run_store);
2128 let task_id = StepId::new();
2129 state
2130 .engine
2131 .with_state("test.seed_agent_ctx_for_bake", {
2132 let task_id = task_id.clone();
2133 move |s| {
2134 s.tasks.insert(
2135 task_id.clone(),
2136 mlua_swarm::core::state::TaskState::new(
2137 task_id,
2138 mlua_swarm::core::state::TaskSpec {
2139 agent: "coder".to_string(),
2140 initial_directive: json!("x"),
2141 step_ctx: None,
2142 check_policy: None,
2143 },
2144 ),
2145 );
2146 }
2147 })
2148 .await
2149 .expect("seed task");
2150 state
2151 .engine
2152 .bake_worker_system_prompt(&task_id, 1, Some("z".repeat(42)))
2153 .await
2154 .expect("bake_worker_system_prompt");
2155
2156 let Json(body) = agent_render_size(
2157 State(state.clone()),
2158 axum::extract::Path("coder".to_string()),
2159 )
2160 .await;
2161 assert_eq!(body.agent, "coder");
2162 assert_eq!(body.last_rendered_bytes, Some(42));
2163 }
2164
2165 #[tokio::test]
2173 async fn worker_artifact_stages_and_204s_for_valid_request() {
2174 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2175 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2176 let state = test_state(data_store, run_store);
2177 let task_id = StepId::new();
2178 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2179
2180 let status = worker_artifact(
2181 State(state.clone()),
2182 bearer_headers(&handle),
2183 Query(ArtifactQuery {
2184 name: "summary".to_string(),
2185 }),
2186 axum::body::Bytes::from_static(b"hello artifact\n"),
2187 )
2188 .await
2189 .expect("worker_artifact");
2190 assert_eq!(status, StatusCode::NO_CONTENT);
2191
2192 let tail = state.engine.output_tail(&task_id, 1).await;
2193 assert_eq!(tail.len(), 1, "tail: {tail:?}");
2194 match &tail[0] {
2195 OutputEvent::Artifact { name, content } => {
2196 assert_eq!(name, "summary");
2197 match content {
2198 ContentRef::Inline { value } => {
2199 assert_eq!(value, &json!("hello artifact"));
2200 }
2201 other => panic!("expected Inline content, got {other:?}"),
2202 }
2203 }
2204 other => panic!("expected Artifact event, got {other:?}"),
2205 }
2206 }
2207
2208 #[tokio::test]
2215 async fn worker_artifact_rejects_blank_name() {
2216 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2217 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2218 let state = test_state(data_store, run_store);
2219 let task_id = StepId::new();
2220 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2221
2222 let result = worker_artifact(
2223 State(state.clone()),
2224 bearer_headers(&handle),
2225 Query(ArtifactQuery {
2226 name: " ".to_string(),
2227 }),
2228 axum::body::Bytes::from_static(b"x"),
2229 )
2230 .await;
2231 let err = match result {
2232 Ok(_) => panic!("expected 400 ApiError for blank name, got Ok"),
2233 Err(e) => e,
2234 };
2235 assert_eq!(err.into_response().status(), StatusCode::BAD_REQUEST);
2236
2237 assert!(state.engine.output_tail(&task_id, 1).await.is_empty());
2239 }
2240
2241 #[tokio::test]
2247 async fn worker_artifact_staging_same_name_twice_appends_both_events_in_order() {
2248 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2249 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2250 let state = test_state(data_store, run_store);
2251 let task_id = StepId::new();
2252 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2253
2254 for body in [b"first".as_slice(), b"second".as_slice()] {
2255 worker_artifact(
2256 State(state.clone()),
2257 bearer_headers(&handle),
2258 Query(ArtifactQuery {
2259 name: "a".to_string(),
2260 }),
2261 axum::body::Bytes::copy_from_slice(body),
2262 )
2263 .await
2264 .expect("worker_artifact");
2265 }
2266
2267 let tail = state.engine.output_tail(&task_id, 1).await;
2268 assert_eq!(tail.len(), 2, "tail: {tail:?}");
2269 let values: Vec<&str> = tail
2270 .iter()
2271 .map(|ev| match ev {
2272 OutputEvent::Artifact {
2273 content: ContentRef::Inline { value },
2274 ..
2275 } => value.as_str().expect("string value"),
2276 other => panic!("expected Artifact/Inline event, got {other:?}"),
2277 })
2278 .collect();
2279 assert_eq!(values, vec!["first", "second"]);
2280 }
2281
2282 async fn link_task_to_run(state: &AppState, task_id: &StepId, attempt: u32, run_id: &RunId) {
2290 let tid = task_id.clone();
2291 let rid_str = run_id.to_string();
2292 state
2293 .engine
2294 .with_state("test.link_task_to_run", move |s| {
2295 let mut entry = mlua_swarm::core::state::AgentCtxEntry::default();
2296 entry.view.run_id = Some(rid_str);
2297 s.agent_ctx.insert((tid, attempt), entry);
2298 })
2299 .await
2300 .expect("link_task_to_run");
2301 }
2302
2303 #[tokio::test]
2308 async fn submit_and_artifact_against_terminal_run_return_410() {
2309 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2310 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2311 let state = test_state(data_store, run_store.clone());
2312 let task_id = StepId::new();
2313 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2314
2315 let owner_task = TaskId::new();
2316 let run_id = RunId::new();
2317 let mut rec = run_record(&owner_task, &run_id, vec![]);
2318 rec.status = RunStatus::Failed;
2319 run_store.create(rec).await.expect("run create");
2320 link_task_to_run(&state, &task_id, 1, &run_id).await;
2321
2322 let err = worker_submit(
2323 State(state.clone()),
2324 bearer_headers(&handle),
2325 Query(SubmitQuery {
2326 ok: None,
2327 verdict: None,
2328 }),
2329 axum::body::Bytes::from_static(b"LATE OUTPUT"),
2330 )
2331 .await
2332 .expect_err("a submit against a Failed run must be rejected");
2333 assert_eq!(err.status, StatusCode::GONE);
2334 assert!(
2335 err.message.contains(&run_id.to_string()),
2336 "the 410 must name the terminal run: {}",
2337 err.message
2338 );
2339
2340 let err = worker_artifact(
2341 State(state.clone()),
2342 bearer_headers(&handle),
2343 Query(ArtifactQuery {
2344 name: "part.md".to_string(),
2345 }),
2346 axum::body::Bytes::from_static(b"LATE PART"),
2347 )
2348 .await
2349 .expect_err("an artifact staged against a Failed run must be rejected");
2350 assert_eq!(err.status, StatusCode::GONE);
2351
2352 let tail = state.engine.output_tail(&task_id, 1).await;
2354 assert!(tail.is_empty(), "rejected submits must not land: {tail:?}");
2355 }
2356
2357 #[tokio::test]
2361 async fn terminal_run_guard_is_fail_open() {
2362 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2363 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2364 let state = test_state(data_store, run_store.clone());
2365 let task_id = StepId::new();
2366 seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2367
2368 reject_if_run_terminal(&state, &task_id, 1)
2370 .await
2371 .expect("no linkage must fail open");
2372
2373 let unknown_run = RunId::new();
2375 link_task_to_run(&state, &task_id, 1, &unknown_run).await;
2376 reject_if_run_terminal(&state, &task_id, 1)
2377 .await
2378 .expect("unknown run must fail open");
2379
2380 let owner_task = TaskId::new();
2382 let live_run = RunId::new();
2383 run_store
2384 .create(run_record(&owner_task, &live_run, vec![]))
2385 .await
2386 .expect("run create");
2387 link_task_to_run(&state, &task_id, 1, &live_run).await;
2388 reject_if_run_terminal(&state, &task_id, 1)
2389 .await
2390 .expect("a Running run must pass the guard");
2391 }
2392
2393 fn degradation_body(tool: &str, note: Option<&str>) -> DegradationBody {
2398 DegradationBody {
2399 tool: tool.to_string(),
2400 error: "boom".to_string(),
2401 fallback: "used cached value".to_string(),
2402 note: note.map(str::to_string),
2403 }
2404 }
2405
2406 async fn link_task_to_run_with_agent(
2412 state: &AppState,
2413 task_id: &StepId,
2414 attempt: u32,
2415 run_id: &RunId,
2416 agent: &str,
2417 ) {
2418 let tid = task_id.clone();
2419 let rid_str = run_id.to_string();
2420 let agent = agent.to_string();
2421 state
2422 .engine
2423 .with_state("test.link_task_to_run_with_agent", move |s| {
2424 let mut entry = mlua_swarm::core::state::AgentCtxEntry::default();
2425 entry.view.run_id = Some(rid_str);
2426 entry.view.agent = agent;
2427 s.agent_ctx.insert((tid, attempt), entry);
2428 })
2429 .await
2430 .expect("link_task_to_run_with_agent");
2431 }
2432
2433 #[tokio::test]
2438 async fn worker_degradation_persists_entry_when_run_tracked() {
2439 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2440 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2441 let state = test_state(data_store, run_store.clone());
2442 let task_id = StepId::new();
2443 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2444
2445 let owner_task = TaskId::new();
2446 let run_id = RunId::new();
2447 run_store
2448 .create(run_record(&owner_task, &run_id, vec![]))
2449 .await
2450 .expect("run create");
2451 link_task_to_run_with_agent(&state, &task_id, 1, &run_id, "planner").await;
2452
2453 let status = worker_degradation(
2454 State(state.clone()),
2455 bearer_headers(&handle),
2456 Json(degradation_body("web_search", Some("rate limited"))),
2457 )
2458 .await
2459 .expect("worker_degradation");
2460 assert_eq!(status, StatusCode::NO_CONTENT);
2461
2462 let rec = run_store.get(&run_id).await.expect("run get");
2463 assert_eq!(
2464 rec.degradations.len(),
2465 1,
2466 "degradations: {:?}",
2467 rec.degradations
2468 );
2469 let entry = &rec.degradations[0];
2470 assert_eq!(entry.tool, "web_search");
2471 assert_eq!(entry.error, "boom");
2472 assert_eq!(entry.fallback, "used cached value");
2473 assert_eq!(entry.note.as_deref(), Some("rate limited"));
2474 assert_eq!(entry.step_ref.as_deref(), Some("planner"));
2475 assert_eq!(entry.attempt, Some(1));
2476 assert!(entry.at > 0, "at must be a real timestamp: {}", entry.at);
2477 }
2478
2479 #[tokio::test]
2481 async fn worker_degradation_appends_in_order() {
2482 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2483 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2484 let state = test_state(data_store, run_store.clone());
2485 let task_id = StepId::new();
2486 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2487
2488 let owner_task = TaskId::new();
2489 let run_id = RunId::new();
2490 run_store
2491 .create(run_record(&owner_task, &run_id, vec![]))
2492 .await
2493 .expect("run create");
2494 link_task_to_run(&state, &task_id, 1, &run_id).await;
2495
2496 for tool in ["first_tool", "second_tool"] {
2497 worker_degradation(
2498 State(state.clone()),
2499 bearer_headers(&handle),
2500 Json(degradation_body(tool, None)),
2501 )
2502 .await
2503 .expect("worker_degradation");
2504 }
2505
2506 let rec = run_store.get(&run_id).await.expect("run get");
2507 let tools: Vec<&str> = rec.degradations.iter().map(|e| e.tool.as_str()).collect();
2508 assert_eq!(tools, vec!["first_tool", "second_tool"]);
2509 }
2510
2511 #[tokio::test]
2515 async fn worker_degradation_silent_ok_when_no_run_tracked() {
2516 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2517 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2518 let state = test_state(data_store, run_store);
2519 let task_id = StepId::new();
2520 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2521
2522 let status = worker_degradation(
2523 State(state.clone()),
2524 bearer_headers(&handle),
2525 Json(degradation_body("some_tool", None)),
2526 )
2527 .await
2528 .expect("worker_degradation must not error on missing run linkage");
2529 assert_eq!(status, StatusCode::NO_CONTENT);
2530 }
2531
2532 #[tokio::test]
2535 async fn worker_degradation_rejects_terminal_run() {
2536 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2537 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2538 let state = test_state(data_store, run_store.clone());
2539 let task_id = StepId::new();
2540 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2541
2542 let owner_task = TaskId::new();
2543 let run_id = RunId::new();
2544 let mut rec = run_record(&owner_task, &run_id, vec![]);
2545 rec.status = RunStatus::Done;
2546 run_store.create(rec).await.expect("run create");
2547 link_task_to_run(&state, &task_id, 1, &run_id).await;
2548
2549 let err = worker_degradation(
2550 State(state.clone()),
2551 bearer_headers(&handle),
2552 Json(degradation_body("some_tool", None)),
2553 )
2554 .await
2555 .expect_err("a degradation against a Done run must be rejected");
2556 assert_eq!(err.status, StatusCode::GONE);
2557
2558 let rec = run_store.get(&run_id).await.expect("run get");
2559 assert!(
2560 rec.degradations.is_empty(),
2561 "rejected degradation must not land: {:?}",
2562 rec.degradations
2563 );
2564 }
2565
2566 async fn seed_work_dir(
2580 state: &AppState,
2581 task_id: &StepId,
2582 attempt: u32,
2583 work_dir: &str,
2584 allow_file_submit: Option<Value>,
2585 ) {
2586 let tid = task_id.clone();
2587 let work_dir = work_dir.to_string();
2588 state
2589 .engine
2590 .with_state("test.seed_work_dir", move |s| {
2591 let mut entry = mlua_swarm::core::state::AgentCtxEntry::default();
2592 entry.view.work_dir = Some(work_dir);
2593 if let Some(v) = allow_file_submit {
2594 entry
2595 .view
2596 .extra
2597 .insert(FILE_SENTINEL_ALLOW_KEY.to_string(), v);
2598 }
2599 s.agent_ctx.insert((tid, attempt), entry);
2600 })
2601 .await
2602 .expect("seed_work_dir");
2603 }
2604
2605 #[tokio::test]
2609 async fn worker_submit_resolves_file_sentinel_under_work_dir() {
2610 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2611 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2612 let state = test_state(data_store.clone(), run_store);
2613 let task_id = StepId::new();
2614 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2615
2616 let tmp = tempfile::tempdir().expect("tempdir");
2617 let work_dir = tmp.path().to_path_buf();
2618 seed_work_dir(
2619 &state,
2620 &task_id,
2621 1,
2622 work_dir.to_str().expect("work_dir utf-8"),
2623 Some(Value::Bool(true)),
2624 )
2625 .await;
2626
2627 let payload_path = work_dir.join("scout.md");
2628 let payload = "## Context Package (broad)\n\nlarge body content\n";
2629 tokio::fs::write(&payload_path, payload)
2630 .await
2631 .expect("write payload");
2632 let body = format!(
2633 "@file:{}",
2634 payload_path.to_str().expect("payload path utf-8")
2635 );
2636
2637 let status = worker_submit(
2638 State(state.clone()),
2639 bearer_headers(&handle),
2640 Query(SubmitQuery {
2641 ok: None,
2642 verdict: None,
2643 }),
2644 axum::body::Bytes::from(body),
2645 )
2646 .await
2647 .expect("worker_submit sentinel");
2648 assert_eq!(status, StatusCode::NO_CONTENT);
2649
2650 let tid = task_id.clone();
2654 let value = state
2655 .engine
2656 .with_state("test.inspect_output_store", move |s| {
2657 s.output_store.get(&(tid.clone(), 1)).and_then(|evs| {
2658 evs.iter().find_map(|ev| match ev {
2659 OutputEvent::Final {
2660 content: ContentRef::Inline { value },
2661 ..
2662 } => Some(value.clone()),
2663 _ => None,
2664 })
2665 })
2666 })
2667 .await
2668 .expect("with_state")
2669 .expect("Final event present");
2670 assert_eq!(value, Value::String(payload.trim_end().to_string()));
2671 }
2672
2673 #[tokio::test]
2676 async fn worker_submit_passes_non_sentinel_body_unchanged() {
2677 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2678 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2679 let state = test_state(data_store.clone(), run_store);
2680 let task_id = StepId::new();
2681 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2682 let status = worker_submit(
2686 State(state.clone()),
2687 bearer_headers(&handle),
2688 Query(SubmitQuery {
2689 ok: None,
2690 verdict: None,
2691 }),
2692 axum::body::Bytes::from_static(b"DONE yes=1 maybe=0 no=0"),
2693 )
2694 .await
2695 .expect("worker_submit inline");
2696 assert_eq!(status, StatusCode::NO_CONTENT);
2697
2698 let tid = task_id.clone();
2699 let value = state
2700 .engine
2701 .with_state("test.inspect_output_store", move |s| {
2702 s.output_store.get(&(tid.clone(), 1)).and_then(|evs| {
2703 evs.iter().find_map(|ev| match ev {
2704 OutputEvent::Final {
2705 content: ContentRef::Inline { value },
2706 ..
2707 } => Some(value.clone()),
2708 _ => None,
2709 })
2710 })
2711 })
2712 .await
2713 .expect("with_state")
2714 .expect("Final event present");
2715 assert_eq!(value, Value::String("DONE yes=1 maybe=0 no=0".to_string()));
2716 }
2717
2718 #[tokio::test]
2723 async fn worker_submit_rejects_sentinel_path_outside_work_dir() {
2724 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2725 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2726 let state = test_state(data_store, run_store);
2727 let task_id = StepId::new();
2728 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2729
2730 let allowed = tempfile::tempdir().expect("allowed tempdir");
2731 let outside = tempfile::tempdir().expect("outside tempdir");
2732 seed_work_dir(
2733 &state,
2734 &task_id,
2735 1,
2736 allowed.path().to_str().expect("utf-8"),
2737 Some(Value::Bool(true)),
2738 )
2739 .await;
2740
2741 let outside_file = outside.path().join("leak.md");
2742 tokio::fs::write(&outside_file, b"outside content")
2743 .await
2744 .expect("write outside");
2745 let body = format!(
2746 "@file:{}",
2747 outside_file.to_str().expect("outside path utf-8")
2748 );
2749
2750 let err = worker_submit(
2751 State(state.clone()),
2752 bearer_headers(&handle),
2753 Query(SubmitQuery {
2754 ok: None,
2755 verdict: None,
2756 }),
2757 axum::body::Bytes::from(body),
2758 )
2759 .await
2760 .expect_err("outside-work_dir sentinel must be rejected");
2761 assert_eq!(err.status, StatusCode::BAD_REQUEST);
2762 }
2763
2764 #[tokio::test]
2766 async fn worker_submit_rejects_sentinel_missing_file() {
2767 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2768 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2769 let state = test_state(data_store, run_store);
2770 let task_id = StepId::new();
2771 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2772
2773 let tmp = tempfile::tempdir().expect("tempdir");
2774 seed_work_dir(
2775 &state,
2776 &task_id,
2777 1,
2778 tmp.path().to_str().expect("utf-8"),
2779 Some(Value::Bool(true)),
2780 )
2781 .await;
2782 let missing = tmp.path().join("does-not-exist.md");
2783 let body = format!("@file:{}", missing.to_str().expect("utf-8"));
2784
2785 let err = worker_submit(
2786 State(state.clone()),
2787 bearer_headers(&handle),
2788 Query(SubmitQuery {
2789 ok: None,
2790 verdict: None,
2791 }),
2792 axum::body::Bytes::from(body),
2793 )
2794 .await
2795 .expect_err("missing-file sentinel must be rejected");
2796 assert_eq!(err.status, StatusCode::NOT_FOUND);
2797 }
2798
2799 #[tokio::test]
2801 async fn worker_submit_rejects_sentinel_relative_path() {
2802 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2803 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2804 let state = test_state(data_store, run_store);
2805 let task_id = StepId::new();
2806 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2807
2808 let err = worker_submit(
2809 State(state.clone()),
2810 bearer_headers(&handle),
2811 Query(SubmitQuery {
2812 ok: None,
2813 verdict: None,
2814 }),
2815 axum::body::Bytes::from_static(b"@file:relative/path.md"),
2816 )
2817 .await
2818 .expect_err("relative-path sentinel must be rejected");
2819 assert_eq!(err.status, StatusCode::BAD_REQUEST);
2820 }
2821
2822 #[tokio::test]
2826 async fn worker_submit_rejects_sentinel_without_agent_context_view() {
2827 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2828 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2829 let state = test_state(data_store, run_store);
2830 let task_id = StepId::new();
2831 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2832 let err = worker_submit(
2835 State(state.clone()),
2836 bearer_headers(&handle),
2837 Query(SubmitQuery {
2838 ok: None,
2839 verdict: None,
2840 }),
2841 axum::body::Bytes::from_static(b"@file:/tmp/anywhere.md"),
2842 )
2843 .await
2844 .expect_err("missing AgentContextView must reject sentinel");
2845 assert_eq!(err.status, StatusCode::BAD_REQUEST);
2846 }
2847
2848 #[tokio::test]
2852 async fn worker_artifact_resolves_file_sentinel_under_work_dir() {
2853 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2854 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2855 let state = test_state(data_store, run_store);
2856 let task_id = StepId::new();
2857 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2858
2859 let tmp = tempfile::tempdir().expect("tempdir");
2860 seed_work_dir(
2861 &state,
2862 &task_id,
2863 1,
2864 tmp.path().to_str().expect("utf-8"),
2865 Some(Value::Bool(true)),
2866 )
2867 .await;
2868
2869 let payload_path = tmp.path().join("part.md");
2870 let payload = "artifact part body\n";
2871 tokio::fs::write(&payload_path, payload)
2872 .await
2873 .expect("write payload");
2874 let body = format!("@file:{}", payload_path.to_str().expect("utf-8"));
2875
2876 let status = worker_artifact(
2877 State(state.clone()),
2878 bearer_headers(&handle),
2879 Query(ArtifactQuery {
2880 name: "scout".to_string(),
2881 }),
2882 axum::body::Bytes::from(body),
2883 )
2884 .await
2885 .expect("worker_artifact sentinel");
2886 assert_eq!(status, StatusCode::NO_CONTENT);
2887 }
2888
2889 #[tokio::test]
2894 async fn worker_submit_rejects_sentinel_without_allow_flag() {
2895 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2896 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2897 let state = test_state(data_store, run_store);
2898 let task_id = StepId::new();
2899 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2900
2901 let tmp = tempfile::tempdir().expect("tempdir");
2902 seed_work_dir(
2903 &state,
2904 &task_id,
2905 1,
2906 tmp.path().to_str().expect("utf-8"),
2907 None,
2908 )
2909 .await;
2910
2911 let payload_path = tmp.path().join("out.md");
2912 tokio::fs::write(&payload_path, b"resolvable body")
2913 .await
2914 .expect("write payload");
2915 let body = format!("@file:{}", payload_path.to_str().expect("utf-8"));
2916
2917 let err = worker_submit(
2918 State(state.clone()),
2919 bearer_headers(&handle),
2920 Query(SubmitQuery {
2921 ok: None,
2922 verdict: None,
2923 }),
2924 axum::body::Bytes::from(body),
2925 )
2926 .await
2927 .expect_err("missing opt-in must reject sentinel");
2928 assert_eq!(err.status, StatusCode::BAD_REQUEST);
2929 assert!(
2930 err.message.contains("not allowed"),
2931 "rejection must name the opt-in guard, got: {}",
2932 err.message
2933 );
2934 }
2935
2936 #[tokio::test]
2939 async fn worker_submit_rejects_sentinel_with_non_true_allow_values() {
2940 for allow in [Value::Bool(false), Value::String("true".to_string())] {
2941 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
2942 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2943 let state = test_state(data_store, run_store);
2944 let task_id = StepId::new();
2945 let handle = seed_task_with_handle(&state, &task_id, "planner", 1, None).await;
2946
2947 let tmp = tempfile::tempdir().expect("tempdir");
2948 seed_work_dir(
2949 &state,
2950 &task_id,
2951 1,
2952 tmp.path().to_str().expect("utf-8"),
2953 Some(allow.clone()),
2954 )
2955 .await;
2956
2957 let payload_path = tmp.path().join("out.md");
2958 tokio::fs::write(&payload_path, b"resolvable body")
2959 .await
2960 .expect("write payload");
2961 let body = format!("@file:{}", payload_path.to_str().expect("utf-8"));
2962
2963 let err = worker_submit(
2964 State(state.clone()),
2965 bearer_headers(&handle),
2966 Query(SubmitQuery {
2967 ok: None,
2968 verdict: None,
2969 }),
2970 axum::body::Bytes::from(body),
2971 )
2972 .await
2973 .expect_err("non-true opt-in value must reject sentinel");
2974 assert_eq!(err.status, StatusCode::BAD_REQUEST, "value: {allow:?}");
2975 }
2976 }
2977
2978 fn body_verdict_contract(values: &[&str]) -> mlua_swarm_schema::VerdictContract {
2989 mlua_swarm_schema::VerdictContract {
2990 channel: VerdictChannel::Body,
2991 values: values.iter().map(|v| v.to_string()).collect(),
2992 }
2993 }
2994
2995 fn part_verdict_contract(values: &[&str]) -> mlua_swarm_schema::VerdictContract {
2996 mlua_swarm_schema::VerdictContract {
2997 channel: VerdictChannel::Part,
2998 values: values.iter().map(|v| v.to_string()).collect(),
2999 }
3000 }
3001
3002 #[tokio::test]
3005 async fn worker_submit_rejects_body_outside_contract_values_with_422() {
3006 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
3007 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3008 let state = test_state(data_store, run_store);
3009 let task_id = StepId::new();
3010 let handle = seed_task_with_handle(&state, &task_id, "gate", 1, None).await;
3011 state.engine.register_verdict_contracts(HashMap::from([(
3012 "gate".to_string(),
3013 body_verdict_contract(&["PASS", "BLOCKED"]),
3014 )]));
3015
3016 let err = worker_submit(
3017 State(state.clone()),
3018 bearer_headers(&handle),
3019 Query(SubmitQuery {
3020 ok: None,
3021 verdict: None,
3022 }),
3023 axum::body::Bytes::from("UNKNOWN"),
3024 )
3025 .await
3026 .expect_err("value outside declared values must reject");
3027 assert_eq!(err.status, StatusCode::UNPROCESSABLE_ENTITY);
3028 assert!(
3029 err.message.contains("PASS") && err.message.contains("BLOCKED"),
3030 "rejection must echo the declared values, got: {}",
3031 err.message
3032 );
3033 }
3034
3035 #[tokio::test]
3038 async fn worker_submit_accepts_body_inside_contract_values() {
3039 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
3040 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3041 let state = test_state(data_store, run_store);
3042 let task_id = StepId::new();
3043 let handle = seed_task_with_handle(&state, &task_id, "gate", 1, None).await;
3044 state.engine.register_verdict_contracts(HashMap::from([(
3045 "gate".to_string(),
3046 body_verdict_contract(&["PASS", "BLOCKED"]),
3047 )]));
3048
3049 let status = worker_submit(
3050 State(state.clone()),
3051 bearer_headers(&handle),
3052 Query(SubmitQuery {
3053 ok: None,
3054 verdict: None,
3055 }),
3056 axum::body::Bytes::from("PASS"),
3057 )
3058 .await
3059 .expect("value inside declared values must succeed");
3060 assert_eq!(status, StatusCode::NO_CONTENT);
3061 }
3062
3063 #[tokio::test]
3067 async fn worker_submit_without_a_declared_contract_is_unaffected() {
3068 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
3069 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3070 let state = test_state(data_store, run_store);
3071 let task_id = StepId::new();
3072 let handle = seed_task_with_handle(&state, &task_id, "undeclared-agent", 1, None).await;
3074
3075 let status = worker_submit(
3076 State(state.clone()),
3077 bearer_headers(&handle),
3078 Query(SubmitQuery {
3079 ok: None,
3080 verdict: None,
3081 }),
3082 axum::body::Bytes::from("anything at all, no contract to violate"),
3083 )
3084 .await
3085 .expect("no contract declared must never reject");
3086 assert_eq!(status, StatusCode::NO_CONTENT);
3087 }
3088
3089 #[tokio::test]
3092 async fn worker_artifact_verdict_part_rejects_value_outside_contract_with_422() {
3093 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
3094 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3095 let state = test_state(data_store, run_store);
3096 let task_id = StepId::new();
3097 let handle = seed_task_with_handle(&state, &task_id, "gate", 1, None).await;
3098 state.engine.register_verdict_contracts(HashMap::from([(
3099 "gate".to_string(),
3100 part_verdict_contract(&["PASS", "BLOCKED"]),
3101 )]));
3102
3103 let err = worker_artifact(
3104 State(state.clone()),
3105 bearer_headers(&handle),
3106 Query(ArtifactQuery {
3107 name: "verdict".to_string(),
3108 }),
3109 axum::body::Bytes::from("UNKNOWN"),
3110 )
3111 .await
3112 .expect_err("value outside declared values must reject");
3113 assert_eq!(err.status, StatusCode::UNPROCESSABLE_ENTITY);
3114 }
3115
3116 #[tokio::test]
3120 async fn worker_artifact_non_verdict_part_skips_the_gate() {
3121 let data_store: Arc<dyn OutputStore> = Arc::new(InMemoryOutputStore::new());
3122 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3123 let state = test_state(data_store, run_store);
3124 let task_id = StepId::new();
3125 let handle = seed_task_with_handle(&state, &task_id, "gate", 1, None).await;
3126 state.engine.register_verdict_contracts(HashMap::from([(
3127 "gate".to_string(),
3128 part_verdict_contract(&["PASS", "BLOCKED"]),
3129 )]));
3130
3131 let status = worker_artifact(
3132 State(state.clone()),
3133 bearer_headers(&handle),
3134 Query(ArtifactQuery {
3135 name: "notes".to_string(),
3136 }),
3137 axum::body::Bytes::from("anything at all"),
3138 )
3139 .await
3140 .expect("non-verdict part name must never be gated");
3141 assert_eq!(status, StatusCode::NO_CONTENT);
3142 }
3143}