1use axum::{
123 extract::{Path, Query, State},
124 http::{header, HeaderMap, HeaderValue, StatusCode},
125 response::IntoResponse,
126 Json,
127};
128use mlua_swarm::core::engine::Engine;
129use mlua_swarm::core::projection::{
130 ProjectionAdapter, ProjectionError, ProjectionKey, ProjectionRef,
131};
132use mlua_swarm::core::projection_placement::ProjectionPlacement;
133use mlua_swarm::core::step_naming::StepNaming;
134use mlua_swarm::store::output::{ContentRef, OutputEvent, OutputStore, OutputStoreError};
135use mlua_swarm::store::run::{RunRecord, RunStore};
136use mlua_swarm::{RunId, StepId, TaskId};
137use serde::{Deserialize, Serialize};
138use serde_json::Value;
139use sha2::Digest as _;
140use std::sync::Arc;
141
142use crate::tasks::map_task_store_err;
143use crate::{ApiError, AppState};
144
145pub struct McpQueryAdapter {
153 data_store: Arc<dyn OutputStore>,
154 run_store: Arc<dyn RunStore>,
155 engine: Engine,
156}
157
158#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, schemars::JsonSchema)]
162#[serde(rename_all = "snake_case")]
163pub enum ProjectionSource {
164 DataPlane,
167 ResultRef,
170}
171
172#[derive(Debug, Clone)]
178pub(crate) struct ResolvedStep {
179 pub(crate) name: String,
182 pub(crate) value: Value,
184 pub(crate) source: ProjectionSource,
186}
187
188fn final_value(event: &OutputEvent) -> Option<Value> {
194 match event {
195 OutputEvent::Final { content, .. } => Some(content_to_value(content)),
196 _ => None,
197 }
198}
199
200fn content_to_value(content: &ContentRef) -> Value {
205 match content {
206 ContentRef::Inline { value } => value.clone(),
207 ContentRef::FileRef {
208 path,
209 mime,
210 size_hint,
211 } => serde_json::json!({
212 "file_ref": path.to_string_lossy(),
213 "mime": mime,
214 "size_hint": size_hint,
215 }),
216 }
217}
218
219fn find_step_id_for_canonical(
228 run: &RunRecord,
229 naming: Option<&StepNaming>,
230 canonical: &str,
231) -> Option<StepId> {
232 run.step_entries
233 .iter()
234 .rev()
235 .find(|entry| {
236 let Some(step_ref) = entry.step_ref.as_deref() else {
237 return false;
238 };
239 match naming {
240 Some(n) => n.canonical_of_producer(step_ref) == Some(canonical),
241 None => step_ref == canonical,
242 }
243 })
244 .map(|entry| entry.step_id.clone())
245}
246
247fn candidate_names<'a>(
256 naming: Option<&'a StepNaming>,
257 canonical: &'a str,
258 raw_step: &'a str,
259) -> Vec<&'a str> {
260 let mut names = vec![canonical];
261 if let Some(entry) = naming.and_then(|n| n.entries().find(|e| e.canonical == canonical)) {
262 for alias in &entry.aliases {
263 if alias != canonical {
264 names.push(alias.as_str());
265 }
266 }
267 }
268 if !names.contains(&raw_step) {
269 names.push(raw_step);
270 }
271 names
272}
273
274impl McpQueryAdapter {
275 pub fn new(
280 data_store: Arc<dyn OutputStore>,
281 run_store: Arc<dyn RunStore>,
282 engine: Engine,
283 ) -> Self {
284 Self {
285 data_store,
286 run_store,
287 engine,
288 }
289 }
290
291 async fn step_naming_for_run(&self, run: &RunRecord) -> Option<Arc<StepNaming>> {
307 resolve_step_naming_for_run(&self.engine, run).await
308 }
309
310 pub(crate) async fn resolve_step_name(&self, run: &RunRecord, raw: &str) -> String {
318 match self.step_naming_for_run(run).await {
319 Some(naming) => naming.resolve(raw).unwrap_or(raw).to_string(),
320 None => raw.to_string(),
321 }
322 }
323
324 async fn resolve_run(
332 &self,
333 task_id: &TaskId,
334 run_id: Option<&str>,
335 ) -> Result<RunRecord, ProjectionError> {
336 match run_id {
337 Some(rid) => {
338 let run_id = RunId::parse(rid.to_string())
339 .map_err(|e| ProjectionError::InvalidKey(format!("run_id: {e}")))?;
340 let run = self.run_store.get(&run_id).await.map_err(|_| {
341 ProjectionError::NotFound(ProjectionKey {
342 task_id: task_id.to_string(),
343 run_id: Some(rid.to_string()),
344 step: None,
345 path: None,
346 })
347 })?;
348 if &run.task_id != task_id {
349 return Err(ProjectionError::NotFound(ProjectionKey {
350 task_id: task_id.to_string(),
351 run_id: Some(rid.to_string()),
352 step: None,
353 path: None,
354 }));
355 }
356 Ok(run)
357 }
358 None => {
359 let mut runs = self.run_store.list_by_task(task_id).await.map_err(|_| {
360 ProjectionError::NotFound(ProjectionKey {
361 task_id: task_id.to_string(),
362 run_id: None,
363 step: None,
364 path: None,
365 })
366 })?;
367 runs.pop().ok_or_else(|| {
368 ProjectionError::NotFound(ProjectionKey {
369 task_id: task_id.to_string(),
370 run_id: None,
371 step: None,
372 path: None,
373 })
374 })
375 }
376 }
377 }
378
379 async fn resolve_async(
399 &self,
400 key: &ProjectionKey,
401 ) -> Result<(RunRecord, Value), ProjectionError> {
402 let task_id = TaskId::parse(key.task_id.clone())
403 .map_err(|e| ProjectionError::InvalidKey(format!("task_id: {e}")))?;
404 let run = self.resolve_run(&task_id, key.run_id.as_deref()).await?;
405
406 let Some(raw_step) = &key.step else {
407 let ctx_data = run.result_ref.clone().unwrap_or(Value::Null);
410 let value = key
411 .resolve(&ctx_data)
412 .cloned()
413 .ok_or_else(|| ProjectionError::NotFound(key.clone()))?;
414 return Ok((run, value));
415 };
416
417 let naming = self.step_naming_for_run(&run).await;
418 let canonical = naming
419 .as_deref()
420 .and_then(|n| n.resolve(raw_step))
421 .unwrap_or(raw_step.as_str())
422 .to_string();
423
424 if let Some(step_id) = find_step_id_for_canonical(&run, naming.as_deref(), &canonical) {
426 match self
427 .data_store
428 .get_latest_by_name_in_run(step_id.as_str(), 1, &canonical)
429 .await
430 {
431 Ok(record) => {
432 if let Some(value) = final_value(&record.event) {
433 let narrowed = match &key.path {
434 None => Some(value),
435 Some(_) => {
436 let path_only = ProjectionKey {
442 task_id: key.task_id.clone(),
443 run_id: key.run_id.clone(),
444 step: None,
445 path: key.path.clone(),
446 };
447 path_only.resolve(&value).cloned()
448 }
449 };
450 if let Some(value) = narrowed {
451 return Ok((run, value));
452 }
453 }
454 }
455 Err(OutputStoreError::NotFound(_)) => {
456 }
459 Err(other) => {
460 return Err(ProjectionError::Io(std::io::Error::other(format!(
461 "OutputStore::get_latest_by_name_in_run: {other}"
462 ))));
463 }
464 }
465 }
466
467 let ctx_data = run.result_ref.clone().unwrap_or(Value::Null);
472 for candidate in candidate_names(naming.as_deref(), &canonical, raw_step) {
473 let candidate_key = ProjectionKey {
474 task_id: key.task_id.clone(),
475 run_id: key.run_id.clone(),
476 step: Some(candidate.to_string()),
477 path: key.path.clone(),
478 };
479 if let Some(value) = candidate_key.resolve(&ctx_data) {
480 return Ok((run, value.clone()));
481 }
482 }
483 Err(ProjectionError::NotFound(key.clone()))
484 }
485
486 pub(crate) async fn list_steps(
492 &self,
493 task_id: &TaskId,
494 run_id: Option<&str>,
495 ) -> Result<(RunRecord, Vec<ResolvedStep>), ProjectionError> {
496 let run = self.resolve_run(task_id, run_id).await?;
497 let steps = self.enumerate_steps(&run).await;
498 Ok((run, steps))
499 }
500
501 pub(crate) async fn list_steps_by_run_id(
509 &self,
510 run_id: &RunId,
511 ) -> Result<(RunRecord, Vec<ResolvedStep>), ProjectionError> {
512 let run = self.run_store.get(run_id).await.map_err(|_| {
513 ProjectionError::NotFound(ProjectionKey {
514 task_id: String::new(),
515 run_id: Some(run_id.to_string()),
516 step: None,
517 path: None,
518 })
519 })?;
520 let steps = self.enumerate_steps(&run).await;
521 Ok((run, steps))
522 }
523
524 async fn enumerate_steps(&self, run: &RunRecord) -> Vec<ResolvedStep> {
531 match self.step_naming_for_run(run).await {
532 Some(naming) => self.enumerate_steps_via_table(run, &naming).await,
533 None => self.enumerate_steps_legacy_union(run).await,
534 }
535 }
536
537 async fn enumerate_steps_via_table(
556 &self,
557 run: &RunRecord,
558 naming: &StepNaming,
559 ) -> Vec<ResolvedStep> {
560 let mut resolved: std::collections::BTreeMap<String, ResolvedStep> =
561 std::collections::BTreeMap::new();
562
563 for entry in &run.step_entries {
564 let Some(step_ref) = entry.step_ref.as_deref() else {
565 continue;
566 };
567 let canonical = naming
568 .canonical_of_producer(step_ref)
569 .unwrap_or(step_ref)
570 .to_string();
571 if let Ok(record) = self
572 .data_store
573 .get_latest_by_name_in_run(entry.step_id.as_str(), 1, &canonical)
574 .await
575 {
576 if let Some(value) = final_value(&record.event) {
577 resolved.insert(
578 canonical.clone(),
579 ResolvedStep {
580 name: canonical,
581 value,
582 source: ProjectionSource::DataPlane,
583 },
584 );
585 }
586 }
587
588 if let Ok(records) = self
604 .data_store
605 .list_for_attempt(entry.step_id.as_str(), 1)
606 .await
607 {
608 for record in records {
609 if let OutputEvent::Artifact { name, content } = &record.event {
610 resolved
611 .entry(name.clone())
612 .or_insert_with(|| ResolvedStep {
613 name: name.clone(),
614 value: content_to_value(content),
615 source: ProjectionSource::DataPlane,
616 });
617 }
618 }
619 }
620 }
621
622 if let Some(Value::Object(map)) = &run.result_ref {
623 for entry in naming.entries() {
624 if resolved.contains_key(&entry.canonical) {
625 continue;
626 }
627 let hit = entry
628 .aliases
629 .iter()
630 .find_map(|alias| map.get(alias))
631 .or_else(|| map.get(&entry.canonical));
632 if let Some(value) = hit {
633 resolved.insert(
634 entry.canonical.clone(),
635 ResolvedStep {
636 name: entry.canonical.clone(),
637 value: value.clone(),
638 source: ProjectionSource::ResultRef,
639 },
640 );
641 }
642 }
643 }
644
645 resolved.into_values().collect()
646 }
647
648 async fn enumerate_steps_legacy_union(&self, run: &RunRecord) -> Vec<ResolvedStep> {
659 let mut out = Vec::new();
660 let mut attempted = std::collections::HashSet::new();
661 let mut resolved_names = std::collections::HashSet::new();
662
663 for entry in &run.step_entries {
664 let Some(name) = &entry.step_ref else {
665 continue;
666 };
667 if !attempted.insert(name.clone()) {
668 continue;
669 }
670 if let Ok(record) = self.data_store.get_latest_by_name(name).await {
671 if let Some(value) = final_value(&record.event) {
672 out.push(ResolvedStep {
673 name: name.clone(),
674 value,
675 source: ProjectionSource::DataPlane,
676 });
677 resolved_names.insert(name.clone());
678 }
679 }
680 }
681
682 if let Some(Value::Object(map)) = &run.result_ref {
683 for (name, value) in map {
684 if resolved_names.contains(name) {
685 continue;
686 }
687 out.push(ResolvedStep {
688 name: name.clone(),
689 value: value.clone(),
690 source: ProjectionSource::ResultRef,
691 });
692 }
693 }
694
695 out
696 }
697}
698
699impl ProjectionAdapter for McpQueryAdapter {
700 fn name(&self) -> &'static str {
701 "mcp-query"
702 }
703
704 fn project(
713 &self,
714 key: &ProjectionKey,
715 ctx_data: &Value,
716 ) -> Result<ProjectionRef, ProjectionError> {
717 if key.task_id.is_empty() {
718 return Err(ProjectionError::InvalidKey(
719 "task_id must not be empty".to_string(),
720 ));
721 }
722 key.resolve(ctx_data)
723 .ok_or_else(|| ProjectionError::NotFound(key.clone()))?;
724 Ok(ProjectionRef::Query {
725 endpoint: format!(
726 "/v1/tasks/{}/runs/{}/steps/{}/content",
727 key.task_id,
728 key.run_id.as_deref().unwrap_or("latest"),
729 key.step.as_deref().unwrap_or("_ctx")
730 ),
731 key: key.clone(),
732 })
733 }
734
735 fn fetch(&self, key: &ProjectionKey) -> Result<Value, ProjectionError> {
736 let handle = tokio::runtime::Handle::try_current().map_err(|e| {
741 ProjectionError::Io(std::io::Error::other(format!(
742 "McpQueryAdapter::fetch requires a Tokio runtime: {e}"
743 )))
744 })?;
745 let (_run, value) =
746 tokio::task::block_in_place(|| handle.block_on(self.resolve_async(key)))?;
747 Ok(value)
748 }
749
750 fn pointer_line(&self, r: &ProjectionRef) -> String {
751 match r {
752 ProjectionRef::Query { endpoint, key } => {
753 format!("projection(mcp-query): {endpoint} task_id={}", key.task_id)
754 }
755 ProjectionRef::File { path } => format!("projection(file): {path}"),
756 }
757 }
758}
759
760#[derive(Debug, Clone, Serialize, schemars::JsonSchema)]
766pub struct StepList {
767 pub task_id: String,
769 pub run_id: String,
772 pub steps: Vec<StepSummary>,
776}
777
778#[derive(Debug, Clone, Serialize, schemars::JsonSchema)]
782pub struct StepSummary {
783 pub name: String,
785 pub size_bytes: u64,
788 pub content_type: String,
793 pub sha256: String,
796 pub source: ProjectionSource,
798 #[serde(default, skip_serializing_if = "Option::is_none")]
806 pub file_path: Option<String>,
807 pub content_url: String,
812 pub preview: String,
815 pub truncated: bool,
819}
820
821#[derive(Debug, Deserialize, Default, schemars::JsonSchema)]
826pub struct StepPathQuery {
827 #[serde(default)]
830 pub path: Option<String>,
831}
832
833fn narrow_step_value(value: &Value, path: Option<&str>) -> Option<Value> {
837 match path {
838 None => Some(value.clone()),
839 Some(p) => {
840 let path_only = ProjectionKey {
841 task_id: String::new(),
842 run_id: None,
843 step: None,
844 path: Some(p.to_string()),
845 };
846 path_only.resolve(value).cloned()
847 }
848 }
849}
850
851fn materialized_file_path(
858 placement: &ProjectionPlacement,
859 root: &str,
860 step_id: &StepId,
861 name: &str,
862) -> std::path::PathBuf {
863 placement.target_path(root, step_id.as_ref(), name)
864}
865
866async fn resolve_materialized_file(
886 state: &AppState,
887 run: &RunRecord,
888 name: &str,
889) -> Option<(std::path::PathBuf, Vec<u8>)> {
890 let naming = resolve_step_naming_for_run(&state.engine, run).await;
891 let step_id = find_step_id_for_canonical(run, naming.as_deref(), name)?;
892 let view = state.engine.agent_context_for(&step_id, 1).await?;
893 let placement = state
894 .engine
895 .projection_placement_for(&step_id)
896 .await
897 .unwrap_or_default();
898 let root = placement.resolve_root(&view)?;
899 let path = materialized_file_path(&placement, &root, &step_id, name);
900 let bytes = std::fs::read(&path).ok()?;
901 Some((path, bytes))
902}
903
904async fn resolve_step_naming_for_run(engine: &Engine, run: &RunRecord) -> Option<Arc<StepNaming>> {
909 for entry in &run.step_entries {
910 if let Some(naming) = engine.step_naming_for(&entry.step_id).await {
911 return Some(naming);
912 }
913 }
914 None
915}
916
917async fn render_step_body(
924 state: &AppState,
925 run: &RunRecord,
926 step: &ResolvedStep,
927 path: Option<&str>,
928) -> Option<(Vec<u8>, &'static str, Option<String>)> {
929 if path.is_none() {
930 if let Some((file_path, bytes)) = resolve_materialized_file(state, run, &step.name).await {
931 return Some((
932 bytes,
933 "text/markdown; charset=utf-8",
934 Some(file_path.to_string_lossy().into_owned()),
935 ));
936 }
937 }
938 let narrowed = narrow_step_value(&step.value, path)?;
939 let body = serde_json::to_vec_pretty(&narrowed).ok()?;
940 Some((body, "application/json", None))
941}
942
943fn build_preview(body: &[u8]) -> (String, bool) {
950 const MAX_PREVIEW_BYTES: usize = 512;
951 if body.len() <= MAX_PREVIEW_BYTES {
952 return (String::from_utf8_lossy(body).into_owned(), false);
953 }
954 let preview = match std::str::from_utf8(body) {
955 Ok(s) => {
956 let mut end = MAX_PREVIEW_BYTES;
957 while end > 0 && !s.is_char_boundary(end) {
958 end -= 1;
959 }
960 s[..end].to_string()
961 }
962 Err(_) => String::from_utf8_lossy(&body[..MAX_PREVIEW_BYTES]).into_owned(),
963 };
964 (format!("{preview}…"), true)
965}
966
967fn build_content_url(
973 base_url: &Option<Arc<str>>,
974 task_id: &TaskId,
975 run_id: &RunId,
976 name: &str,
977 path: Option<&str>,
978) -> String {
979 let mut url = format!("/v1/tasks/{task_id}/runs/{run_id}/steps/{name}/content");
980 if let Some(p) = path {
981 url.push_str("?path=");
982 url.push_str(p);
983 }
984 match base_url {
985 Some(base) => format!("{}{}", base.trim_end_matches('/'), url),
986 None => url,
987 }
988}
989
990async fn build_step_summary(
993 state: &AppState,
994 run: &RunRecord,
995 step: &ResolvedStep,
996 path: Option<&str>,
997) -> Option<StepSummary> {
998 let (body, content_type, file_path) = render_step_body(state, run, step, path).await?;
999 let sha256 = hex::encode(sha2::Sha256::digest(&body));
1000 let size_bytes = body.len() as u64;
1001 let (preview, truncated) = build_preview(&body);
1002 let content_url = build_content_url(&state.base_url, &run.task_id, &run.id, &step.name, path);
1003 Some(StepSummary {
1004 name: step.name.clone(),
1005 size_bytes,
1006 content_type: content_type.to_string(),
1007 sha256,
1008 source: step.source,
1009 file_path,
1010 content_url,
1011 preview,
1012 truncated,
1013 })
1014}
1015
1016pub(crate) async fn resolve_step_pointer_fields(
1027 state: &AppState,
1028 run: &RunRecord,
1029 step: &ResolvedStep,
1030) -> Option<(u64, Option<String>, String, String)> {
1031 let (body, _content_type, file_path) = render_step_body(state, run, step, None).await?;
1032 let sha256 = hex::encode(sha2::Sha256::digest(&body));
1033 let size_bytes = body.len() as u64;
1034 let content_url = build_content_url(&state.base_url, &run.task_id, &run.id, &step.name, None);
1035 Some((size_bytes, file_path, content_url, sha256))
1036}
1037
1038async fn resolve_run_and_steps(
1047 state: &AppState,
1048 id: &str,
1049 run: &str,
1050) -> Result<(McpQueryAdapter, RunRecord, Vec<ResolvedStep>), ApiError> {
1051 let task_id = TaskId::parse(id.to_string())
1052 .map_err(|e| ApiError::bad_request(format!("invalid task id: {e}")))?;
1053 state
1054 .task_store
1055 .get(&task_id)
1056 .await
1057 .map_err(map_task_store_err)?;
1058 let adapter = McpQueryAdapter::new(
1059 state.data_store.clone(),
1060 state.run_store.clone(),
1061 state.engine.clone(),
1062 );
1063 let run_sel = if run == "latest" { None } else { Some(run) };
1064 let (run_record, steps) = adapter
1065 .list_steps(&task_id, run_sel)
1066 .await
1067 .map_err(map_projection_err)?;
1068 Ok((adapter, run_record, steps))
1069}
1070
1071pub async fn steps_list(
1074 State(state): State<AppState>,
1075 Path((id, run)): Path<(String, String)>,
1076) -> Result<Json<StepList>, ApiError> {
1077 let (_adapter, run_record, steps) = resolve_run_and_steps(&state, &id, &run).await?;
1078 let mut summaries = Vec::with_capacity(steps.len());
1079 for step in &steps {
1080 if let Some(summary) = build_step_summary(&state, &run_record, step, None).await {
1081 summaries.push(summary);
1082 }
1083 }
1084 Ok(Json(StepList {
1085 task_id: run_record.task_id.to_string(),
1086 run_id: run_record.id.to_string(),
1087 steps: summaries,
1088 }))
1089}
1090
1091pub async fn step_get(
1097 State(state): State<AppState>,
1098 Path((id, run, step)): Path<(String, String, String)>,
1099 Query(q): Query<StepPathQuery>,
1100) -> Result<Json<StepSummary>, ApiError> {
1101 let (adapter, run_record, steps) = resolve_run_and_steps(&state, &id, &run).await?;
1102 let canonical = adapter.resolve_step_name(&run_record, &step).await;
1103 let resolved = steps
1104 .into_iter()
1105 .find(|s| s.name == canonical)
1106 .ok_or_else(|| ApiError::not_found(format!("step not found: {step}")))?;
1107 let summary = build_step_summary(&state, &run_record, &resolved, q.path.as_deref())
1108 .await
1109 .ok_or_else(|| ApiError::not_found(format!("path not found: {:?}", q.path)))?;
1110 Ok(Json(summary))
1111}
1112
1113pub async fn step_content(
1119 State(state): State<AppState>,
1120 Path((id, run, step)): Path<(String, String, String)>,
1121 Query(q): Query<StepPathQuery>,
1122) -> Result<impl IntoResponse, ApiError> {
1123 let (adapter, run_record, steps) = resolve_run_and_steps(&state, &id, &run).await?;
1124 let canonical = adapter.resolve_step_name(&run_record, &step).await;
1125 let resolved = steps
1126 .into_iter()
1127 .find(|s| s.name == canonical)
1128 .ok_or_else(|| ApiError::not_found(format!("step not found: {step}")))?;
1129 let (body, content_type, _file_path) =
1130 render_step_body(&state, &run_record, &resolved, q.path.as_deref())
1131 .await
1132 .ok_or_else(|| ApiError::not_found(format!("path not found: {:?}", q.path)))?;
1133 let sha256 = hex::encode(sha2::Sha256::digest(&body));
1134 let mut headers = HeaderMap::new();
1135 headers.insert(
1136 header::CONTENT_TYPE,
1137 HeaderValue::from_str(content_type).expect("content_type is a static ASCII literal"),
1138 );
1139 headers.insert(
1140 header::ETAG,
1141 HeaderValue::from_str(&format!("\"sha256:{sha256}\""))
1142 .expect("hex digest is ASCII-safe for a header value"),
1143 );
1144 Ok((StatusCode::OK, headers, body))
1145}
1146
1147fn map_projection_err(e: ProjectionError) -> ApiError {
1148 match e {
1149 ProjectionError::NotFound(key) => {
1150 ApiError::not_found(format!("projection not found for key {key:?}"))
1151 }
1152 ProjectionError::InvalidKey(msg) => ApiError::bad_request(msg),
1153 other => ApiError::engine(other),
1154 }
1155}
1156
1157#[cfg(test)]
1162mod tests {
1163 use super::*;
1164 use crate::TaskLaunchRequest;
1165 use axum::http::StatusCode;
1166 use mlua_swarm::application::BlueprintRef;
1167 use mlua_swarm::blueprint::{
1168 current_schema_version, AgentDef, AgentKind, AgentMeta, Blueprint, BlueprintMetadata,
1169 CompilerHints, CompilerStrategy, ProjectionPlacementSpec,
1170 };
1171 use mlua_swarm::core::config::{CheckPolicy, EngineCfg};
1172 use mlua_swarm::core::engine::Engine;
1173 use mlua_swarm::store::output::InMemoryOutputStore;
1174 use mlua_swarm::store::run::InMemoryRunStore;
1175 use mlua_swarm::store::task::InMemoryTaskStore;
1176 use serde_json::json;
1177 use std::collections::HashMap;
1178 use tokio::sync::Mutex;
1179
1180 fn greeting_blueprint() -> Blueprint {
1188 Blueprint {
1189 schema_version: current_schema_version(),
1190 id: "projection-test-greeting-bp".into(),
1191 flow: serde_json::from_value(json!({
1192 "kind": "step",
1193 "ref": mlua_swarm::worker::baseline::AG_IDENTITY,
1194 "in": {"op": "path", "at": "$.greeting"},
1195 "out": {"op": "path", "at": "$.out"},
1196 }))
1197 .expect("flow parse"),
1198 agents: vec![AgentDef {
1199 name: mlua_swarm::worker::baseline::AG_IDENTITY.into(),
1200 kind: AgentKind::RustFn,
1201 spec: json!({"fn_id": mlua_swarm::worker::baseline::AG_IDENTITY}),
1202 profile: None,
1203 meta: None,
1204 runner: None,
1205 runner_ref: None,
1206 verdict: None,
1207 lints: None,
1208 }],
1209 operators: vec![],
1210 metas: vec![],
1211 hints: CompilerHints::default(),
1212 strategy: CompilerStrategy::default(),
1213 metadata: BlueprintMetadata::default(),
1214 spawner_hints: Default::default(),
1215 default_agent_kind: AgentKind::Operator,
1216 default_operator_kind: None,
1217 default_init_ctx: None,
1218 default_agent_ctx: None,
1219 default_context_policy: None,
1220 projection_placement: None,
1221 audits: vec![],
1222 degradation_policy: None,
1223 runners: vec![],
1224 default_runner: None,
1225 subprocesses: vec![],
1226 check_policy: None,
1227 blueprint_ref_includes: Vec::new(),
1228 }
1229 }
1230
1231 fn test_state() -> AppState {
1232 let engine = Engine::new_with_layers(EngineCfg::default(), crate::default_layer_registry());
1233 let compiler = mlua_swarm::Compiler::new(crate::default_registry());
1234 let launch = Arc::new(mlua_swarm::TaskLaunchService::new(engine.clone(), compiler));
1235 let data_store: Arc<dyn mlua_swarm::store::output::OutputStore> =
1236 Arc::new(InMemoryOutputStore::new());
1237 engine.set_output_store(data_store.clone());
1243 AppState {
1244 engine,
1245 sessions: Arc::new(Mutex::new(crate::SessionStore::default())),
1246 task_app: Arc::new(mlua_swarm::TaskApplication::new_inline_only(launch)),
1247 ws_operator_factory: None,
1248 data_store,
1249 operator_sessions: Arc::new(Mutex::new(HashMap::new())),
1250 roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
1251 task_store: Arc::new(InMemoryTaskStore::new()),
1252 run_store: Arc::new(InMemoryRunStore::new()),
1253 replay_store: Arc::new(mlua_swarm::store::replay::InMemoryReplayStore::new()),
1254 run_trace_store: Arc::new(mlua_swarm::store::trace::InMemoryRunTraceStore::new()),
1255 base_url: None,
1256 sync_timeout_secs: 300,
1257 }
1258 }
1259
1260 fn greeting_task_req(greeting: &str) -> TaskLaunchRequest {
1261 TaskLaunchRequest {
1262 blueprint: BlueprintRef::Inline {
1263 value: Box::new(greeting_blueprint()),
1264 },
1265 init_ctx: json!({ "greeting": greeting }),
1266 project_root: None,
1267 work_dir: None,
1268 task_metadata: None,
1269 ttl_secs: None,
1270 operator: None,
1271 operator_sid: None,
1272 timeout_secs: None,
1273 goal: Some("projection test goal".to_string()),
1274 detach: false,
1275 check_policy: None,
1276 }
1277 }
1278
1279 fn declared_projection_name_blueprint(projection_name: &str) -> Blueprint {
1285 Blueprint {
1286 schema_version: current_schema_version(),
1287 id: "projection-test-declared-name-bp".into(),
1288 flow: serde_json::from_value(json!({
1289 "kind": "step",
1290 "ref": mlua_swarm::worker::baseline::AG_IDENTITY,
1291 "in": {"op": "path", "at": "$.greeting"},
1292 "out": {"op": "path", "at": "$.out"},
1293 }))
1294 .expect("flow parse"),
1295 agents: vec![AgentDef {
1296 name: mlua_swarm::worker::baseline::AG_IDENTITY.into(),
1297 kind: AgentKind::RustFn,
1298 spec: json!({"fn_id": mlua_swarm::worker::baseline::AG_IDENTITY}),
1299 profile: None,
1300 meta: Some(AgentMeta {
1301 projection_name: Some(projection_name.to_string()),
1302 ..Default::default()
1303 }),
1304 runner: None,
1305 runner_ref: None,
1306 verdict: None,
1307 lints: None,
1308 }],
1309 operators: vec![],
1310 metas: vec![],
1311 hints: CompilerHints::default(),
1312 strategy: CompilerStrategy::default(),
1313 metadata: BlueprintMetadata::default(),
1314 spawner_hints: Default::default(),
1315 default_agent_kind: AgentKind::Operator,
1316 default_operator_kind: None,
1317 default_init_ctx: None,
1318 default_agent_ctx: None,
1319 default_context_policy: None,
1320 projection_placement: None,
1321 audits: vec![],
1322 degradation_policy: None,
1323 runners: vec![],
1324 default_runner: None,
1325 subprocesses: vec![],
1326 check_policy: None,
1327 blueprint_ref_includes: Vec::new(),
1328 }
1329 }
1330
1331 fn declared_task_req(greeting: &str, projection_name: &str) -> TaskLaunchRequest {
1332 TaskLaunchRequest {
1333 blueprint: BlueprintRef::Inline {
1334 value: Box::new(declared_projection_name_blueprint(projection_name)),
1335 },
1336 init_ctx: json!({ "greeting": greeting }),
1337 project_root: None,
1338 work_dir: None,
1339 task_metadata: None,
1340 ttl_secs: None,
1341 operator: None,
1342 operator_sid: None,
1343 timeout_secs: None,
1344 goal: Some("projection test goal (declared name)".to_string()),
1345 detach: false,
1346 check_policy: None,
1347 }
1348 }
1349
1350 fn strict_greeting_blueprint() -> Blueprint {
1359 Blueprint {
1360 check_policy: Some(CheckPolicy::Strict),
1361 ..greeting_blueprint()
1362 }
1363 }
1364
1365 fn strict_greeting_task_req(greeting: &str, project_root: Option<&str>) -> TaskLaunchRequest {
1366 TaskLaunchRequest {
1367 blueprint: BlueprintRef::Inline {
1368 value: Box::new(strict_greeting_blueprint()),
1369 },
1370 init_ctx: json!({ "greeting": greeting }),
1371 project_root: project_root.map(String::from),
1372 work_dir: None,
1373 task_metadata: None,
1374 ttl_secs: None,
1375 operator: None,
1376 operator_sid: None,
1377 timeout_secs: None,
1378 goal: Some("pre-dispatch guard test goal".to_string()),
1379 detach: false,
1380 check_policy: None,
1381 }
1382 }
1383
1384 #[tokio::test]
1391 async fn tasks_start_rejects_strict_launch_with_no_roots_as_400() {
1392 let state = test_state();
1393 let result =
1397 crate::tasks_start(State(state), Json(strict_greeting_task_req("hello", None))).await;
1398 let err = match result {
1399 Err(e) => e,
1400 Ok(_) => {
1401 panic!("strict check_policy + no roots must be rejected before dispatch, got Ok")
1402 }
1403 };
1404 assert_eq!(err.status, StatusCode::BAD_REQUEST);
1405 assert!(
1406 err.message.contains("pre-dispatch"),
1407 "expected a pre-dispatch guard message, got: {}",
1408 err.message
1409 );
1410 }
1411
1412 #[tokio::test]
1417 async fn tasks_start_strict_launch_with_project_root_succeeds() {
1418 let state = test_state();
1419 let reply = crate::tasks_start(
1420 State(state),
1421 Json(strict_greeting_task_req("hello", Some("/repo"))),
1422 )
1423 .await
1424 .expect("strict check_policy + project_root supplied must pass the guard");
1425 assert_eq!(reply.1, StatusCode::OK);
1426 assert_eq!(reply.0.final_ctx["out"]["echoed"], "hello");
1427 }
1428
1429 #[tokio::test]
1439 async fn steps_list_undeclared_step_resolves_to_single_canonical_entry() {
1440 let state = test_state();
1441 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("hello")))
1442 .await
1443 .expect("tasks_start")
1444 .0;
1445
1446 let resp = steps_list(
1447 State(state.clone()),
1448 Path((posted.task_id.to_string(), "latest".to_string())),
1449 )
1450 .await
1451 .expect("steps_list")
1452 .0;
1453
1454 assert_eq!(resp.task_id, posted.task_id.to_string());
1455 assert_eq!(resp.run_id, posted.run_id.to_string());
1456 let identity_name = mlua_swarm::worker::baseline::AG_IDENTITY;
1457 assert_eq!(resp.steps.len(), 1, "steps: {:?}", resp.steps);
1458 let entry = &resp.steps[0];
1459 assert_eq!(entry.name, identity_name);
1460 assert_eq!(entry.source, ProjectionSource::DataPlane);
1461 }
1462
1463 #[tokio::test]
1468 async fn step_get_resolves_alias_name_to_canonical_entry() {
1469 let state = test_state();
1470 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("hi")))
1471 .await
1472 .expect("tasks_start")
1473 .0;
1474
1475 let identity_name = mlua_swarm::worker::baseline::AG_IDENTITY;
1476 let via_ref = step_get(
1477 State(state.clone()),
1478 Path((
1479 posted.task_id.to_string(),
1480 "latest".to_string(),
1481 identity_name.to_string(),
1482 )),
1483 Query(StepPathQuery::default()),
1484 )
1485 .await
1486 .expect("step_get via own ref name")
1487 .0;
1488 let via_alias = step_get(
1489 State(state.clone()),
1490 Path((
1491 posted.task_id.to_string(),
1492 "latest".to_string(),
1493 "out".to_string(),
1494 )),
1495 Query(StepPathQuery::default()),
1496 )
1497 .await
1498 .expect("step_get via out-top alias")
1499 .0;
1500
1501 assert_eq!(via_ref.name, identity_name);
1502 assert_eq!(
1503 via_alias.name, identity_name,
1504 "alias lookup must report the canonical name"
1505 );
1506 assert_eq!(
1507 via_ref.sha256, via_alias.sha256,
1508 "same OUTPUT regardless of which name was queried"
1509 );
1510 }
1511
1512 #[tokio::test]
1517 async fn declared_projection_name_e2e_resolves_via_canonical_and_alias() {
1518 let state = test_state();
1519 let posted = crate::tasks_start(
1520 State(state.clone()),
1521 Json(declared_task_req("hi", "plan-out")),
1522 )
1523 .await
1524 .expect("tasks_start")
1525 .0;
1526
1527 let list = steps_list(
1528 State(state.clone()),
1529 Path((posted.task_id.to_string(), "latest".to_string())),
1530 )
1531 .await
1532 .expect("steps_list")
1533 .0;
1534 assert_eq!(list.steps.len(), 1, "steps: {:?}", list.steps);
1535 assert_eq!(list.steps[0].name, "plan-out");
1536 assert_eq!(list.steps[0].source, ProjectionSource::DataPlane);
1537
1538 let by_canonical = step_get(
1539 State(state.clone()),
1540 Path((
1541 posted.task_id.to_string(),
1542 "latest".to_string(),
1543 "plan-out".to_string(),
1544 )),
1545 Query(StepPathQuery::default()),
1546 )
1547 .await
1548 .expect("step_get canonical")
1549 .0;
1550 assert_eq!(by_canonical.name, "plan-out");
1551
1552 let identity_name = mlua_swarm::worker::baseline::AG_IDENTITY;
1553 let by_ref_alias = step_get(
1554 State(state.clone()),
1555 Path((
1556 posted.task_id.to_string(),
1557 "latest".to_string(),
1558 identity_name.to_string(),
1559 )),
1560 Query(StepPathQuery::default()),
1561 )
1562 .await
1563 .expect("step_get ref alias")
1564 .0;
1565 assert_eq!(by_ref_alias.name, "plan-out");
1566 assert_eq!(by_ref_alias.sha256, by_canonical.sha256);
1567
1568 let by_out_alias = step_get(
1569 State(state.clone()),
1570 Path((
1571 posted.task_id.to_string(),
1572 "latest".to_string(),
1573 "out".to_string(),
1574 )),
1575 Query(StepPathQuery::default()),
1576 )
1577 .await
1578 .expect("step_get out-top alias")
1579 .0;
1580 assert_eq!(by_out_alias.name, "plan-out");
1581 assert_eq!(by_out_alias.sha256, by_canonical.sha256);
1582 }
1583
1584 #[tokio::test]
1591 async fn declared_projection_name_materialized_file_stem_is_canonical() {
1592 let dir = tempfile::TempDir::new().unwrap();
1593 let state = test_state();
1594 let mut req = declared_task_req("materialized-declared", "plan-out");
1595 req.work_dir = Some(dir.path().to_string_lossy().into_owned());
1596 let posted = crate::tasks_start(State(state.clone()), Json(req))
1597 .await
1598 .expect("tasks_start")
1599 .0;
1600
1601 let summary = step_get(
1602 State(state.clone()),
1603 Path((
1604 posted.task_id.to_string(),
1605 "latest".to_string(),
1606 "plan-out".to_string(),
1607 )),
1608 Query(StepPathQuery::default()),
1609 )
1610 .await
1611 .expect("step_get")
1612 .0;
1613
1614 let file_path = summary.file_path.expect("materialized file_path present");
1615 assert!(
1616 file_path.ends_with("plan-out.md"),
1617 "materialized file stem must be the canonical name: {file_path}"
1618 );
1619 }
1620
1621 #[tokio::test]
1633 async fn declared_projection_placement_e2e_write_and_read_back_converge() {
1634 let project_root_dir = tempfile::TempDir::new().unwrap();
1635 let state = test_state();
1636 let mut bp = declared_projection_name_blueprint("plan-out");
1637 bp.projection_placement = Some(ProjectionPlacementSpec {
1638 root: Some("project_root".to_string()),
1639 dir_template: Some("custom/{task_id}/out".to_string()),
1640 });
1641 let req = TaskLaunchRequest {
1642 blueprint: BlueprintRef::Inline {
1643 value: Box::new(bp),
1644 },
1645 init_ctx: json!({ "greeting": "materialized-custom-placement" }),
1646 project_root: Some(project_root_dir.path().to_string_lossy().into_owned()),
1647 work_dir: None,
1648 task_metadata: None,
1649 ttl_secs: None,
1650 operator: None,
1651 operator_sid: None,
1652 timeout_secs: None,
1653 goal: Some("projection placement test goal".to_string()),
1654 detach: false,
1655 check_policy: None,
1656 };
1657 let posted = crate::tasks_start(State(state.clone()), Json(req))
1658 .await
1659 .expect("tasks_start")
1660 .0;
1661
1662 let summary = step_get(
1663 State(state.clone()),
1664 Path((
1665 posted.task_id.to_string(),
1666 "latest".to_string(),
1667 "plan-out".to_string(),
1668 )),
1669 Query(StepPathQuery::default()),
1670 )
1671 .await
1672 .expect("step_get")
1673 .0;
1674
1675 let file_path = summary.file_path.expect("materialized file_path present");
1684 let path = std::path::Path::new(&file_path);
1685 assert!(
1686 path.starts_with(project_root_dir.path()),
1687 "file must be rooted at project_root (root_preference=ProjectRoot): {file_path}"
1688 );
1689 assert!(
1690 file_path.ends_with("out/plan-out.md"),
1691 "file must follow the custom dir_template's tail: {file_path}"
1692 );
1693 assert!(
1694 file_path.contains("/custom/"),
1695 "file must follow the custom dir_template's prefix segment: {file_path}"
1696 );
1697 assert!(
1698 path.exists(),
1699 "the write side must have materialized the file the read-back reports: {file_path}"
1700 );
1701 }
1702
1703 #[tokio::test]
1710 async fn declared_projection_name_colliding_with_another_steps_ref_is_rejected_at_register_time(
1711 ) {
1712 use mlua_flow_ir::{Expr, Node as FlowNode};
1713 use mlua_swarm::worker::adapter::WorkerResult;
1714 use mlua_swarm::{RustFnInProcessSpawnerFactory, SpawnerRegistry};
1715
1716 let factory = RustFnInProcessSpawnerFactory::new()
1717 .register_fn("step-a", |inv| async move {
1718 Ok(WorkerResult {
1719 value: json!(inv.prompt),
1720 ok: true,
1721 stats: None,
1722 })
1723 })
1724 .register_fn("step-b", |inv| async move {
1725 Ok(WorkerResult {
1726 value: json!(inv.prompt),
1727 ok: true,
1728 stats: None,
1729 })
1730 });
1731 let mut reg = SpawnerRegistry::new();
1732 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(factory));
1733
1734 let engine = Engine::new_with_layers(EngineCfg::default(), crate::default_layer_registry());
1735 let data_store: Arc<dyn mlua_swarm::store::output::OutputStore> =
1736 Arc::new(InMemoryOutputStore::new());
1737 engine.set_output_store(data_store.clone());
1738 let compiler = mlua_swarm::Compiler::new(reg);
1739 let launch = Arc::new(mlua_swarm::TaskLaunchService::new(engine.clone(), compiler));
1740 let state = AppState {
1741 engine,
1742 sessions: Arc::new(Mutex::new(crate::SessionStore::default())),
1743 task_app: Arc::new(mlua_swarm::TaskApplication::new_inline_only(launch)),
1744 ws_operator_factory: None,
1745 data_store,
1746 operator_sessions: Arc::new(Mutex::new(HashMap::new())),
1747 roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
1748 task_store: Arc::new(InMemoryTaskStore::new()),
1749 run_store: Arc::new(InMemoryRunStore::new()),
1750 replay_store: Arc::new(mlua_swarm::store::replay::InMemoryReplayStore::new()),
1751 run_trace_store: Arc::new(mlua_swarm::store::trace::InMemoryRunTraceStore::new()),
1752 base_url: None,
1753 sync_timeout_secs: 300,
1754 };
1755
1756 let flow = FlowNode::Seq {
1757 children: vec![
1758 FlowNode::Step {
1759 ref_: "step-a".to_string(),
1760 in_: Expr::Path {
1761 at: "$.greeting".parse().expect("literal test path: $.greeting"),
1762 },
1763 out: Expr::Path {
1764 at: "$.a_out".parse().expect("literal test path: $.a_out"),
1765 },
1766 },
1767 FlowNode::Step {
1768 ref_: "step-b".to_string(),
1769 in_: Expr::Path {
1770 at: "$.greeting".parse().expect("literal test path: $.greeting"),
1771 },
1772 out: Expr::Path {
1773 at: "$.b_out".parse().expect("literal test path: $.b_out"),
1774 },
1775 },
1776 ],
1777 };
1778 let blueprint = Blueprint {
1779 schema_version: current_schema_version(),
1780 id: "projection-test-collision-bp".into(),
1781 flow,
1782 agents: vec![
1783 AgentDef {
1784 name: "step-a".into(),
1785 kind: AgentKind::RustFn,
1786 spec: json!({"fn_id": "step-a"}),
1787 profile: None,
1788 meta: Some(AgentMeta {
1791 projection_name: Some("step-b".to_string()),
1792 ..Default::default()
1793 }),
1794 runner: None,
1795 runner_ref: None,
1796 verdict: None,
1797 lints: None,
1798 },
1799 AgentDef {
1800 name: "step-b".into(),
1801 kind: AgentKind::RustFn,
1802 spec: json!({"fn_id": "step-b"}),
1803 profile: None,
1804 meta: None,
1805 runner: None,
1806 runner_ref: None,
1807 verdict: None,
1808 lints: None,
1809 },
1810 ],
1811 operators: vec![],
1812 metas: vec![],
1813 hints: CompilerHints::default(),
1814 strategy: CompilerStrategy::default(),
1815 metadata: BlueprintMetadata::default(),
1816 spawner_hints: Default::default(),
1817 default_agent_kind: AgentKind::Operator,
1818 default_operator_kind: None,
1819 default_init_ctx: None,
1820 default_agent_ctx: None,
1821 default_context_policy: None,
1822 projection_placement: None,
1823 audits: vec![],
1824 degradation_policy: None,
1825 runners: vec![],
1826 default_runner: None,
1827 subprocesses: vec![],
1828 check_policy: None,
1829 blueprint_ref_includes: Vec::new(),
1830 };
1831
1832 let req = TaskLaunchRequest {
1833 blueprint: BlueprintRef::Inline {
1834 value: Box::new(blueprint),
1835 },
1836 init_ctx: json!({ "greeting": "hi" }),
1837 project_root: None,
1838 work_dir: None,
1839 task_metadata: None,
1840 ttl_secs: None,
1841 operator: None,
1842 operator_sid: None,
1843 timeout_secs: None,
1844 goal: None,
1845 detach: false,
1846 check_policy: None,
1847 };
1848
1849 let result = crate::tasks_start(State(state), Json(req)).await;
1853 let err = match result {
1854 Err(e) => e,
1855 Ok(_) => {
1856 panic!("declared projection_name colliding with another step's own ref must reject")
1857 }
1858 };
1859 assert_eq!(err.status, StatusCode::BAD_REQUEST);
1860 }
1861
1862 #[tokio::test]
1871 async fn steps_list_run_scoped_lookup_does_not_bleed_across_tasks_sharing_a_producer_name() {
1872 let state = test_state();
1873 let first = crate::tasks_start(State(state.clone()), Json(greeting_task_req("first-task")))
1874 .await
1875 .expect("first tasks_start")
1876 .0;
1877 let second =
1878 crate::tasks_start(State(state.clone()), Json(greeting_task_req("second-task")))
1879 .await
1880 .expect("second tasks_start")
1881 .0;
1882
1883 let first_steps = steps_list(
1884 State(state.clone()),
1885 Path((first.task_id.to_string(), "latest".to_string())),
1886 )
1887 .await
1888 .expect("first steps_list")
1889 .0;
1890 let second_steps = steps_list(
1891 State(state.clone()),
1892 Path((second.task_id.to_string(), "latest".to_string())),
1893 )
1894 .await
1895 .expect("second steps_list")
1896 .0;
1897
1898 let identity_name = mlua_swarm::worker::baseline::AG_IDENTITY;
1899 let first_entry = first_steps
1900 .steps
1901 .iter()
1902 .find(|s| s.name == identity_name)
1903 .expect("first entry present");
1904 let second_entry = second_steps
1905 .steps
1906 .iter()
1907 .find(|s| s.name == identity_name)
1908 .expect("second entry present");
1909 assert_eq!(first_entry.source, ProjectionSource::DataPlane);
1910 assert_eq!(second_entry.source, ProjectionSource::DataPlane);
1911 assert_ne!(
1912 first_entry.sha256, second_entry.sha256,
1913 "each Task's own greeting must resolve, not the globally-latest submission"
1914 );
1915 }
1916
1917 #[tokio::test]
1920 async fn steps_list_latest_resolves_newest_run_explicit_pin_still_works() {
1921 let state = test_state();
1922 let first = crate::tasks_start(State(state.clone()), Json(greeting_task_req("first")))
1923 .await
1924 .expect("tasks_start")
1925 .0;
1926 let (status, rekicked) = crate::tasks::task_rekick(
1927 State(state.clone()),
1928 Path(first.task_id.to_string()),
1929 Some(Json(crate::tasks::RunKickRequest {
1930 init_ctx_override: Some(json!({ "greeting": "second" })),
1931 task_input_override: None,
1932 timeout_secs: None,
1933 detach: false,
1934 operator_sid: None,
1935 })),
1936 )
1937 .await
1938 .expect("task_rekick");
1939 assert_eq!(status, StatusCode::CREATED);
1940
1941 let latest = steps_list(
1942 State(state.clone()),
1943 Path((first.task_id.to_string(), "latest".to_string())),
1944 )
1945 .await
1946 .expect("steps_list latest")
1947 .0;
1948 assert_eq!(latest.run_id, rekicked.0.run_id.to_string());
1949
1950 let pinned = steps_list(
1951 State(state.clone()),
1952 Path((first.task_id.to_string(), first.run_id.to_string())),
1953 )
1954 .await
1955 .expect("steps_list pinned")
1956 .0;
1957 assert_eq!(pinned.run_id, first.run_id.to_string());
1958 }
1959
1960 #[tokio::test]
1963 async fn step_get_preview_is_utf8_boundary_safe_and_truncated_flag_is_correct() {
1964 let state = test_state();
1965 let long_value = "あ".repeat(300); let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req(&long_value)))
1970 .await
1971 .expect("tasks_start")
1972 .0;
1973
1974 let identity_name = mlua_swarm::worker::baseline::AG_IDENTITY;
1975 let summary = step_get(
1976 State(state.clone()),
1977 Path((
1978 posted.task_id.to_string(),
1979 "latest".to_string(),
1980 identity_name.to_string(),
1981 )),
1982 Query(StepPathQuery::default()),
1983 )
1984 .await
1985 .expect("step_get")
1986 .0;
1987
1988 assert!(
1989 summary.preview.len() <= 512 + "…".len(),
1990 "preview must stay near the 512-byte cap: {} bytes",
1991 summary.preview.len()
1992 );
1993 assert!(
1994 summary.truncated,
1995 "a 900-byte body must be reported truncated"
1996 );
1997 assert!(
1998 summary.preview.ends_with('…'),
1999 "truncated preview must end with an ellipsis: {}",
2000 summary.preview
2001 );
2002 assert!(summary.preview.chars().all(|c| c != '\u{FFFD}'));
2008 }
2009
2010 #[tokio::test]
2013 async fn step_content_in_memory_fallback_is_json_with_matching_etag() {
2014 let state = test_state();
2015 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("hi")))
2016 .await
2017 .expect("tasks_start")
2018 .0;
2019
2020 let resp = step_content(
2021 State(state.clone()),
2022 Path((
2023 posted.task_id.to_string(),
2024 "latest".to_string(),
2025 "out".to_string(),
2026 )),
2027 Query(StepPathQuery::default()),
2028 )
2029 .await
2030 .expect("step_content")
2031 .into_response();
2032
2033 assert_eq!(resp.status(), StatusCode::OK);
2034 let content_type = resp
2035 .headers()
2036 .get(header::CONTENT_TYPE)
2037 .expect("content-type header")
2038 .to_str()
2039 .expect("ascii");
2040 assert_eq!(content_type, "application/json");
2041 let etag = resp
2042 .headers()
2043 .get(header::ETAG)
2044 .expect("etag header")
2045 .to_str()
2046 .expect("ascii")
2047 .to_string();
2048 let body_bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
2049 .await
2050 .expect("body bytes");
2051 let expected_sha = hex::encode(sha2::Sha256::digest(&body_bytes));
2052 assert_eq!(etag, format!("\"sha256:{expected_sha}\""));
2053 let parsed: Value = serde_json::from_slice(&body_bytes).expect("valid json body");
2054 assert_eq!(parsed["echoed"], json!("hi"));
2055 }
2056
2057 #[tokio::test]
2062 async fn step_content_materialized_file_is_served_as_markdown() {
2063 let dir = tempfile::TempDir::new().unwrap();
2064 let state = test_state();
2065 let mut req = greeting_task_req("materialized");
2066 req.work_dir = Some(dir.path().to_string_lossy().into_owned());
2067 let posted = crate::tasks_start(State(state.clone()), Json(req))
2068 .await
2069 .expect("tasks_start")
2070 .0;
2071
2072 let identity_name = mlua_swarm::worker::baseline::AG_IDENTITY;
2073 let resp = step_content(
2074 State(state.clone()),
2075 Path((
2076 posted.task_id.to_string(),
2077 "latest".to_string(),
2078 identity_name.to_string(),
2079 )),
2080 Query(StepPathQuery::default()),
2081 )
2082 .await
2083 .expect("step_content")
2084 .into_response();
2085
2086 assert_eq!(resp.status(), StatusCode::OK);
2087 let content_type = resp
2088 .headers()
2089 .get(header::CONTENT_TYPE)
2090 .expect("content-type header")
2091 .to_str()
2092 .expect("ascii");
2093 assert_eq!(content_type, "text/markdown; charset=utf-8");
2094 let body_bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
2095 .await
2096 .expect("body bytes");
2097 let body_str = String::from_utf8(body_bytes.to_vec()).expect("utf8 body");
2098 assert!(
2099 body_str.contains("```json"),
2100 "materialized file must carry the fenced json block: {body_str}"
2101 );
2102 }
2103
2104 #[tokio::test]
2107 async fn step_content_path_narrow_returns_json_fragment() {
2108 let state = test_state();
2109 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("narrowed")))
2110 .await
2111 .expect("tasks_start")
2112 .0;
2113
2114 let resp = step_content(
2115 State(state.clone()),
2116 Path((
2117 posted.task_id.to_string(),
2118 "latest".to_string(),
2119 "out".to_string(),
2120 )),
2121 Query(StepPathQuery {
2122 path: Some("echoed".to_string()),
2123 }),
2124 )
2125 .await
2126 .expect("step_content narrowed")
2127 .into_response();
2128
2129 assert_eq!(resp.status(), StatusCode::OK);
2130 let content_type = resp
2131 .headers()
2132 .get(header::CONTENT_TYPE)
2133 .expect("content-type header")
2134 .to_str()
2135 .expect("ascii");
2136 assert_eq!(content_type, "application/json");
2137 let body_bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
2138 .await
2139 .expect("body bytes");
2140 let parsed: Value = serde_json::from_slice(&body_bytes).expect("valid json body");
2141 assert_eq!(parsed, json!("narrowed"));
2142 }
2143
2144 #[tokio::test]
2147 async fn steps_list_unknown_task_returns_404() {
2148 let state = test_state();
2149 let err = steps_list(
2150 State(state),
2151 Path(("T-does-not-exist".to_string(), "latest".to_string())),
2152 )
2153 .await
2154 .expect_err("unknown task must 404");
2155 assert_eq!(err.status, StatusCode::NOT_FOUND);
2156 }
2157
2158 #[tokio::test]
2159 async fn steps_list_unknown_run_returns_404() {
2160 let state = test_state();
2161 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("hi")))
2162 .await
2163 .expect("tasks_start")
2164 .0;
2165 let err = steps_list(
2166 State(state),
2167 Path((posted.task_id.to_string(), "R-does-not-exist".to_string())),
2168 )
2169 .await
2170 .expect_err("unknown run must 404");
2171 assert_eq!(err.status, StatusCode::NOT_FOUND);
2172 }
2173
2174 #[tokio::test]
2175 async fn step_get_unknown_step_returns_404() {
2176 let state = test_state();
2177 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("hi")))
2178 .await
2179 .expect("tasks_start")
2180 .0;
2181 let err = step_get(
2182 State(state),
2183 Path((
2184 posted.task_id.to_string(),
2185 "latest".to_string(),
2186 "does-not-exist".to_string(),
2187 )),
2188 Query(StepPathQuery::default()),
2189 )
2190 .await
2191 .expect_err("unknown step must 404");
2192 assert_eq!(err.status, StatusCode::NOT_FOUND);
2193 }
2194
2195 #[tokio::test]
2198 async fn old_ctx_route_returns_404_not_found_by_router() {
2199 let engine = Engine::new(EngineCfg::default());
2200 let router = mlua_swarm_server_router_for_test(engine);
2201 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
2202 .await
2203 .expect("bind ephemeral port");
2204 let addr = listener.local_addr().expect("local addr");
2205 tokio::spawn(async move {
2206 let _ = axum::serve(listener, router).await;
2207 });
2208 let client = reqwest::Client::new();
2209 let resp = client
2210 .get(format!("http://{addr}/v1/tasks/T-anything/ctx"))
2211 .send()
2212 .await
2213 .expect("request");
2214 assert_eq!(resp.status(), reqwest::StatusCode::NOT_FOUND);
2215 }
2216
2217 fn mlua_swarm_server_router_for_test(engine: Engine) -> axum::Router {
2221 crate::build_router(engine)
2222 }
2223
2224 #[test]
2227 fn mcp_query_adapter_project_builds_query_ref() {
2228 let adapter = McpQueryAdapter::new(
2229 Arc::new(InMemoryOutputStore::new()),
2230 Arc::new(InMemoryRunStore::new()),
2231 Engine::new(EngineCfg::default()),
2232 );
2233 let key = ProjectionKey {
2234 task_id: "T-abc".to_string(),
2235 run_id: None,
2236 step: Some("planner".to_string()),
2237 path: None,
2238 };
2239 let ctx_data = json!({"planner": {"plan": "do it"}});
2240 let reference = adapter.project(&key, &ctx_data).expect("project");
2241 match &reference {
2242 ProjectionRef::Query { endpoint, key: k } => {
2243 assert!(endpoint.contains("/steps/planner/content"));
2244 assert_eq!(k, &key);
2245 }
2246 other => panic!("expected Query ref, got {other:?}"),
2247 }
2248 let line = adapter.pointer_line(&reference);
2249 assert!(line.contains("T-abc"));
2250 }
2251
2252 #[test]
2253 fn mcp_query_adapter_project_rejects_key_not_present_in_ctx_data() {
2254 let adapter = McpQueryAdapter::new(
2255 Arc::new(InMemoryOutputStore::new()),
2256 Arc::new(InMemoryRunStore::new()),
2257 Engine::new(EngineCfg::default()),
2258 );
2259 let key = ProjectionKey {
2260 task_id: "T-abc".to_string(),
2261 run_id: None,
2262 step: Some("missing".to_string()),
2263 path: None,
2264 };
2265 let err = adapter.project(&key, &json!({"planner": {}})).unwrap_err();
2266 assert!(matches!(err, ProjectionError::NotFound(_)));
2267 }
2268
2269 #[tokio::test(flavor = "multi_thread")]
2270 async fn mcp_query_adapter_fetch_bridges_to_resolve_async() {
2271 let state = test_state();
2272 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("bridged")))
2273 .await
2274 .expect("tasks_start")
2275 .0;
2276
2277 let adapter = McpQueryAdapter::new(
2278 state.data_store.clone(),
2279 state.run_store.clone(),
2280 state.engine.clone(),
2281 );
2282 let key = ProjectionKey {
2283 task_id: posted.task_id.to_string(),
2284 run_id: None,
2285 step: Some("out".to_string()),
2286 path: Some("echoed".to_string()),
2287 };
2288 let value = adapter.fetch(&key).expect("fetch");
2296 assert_eq!(value, json!("bridged"));
2297 }
2298
2299 #[tokio::test]
2312 async fn resolve_async_path_narrows_within_data_plane_final_content() {
2313 let state = test_state();
2314 let posted = crate::tasks_start(State(state.clone()), Json(greeting_task_req("hi")))
2315 .await
2316 .expect("tasks_start")
2317 .0;
2318
2319 let adapter = McpQueryAdapter::new(
2320 state.data_store.clone(),
2321 state.run_store.clone(),
2322 state.engine.clone(),
2323 );
2324 let key = ProjectionKey {
2325 task_id: posted.task_id.to_string(),
2326 run_id: None,
2327 step: Some(mlua_swarm::worker::baseline::AG_IDENTITY.to_string()),
2328 path: Some("echoed".to_string()),
2329 };
2330 let (_run, value) = adapter.resolve_async(&key).await.expect("resolve_async");
2331 assert_eq!(value, json!("hi"));
2332 }
2333
2334 #[tokio::test(flavor = "multi_thread")]
2345 async fn steps_list_returns_in_flight_step_output_before_run_completes() {
2346 use mlua_flow_ir::{Expr, Node as FlowNode};
2347 use mlua_swarm::worker::adapter::WorkerResult;
2348 use mlua_swarm::{RustFnInProcessSpawnerFactory, SpawnerRegistry};
2349
2350 let started = Arc::new(tokio::sync::Notify::new());
2351 let gate = Arc::new(tokio::sync::Notify::new());
2352 let started_bg = started.clone();
2353 let gate_bg = gate.clone();
2354
2355 let factory = RustFnInProcessSpawnerFactory::new()
2356 .register_fn("step1", |inv| async move {
2357 Ok(WorkerResult {
2358 value: json!({ "step1_out": inv.prompt }),
2359 ok: true,
2360 stats: None,
2361 })
2362 })
2363 .register_fn("step2", move |_inv| {
2364 let started = started_bg.clone();
2365 let gate = gate_bg.clone();
2366 async move {
2367 started.notify_one();
2368 gate.notified().await;
2369 Ok(WorkerResult {
2370 value: json!("step2 done"),
2371 ok: true,
2372 stats: None,
2373 })
2374 }
2375 });
2376 let mut reg = SpawnerRegistry::new();
2377 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(factory));
2378
2379 let engine = Engine::new_with_layers(EngineCfg::default(), crate::default_layer_registry());
2380 let data_store: Arc<dyn mlua_swarm::store::output::OutputStore> =
2381 Arc::new(InMemoryOutputStore::new());
2382 engine.set_output_store(data_store.clone());
2383 let compiler = mlua_swarm::Compiler::new(reg);
2384 let launch = Arc::new(mlua_swarm::TaskLaunchService::new(engine.clone(), compiler));
2385 let state = AppState {
2386 engine,
2387 sessions: Arc::new(Mutex::new(crate::SessionStore::default())),
2388 task_app: Arc::new(mlua_swarm::TaskApplication::new_inline_only(launch)),
2389 ws_operator_factory: None,
2390 data_store,
2391 operator_sessions: Arc::new(Mutex::new(HashMap::new())),
2392 roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
2393 task_store: Arc::new(InMemoryTaskStore::new()),
2394 run_store: Arc::new(InMemoryRunStore::new()),
2395 replay_store: Arc::new(mlua_swarm::store::replay::InMemoryReplayStore::new()),
2396 run_trace_store: Arc::new(mlua_swarm::store::trace::InMemoryRunTraceStore::new()),
2397 base_url: None,
2398 sync_timeout_secs: 300,
2399 };
2400
2401 let flow = FlowNode::Seq {
2402 children: vec![
2403 FlowNode::Step {
2404 ref_: "step1".to_string(),
2405 in_: Expr::Path {
2406 at: "$.greeting".parse().expect("literal test path: $.greeting"),
2407 },
2408 out: Expr::Path {
2409 at: "$.step1".parse().expect("literal test path: $.step1"),
2410 },
2411 },
2412 FlowNode::Step {
2413 ref_: "step2".to_string(),
2414 in_: Expr::Path {
2415 at: "$.step1".parse().expect("literal test path: $.step1"),
2416 },
2417 out: Expr::Path {
2418 at: "$.step2".parse().expect("literal test path: $.step2"),
2419 },
2420 },
2421 ],
2422 };
2423 let blueprint = Blueprint {
2424 schema_version: current_schema_version(),
2425 id: "projection-test-in-flight-bp".into(),
2426 flow,
2427 agents: vec![
2428 AgentDef {
2429 name: "step1".into(),
2430 kind: AgentKind::RustFn,
2431 spec: json!({"fn_id": "step1"}),
2432 profile: None,
2433 meta: None,
2434 runner: None,
2435 runner_ref: None,
2436 verdict: None,
2437 lints: None,
2438 },
2439 AgentDef {
2440 name: "step2".into(),
2441 kind: AgentKind::RustFn,
2442 spec: json!({"fn_id": "step2"}),
2443 profile: None,
2444 meta: None,
2445 runner: None,
2446 runner_ref: None,
2447 verdict: None,
2448 lints: None,
2449 },
2450 ],
2451 operators: vec![],
2452 metas: vec![],
2453 hints: CompilerHints::default(),
2454 strategy: CompilerStrategy::default(),
2455 metadata: BlueprintMetadata::default(),
2456 spawner_hints: Default::default(),
2457 default_agent_kind: AgentKind::Operator,
2458 default_operator_kind: None,
2459 default_init_ctx: None,
2460 default_agent_ctx: None,
2461 default_context_policy: None,
2462 projection_placement: None,
2463 audits: vec![],
2464 degradation_policy: None,
2465 runners: vec![],
2466 default_runner: None,
2467 subprocesses: vec![],
2468 check_policy: None,
2469 blueprint_ref_includes: Vec::new(),
2470 };
2471
2472 let req = TaskLaunchRequest {
2473 blueprint: BlueprintRef::Inline {
2474 value: Box::new(blueprint),
2475 },
2476 init_ctx: json!({ "greeting": "hi" }),
2477 project_root: None,
2478 work_dir: None,
2479 task_metadata: None,
2480 ttl_secs: None,
2481 operator: None,
2482 operator_sid: None,
2483 timeout_secs: None,
2484 goal: None,
2485 detach: false,
2486 check_policy: None,
2487 };
2488
2489 let state_bg = state.clone();
2490 let launch_handle =
2491 tokio::spawn(async move { crate::tasks_start(State(state_bg), Json(req)).await });
2492
2493 started.notified().await;
2497
2498 let in_flight_tasks = state.task_store.list().await.expect("task_store list");
2499 assert_eq!(in_flight_tasks.len(), 1, "exactly one Task minted");
2500 let task_id = in_flight_tasks[0].id.clone();
2501
2502 let resp = steps_list(
2503 State(state.clone()),
2504 Path((task_id.to_string(), "latest".to_string())),
2505 )
2506 .await
2507 .expect("steps_list while step2 is still in flight");
2508 let step1_entry = resp
2509 .steps
2510 .iter()
2511 .find(|s| s.name == "step1")
2512 .expect("step1 must already be visible");
2513 assert_eq!(step1_entry.source, ProjectionSource::DataPlane);
2514
2515 gate.notify_one();
2518 let posted = launch_handle.await.expect("join").expect("tasks_start").0;
2519 assert_eq!(posted.final_ctx["step2"], json!("step2 done"));
2520 }
2521}