mlua_swarm_server/lib.rs
1//! the server lib: axum Router + handler set. Split out as a library so it can
2//! be used from both `main.rs` (CLI) and integration tests.
3//!
4//! # Endpoints
5//!
6//! - `GET /v1/healthz`
7//! - `POST /v1/sessions` / `DELETE /v1/sessions` (= operator attach / detach, Bearer sid)
8//! - `POST /v1/tasks` (= unified Flow-form entry, Operator inject supported;
9//! `operator_sid` explicitly pins the task to a registered Operator session, S2).
10//! Also creates a `TaskRecord` + `RunRecord` (issue #13 ID-hierarchy persistence)
11//! and echoes their ids in the response; see the `tasks` module doc. Always
12//! synchronous, guarded against hanging (GH #33) by a readiness precheck
13//! (`503` when the launch resolves to an operator-delegate path with zero
14//! attached operators) and a `tokio::time::timeout` ceiling around the
15//! dispatch await (`504` on expiry) — see `run_flow_form`'s doc comment.
16//! - `GET /v1/tasks` — list every persisted `TaskRecord` (newest first).
17//! - `GET /v1/tasks/:id` — a `TaskRecord` plus every `RunRecord` kicked from it.
18//! - `POST /v1/tasks/:id/runs` — re-kick an existing Task (new `RunId`, same
19//! `blueprint_ref` / `input_ctx`).
20//! - `GET /v1/tasks/:id/runs/:run/steps` / `.../steps/:step` /
21//! `.../steps/:step/content` — the metadata + content debug plane over a
22//! Run's step OUTPUT (`:run` accepts `latest` or an explicit `R-<hex>`,
23//! `projection::McpQueryAdapter`); see the `projection` module doc. This
24//! is the operator / human-debug counterpart to the Worker axis's
25//! `context.steps` pointer list on `GET /v1/worker/prompt`
26//! (`projection-adapter` ST5 — replaces the ST2/ST4 single-value `GET
27//! /v1/tasks/:id/ctx`).
28//! - `GET /v1/runs/:id` — a single `RunRecord` (its `step_entries` trace included).
29//! - `GET /v1/runs/:id/bindings` — immutable requested/effective AgentProvider
30//! binding explain for that Run.
31//! - `POST /v1/runs/:id/resume` — resume an `Interrupted` Run under the same `run_id`.
32//! - `POST /v1/runs/:id/rerun-from` — GH #71 Layer A. Rerun a terminal Run from
33//! a caller-specified step under the same `run_id` (physically truncates the
34//! replay log at the cut point). See `tasks::run_rerun_from`.
35//! - `POST /v1/operators` / `GET /v1/operators/:sid` / `DELETE /v1/operators/:sid` /
36//! `GET /v1/operators/:sid/ws` (WS upgrade) — REST-like Operator login flow,
37//! Bearer-mandatory; the sole WS Operator session route. See `operator_ws::login`
38//! module doc.
39//!
40//! The Enhance issue axis (`/issues`) lives in the `issues` module; callers merge
41//! `build_issues_router` to integrate it into the same server.
42//!
43//! # The 3 faces of the Operator role (= registered directly on the engine SoT)
44//!
45//! The engine stateless-executor refactor removed the three
46//! `AppState` registries (former `HookRegistry` / `BridgeRegistry` / `OperatorRegistry`);
47//! all registration now goes directly to the engine SoT via
48//! `engine.register_spawn_hook` / `register_senior_bridge` / `register_operator`.
49//! `WSOperatorSession` (in the `operator_ws` module) registers all three traits
50//! simultaneously under a single sid — one WS connection covers all 3 faces of
51//! the Operator role, the canonical pattern.
52//!
53//! # `build_*` family
54//!
55//! - [`build_router`] — minimal entry (= `default_registry()`)
56//! - [`build_router_with`] — caller provides a `SpawnerRegistry` and optional `BlueprintStore`
57//!
58//! The engine should be started with [`default_layer_registry`] (= `Engine::new_with_layers`);
59//! otherwise `Blueprint.spawner_hints` is ignored.
60
61#![warn(missing_docs)]
62
63pub mod binding;
64/// HTTP surface for inspecting/registering Blueprint state (`/v1/blueprints/*`).
65pub mod blueprints;
66/// Server config file support (`~/.mse/config.toml`, CLI > file > default merge).
67pub mod config;
68/// `/v1/data/*` endpoints (v9 Big Response handling, Store-owner direct path).
69pub mod data;
70/// `GET /v1/doctor` — read-only startup config / Store snapshot.
71pub mod doctor;
72/// HTTP surface for the `/v1/enhance/log` axis.
73pub mod enhance_log;
74/// `EnhanceSetting` HTTP CRUD (`/v1/enhance-settings*`).
75pub mod enhance_settings;
76/// HTTP surface for the Enhance issue axis (`/v1/issues*`).
77pub mod issues;
78/// WebSocket Operator Callback IF (`/v1/operators*`).
79pub mod operator_ws;
80/// `GET /v1/tasks/:id/runs/:run/steps*` (the metadata + content debug
81/// plane over a Run's step OUTPUT — `McpQueryAdapter`, a server-side
82/// `mlua_swarm::core::projection::ProjectionAdapter` impl reading through
83/// the Data-plane `OutputStore` with a persisted `RunRecord.result_ref`
84/// fallback). See the module doc for how this relates to
85/// `operator_ws::session`'s in-flight `FileProjectionAdapter` hook and
86/// `worker`'s Worker-axis `context.steps` pointer assembly.
87pub mod projection;
88/// HTTP surface for the Task/Run persistence axis (issue #13 ID hierarchy;
89/// `GET /v1/tasks`, `GET /v1/tasks/:id`, `POST /v1/tasks/:id/runs`,
90/// `GET /v1/runs/:id`). `POST /v1/tasks` itself stays in this module (it is
91/// the entry point `tasks_start` shares with the flow-eval path) — see the
92/// `tasks` module doc for the split rationale.
93pub mod tasks;
94/// `/v1/worker/*` endpoints (SubAgent self-fetch path).
95pub mod worker;
96pub use blueprints::{
97 build_blueprints_router, build_blueprints_router_with_refs, BindingRequirementsResponse,
98};
99pub use enhance_log::build_enhance_log_router;
100pub use enhance_settings::build_enhance_settings_router;
101pub use issues::{build_issues_router, GetIssueResponse, PostIssueRequest, PostIssueResponse};
102pub use operator_ws::{
103 operators_create, operators_delete, operators_delete_by_role, operators_info, operators_list,
104 operators_ws_connect, ClientMsg, OperatorSessionEntry, OperatorsListEntry, OperatorsListResp,
105 ServerMsg, WSOperatorSession,
106};
107pub use projection::{McpQueryAdapter, ProjectionSource, StepList, StepPathQuery, StepSummary};
108pub use tasks::{
109 RunBindingDifference, RunBindingExplainEntry, RunBindingStatus, RunBindingsExplainResponse,
110 RunKickRequest, RunKickResponse, RunResumeResponse, RunStepsResponse, TaskDetailResponse,
111};
112pub use worker::{
113 worker_artifact, worker_prompt, worker_result, ArtifactQuery, DegradationBody, PromptQuery,
114 StatsBody, WorkerResultReq,
115};
116
117use axum::{
118 extract::{DefaultBodyLimit, State},
119 http::{header::AUTHORIZATION, HeaderMap, StatusCode},
120 response::{IntoResponse, Response},
121 routing::{get, post},
122 Json, Router,
123};
124use mlua_swarm::application::{BlueprintRef, TaskApplication, TaskApplicationError};
125use mlua_swarm::blueprint::store::BlueprintStore;
126use mlua_swarm::core::config::CheckPolicy;
127use mlua_swarm::service::{TaskLaunchError, TaskLaunchService};
128use mlua_swarm::store::replay::{InMemoryReplayStore, ReplayStore};
129use mlua_swarm::store::run::{RunContext, RunRecord, RunStatus, RunStore};
130use mlua_swarm::store::task::{TaskRecord, TaskRecordStatus, TaskStore};
131use mlua_swarm::{
132 AgentBlockInProcessSpawnerFactory, CapToken, Compiler, Engine, LayerRegistry,
133 LongHoldMiddleware, LuaInProcessSpawnerFactory, MainAIMiddleware, OperatorDelegateMiddleware,
134 OperatorSpawnerFactory, Role, RunId, RustFnInProcessSpawnerFactory, SeniorEscalationMiddleware,
135 SessionId, SpawnerRegistry, SubprocessProcessSpawnerFactory, TaskId,
136};
137use serde::{Deserialize, Serialize};
138use serde_json::{json, Value};
139use std::collections::HashMap;
140use std::sync::Arc;
141use std::time::Duration;
142use tokio::sync::Mutex;
143
144/// In-memory session map backing `/v1/sessions` attach/detach.
145///
146/// The `sid` handed to the client on this REST path is the token nonce
147/// itself (a bearer secret), so the server never uses it as a map key —
148/// entries are keyed by its fingerprint
149/// (`mlua_swarm::types::token_fingerprint`; issue #14).
150#[derive(Default)]
151pub struct SessionStore {
152 /// Live session tokens keyed by the sid's fingerprint.
153 pub map: HashMap<String, CapToken>,
154}
155
156/// Shared axum handler state for the whole router. Cloned per-request (all
157/// fields are `Arc`/cheap-clone), constructed once in [`build_router_with_ws_factory`].
158#[derive(Clone)]
159pub struct AppState {
160 /// The engine SoT (attach/detach, dispatch, registries).
161 pub engine: Engine,
162 /// Live `/v1/sessions` attach records (Operator/Worker/etc session tokens).
163 pub sessions: Arc<Mutex<SessionStore>>,
164 /// Application used at the task entry to resolve `BlueprintRef`. Without a Store, runs in Inline-only mode.
165 pub task_app: Arc<TaskApplication>,
166 /// When `Some`, on WS connect a new `WSOperatorSession` is automatically registered
167 /// with this factory under the sid name (= a `kind=operator` + `operator_ref=<sid>` AgentDef
168 /// binds to the `WSOperatorSession` backend).
169 /// When `None`, no auto-registration happens; the session is only registered on
170 /// `engine.OperatorRegistry` (= only the `OperatorDelegateMiddleware` path is effective;
171 /// the `OperatorSpawnerFactory` path is dead).
172 pub ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
173 /// Owner of the Store on the Data path (Big Response handling). Added in v9.
174 /// Independent layer — the Engine core and the Domain path (`/v1/worker/result`)
175 /// are not involved.
176 /// Default = `InMemoryOutputStore` (constructed inside `build_router_with_ws_factory`);
177 /// callers can swap in an sqlite/fs backend later (future carry).
178 pub data_store: Arc<dyn mlua_swarm::store::output::OutputStore>,
179 /// Login-flow session store (`POST /v1/operators` mint records). `sid` →
180 /// `OperatorSessionEntry`. This is the sole session store for the WS
181 /// Operator role. See `operator_ws::login` module doc.
182 pub operator_sessions:
183 Arc<Mutex<HashMap<SessionId, Arc<crate::operator_ws::login::OperatorSessionEntry>>>>,
184 /// S1 login-flow roles-exclusivity map. Role name → owning `sid`. Checked
185 /// (and updated) atomically under a single lock in
186 /// `operator_ws::login::operators_create` — a role already present here
187 /// causes `POST /v1/operators` to return `409 CONFLICT`. Entries are
188 /// released on `DELETE /v1/operators/:sid`.
189 pub roles_to_sid: Arc<Mutex<HashMap<String, SessionId>>>,
190 /// Persistence for `Task` records (issue #13 ID-hierarchy work-item
191 /// identity; see `mlua_swarm::store::task` module doc). Default =
192 /// `InMemoryTaskStore` (constructed inside `build_router_full`); callers
193 /// can swap in a `SqliteTaskStore` via the `task_store` argument.
194 pub task_store: Arc<dyn TaskStore>,
195 /// Persistence for `Run` records (one kick of a Task; see
196 /// `mlua_swarm::store::run` module doc). Default = `InMemoryRunStore`;
197 /// callers can swap in a `SqliteRunStore` via the `run_store` argument.
198 pub run_store: Arc<dyn RunStore>,
199 /// Per-run replay log — the Ctx-snapshot + step-output store the engine
200 /// appends to after every completed step (see `mlua_swarm::store::replay`
201 /// module doc). Threaded into `RunContext` at every dispatch site so a
202 /// later restart-equivalent recovery can reconstruct the run. Default =
203 /// `InMemoryReplayStore` (process-volatile); callers can swap in a
204 /// `SqliteReplayStore` via the `replay_store` argument.
205 pub replay_store: Arc<dyn ReplayStore>,
206 /// Per-Run trace stream (the RunTrace rail — see
207 /// `mlua_swarm::store::trace` module doc). A `TraceHandle` bound to
208 /// this store is threaded into `RunContext` at every dispatch site
209 /// (`core.*` events + middleware/worker insertion) and read back via
210 /// `GET /v1/runs/:id/trace`. Default = `InMemoryRunTraceStore`;
211 /// callers can swap in a `SqliteRunTraceStore` (typically sharing
212 /// the `SqliteRunStore` file) via the terminal builder's
213 /// `run_trace_store` argument.
214 pub run_trace_store: Arc<dyn mlua_swarm::store::trace::RunTraceStore>,
215 /// Public HTTP base URL the server is reachable at (e.g.
216 /// `"http://127.0.0.1:7777"`), sourced from the binary at boot time.
217 /// When `Some`, `WSOperatorSession` renders it literally into the
218 /// Spawn `directive`'s `base_url` line so the receiving operator can
219 /// paste the frame into a SubAgent prompt without a `mse_doctor`
220 /// detour (issue #8). `None` preserves the historical fallback
221 /// (a placeholder that points at `mse_doctor`).
222 pub base_url: Option<Arc<str>>,
223 /// Server-wide fallback ceiling (seconds) for the `POST /v1/tasks`
224 /// synchronous launch await (GH #33 Guard 2; see `run_flow_form`'s doc
225 /// comment). Sourced from `config::ResolvedConfig::sync_timeout_secs`.
226 /// A per-request `TaskLaunchRequest.timeout_secs` override, when
227 /// present, takes priority over this value.
228 pub sync_timeout_secs: u64,
229}
230
231/// Minimal entry point: builds a router with [`default_registry`] and no
232/// `BlueprintStore` (Inline-only mode) or `ws_operator_factory`.
233pub fn build_router(engine: Engine) -> Router {
234 build_router_with(engine, default_registry(), None)
235}
236
237/// Default `LayerRegistry` for the server. Hint keys:
238/// - `"main_ai"` → `MainAIMiddleware` (= fires SpawnHook before/after)
239/// - `"senior_escalation"` → `SeniorEscalationMiddleware` (= on `ok=false`, escalates via `SeniorBridge.ask`)
240/// - `"operator_delegate"` → `OperatorDelegateMiddleware` (= when an operator backend is registered, delegates the entire spawn)
241///
242/// Including any of these keys in `Blueprint.spawner_hints.layers` causes them to
243/// be wrapped into a `SpawnerStack` at `service::linker::link` time (= per-launch;
244/// the old `engine.bind` global-state path is retired).
245/// Callers (the engine builder side) receive it via
246/// `Engine::new_with_layers(cfg, mse_server::default_layer_registry())`.
247pub fn default_layer_registry() -> LayerRegistry {
248 default_layer_registry_with(LayerOptions::default())
249}
250
251/// Optional knobs the terminal `default_layer_registry_with` builder
252/// consumes — the only currently-tunable knob is the LongHold threshold.
253#[derive(Debug, Default, Clone, Copy)]
254pub struct LayerOptions {
255 /// When `Some(ms)`, wire [`LongHoldMiddleware`] as a **base layer**
256 /// (applied to every dispatched step) with `default_hold = ms`.
257 /// The layer stays observational: on threshold breach it broadcasts
258 /// `Event::TaskAttemptCompleted { long_hold_warn: true, .. }` and
259 /// (when the dispatcher registered a `TraceHandle` for the step)
260 /// appends a `mw.long_hold_warn` event to the persistent
261 /// `RunTraceStore`. `None` = the layer is not installed — the same
262 /// no-op default the pre-config shape had.
263 pub long_hold_warn_ms: Option<u64>,
264}
265
266/// Variant of [`default_layer_registry`] that also honours per-server
267/// [`LayerOptions`] (currently: the LongHold threshold). Called by
268/// `mse serve` with the resolved config value; every other caller
269/// (tests, in-tree bins that don't tune the LongHold knob) can keep
270/// using the zero-arg [`default_layer_registry`].
271pub fn default_layer_registry_with(options: LayerOptions) -> LayerRegistry {
272 let mut reg = LayerRegistry::new()
273 .with_hint("main_ai", |_engine| Arc::new(MainAIMiddleware::new()))
274 .with_hint("senior_escalation", |_engine| {
275 Arc::new(SeniorEscalationMiddleware::new())
276 })
277 .with_hint("operator_delegate", |_engine| {
278 Arc::new(OperatorDelegateMiddleware::new())
279 });
280 if let Some(ms) = options.long_hold_warn_ms {
281 // Bake the millisecond threshold into the factory closure so it
282 // rides into every per-launch stack build without any per-BP
283 // state. The `Engine::event_tx()` sender is captured at bind
284 // time (fresh per Engine — the factory takes `&Engine`).
285 reg = reg.with_base(move |engine| {
286 Arc::new(LongHoldMiddleware::new(
287 std::time::Duration::from_millis(ms),
288 engine.event_tx(),
289 ))
290 });
291 }
292 reg
293}
294
295/// Build form where the caller supplies a registry and an optional `BlueprintStore`.
296/// The Operator callback path (= external HTTP / WS callers acting as an Operator)
297/// must be pre-registered via `engine.register_*` (= the engine is the SoT).
298/// See the `operator_ws` module doc and `OperatorInfo` (engine-side `ctx.rs`) for details.
299pub fn build_router_with(
300 engine: Engine,
301 registry: SpawnerRegistry,
302 store: Option<Arc<dyn BlueprintStore>>,
303) -> Router {
304 build_router_with_ws_factory(engine, registry, store, None)
305}
306
307/// 4-argument variant of `build_router_with`. Passing `ws_operator_factory = Some(arc)`
308/// causes each WS connect to auto-register a new `WSOperatorSession` under its sid
309/// name with the factory (= a `kind=operator` AgentDef with `operator_ref: <sid>`
310/// can then bind to the WS client backend). Callers are expected to also install
311/// the same `Arc` into the `SpawnerRegistry` via
312/// `reg.register::<OperatorSpawnerFactory>(arc.clone())`.
313pub fn build_router_with_ws_factory(
314 engine: Engine,
315 registry: SpawnerRegistry,
316 store: Option<Arc<dyn BlueprintStore>>,
317 ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
318) -> Router {
319 build_router_with_ws_factory_and_output(engine, registry, store, ws_operator_factory, None)
320}
321
322/// 5-argument variant of [`build_router_with_ws_factory`]. Passing
323/// `output_store = Some(arc)` swaps the default `InMemoryOutputStore` for a
324/// caller-supplied backend (a `SqliteOutputStore`, for instance). `None`
325/// preserves the historical behaviour (fresh in-memory store per call).
326pub fn build_router_with_ws_factory_and_output(
327 engine: Engine,
328 registry: SpawnerRegistry,
329 store: Option<Arc<dyn BlueprintStore>>,
330 ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
331 output_store: Option<Arc<dyn mlua_swarm::store::output::OutputStore>>,
332) -> Router {
333 build_router_full(
334 engine,
335 registry,
336 store,
337 ws_operator_factory,
338 output_store,
339 None,
340 None,
341 None,
342 None,
343 crate::config::default_sync_timeout_secs(),
344 )
345}
346
347// Backend-availability note for the trace rail: `build_router_full`
348// keeps its pre-trace signature (every existing caller gets the
349// in-memory default); a persistent `RunTraceStore` is injected via the
350// terminal `build_router_full_with_legacy_worker_binding_policy`'s
351// `run_trace_store` argument (the CLI `serve` path does this, sharing
352// the `SqliteRunStore` file).
353
354/// 8-argument variant of [`build_router_with_ws_factory_and_output`].
355/// Passing `base_url = Some(...)` (e.g. `"http://127.0.0.1:7777"`) makes
356/// `WSOperatorSession` render the actual server bind into the Spawn
357/// directive's `base_url` line, so the receiving operator can copy the
358/// frame straight into a SubAgent prompt (issue #8). `None` preserves
359/// the historical fallback (`<check with mse_doctor>` placeholder).
360/// `task_store` / `run_store` swap the default `InMemoryTaskStore` /
361/// `InMemoryRunStore` (issue #13 ID-hierarchy persistence) for a
362/// caller-supplied backend (`SqliteTaskStore` / `SqliteRunStore`, for
363/// instance); `None` preserves the process-volatile default.
364/// `sync_timeout_secs` is the server-wide fallback ceiling for the `POST
365/// /v1/tasks` synchronous launch await (GH #33 Guard 2) — see
366/// `AppState::sync_timeout_secs` / `run_flow_form`'s doc comment.
367// This is the terminal builder in the `build_router*` delegation chain
368// (each variant adds one more caller-overridable store/factory); the
369// argument count grows with the number of pluggable backends, not with
370// unrelated responsibilities, so a plain allow is preferable to bundling
371// them into a config struct only this one function would consume.
372#[allow(clippy::too_many_arguments)]
373pub fn build_router_full(
374 engine: Engine,
375 registry: SpawnerRegistry,
376 store: Option<Arc<dyn BlueprintStore>>,
377 ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
378 output_store: Option<Arc<dyn mlua_swarm::store::output::OutputStore>>,
379 base_url: Option<Arc<str>>,
380 task_store: Option<Arc<dyn TaskStore>>,
381 run_store: Option<Arc<dyn RunStore>>,
382 replay_store: Option<Arc<dyn ReplayStore>>,
383 sync_timeout_secs: u64,
384) -> Router {
385 build_router_full_with_legacy_worker_binding_policy(
386 engine,
387 registry,
388 store,
389 ws_operator_factory,
390 output_store,
391 base_url,
392 task_store,
393 run_store,
394 replay_store,
395 None,
396 sync_timeout_secs,
397 mlua_swarm::LegacyWorkerBindingPolicy::Allow,
398 )
399}
400
401/// Full router builder with an explicit migration gate for deprecated
402/// `AgentProfile.worker_binding` Runner fallback. The existing
403/// [`build_router_full`] remains compatibility-defaulted to `Allow`.
404#[allow(clippy::too_many_arguments)]
405pub fn build_router_full_with_legacy_worker_binding_policy(
406 engine: Engine,
407 registry: SpawnerRegistry,
408 store: Option<Arc<dyn BlueprintStore>>,
409 ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
410 output_store: Option<Arc<dyn mlua_swarm::store::output::OutputStore>>,
411 base_url: Option<Arc<str>>,
412 task_store: Option<Arc<dyn TaskStore>>,
413 run_store: Option<Arc<dyn RunStore>>,
414 replay_store: Option<Arc<dyn ReplayStore>>,
415 run_trace_store: Option<Arc<dyn mlua_swarm::store::trace::RunTraceStore>>,
416 sync_timeout_secs: u64,
417 legacy_worker_binding_policy: mlua_swarm::LegacyWorkerBindingPolicy,
418) -> Router {
419 let operator_sessions = Arc::new(Mutex::new(HashMap::new()));
420 let roles_to_sid = Arc::new(Mutex::new(HashMap::new()));
421 let compiler = Compiler::new(registry);
422 let binding_provider = Arc::new(binding::OperatorSessionBindingProvider::new(
423 operator_sessions.clone(),
424 roles_to_sid.clone(),
425 ));
426 let launch = Arc::new(
427 TaskLaunchService::new(engine.clone(), compiler)
428 .with_binding_provider(binding_provider)
429 .with_legacy_worker_binding_policy(legacy_worker_binding_policy),
430 );
431 let task_app = Arc::new(match store {
432 Some(s) => TaskApplication::new(launch, s),
433 None => TaskApplication::new_inline_only(launch),
434 });
435 let data_store: Arc<dyn mlua_swarm::store::output::OutputStore> = match output_store {
436 Some(s) => s,
437 None => Arc::new(mlua_swarm::store::output::InMemoryOutputStore::new()),
438 };
439 // subtask-4 / ST2 rework: wire the SAME `data_store` instance into the
440 // engine's submit-time projection sink (`Engine::submit_output` /
441 // `submit_worker_result_trusted`), so an ordinary worker
442 // `/v1/worker/submit` — not just the explicit `POST /v1/data/emit` —
443 // lands in this store too. `projection::McpQueryAdapter` (`GET
444 // /v1/tasks/:id/runs/:run/steps*`) reads through this same `Arc`,
445 // which is what makes an in-flight run's already-submitted step
446 // OUTPUT queryable.
447 engine.set_output_store(data_store.clone());
448 let task_store: Arc<dyn TaskStore> = match task_store {
449 Some(s) => s,
450 None => Arc::new(mlua_swarm::store::task::InMemoryTaskStore::new()),
451 };
452 let run_store: Arc<dyn RunStore> = match run_store {
453 Some(s) => s,
454 None => Arc::new(mlua_swarm::store::run::InMemoryRunStore::new()),
455 };
456 let replay_store: Arc<dyn ReplayStore> = match replay_store {
457 Some(s) => s,
458 None => Arc::new(InMemoryReplayStore::new()),
459 };
460 let run_trace_store: Arc<dyn mlua_swarm::store::trace::RunTraceStore> = match run_trace_store {
461 Some(s) => s,
462 None => Arc::new(mlua_swarm::store::trace::InMemoryRunTraceStore::new()),
463 };
464 let state = AppState {
465 engine,
466 sessions: Arc::new(Mutex::new(SessionStore::default())),
467 task_app,
468 ws_operator_factory,
469 data_store,
470 operator_sessions,
471 roles_to_sid,
472 task_store,
473 run_store,
474 replay_store,
475 run_trace_store,
476 base_url,
477 sync_timeout_secs,
478 };
479 Router::new()
480 .route("/v1/healthz", get(healthz))
481 .route("/v1/status", get(status_get))
482 // session = collection (POST = attach, DELETE = detach, sid via Authorization)
483 .route(
484 "/v1/sessions",
485 post(sessions_attach).delete(sessions_detach),
486 )
487 // task = flat, single level; authz resolved via Authorization: Bearer <sid>
488 .route("/v1/tasks", post(tasks_start).get(tasks::tasks_list))
489 .route("/v1/tasks/:id", get(tasks::task_get))
490 .route("/v1/tasks/:id/runs", post(tasks::task_rekick))
491 .route("/v1/tasks/:id/runs/:run/steps", get(projection::steps_list))
492 .route(
493 "/v1/tasks/:id/runs/:run/steps/:step",
494 get(projection::step_get),
495 )
496 .route(
497 "/v1/tasks/:id/runs/:run/steps/:step/content",
498 get(projection::step_content),
499 )
500 // Run collection + sub-resources (per-step run stats / trace rail):
501 // `GET /v1/runs` = filtered list, `DELETE /v1/runs/:id` = retention
502 // prune (run row + trace stream together), `:id/steps` = the
503 // terminal per-step stats, `:id/trace` = the TraceEvent stream.
504 .route("/v1/runs", get(tasks::runs_list))
505 .route(
506 "/v1/runs/:id",
507 get(tasks::run_get).delete(tasks::run_delete),
508 )
509 .route("/v1/runs/:id/cancel", post(tasks::run_cancel))
510 .route("/v1/runs/:id/steps", get(tasks::run_steps))
511 .route("/v1/runs/:id/trace", get(tasks::run_trace))
512 .route("/v1/runs/:id/bindings", get(tasks::run_bindings_explain))
513 // Resume an Interrupted Run under the SAME run_id (replay cursor +
514 // stored launch-input snapshot); see `tasks::run_resume`.
515 .route("/v1/runs/:id/resume", post(tasks::run_resume))
516 // Rerun-from-step on a terminal Run under the SAME run_id (physically
517 // truncates the replay log at the cut point); see
518 // `tasks::run_rerun_from` for the full contract (GH #71 Layer A).
519 .route("/v1/runs/:id/rerun-from", post(tasks::run_rerun_from))
520 // REST-like Operator login flow (Bearer-mandatory, roles exclusivity).
521 // Sole WS Operator session route; see `operator_ws::login` module doc.
522 // GH #81 Layer 2: `GET /v1/operators` (list, read-only observability
523 // — no Bearer, same trust tier as `GET /v1/status`) and
524 // `DELETE /v1/operators/by-role/:role` (stale-session recovery
525 // without knowing the sid — same trust tier as
526 // `mlua_swarm_server_shutdown`) close the pre-#81 recovery gap
527 // where a stale session was only clearable via a full server
528 // restart. Order matters: the `by-role` route is declared BEFORE
529 // the `:sid` route so `axum` matches `by-role/:role` as its own
530 // path, not as a `:sid` extract of literal `by-role`.
531 .route("/v1/operators", post(operators_create).get(operators_list))
532 .route("/v1/operators/:sid/ws", get(operators_ws_connect))
533 .route(
534 "/v1/operators/by-role/:role",
535 axum::routing::delete(operators_delete_by_role),
536 )
537 .route(
538 "/v1/operators/:sid",
539 get(operators_info).delete(operators_delete),
540 )
541 // SubAgent self-fetch path (the SubAgent self-fetch design). The SubAgent puts the
542 // CapToken handed over via WS Spawn into Bearer and hits the prompt / result
543 // endpoints directly over HTTP. See the `worker` module doc for details.
544 .route("/v1/worker/prompt", get(worker::worker_prompt))
545 .route("/v1/worker/result", post(worker::worker_result))
546 // Simplified endpoint (= worker POSTs with just token + raw body; task_id is auto-looked-up).
547 // `DefaultBodyLimit::max` is applied explicitly here (and on the sibling
548 // `/v1/worker/artifact` below) — same 2MB axum ships as its implicit
549 // global default, made visible rather than relied on silently.
550 .route(
551 "/v1/worker/submit",
552 post(worker::worker_submit).layer(DefaultBodyLimit::max(2 * 1024 * 1024)),
553 )
554 // GH #36 ST1: named multi-part worker output. A worker stages one
555 // named part per POST here, then completes the attempt with the
556 // ordinary `/v1/worker/submit` above — see the `worker` module doc.
557 .route(
558 "/v1/worker/artifact",
559 post(worker::worker_artifact).layer(DefaultBodyLimit::max(2 * 1024 * 1024)),
560 )
561 // GH #31: `Http`-mode fetch target for `system_ref.uri` (raw baked system
562 // bytes, same Bearer flow as `/v1/worker/prompt`) + live per-agent render-size
563 // lookup for `bp_doctor` (no Bearer, same trust tier as blueprints `get_head`).
564 .route(
565 "/v1/worker/prompt/system",
566 get(worker::worker_prompt_system),
567 )
568 .route(
569 "/v1/agents/:name/render-size",
570 get(worker::agent_render_size),
571 )
572 // GH #32: structured worker degradation reporting — independent channel,
573 // never touches OutputStore / the fold path. See the `worker` module doc.
574 .route("/v1/worker/degradation", post(worker::worker_degradation))
575 .route("/v1/worker/stats", post(worker::worker_stats))
576 // Data path (v9 Big Response handling, independent from Domain / verdict flow)
577 .route("/v1/data/emit", post(data::data_emit))
578 .route(
579 "/v1/data/:key",
580 get(data::data_get).post(data::data_emit_named),
581 )
582 .with_state(state)
583}
584
585/// Default registry = Subprocess + RustFn (baseline `identity` worker pre-baked) + Lua + AgentBlock + empty Operator factory.
586///
587/// `RustFnInProcessSpawnerFactory` gets one baseline entry (`fn_id = "identity"`)
588/// baked in via [`mlua_swarm::worker::baseline::extend_with_baseline`]. This
589/// is the shared bootstrap / smoke worker SoT across each binary (the server / MCP adapter /
590/// one-shot runner) — it structurally replaces the old per-binary inline echo injection.
591///
592/// Usage: default Task path at server startup. If production needs additional
593/// backends, callers bring in a different registry via
594/// `build_router_with(engine, custom_registry)`. The enhance flow
595/// (= patch-spawner / patch-applier / verifier-router / committer axes) uses
596/// [`default_registry_with_enhance_flow`].
597///
598/// The Operator factory is an empty shell with zero registrations (= sids are
599/// dynamically registered per WS connect; see the `operator_ws` module).
600pub fn default_registry() -> SpawnerRegistry {
601 let rustfn_factory =
602 mlua_swarm::worker::baseline::extend_with_baseline(RustFnInProcessSpawnerFactory::new());
603
604 let mut reg = SpawnerRegistry::new();
605 reg.register::<SubprocessProcessSpawnerFactory>(Arc::new(SubprocessProcessSpawnerFactory));
606 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(rustfn_factory));
607 // Empty `LuaInProcessSpawnerFactory`: no `fn_id` is pre-registered here,
608 // but BP agents can still declare `kind: lua` by carrying an inline
609 // `spec.source` (or a `$file`-expanded Lua chunk). This lets a BP ship
610 // deterministic Lua gates on the vanilla registry, without opting into
611 // the enhance flow. See `LuaInProcessSpawnerFactory` docs for the spec
612 // shape.
613 reg.register::<LuaInProcessSpawnerFactory>(Arc::new(LuaInProcessSpawnerFactory::new()));
614 // GH #86: the stateless AgentBlock factory belongs on the vanilla path
615 // too — every per-agent specialization lives in `AgentDef.spec` /
616 // `.profile` / `.runner`, so registering it here grants no enhance-flow
617 // capability, it only makes the first-class `AgentKind::AgentBlock`
618 // dispatchable. Before this, a BP declaring `kind = "agent_block"`
619 // compiled only under `--enable-enhance-flow`; the enhance branch below
620 // still differs by baking the enhance-flow Lua `fn_id`s.
621 reg.register::<AgentBlockInProcessSpawnerFactory>(Arc::new(
622 AgentBlockInProcessSpawnerFactory::new(),
623 ));
624 reg.register::<OperatorSpawnerFactory>(Arc::new(OperatorSpawnerFactory::new()));
625 reg
626}
627
628/// Opt-in registry that merges [`default_registry`] with the enhance flow
629/// (Lua factory + AgentBlock factory).
630///
631/// Selected via the server CLI flag `--enable-enhance-flow`. The enhance flow
632/// is a separate-axis wrapper, and this registry bakes both of its halves in as
633/// pipeline defaults:
634///
635/// - the **Lua factory** — the three enhance workers (`patch-applier` /
636/// `verifier-router` / `committer`) plus the three host bridges they call.
637/// Their Lua bodies are `include_str!`-embedded, so no file has to exist on
638/// disk for them to dispatch.
639/// - the **AgentBlock factory** — the `patch-spawner` axis. The bundled default
640/// declares no `spec.script_path`, so it runs in PromptBasedAgent mode with
641/// its `profile.system_prompt` carrying the whole `ops` / `bump` /
642/// `rationale` output contract; what it needs at run time is a credential for
643/// the provider behind its declared `profile.model` (`ANTHROPIC_API_KEY` for
644/// the model the bundled default declares). Driving the spawner on a
645/// different backend is a setting-level swap (`EnhanceSetting.spawner`), not
646/// a Blueprint rewrite.
647///
648/// The baseline RustFn (`identity`) is pre-baked the same way as in
649/// [`default_registry`]. End-to-end walkthrough of the flow (prerequisites,
650/// HTTP surface, spawner contract): the bundled `mse://guides/enhance-flow`.
651pub fn default_registry_with_enhance_flow() -> SpawnerRegistry {
652 let lua_factory =
653 mlua_swarm::enhance::blueprint::extend_factory(LuaInProcessSpawnerFactory::new());
654 // The Factory is stateless (= 1 process → 1 factory shared by all AgentDefs).
655 // Per-agent specialization (script_path / project_root, etc.) goes through AgentDef.spec.
656 // The enhance-flow patch-spawner is declared literally in agents[].spec of `default_blueprint.yaml`.
657 let agent_block_factory = AgentBlockInProcessSpawnerFactory::new();
658 let rustfn_factory =
659 mlua_swarm::worker::baseline::extend_with_baseline(RustFnInProcessSpawnerFactory::new());
660
661 let mut reg = SpawnerRegistry::new();
662 reg.register::<SubprocessProcessSpawnerFactory>(Arc::new(SubprocessProcessSpawnerFactory));
663 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(rustfn_factory));
664 reg.register::<LuaInProcessSpawnerFactory>(Arc::new(lua_factory));
665 reg.register::<AgentBlockInProcessSpawnerFactory>(Arc::new(agent_block_factory));
666 reg.register::<OperatorSpawnerFactory>(Arc::new(OperatorSpawnerFactory::new()));
667 reg
668}
669
670// ─── handlers ────────────────────────────────────────────────────────────
671
672async fn healthz() -> &'static str {
673 "ok"
674}
675
676/// Response body for `GET /v1/status` (issue #35 ST4 — lifecycle
677/// occupancy guard). Cheap-to-poll summary of "is it safe to kill this
678/// server right now".
679#[derive(Debug, Clone, Serialize, schemars::JsonSchema)]
680pub struct StatusResponse {
681 /// Count of `Run`s currently `Running` (`RunStore::list_running`).
682 /// Degrades to `0` on a store error rather than 500ing — see
683 /// module doc rationale.
684 pub running_runs: usize,
685 /// Count of attached Operator ids (`engine.list_operator_ids()`,
686 /// same idiom as `run_flow_form`'s Guard 1).
687 pub attached_operators: usize,
688}
689
690/// `GET /v1/status`. Infallible summary for the ST4 occupancy guard —
691/// store/engine query failures degrade the corresponding count to `0`
692/// (logged via `tracing::warn!`) rather than 500ing, since this
693/// endpoint may be polled frequently by a lifecycle-check caller that
694/// should not itself become a hang/error surface.
695async fn status_get(State(state): State<AppState>) -> Json<StatusResponse> {
696 let running_runs = state
697 .run_store
698 .list_running()
699 .await
700 .map(|v| v.len())
701 .unwrap_or_else(|e| {
702 tracing::warn!(error = %e, "status_get: list_running failed");
703 0
704 });
705 let attached_operators = state.engine.list_operator_ids().await.len();
706 Json(StatusResponse {
707 running_runs,
708 attached_operators,
709 })
710}
711
712#[derive(Deserialize)]
713struct AttachReq {
714 agent_id: String,
715 role: String,
716 ttl_secs: u64,
717}
718
719#[derive(Serialize)]
720struct AttachResp {
721 session_id: String,
722 role: String,
723}
724
725async fn sessions_attach(
726 State(state): State<AppState>,
727 Json(req): Json<AttachReq>,
728) -> Result<Json<AttachResp>, ApiError> {
729 let role = parse_role(&req.role)?;
730 let token = state
731 .engine
732 .attach(req.agent_id, role, Duration::from_secs(req.ttl_secs))
733 .await
734 .map_err(ApiError::engine)?;
735 // The wire `session_id` stays the nonce (Bearer credential contract);
736 // the server-side map key is its fingerprint (issue #14).
737 let sid = token.nonce.clone();
738 let key = token.fingerprint();
739 state.sessions.lock().await.map.insert(key, token);
740 Ok(Json(AttachResp {
741 session_id: sid,
742 role: req.role,
743 }))
744}
745
746async fn sessions_detach(
747 State(state): State<AppState>,
748 headers: HeaderMap,
749) -> Result<StatusCode, ApiError> {
750 let sid = extract_bearer(&headers)?;
751 let token = take_session_token(&state, &sid).await?;
752 state
753 .engine
754 .detach(&token)
755 .await
756 .map_err(ApiError::engine)?;
757 Ok(StatusCode::NO_CONTENT)
758}
759
760// ─── Unified /v1/tasks schema (= flow-eval path, Operator inject supported) ───────
761
762/// `/v1/tasks` POST schema. Uses the flow-eval path and supports Operator inject
763/// (kind / spawn_hook / senior_bridge). Expressing a one-shot task as a 1-Step
764/// Blueprint is the only correct model.
765///
766/// `pub` (issue #19 ST5) so its `schemars`-derived JSON Schema can be
767/// generated cross-crate by `mlua-swarm-cli`'s `mse://api/http-endpoints`
768/// MCP resource; fields stay module-private (no public field-level API
769/// surface is intended).
770#[derive(Deserialize, schemars::JsonSchema)]
771pub struct TaskLaunchRequest {
772 /// `BlueprintRef` selects Inline (a full Blueprint value) or Id (a
773 /// store lookup). Left opaque here — its own schema nests the full
774 /// `Blueprint` schema (owned by `mse://api/blueprint-schema`), and
775 /// mixing the two into this HTTP-endpoint resource would violate
776 /// their separation of concerns (see the resource's module doc).
777 #[schemars(with = "Value")]
778 blueprint: BlueprintRef,
779 /// flow.ir's initial `ctx` — every `Step.in` `$.<path>` reads from
780 /// here. This field's role is limited to the flow-ir eval seed
781 /// (issue #19); the Task-level execution context lives in the
782 /// sibling top-level fields below (`project_root` / `work_dir` /
783 /// `task_metadata`), promoted out of `init_ctx` to remove the
784 /// prior "free bag nested in free JSON" duplication.
785 ///
786 /// Backward compat: the pre-#19 shape — the same three keys nested
787 /// directly inside this object — is still honored as a fallback
788 /// when the sibling field is absent; see `run_flow_form`'s 2-stage
789 /// resolution and `TaskInputMiddleware::from_init_ctx`.
790 #[schemars(with = "Value")]
791 init_ctx: Value,
792 /// Task-level project root (issue #19 canonical Task IF field —
793 /// promoted out of `init_ctx`). Takes priority over a same-named
794 /// key nested inside `init_ctx` (backward-compat fallback).
795 #[serde(default)]
796 project_root: Option<String>,
797 /// Task-level working directory (issue #19), same priority rule as
798 /// `project_root`.
799 #[serde(default)]
800 work_dir: Option<String>,
801 /// Task-level arbitrary metadata bag (issue #19), same priority
802 /// rule as `project_root`.
803 #[serde(default)]
804 #[schemars(with = "Option<Value>")]
805 task_metadata: Option<Value>,
806 /// TTL in seconds. When unspecified (`None`), falls back in this order:
807 /// (1) `metadata.default_run_ttl_secs` from the resolved BP,
808 /// (2) if absent, the server global `default_run_ttl()` (1800s).
809 #[serde(default)]
810 ttl_secs: Option<u64>,
811 #[serde(default)]
812 operator: Option<OperatorReq>,
813 /// Explicit Operator session sid (or role alias) this task's entire Spawn
814 /// stream should be routed to (runtime Operator match stage 1).
815 ///
816 /// When `Some`, it is validated at request time against
817 /// `state.engine.list_operator_ids()` (the live `engine.operators`
818 /// registry key set): an unknown/never-registered id returns `400`
819 /// immediately — this is a deliberate hard-fail, in contrast to
820 /// `OperatorDelegateWrapped::spawn`, which silently falls through to
821 /// `inner.spawn` on a registry miss. A sid that *was* registered but has
822 /// since disconnected (WS `tx` cleared, session entry retained for
823 /// reconnect) passes this check and surfaces as an explicit dispatch-time
824 /// error instead (`WSOperatorSession::send_and_await` returns `Err` when
825 /// `tx` is `None`), which also propagates as a request failure rather
826 /// than a silent fallback.
827 ///
828 /// On success this value **overrides** `operator.operator_backend_id`
829 /// (last-write-wins, `operator_sid` takes priority) before the flow is
830 /// dispatched — see `run_flow_form`. Dispatch still only delegates if the
831 /// Blueprint opts into `spawner_hints.layers = ["operator_delegate"]`
832 /// (unchanged precondition, same as the existing `operator_backend_id`
833 /// field).
834 ///
835 /// The field also pins the **AgentSpec axis** (the per-agent
836 /// `spec.operator_ref` route every Blueprint with `kind = Operator`
837 /// agents uses, whether or not it declares the delegate layer):
838 /// `TaskApplicationInput.operator_pin` carries the sid down to the
839 /// compiler, which resolves those agents against the pinned session
840 /// instead of the role's current process-global holder, and to the
841 /// binding provider, which attests their manifests through the same
842 /// session. Blueprints keep declaring the logical role; which session
843 /// that role means for this run becomes a launch-time fact, recorded on
844 /// `RunRecord.operator_sid`. A pin naming no live session fails the
845 /// launch — there is no fallback to the role, because that fallback is
846 /// exactly how a run ends up on another driver's session.
847 ///
848 /// When unset, behavior is unchanged: whatever
849 /// `operator.operator_backend_id` / BP-level `operator_ref` alias
850 /// resolution already does still applies.
851 #[serde(default)]
852 operator_sid: Option<String>,
853 /// Per-request override for the sync launch's timeout ceiling (GH #33
854 /// Guard 2, see `run_flow_form`'s doc comment). `None` (the default;
855 /// existing clients are unaffected) falls back to
856 /// `AppState::sync_timeout_secs` (server config), then the built-in
857 /// default (300s). `Some(0)` is rejected with `400` — omit the field
858 /// to defer to the server default rather than sending an explicit
859 /// zero.
860 #[serde(default)]
861 timeout_secs: Option<u64>,
862 /// Human-facing description of the work item (e.g. "resolve issue #10"),
863 /// stashed verbatim into the minted `TaskRecord.goal`. Omitted / `None`
864 /// stores an empty string — the flow-eval path itself never reads it.
865 #[serde(default)]
866 goal: Option<String>,
867 /// The "launch request" tier (tier 1, highest
868 /// priority) of the `check_policy` cascade
869 /// (`launch request > blueprint > server config`). `None` (the default;
870 /// existing clients are unaffected) leaves the tier unspecified so the
871 /// Blueprint-declared `check_policy` and, failing that, the server-wide
872 /// `EngineCfg.check_policy` default decide. Wire form is snake_case
873 /// (`"silent"` / `"warn"` / `"strict"`). Threaded verbatim into
874 /// `TaskApplicationInput.check_policy`.
875 #[serde(default)]
876 check_policy: Option<CheckPolicy>,
877 /// GH #37: opt into the detached (asynchronous) launch. `false` (the
878 /// default; existing clients are unaffected) keeps the synchronous
879 /// launch: the handler drives the flow eval inline and returns the
880 /// `final_ctx` on completion. `true` spawns the flow eval as a
881 /// detached background task and returns `202 Accepted` immediately
882 /// with `{task_id, run_id, status: "running"}` (`final_ctx` is
883 /// `null`) — the run's only lifetime bound is `ttl_secs`, and its
884 /// outcome is observed via `GET /v1/runs/:id` (or the `swarm_status`
885 /// MCP tool). Mutually exclusive with `timeout_secs` (the sync-launch
886 /// ceiling has no meaning for a detached run; combining them is a
887 /// `400`).
888 #[serde(default)]
889 detach: bool,
890}
891
892/// Operator inject sub-schema of [`TaskLaunchRequest`] (`kind` / `id` /
893/// `spawn_hook_id` / `senior_bridge_id` / `operator_backend_id` /
894/// `per_agent_kinds`). `pub` for the same cross-crate schema-generation
895/// reason as `TaskLaunchRequest`.
896#[derive(Deserialize, Default, schemars::JsonSchema)]
897pub struct OperatorReq {
898 /// `main_ai` / `automate` / `composite`. This is the "Runtime Global"
899 /// tier of the 4-tier `OperatorKind` cascade (see `mlua_swarm
900 /// ::ctx::collapse_operator_kind`); when unspecified, falls through to
901 /// the BP-level tiers (`OperatorDef.kind` / `Blueprint
902 /// .default_operator_kind`) instead of eagerly defaulting to `automate`.
903 #[serde(default)]
904 kind: Option<String>,
905 /// Operator id at attach time (= sessions tracking key in the EventLog); unspecified defaults to `"http-run"`.
906 #[serde(default)]
907 id: Option<String>,
908 /// Name of a hook pre-registered via `engine.register_spawn_hook`; `None` if unspecified.
909 #[serde(default)]
910 spawn_hook_id: Option<String>,
911 /// Name of a bridge pre-registered via `engine.register_senior_bridge`; `None` if unspecified.
912 #[serde(default)]
913 senior_bridge_id: Option<String>,
914 /// Name of an Operator backend pre-registered via `engine.register_operator`
915 /// (= the path that delegates the entire spawn to an external Operator);
916 /// `None` if unspecified. When `kind == MainAi/Composite` and this id is `Some`,
917 /// `OperatorDelegateMiddleware` bypasses `inner.spawn` and calls `operator.execute` instead.
918 /// This is a different axis from `operator.id` (= session tracking label);
919 /// `operator_backend_id` is the registry lookup key.
920 #[serde(default)]
921 operator_backend_id: Option<String>,
922 /// "Runtime Agent-level" tier (highest priority) of the `OperatorKind`
923 /// cascade — per-agent override, keyed by `AgentDef.name`, value is
924 /// `main_ai` / `automate` / `composite` (same parsing as `kind`).
925 /// `None` / absent means no per-agent override.
926 #[serde(default)]
927 per_agent_kinds: Option<HashMap<String, String>>,
928}
929
930/// Parse a wire-level kind string (`"main_ai"` / `"automate"` / `"composite"`)
931/// into `OperatorKind`. Shared by `OperatorReq.kind` and
932/// `OperatorReq.per_agent_kinds` values.
933fn parse_operator_kind_str(s: &str) -> Result<mlua_swarm::OperatorKind, ApiError> {
934 use mlua_swarm::OperatorKind;
935 match s {
936 "main_ai" => Ok(OperatorKind::MainAi),
937 "composite" => Ok(OperatorKind::Composite),
938 "automate" => Ok(OperatorKind::Automate),
939 other => Err(ApiError::bad_request(format!(
940 "operator kind: unknown value '{other}' (expected main_ai|automate|composite)"
941 ))),
942 }
943}
944
945/// `/v1/tasks` POST response body. `pub` for the same cross-crate
946/// schema-generation reason as [`TaskLaunchRequest`].
947#[derive(Serialize, schemars::JsonSchema)]
948pub struct TaskLaunchResponse {
949 /// The final flow.ir `ctx` after every `Step.out` has been written.
950 #[schemars(with = "Value")]
951 final_ctx: Value,
952 /// Debug-formatted `BlueprintVersion` the run resolved against, when
953 /// the Blueprint came from a store lookup (`None` for `Inline` refs).
954 bound_version: Option<String>,
955 /// Resolved TTL (seconds) actually applied to the run. Exposes the
956 /// 3-layer cascade (request body → BP metadata → server default) so
957 /// clients can verify which value took effect without re-deriving it.
958 effective_ttl_secs: u64,
959 /// Which layer of the TTL cascade won.
960 ttl_source: TtlSource,
961 /// The `TaskRecord` minted for this request (issue #13 ID-hierarchy
962 /// persistence). `GET /v1/tasks/:id` re-fetches it; `POST
963 /// /v1/tasks/:id/runs` re-kicks it under a fresh `RunId`.
964 #[schemars(with = "String")]
965 task_id: TaskId,
966 /// The `RunRecord` minted for this specific kick. `GET /v1/runs/:id`
967 /// re-fetches it (`step_entries` included).
968 #[schemars(with = "String")]
969 run_id: RunId,
970 /// Launch outcome at response time (GH #37). The synchronous path
971 /// (default) reports `done` — the flow eval completed before this
972 /// response was built. A detached launch (`detach: true`) reports
973 /// `running` — the eval continues in the background; poll `GET
974 /// /v1/runs/:id` for the terminal status and result.
975 status: RunStatus,
976}
977
978/// `tasks_start`'s reply — a [`TaskLaunchResponse`] plus the HTTP status
979/// it rides out on (`200 OK` for the synchronous path, `202 Accepted` for
980/// a detached launch, GH #37). A tuple struct with the body first so
981/// handler-level tests keep their established `.0` access to the response
982/// body regardless of which path produced it.
983pub struct TaskLaunchReply(pub TaskLaunchResponse, pub StatusCode);
984
985impl IntoResponse for TaskLaunchReply {
986 fn into_response(self) -> Response {
987 (self.1, Json(self.0)).into_response()
988 }
989}
990
991/// Which layer of the TTL cascade (request body → BP metadata → server
992/// default) resolved [`TaskLaunchResponse::effective_ttl_secs`]. `pub` for
993/// the same cross-crate schema-generation reason as `TaskLaunchRequest`.
994#[derive(Serialize, Clone, Copy, Debug, PartialEq, Eq, schemars::JsonSchema)]
995#[serde(rename_all = "snake_case")]
996pub enum TtlSource {
997 /// The request body's `ttl_secs` was set explicitly.
998 RequestBody,
999 /// The request body omitted `ttl_secs`; the resolved Blueprint's
1000 /// `metadata.default_run_ttl_secs` was set.
1001 BpMetadata,
1002 /// Both the request body and the Blueprint metadata omitted a TTL;
1003 /// the server-global `default_run_ttl()` (1800s) applied.
1004 ServerDefault,
1005}
1006
1007/// Unified `/v1/tasks` POST entry (= Flow form only).
1008/// Runs `Blueprint.flow` to completion via flow eval in a single round-trip.
1009/// One-shot tasks are also expressed as a 1-Step Blueprint. Operator
1010/// (kind / spawn_hook / senior_bridge) can be injected per request body.
1011/// `operator_sid` (S2, runtime Operator match stage 1) additionally
1012/// lets the caller pin the task to a specific already-registered Operator
1013/// session sid, bypassing BP-level alias lookup — see `TaskLaunchRequest` doc.
1014async fn tasks_start(
1015 State(state): State<AppState>,
1016 Json(req): Json<TaskLaunchRequest>,
1017) -> Result<TaskLaunchReply, ApiError> {
1018 run_flow_form(&state, req).await
1019}
1020
1021/// Flow-form path (= via `TaskApplication::handle_with_run`).
1022/// Core handler behind the `/v1/tasks` entry (`tasks_start`).
1023///
1024/// Engine stateless-executor refactor: the per-request
1025/// sub_engine + 3-registry propagate loop is retired; the startup-built
1026/// `state.task_app` (= a `TaskLaunchService` wrap around `state.engine`) is
1027/// used directly. The Operator callback IF (`spawn_hook_id` /
1028/// `senior_bridge_id` / `operator_backend_id`) is registered on
1029/// `state.engine.register_*` at WS connect time — the engine is the SoT.
1030/// See the `operator_ws` module doc for details.
1031///
1032/// # GH #33 — sync-hang guards
1033///
1034/// This handler is always synchronous end-to-end (no sync/async branch);
1035/// two fail-loud guards keep a bad launch from hanging the HTTP request
1036/// forever:
1037///
1038/// - **Guard 1 (readiness precheck, `503`)**: when the request/BP
1039/// references an operator backend (`operator.operator_backend_id`, set
1040/// directly or via `operator_sid`) and `state.engine.list_operator_ids()`
1041/// is empty, the request fails immediately rather than dispatching into
1042/// a session with nothing attached to serve it. Coarse by design — a
1043/// launch that cannot be cheaply determined to route through an operator
1044/// is never rejected here (Guard 2 still covers the hang in that case).
1045/// - **Guard 2 (sync timeout, `504`)**: the `handle_with_run` driver is
1046/// wrapped in `tokio::time::timeout`. Ceiling cascade, highest priority
1047/// first: request `timeout_secs` (rejecting `Some(0)` with `400`), then
1048/// `AppState::sync_timeout_secs` (server config), then the built-in
1049/// default (300s). On expiry the timed-out future is dropped — this
1050/// cancels the in-process flow eval (the flow is abandoned, not
1051/// resumed; intended v1 semantics) — and the Task/Run records are
1052/// best-effort marked `Failed` so they do not stay `Running` forever.
1053///
1054/// # Driver lifetime (survives client disconnect)
1055///
1056/// The synchronous path spawns its driver exactly like the detached one
1057/// below: the eval, the Guard 2 ceiling, the panic guard and
1058/// `finalize_run` all live in a `tokio::spawn`ed task, and the handler
1059/// only awaits that task's verdict over a `oneshot`. A client disconnect
1060/// therefore drops the *wait*, not the run — before this, dropping the
1061/// request future dropped the driver with it, and any `/v1/worker/submit`
1062/// that arrived afterwards was written into `EngineState` with no reader
1063/// left to fold it (the Run then sat `Running` until the stale-run
1064/// sweeper reaped it).
1065///
1066/// # GH #37 — detached launch (`detach: true`)
1067///
1068/// `detach: true` additionally decouples the *response* from the run: the
1069/// driver's only lifetime bound is the resolved `ttl_secs` (marked
1070/// `Failed` on expiry, same best-effort persistence as Guard 2), and the
1071/// handler returns `202 Accepted` with `status: "running"` immediately
1072/// instead of waiting for the terminal outcome. Guard 1 still applies
1073/// (checked before any store write); Guard 2's ceiling does not
1074/// (`timeout_secs` + `detach` together is a `400`).
1075async fn run_flow_form(
1076 state: &AppState,
1077 req: TaskLaunchRequest,
1078) -> Result<TaskLaunchReply, ApiError> {
1079 use mlua_swarm::application::{
1080 BlueprintRef as AppBlueprintRef, TaskApplicationInput, TaskApplicationOutput,
1081 };
1082 use mlua_swarm::OperatorKind;
1083
1084 // Snapshot everything the TaskRecord needs before `req.blueprint` /
1085 // `req.init_ctx` are moved into the dispatch path below.
1086 let blueprint_ref_json = serde_json::to_value(&req.blueprint)
1087 .map_err(|e| ApiError::bad_request(format!("blueprint snapshot: {e}")))?;
1088 let input_ctx_snapshot = req.init_ctx.clone();
1089 let goal = req.goal.clone().unwrap_or_default();
1090
1091 // issue #19 ST2: resolve the Task-level canonical fields
1092 // (`project_root` / `work_dir` / `task_metadata`) once, at the wire
1093 // boundary. Sibling top-level fields on the request body take
1094 // priority; the pre-#19 shape (same key nested inside `init_ctx`) is
1095 // only a fallback for legacy callers. The result is threaded straight
1096 // through as `TaskApplicationInput.task_input` — `init_ctx` itself is
1097 // NOT mutated, so it stays a pure flow-ir eval seed identical to
1098 // whatever the caller sent.
1099 let task_input_spec = build_task_input_spec_from_request(&req);
1100 // Issue #19 ST4: snapshot the resolved spec into the `TaskRecord` (JSON,
1101 // same "bare `Value`" rationale as `blueprint_ref_json` /
1102 // `input_ctx_snapshot` above) so `POST /v1/tasks/:id/runs` can resolve
1103 // it back out on rekick without re-deriving it from a since-stale
1104 // request body. Cloned rather than computed from `task_input_spec`
1105 // after the fact — the original is still moved into
1106 // `TaskApplicationInput.task_input` below.
1107 let task_input_spec_snapshot = task_input_spec
1108 .clone()
1109 .map(|spec| serde_json::to_value(&spec))
1110 .transpose()
1111 .map_err(|e| ApiError::bad_request(format!("task_input_spec snapshot: {e}")))?;
1112 let init_ctx = req.init_ctx.clone();
1113
1114 let mut op_req = req.operator.unwrap_or_default();
1115
1116 // S2: explicit `operator_sid` override (runtime Operator match stage 1).
1117 // Resolved *before* building `operator_kind` / dispatching so an
1118 // unknown sid fails fast with a 400, never silently falling back to the
1119 // BP-level alias lookup. See `TaskLaunchRequest::operator_sid` doc for the
1120 // disconnected-vs-unknown distinction.
1121 if let Some(sid) = &req.operator_sid {
1122 let known_ids = state.engine.list_operator_ids().await;
1123 if !known_ids.iter().any(|id| id == sid) {
1124 return Err(ApiError::bad_request(format!(
1125 "operator_sid: no such registered operator session '{sid}'"
1126 )));
1127 }
1128 op_req.operator_backend_id = Some(sid.clone());
1129 }
1130
1131 // GH #33 Guard 2 ceiling resolution: request field > server config >
1132 // built-in default (300s, `config::default_sync_timeout_secs`).
1133 // Validated up front — before any TaskRecord/RunRecord side effects —
1134 // so a caller-supplied `Some(0)` fails fast with `400` rather than
1135 // minting records for a launch that was never going to dispatch.
1136 // GH #37: `detach: true` makes the sync ceiling meaningless (the
1137 // detached run is bounded by `ttl_secs` alone) — combining the two
1138 // is rejected here, same fail-fast-before-side-effects ordering.
1139 let detach = req.detach;
1140 let sync_timeout_secs = match (detach, req.timeout_secs) {
1141 (true, Some(_)) => {
1142 return Err(ApiError::bad_request(
1143 "timeout_secs is the synchronous launch ceiling and does not apply to a \
1144 detached launch (detach: true), whose lifetime bound is ttl_secs — omit \
1145 timeout_secs"
1146 .into(),
1147 ));
1148 }
1149 (false, Some(0)) => {
1150 return Err(ApiError::bad_request(
1151 "timeout_secs: 0 is invalid; omit the field to use the server default".into(),
1152 ));
1153 }
1154 (false, Some(v)) => v,
1155 (_, None) => state.sync_timeout_secs,
1156 };
1157
1158 // GH #33 Guard 1: operator readiness precheck. Coarse signal — this
1159 // handler can cheaply see whether the request/BP references an
1160 // operator backend (`operator.operator_backend_id`, set directly or
1161 // resolved above from `operator_sid`), but not the full
1162 // `OperatorDelegateMiddleware` routing decision (that also considers
1163 // BP-level `kind` tiers, resolved only at dispatch time). When a
1164 // backend is referenced and *zero* operators are attached at all,
1165 // fail fast rather than dispatching into a session nothing can serve.
1166 // A launch this coarse check cannot positively identify as
1167 // operator-delegate is never rejected here — Guard 2 (the timeout
1168 // wrap below) still covers the hang in that case.
1169 if let Some(backend_id) = op_req.operator_backend_id.as_deref() {
1170 let attached = state.engine.list_operator_ids().await;
1171 if attached.is_empty() {
1172 return Err(ApiError::unavailable(format!(
1173 "no operator attached to serve this launch (operator backend '{backend_id}' \
1174 requested): attach an operator via POST /v1/operators + WS, or use the \
1175 poll-style flow (GET /v1/worker/prompt + POST /v1/worker/submit)"
1176 )));
1177 }
1178 }
1179
1180 // "Runtime Global" tier: `Some(_)` — including `Some(Automate)` — is
1181 // always an explicit request that outranks the BP-level tiers; an
1182 // absent/unset `kind` in the request body stays `None`, leaving the
1183 // BP-level tiers (`OperatorDef.kind` / `Blueprint.default_operator_kind`)
1184 // to decide instead of eagerly defaulting to `Automate`.
1185 let operator_kind = op_req
1186 .kind
1187 .as_deref()
1188 .map(parse_operator_kind_str)
1189 .transpose()?;
1190 let operator_id = op_req.id.unwrap_or_else(|| "http-run".to_string());
1191 // "Runtime Agent-level" tier: per-agent overrides. Absent/empty = no
1192 // override for any agent, letting the BP-level tiers decide per agent.
1193 let mut operator_kind_overrides: HashMap<String, OperatorKind> = HashMap::new();
1194 for (agent, kind_str) in op_req.per_agent_kinds.take().unwrap_or_default() {
1195 operator_kind_overrides.insert(agent, parse_operator_kind_str(&kind_str)?);
1196 }
1197
1198 let blueprint: AppBlueprintRef = match req.blueprint {
1199 AppBlueprintRef::Inline { value } => AppBlueprintRef::Inline { value },
1200 AppBlueprintRef::Id { id, version } => AppBlueprintRef::Id { id, version },
1201 };
1202
1203 // TTL resolution cascade: (1) request body value, (2) BP metadata `default_run_ttl_secs`,
1204 // (3) server global default (`default_run_ttl()`, 1800s).
1205 let (ttl_secs, ttl_source) = match req.ttl_secs {
1206 Some(v) => (v, TtlSource::RequestBody),
1207 None => {
1208 let (resolved_bp, _ver) = state
1209 .task_app
1210 .resolve(&blueprint)
1211 .await
1212 .map_err(|e| ApiError::from_task_resolve(&e, "bp resolve"))?;
1213 match resolved_bp.metadata.default_run_ttl_secs {
1214 Some(v) => (v, TtlSource::BpMetadata),
1215 None => (default_run_ttl(), TtlSource::ServerDefault),
1216 }
1217 }
1218 };
1219
1220 // Build the launch input up front so a snapshot of it can be persisted
1221 // into the RunRecord below — an Interrupted Run is resumed from that
1222 // snapshot (`POST /v1/runs/:id/resume`) under the same run_id.
1223 let input = TaskApplicationInput {
1224 blueprint,
1225 operator_id: operator_id.clone(),
1226 role: Role::Operator,
1227 ttl: Duration::from_secs(ttl_secs),
1228 init_ctx,
1229 operator_kind,
1230 bridge_id: op_req.senior_bridge_id,
1231 hook_id: op_req.spawn_hook_id,
1232 operator_backend_id: op_req.operator_backend_id,
1233 // Axis-independent half of `operator_sid` (see its doc on
1234 // `TaskLaunchRequest`): the same sid binds this launch's AgentSpec
1235 // axis — Operator agents compile against the pinned session and
1236 // their manifests are attested through it — while the field above
1237 // keeps feeding the opt-in delegate layer unchanged.
1238 operator_pin: req.operator_sid.clone(),
1239 operator_kind_overrides,
1240 task_input: task_input_spec,
1241 // The request-body top-level `check_policy` (tier 1)
1242 // flows straight into the cascade resolved once in
1243 // `TaskLaunchService::launch`.
1244 check_policy: req.check_policy,
1245 };
1246 let input_json = Some(tasks::snapshot_launch_input(&input)?);
1247
1248 // issue #13 ID-hierarchy persistence: mint the work-item identity (Task)
1249 // and this kick's identity (Run) *before* dispatching, so a Task/Run
1250 // pair always exists even if the flow itself fails mid-way (the
1251 // Failed-status paths below still have a row to update).
1252 let task_id = TaskId::new();
1253 let run_id = RunId::new();
1254 let now = tasks::now_secs();
1255 state
1256 .task_store
1257 .create(TaskRecord {
1258 id: task_id.clone(),
1259 goal,
1260 blueprint_ref: blueprint_ref_json,
1261 input_ctx: input_ctx_snapshot,
1262 task_input_spec: task_input_spec_snapshot,
1263 status: TaskRecordStatus::Running,
1264 created_at: now,
1265 updated_at: now,
1266 })
1267 .await
1268 .map_err(ApiError::engine)?;
1269 state
1270 .run_store
1271 .create(RunRecord {
1272 id: run_id.clone(),
1273 task_id: task_id.clone(),
1274 status: RunStatus::Running,
1275 step_entries: Vec::new(),
1276 degradations: Vec::new(),
1277 operator_sid: req.operator_sid.clone(),
1278 result_ref: None,
1279 input_json,
1280 created_at: now,
1281 updated_at: now,
1282 })
1283 .await
1284 .map_err(ApiError::engine)?;
1285
1286 let trace =
1287 mlua_swarm::store::trace::TraceHandle::new(run_id.clone(), state.run_trace_store.clone());
1288 trace
1289 .append(
1290 mlua_swarm::store::trace::kind::RUN_STARTED,
1291 None,
1292 None,
1293 json!({"mode": "launch"}),
1294 )
1295 .await;
1296 let run_ctx = RunContext::new(run_id.clone(), state.run_store.clone())
1297 .with_replay_store(state.replay_store.clone())
1298 .with_trace(trace);
1299
1300 // GH #37 detached launch: the eval driver runs in its own spawned
1301 // task — its lifetime is bound to `ttl_secs`, not to this request's
1302 // future (client disconnect / handler completion cannot cancel it).
1303 // The spawned task owns the run to its terminal status: `finalize_run`
1304 // on completion, or the same best-effort `Failed` marking as Guard 2
1305 // if the ttl ceiling expires first.
1306 if detach {
1307 let bg_state = state.clone();
1308 let bg_task_id = task_id.clone();
1309 let bg_run_id = run_id.clone();
1310 // Panic guard (see `tasks::catch_run_panic`): without it a panic in
1311 // the driver unwinds this whole spawned task — timeout combinator
1312 // included — and strands the Run in `Running`.
1313 let guard_state = state.clone();
1314 let guard_task_id = task_id.clone();
1315 let guard_run_id = run_id.clone();
1316 tokio::spawn(async move {
1317 let driver = async move {
1318 let outcome = match tokio::time::timeout(
1319 Duration::from_secs(ttl_secs),
1320 bg_state.task_app.handle_with_run(input, Some(run_ctx)),
1321 )
1322 .await
1323 {
1324 Ok(outcome) => outcome,
1325 Err(_elapsed) => {
1326 let reason = json!({
1327 "error": format!("detached run exceeded {ttl_secs}s ttl ceiling"),
1328 });
1329 if let Err(e) = bg_state.run_store.set_result(&bg_run_id, reason).await {
1330 tracing::warn!(%bg_run_id, error = %e, "run_flow_form: detached ttl set_result failed");
1331 }
1332 if let Err(e) = bg_state
1333 .run_store
1334 .update_status(&bg_run_id, RunStatus::Failed)
1335 .await
1336 {
1337 tracing::warn!(%bg_run_id, error = %e, "run_flow_form: detached ttl run update_status(Failed) failed");
1338 }
1339 if let Err(e) = bg_state
1340 .task_store
1341 .update_status(&bg_task_id, TaskRecordStatus::Failed)
1342 .await
1343 {
1344 tracing::warn!(%bg_task_id, error = %e, "run_flow_form: detached ttl task update_status(Failed) failed");
1345 }
1346 // This arm never reaches `finalize_run`, so the trace
1347 // stream gets its terminal marker here.
1348 mlua_swarm::store::trace::TraceHandle::new(
1349 bg_run_id.clone(),
1350 bg_state.run_trace_store.clone(),
1351 )
1352 .append(
1353 mlua_swarm::store::trace::kind::RUN_FINISHED,
1354 None,
1355 None,
1356 json!({ "status": "failed", "reason": format!("ttl {ttl_secs}s exceeded") }),
1357 )
1358 .await;
1359 return;
1360 }
1361 };
1362 // `finalize_run` persists both the Ok and Err outcomes itself;
1363 // the passthrough return value has no consumer here.
1364 let _ = tasks::finalize_run(&bg_state, &bg_task_id, &bg_run_id, outcome).await;
1365 };
1366 let _ = tasks::catch_run_panic(
1367 &guard_state,
1368 &guard_task_id,
1369 &guard_run_id,
1370 "launch.detach",
1371 driver,
1372 )
1373 .await;
1374 });
1375 return Ok(TaskLaunchReply(
1376 TaskLaunchResponse {
1377 final_ctx: Value::Null,
1378 bound_version: None,
1379 effective_ttl_secs: ttl_secs,
1380 ttl_source,
1381 task_id,
1382 run_id,
1383 status: RunStatus::Running,
1384 },
1385 StatusCode::ACCEPTED,
1386 ));
1387 }
1388
1389 // GH #33 Guard 2 + driver-lifetime fix: the driver runs in its own spawned
1390 // task — the same shape as the detached branch above — and this
1391 // handler only awaits its verdict over a `oneshot`. What the client's
1392 // connection owns is therefore the *wait*, not the run: a disconnect
1393 // drops the receiver while the spawned driver keeps going to its own
1394 // terminal step (`finalize_run`, or the ceiling's `Failed` marking),
1395 // so a `/v1/worker/submit` arriving after the disconnect still has a
1396 // driver to fold it into.
1397 //
1398 // Guard 2's ceiling still bounds the driver itself: on expiry the
1399 // timed-out future is dropped, cancelling the in-process flow eval —
1400 // the flow is abandoned, not resumed (intended v1 semantics;
1401 // stage-granularity resume is a coarser guarantee than this handler
1402 // makes, out of scope here). No second handler-side timeout exists:
1403 // the driver self-bounds, so this await ends when the driver ends.
1404 //
1405 // Wrapped in the panic guard (`tasks::catch_run_panic`) so a panicking
1406 // driver returns a structured 500 with a resumable `Interrupted` Run
1407 // instead of unwinding the spawned task with the Run stuck `Running`.
1408 let (tx, rx) = tokio::sync::oneshot::channel::<Result<TaskApplicationOutput, ApiError>>();
1409 let bg_state = state.clone();
1410 let bg_task_id = task_id.clone();
1411 let bg_run_id = run_id.clone();
1412 let guard_state = state.clone();
1413 let guard_task_id = task_id.clone();
1414 let guard_run_id = run_id.clone();
1415 tokio::spawn(async move {
1416 let driver = async move {
1417 let outcome = match tokio::time::timeout(
1418 Duration::from_secs(sync_timeout_secs),
1419 bg_state.task_app.handle_with_run(input, Some(run_ctx)),
1420 )
1421 .await
1422 {
1423 Ok(outcome) => outcome,
1424 Err(_elapsed) => {
1425 // Best effort: mark the Task/Run so they do not stay
1426 // `Running` forever. Reuses the existing `Failed`
1427 // variant (no new schema-crate enum additions) and
1428 // stashes a reason string into `RunRecord.result_ref`
1429 // — the only free-form field the Run schema carries;
1430 // secondary persistence failures here are logged and
1431 // swallowed, mirroring `tasks::finalize_run`'s
1432 // error-path convention.
1433 let reason = json!({
1434 "error": format!("sync launch exceeded {sync_timeout_secs}s timeout ceiling"),
1435 });
1436 if let Err(e) = bg_state.run_store.set_result(&bg_run_id, reason).await {
1437 tracing::warn!(%bg_run_id, error = %e, "run_flow_form: timeout run set_result failed");
1438 }
1439 if let Err(e) = bg_state
1440 .run_store
1441 .update_status(&bg_run_id, RunStatus::Failed)
1442 .await
1443 {
1444 tracing::warn!(%bg_run_id, error = %e, "run_flow_form: timeout run update_status(Failed) failed");
1445 }
1446 if let Err(e) = bg_state
1447 .task_store
1448 .update_status(&bg_task_id, TaskRecordStatus::Failed)
1449 .await
1450 {
1451 tracing::warn!(%bg_task_id, error = %e, "run_flow_form: timeout task update_status(Failed) failed");
1452 }
1453 return Err(ApiError::timeout(format!(
1454 "sync launch exceeded {sync_timeout_secs}s timeout ceiling: the in-process flow \
1455 eval was abandoned (dropping the future cancels it); attach an operator that \
1456 acks promptly (POST /v1/operators + WS), or raise timeout_secs / sync_timeout_secs"
1457 )));
1458 }
1459 };
1460 tasks::finalize_run(&bg_state, &bg_task_id, &bg_run_id, outcome)
1461 .await
1462 .map_err(flow_eval_error_to_api_error)
1463 };
1464 let reply = match tasks::catch_run_panic(
1465 &guard_state,
1466 &guard_task_id,
1467 &guard_run_id,
1468 "launch.sync",
1469 driver,
1470 )
1471 .await
1472 {
1473 Ok(reply) => reply,
1474 Err(msg) => Err(ApiError::engine(format!(
1475 "run driver panicked: {msg}; the run was marked Interrupted and can be resumed \
1476 via POST /v1/runs/{guard_run_id}/resume"
1477 ))),
1478 };
1479 // A disconnected client leaves no receiver; the run is already
1480 // persisted, so the undeliverable reply is dropped.
1481 let _ = tx.send(reply);
1482 });
1483
1484 // Only reachable if the spawned task died without sending — a panic
1485 // outside the guard, or a runtime shutdown.
1486 let out = rx.await.map_err(|_| {
1487 ApiError::engine(format!(
1488 "run driver task ended without reporting an outcome; see GET /v1/runs/{run_id} \
1489 for the run's persisted status"
1490 ))
1491 })??;
1492
1493 Ok(TaskLaunchReply(
1494 TaskLaunchResponse {
1495 final_ctx: out.final_ctx,
1496 bound_version: out.bound_version.map(|v| format!("{:?}", v)),
1497 effective_ttl_secs: ttl_secs,
1498 ttl_source,
1499 task_id,
1500 run_id,
1501 status: RunStatus::Done,
1502 },
1503 StatusCode::OK,
1504 ))
1505}
1506
1507/// issue #19 ST2 direct sibling-field resolver — extracts the three
1508/// Task-level canonical fields (`project_root` / `work_dir` /
1509/// `task_metadata`) once at the wire boundary. Sibling top-level body
1510/// fields take priority; the pre-#19 shape (same key nested inside
1511/// `init_ctx`) is only a fallback for legacy callers. Unlike the ST1
1512/// `resolve_task_level_init_ctx` bridge this replaced, `init_ctx` is
1513/// NOT mutated — the resolved values are handed straight to
1514/// [`mlua_swarm::service::TaskLaunchInput::task_input`], keeping
1515/// `init_ctx` a pure flow-ir eval seed.
1516///
1517/// Returns `None` when all three fields resolve to `None` (no
1518/// middleware is layered onto the spawner stack downstream — the
1519/// [`mlua_swarm::middleware::task_input::TaskInputMiddleware::new_from_fields`]
1520/// contract).
1521fn build_task_input_spec_from_request(
1522 req: &TaskLaunchRequest,
1523) -> Option<mlua_swarm::service::TaskInputSpec> {
1524 let project_root = req.project_root.clone().or_else(|| {
1525 req.init_ctx
1526 .get("project_root")
1527 .and_then(Value::as_str)
1528 .map(String::from)
1529 });
1530 let work_dir = req.work_dir.clone().or_else(|| {
1531 req.init_ctx
1532 .get("work_dir")
1533 .and_then(Value::as_str)
1534 .map(String::from)
1535 });
1536 let task_metadata = req.task_metadata.clone().or_else(|| {
1537 req.init_ctx
1538 .get("task_metadata")
1539 .filter(|v| v.is_object())
1540 .cloned()
1541 });
1542
1543 if project_root.is_none() && work_dir.is_none() && task_metadata.is_none() {
1544 None
1545 } else {
1546 Some(mlua_swarm::service::TaskInputSpec {
1547 project_root,
1548 work_dir,
1549 task_metadata,
1550 })
1551 }
1552}
1553
1554// ─── helpers ─────────────────────────────────────────────────────────────
1555
1556async fn take_session_token(state: &AppState, sid: &str) -> Result<CapToken, ApiError> {
1557 // `sid` on this path is the token nonce itself (a bearer secret), so
1558 // both the map key and the not-found diagnostic use its fingerprint
1559 // (issue #14 — never echo the nonce back in an error body).
1560 let key = mlua_swarm::types::token_fingerprint(sid);
1561 state
1562 .sessions
1563 .lock()
1564 .await
1565 .map
1566 .remove(&key)
1567 .ok_or_else(|| ApiError::not_found(format!("session: fp={key}")))
1568}
1569
1570/// Extracts sid from `Authorization: Bearer <sid>`. Strict — does not accept any other scheme prefix.
1571fn extract_bearer(headers: &HeaderMap) -> Result<String, ApiError> {
1572 let v = headers
1573 .get(AUTHORIZATION)
1574 .ok_or_else(|| ApiError::bad_request("missing Authorization header".into()))?
1575 .to_str()
1576 .map_err(|_| ApiError::bad_request("invalid Authorization header encoding".into()))?;
1577 let sid = v
1578 .strip_prefix("Bearer ")
1579 .ok_or_else(|| ApiError::bad_request("Authorization must be 'Bearer <sid>'".into()))?
1580 .trim();
1581 if sid.is_empty() {
1582 return Err(ApiError::bad_request("Bearer sid is empty".into()));
1583 }
1584 Ok(sid.to_string())
1585}
1586
1587fn parse_role(s: &str) -> Result<Role, ApiError> {
1588 match s.to_ascii_lowercase().as_str() {
1589 "operator" => Ok(Role::Operator),
1590 "worker" => Ok(Role::Worker),
1591 "observer" => Ok(Role::Observer),
1592 "senior" => Ok(Role::Senior),
1593 other => Err(ApiError::bad_request(format!("unknown role: {other}"))),
1594 }
1595}
1596
1597// ─── error type ──────────────────────────────────────────────────────────
1598
1599/// GH #76 error surface: adapter that lifts a [`TaskApplicationError`] into an
1600/// [`ApiError`], surfacing the structured
1601/// [`TaskLaunchError::FlowEval`] fields
1602/// (`failed_step` / `verdict_value` / `partial_ctx`) into the response
1603/// body's `details` object when the abort originated from a Blueprint
1604/// step. Every other error variant collapses to the pre-#76
1605/// `bad_request(format!("run: {e}"))` shape byte-for-byte, so callers
1606/// that only match on the `{"error": message}` message keep working.
1607fn flow_eval_error_to_api_error(e: TaskApplicationError) -> ApiError {
1608 if let TaskApplicationError::Launch(TaskLaunchError::FlowEval {
1609 message,
1610 failed_step,
1611 verdict_value,
1612 partial_ctx,
1613 }) = &e
1614 {
1615 let details = json!({
1616 "failed_step": failed_step,
1617 "verdict_value": verdict_value,
1618 "partial_ctx": partial_ctx,
1619 });
1620 return ApiError::bad_request(format!("run: flow eval: {message}")).with_details(details);
1621 }
1622 ApiError::bad_request(format!("run: {e}"))
1623}
1624
1625/// Uniform error response type for the handlers in this module. Converts to
1626/// a JSON `{"error": message}` body with the given status via [`IntoResponse`].
1627///
1628/// GH #76 error surface: an optional `details` field carries the structured
1629/// [`mlua_swarm::service::TaskLaunchError::FlowEval`] envelope
1630/// (`failed_step` / `verdict_value` / `partial_ctx`) when the abort
1631/// originated from a Blueprint step. When present, the JSON body becomes
1632/// `{"error": message, "details": {...}}` — a pure ADDITIVE schema change
1633/// for consumers that already ignore unknown keys. Absent (the default)
1634/// preserves the pre-#76 `{"error": message}` shape byte-for-byte for
1635/// every other error site.
1636#[derive(Debug)]
1637pub struct ApiError {
1638 status: StatusCode,
1639 message: String,
1640 details: Option<Value>,
1641}
1642
1643impl ApiError {
1644 /// Wraps an engine-side error as `500 Internal Server Error`.
1645 pub fn engine(e: impl std::fmt::Display) -> Self {
1646 Self {
1647 status: StatusCode::INTERNAL_SERVER_ERROR,
1648 message: format!("engine: {e}"),
1649 details: None,
1650 }
1651 }
1652 /// Builds a `404 Not Found` with the given message.
1653 pub fn not_found(m: String) -> Self {
1654 Self {
1655 status: StatusCode::NOT_FOUND,
1656 message: m,
1657 details: None,
1658 }
1659 }
1660 /// Builds a `400 Bad Request` with the given message.
1661 pub fn bad_request(m: String) -> Self {
1662 Self {
1663 status: StatusCode::BAD_REQUEST,
1664 message: m,
1665 details: None,
1666 }
1667 }
1668 /// Builds a `409 Conflict` with the given message (`POST
1669 /// /v1/runs/:id/resume` — the Run is not `Interrupted`, or a concurrent
1670 /// resume already won the `Interrupted -> Running` compare-and-set).
1671 pub fn conflict(m: String) -> Self {
1672 Self {
1673 status: StatusCode::CONFLICT,
1674 message: m,
1675 details: None,
1676 }
1677 }
1678 /// GH #81 Layer 1: map a `TaskApplication::resolve` failure into the
1679 /// same recovery wording the register path already emits when the
1680 /// underlying condition is that the Blueprint is archived. On the
1681 /// archived branch the caller gets a `409 CONFLICT` with the exact
1682 /// hint `blueprint {id} is archived; POST /v1/blueprints/{id}/unarchive
1683 /// first` — byte-identical to `blueprints::seed_blueprint`'s message
1684 /// so downstream tooling can grep either surface. Every other
1685 /// `TaskApplicationError` falls through to the pre-#81 `400 bad
1686 /// request` shape, preserving the launch / rekick error surface for
1687 /// callers that don't distinguish store errors from other resolve
1688 /// failures.
1689 pub fn from_task_resolve(err: &mlua_swarm::TaskApplicationError, prefix: &str) -> Self {
1690 use mlua_swarm::blueprint::store::BlueprintStoreError;
1691 use mlua_swarm::TaskApplicationError as E;
1692 if let E::Store(BlueprintStoreError::Archived(bp_id)) = err {
1693 return Self {
1694 status: StatusCode::CONFLICT,
1695 message: format!(
1696 "{prefix}: blueprint {bp_id} is archived; \
1697 POST /v1/blueprints/{bp_id}/unarchive first"
1698 ),
1699 details: None,
1700 };
1701 }
1702 Self::bad_request(format!("{prefix}: {err}"))
1703 }
1704 /// Builds a `503 Service Unavailable` with the given message (GH #33
1705 /// Guard 1 — operator readiness precheck).
1706 pub fn unavailable(m: String) -> Self {
1707 Self {
1708 status: StatusCode::SERVICE_UNAVAILABLE,
1709 message: m,
1710 details: None,
1711 }
1712 }
1713 /// Builds a `504 Gateway Timeout` with the given message (GH #33
1714 /// Guard 2 — sync launch timeout ceiling).
1715 pub fn timeout(m: String) -> Self {
1716 Self {
1717 status: StatusCode::GATEWAY_TIMEOUT,
1718 message: m,
1719 details: None,
1720 }
1721 }
1722 /// Builds a `410 Gone` with the given message (GH #37 — worker
1723 /// submit/artifact addressed at a Run that already reached a terminal
1724 /// status; the silent-`204`-then-orphan alternative is the failure
1725 /// shape this replaces).
1726 pub fn gone(m: String) -> Self {
1727 Self {
1728 status: StatusCode::GONE,
1729 message: m,
1730 details: None,
1731 }
1732 }
1733 /// Builds a `413 Payload Too Large` with the given message (GH #42 —
1734 /// `@file:` sentinel resolves to a file larger than the shared
1735 /// `DefaultBodyLimit`; same size ceiling as the inline body path).
1736 pub fn payload_too_large(m: String) -> Self {
1737 Self {
1738 status: StatusCode::PAYLOAD_TOO_LARGE,
1739 message: m,
1740 details: None,
1741 }
1742 }
1743 /// Builds a `422 Unprocessable Entity` with the given message (GH #50
1744 /// — a `worker_submit` / `worker_artifact` value violates the
1745 /// dispatching agent's declared `VerdictContract`: rejected before it
1746 /// reaches `submit_worker_result_trusted` / `stage_worker_artifact_trusted`,
1747 /// i.e. before it can land in the flow ctx).
1748 pub fn unprocessable(m: impl Into<String>) -> Self {
1749 Self {
1750 status: StatusCode::UNPROCESSABLE_ENTITY,
1751 message: m.into(),
1752 details: None,
1753 }
1754 }
1755 /// GH #76 error surface: attach a structured details payload alongside the
1756 /// string message. Consumed by [`IntoResponse`] to emit
1757 /// `{"error": message, "details": {...}}`. The current sole caller
1758 /// is `run_flow_form`'s `TaskApplicationError::Launch(FlowEval)` arm,
1759 /// which lifts `failed_step` / `verdict_value` / `partial_ctx` out of
1760 /// the structured [`mlua_swarm::service::TaskLaunchError::FlowEval`]
1761 /// variant into this field.
1762 pub fn with_details(mut self, details: Value) -> Self {
1763 self.details = Some(details);
1764 self
1765 }
1766}
1767
1768impl IntoResponse for ApiError {
1769 fn into_response(self) -> Response {
1770 let body = match self.details {
1771 Some(details) => json!({"error": self.message, "details": details}),
1772 None => json!({"error": self.message}),
1773 };
1774 (self.status, Json(body)).into_response()
1775 }
1776}
1777
1778fn default_run_ttl() -> u64 {
1779 // 1800s (= 30 min). Prevents op_token expiry across a flow.ir multi-step chain
1780 // (= 5+ SubAgent dispatches at 30–60s each). Origin: the observed fvloop smoke
1781 // where a post-gate mock-commit dispatch blew past 300s and expired — sibling of worker_token TTL.
1782 1800
1783}
1784
1785/// TTL cascade resolve helper (Blueprint metadata → server default fallback).
1786/// Second-stage fallback, called when the POST `/v1/tasks` body does not set `ttl_secs`.
1787/// (1) If BP metadata `default_run_ttl_secs` is `Some`, use it.
1788/// (2) If `None`, fall back to the server global `default_run_ttl()` (1800s).
1789///
1790/// # Full cascade (combined in `run_flow_form`)
1791///
1792/// - request body `ttl_secs=Some(v)` → v (this helper is not called)
1793/// - request body `None` + metadata `Some(v)` → v
1794/// - request body `None` + metadata `None` → `default_run_ttl()` = 1800s
1795#[cfg(test)]
1796fn resolve_ttl_from_metadata(metadata_ttl: Option<u64>) -> u64 {
1797 metadata_ttl.unwrap_or_else(default_run_ttl)
1798}
1799
1800#[cfg(test)]
1801mod tests {
1802 use super::*;
1803
1804 /// TTL cascade case 1: when the request body sets it, that value is used as-is
1805 /// (upper branch that does not go through the helper; semantic verify of the
1806 /// `Some(v) => v` direct-return path in `run_flow_form`).
1807 #[test]
1808 fn ttl_cascade_request_body_wins_over_metadata() {
1809 let req_ttl: Option<u64> = Some(100);
1810 let metadata_ttl: Option<u64> = Some(3600);
1811 let effective = match req_ttl {
1812 Some(v) => v,
1813 None => resolve_ttl_from_metadata(metadata_ttl),
1814 };
1815 assert_eq!(
1816 effective, 100,
1817 "request body ttl_secs=100 must win over metadata=3600 (cascade priority (1) > (2))"
1818 );
1819 }
1820
1821 /// TTL cascade case 2: request body omitted + BP metadata `Some(N)` → `N` is effective.
1822 #[test]
1823 fn ttl_cascade_metadata_used_when_body_missing() {
1824 let req_ttl: Option<u64> = None;
1825 let metadata_ttl: Option<u64> = Some(3600);
1826 let effective = match req_ttl {
1827 Some(v) => v,
1828 None => resolve_ttl_from_metadata(metadata_ttl),
1829 };
1830 assert_eq!(
1831 effective, 3600,
1832 "body None + metadata=3600 must resolve to 3600 (cascade (2))"
1833 );
1834 }
1835
1836 /// TTL cascade case 3: request body omitted + BP metadata `None` → server default (1800s).
1837 #[test]
1838 fn ttl_cascade_server_default_when_both_missing() {
1839 let req_ttl: Option<u64> = None;
1840 let metadata_ttl: Option<u64> = None;
1841 let effective = match req_ttl {
1842 Some(v) => v,
1843 None => resolve_ttl_from_metadata(metadata_ttl),
1844 };
1845 assert_eq!(
1846 effective,
1847 default_run_ttl(),
1848 "body None + metadata None must fall back to default_run_ttl() = 1800s"
1849 );
1850 assert_eq!(effective, 1800, "default_run_ttl() literal = 1800s");
1851 }
1852
1853 /// Helper unit: metadata `None` → 1800 (server default expansion).
1854 #[test]
1855 fn resolve_ttl_from_metadata_none_returns_server_default() {
1856 assert_eq!(resolve_ttl_from_metadata(None), 1800);
1857 }
1858
1859 /// Helper unit: metadata `Some(N)` → `N` (server default ignored).
1860 #[test]
1861 fn resolve_ttl_from_metadata_some_returns_value() {
1862 assert_eq!(resolve_ttl_from_metadata(Some(7200)), 7200);
1863 assert_eq!(resolve_ttl_from_metadata(Some(60)), 60);
1864 }
1865
1866 // ──────────────────────────────────────────────────────────────────
1867 // `TaskLaunchRequest.check_policy` wire field (T5)
1868 // ──────────────────────────────────────────────────────────────────
1869
1870 /// T5: a `POST /v1/tasks` body carrying a top-level `check_policy`
1871 /// deserializes into `TaskLaunchRequest.check_policy` using the
1872 /// snake_case wire form.
1873 #[test]
1874 fn task_launch_request_parses_check_policy_wire_field() {
1875 let body = json!({
1876 "blueprint": { "kind": "id", "id": "some-bp" },
1877 "init_ctx": {},
1878 "check_policy": "silent",
1879 });
1880 let req: TaskLaunchRequest =
1881 serde_json::from_value(body).expect("request must deserialize");
1882 assert_eq!(req.check_policy, Some(CheckPolicy::Silent));
1883 }
1884
1885 /// A body that omits `check_policy` leaves the field `None` (existing
1886 /// clients are unaffected — `#[serde(default)]`).
1887 #[test]
1888 fn task_launch_request_check_policy_defaults_to_none_when_omitted() {
1889 let body = json!({
1890 "blueprint": { "kind": "id", "id": "some-bp" },
1891 "init_ctx": {},
1892 });
1893 let req: TaskLaunchRequest =
1894 serde_json::from_value(body).expect("request must deserialize");
1895 assert_eq!(req.check_policy, None);
1896 }
1897
1898 // ──────────────────────────────────────────────────────────────────
1899 // issue #19 ST2: `build_task_input_spec_from_request` direct resolver
1900 // ──────────────────────────────────────────────────────────────────
1901
1902 fn task_req(
1903 init_ctx: Value,
1904 project_root: Option<&str>,
1905 work_dir: Option<&str>,
1906 task_metadata: Option<Value>,
1907 ) -> TaskLaunchRequest {
1908 TaskLaunchRequest {
1909 blueprint: BlueprintRef::Id {
1910 id: mlua_swarm::blueprint::store::BlueprintId::new("ut"),
1911 version: Default::default(),
1912 },
1913 init_ctx,
1914 project_root: project_root.map(String::from),
1915 work_dir: work_dir.map(String::from),
1916 task_metadata,
1917 ttl_secs: None,
1918 operator: None,
1919 operator_sid: None,
1920 timeout_secs: None,
1921 goal: None,
1922 detach: false,
1923 check_policy: None,
1924 }
1925 }
1926
1927 /// (a) Sibling fields only — no legacy keys in `init_ctx` — are
1928 /// returned in the `TaskInputSpec` unchanged. `init_ctx` itself is
1929 /// untouched by this resolver (checked separately at the call site).
1930 #[test]
1931 fn build_task_input_spec_from_request_returns_sibling_fields_when_present() {
1932 let req = task_req(
1933 json!({"free": "form"}),
1934 Some("/repo/sibling"),
1935 Some("/repo/sibling/work"),
1936 Some(json!({"issue": 19})),
1937 );
1938 let spec = build_task_input_spec_from_request(&req).expect("spec must be Some");
1939 assert_eq!(spec.project_root.as_deref(), Some("/repo/sibling"));
1940 assert_eq!(spec.work_dir.as_deref(), Some("/repo/sibling/work"));
1941 assert_eq!(spec.task_metadata, Some(json!({"issue": 19})));
1942 }
1943
1944 /// (b) No sibling fields — the pre-#19 shape (same 3 keys nested
1945 /// inside `init_ctx`) is used as the fallback source.
1946 #[test]
1947 fn build_task_input_spec_from_request_falls_back_to_legacy_init_ctx_shape() {
1948 let req = task_req(
1949 json!({
1950 "project_root": "/repo/legacy",
1951 "work_dir": "/repo/legacy/work",
1952 "task_metadata": {"issue": 17},
1953 }),
1954 None,
1955 None,
1956 None,
1957 );
1958 let spec = build_task_input_spec_from_request(&req).expect("spec must be Some");
1959 assert_eq!(spec.project_root.as_deref(), Some("/repo/legacy"));
1960 assert_eq!(spec.work_dir.as_deref(), Some("/repo/legacy/work"));
1961 assert_eq!(spec.task_metadata, Some(json!({"issue": 17})));
1962 }
1963
1964 /// (c) Both present — the sibling field must win over the legacy
1965 /// `init_ctx`-nested value.
1966 #[test]
1967 fn build_task_input_spec_from_request_sibling_wins_over_legacy_shape() {
1968 let req = task_req(
1969 json!({
1970 "project_root": "/repo/legacy",
1971 "work_dir": "/repo/legacy/work",
1972 "task_metadata": {"issue": 17},
1973 }),
1974 Some("/repo/sibling"),
1975 Some("/repo/sibling/work"),
1976 Some(json!({"issue": 19})),
1977 );
1978 let spec = build_task_input_spec_from_request(&req).expect("spec must be Some");
1979 assert_eq!(
1980 spec.project_root.as_deref(),
1981 Some("/repo/sibling"),
1982 "sibling field must win over the legacy init_ctx-nested value"
1983 );
1984 assert_eq!(spec.work_dir.as_deref(), Some("/repo/sibling/work"));
1985 assert_eq!(spec.task_metadata, Some(json!({"issue": 19})));
1986 }
1987
1988 /// (d) All three fields absent from both sibling and legacy shapes —
1989 /// resolver returns `None`, and no middleware is layered downstream.
1990 #[test]
1991 fn build_task_input_spec_from_request_returns_none_when_no_fields_present() {
1992 let req = task_req(json!({"unrelated": "value"}), None, None, None);
1993 assert!(build_task_input_spec_from_request(&req).is_none());
1994 }
1995
1996 /// Minimal `AppState` for the `status_get` handler-fn-direct-call test
1997 /// below — same construction shape as `tasks.rs::test_state()`
1998 /// (mirrors what `build_router_full` does internally, skipping the
1999 /// `Router` wrapper).
2000 fn status_test_state() -> AppState {
2001 let engine = Engine::new(mlua_swarm::EngineCfg::default());
2002 let compiler = mlua_swarm::Compiler::new(default_registry());
2003 let launch = Arc::new(mlua_swarm::TaskLaunchService::new(engine.clone(), compiler));
2004 AppState {
2005 engine,
2006 sessions: Arc::new(Mutex::new(SessionStore::default())),
2007 task_app: Arc::new(mlua_swarm::TaskApplication::new_inline_only(launch)),
2008 ws_operator_factory: None,
2009 data_store: Arc::new(mlua_swarm::store::output::InMemoryOutputStore::new()),
2010 operator_sessions: Arc::new(Mutex::new(HashMap::new())),
2011 roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
2012 task_store: Arc::new(mlua_swarm::store::task::InMemoryTaskStore::new()),
2013 run_store: Arc::new(mlua_swarm::store::run::InMemoryRunStore::new()),
2014 replay_store: Arc::new(mlua_swarm::store::replay::InMemoryReplayStore::new()),
2015 run_trace_store: Arc::new(mlua_swarm::store::trace::InMemoryRunTraceStore::new()),
2016 base_url: None,
2017 sync_timeout_secs: 300,
2018 }
2019 }
2020
2021 /// issue #35 ST4 Acceptance Criteria: `GET /v1/status` reports the
2022 /// count of `Running` `Run`s (`RunStore::list_running`) and attached
2023 /// Operator ids (`engine.list_operator_ids()`), called directly as a
2024 /// handler fn (no `Router` wrapper — this crate's established
2025 /// unit-test convention).
2026 #[tokio::test]
2027 async fn status_get_reports_running_runs_and_operators() {
2028 let state = status_test_state();
2029
2030 let now = std::time::SystemTime::now()
2031 .duration_since(std::time::UNIX_EPOCH)
2032 .map(|d| d.as_secs())
2033 .unwrap_or(0);
2034 state
2035 .run_store
2036 .create(RunRecord {
2037 id: RunId::new(),
2038 task_id: TaskId::new(),
2039 status: RunStatus::Running,
2040 step_entries: Vec::new(),
2041 degradations: Vec::new(),
2042 operator_sid: None,
2043 result_ref: None,
2044 input_json: None,
2045 created_at: now,
2046 updated_at: now,
2047 })
2048 .await
2049 .expect("seed running RunRecord");
2050
2051 // Throwaway `Operator` impl — only registration/list-count matters
2052 // for this test, `execute` is never dispatched (same idiom as
2053 // `tasks.rs::StallingOperator`).
2054 struct NoopOperator;
2055 #[async_trait::async_trait]
2056 impl mlua_swarm::Operator for NoopOperator {
2057 async fn execute(
2058 &self,
2059 _ctx: &mlua_swarm::Ctx,
2060 _system: Option<String>,
2061 _prompt: Value,
2062 _worker: Option<mlua_swarm::WorkerBinding>,
2063 _worker_token: mlua_swarm::CapToken,
2064 ) -> Result<mlua_swarm::WorkerResult, mlua_swarm::WorkerError> {
2065 unimplemented!("not exercised by this test — only registration/list matters")
2066 }
2067 }
2068 state
2069 .engine
2070 .register_operator("test-op", Arc::new(NoopOperator))
2071 .await;
2072
2073 let Json(resp) = status_get(State(state)).await;
2074 assert_eq!(resp.running_runs, 1);
2075 assert_eq!(resp.attached_operators, 1);
2076 }
2077
2078 // ──────────────────────────────────────────────────────────────────
2079 // GH #76 error surface: ApiError.details + flow_eval_error_to_api_error mapper
2080 // ──────────────────────────────────────────────────────────────────
2081
2082 /// `ApiError::with_details` populates the optional `details` field, and
2083 /// [`IntoResponse`] renders it into the JSON body as a sibling of
2084 /// `error`. Pre-#76 shape (no details) stays byte-for-byte
2085 /// `{"error": message}`.
2086 #[tokio::test]
2087 async fn api_error_details_render_into_response_body() {
2088 use axum::body::to_bytes;
2089 use axum::response::IntoResponse;
2090
2091 // Baseline: no details → pre-#76 shape.
2092 let bare = ApiError::bad_request("something".to_string()).into_response();
2093 let (parts, body) = bare.into_parts();
2094 assert_eq!(parts.status, StatusCode::BAD_REQUEST);
2095 let bytes = to_bytes(body, 1024).await.expect("bare body");
2096 let parsed: serde_json::Value = serde_json::from_slice(&bytes).expect("parse bare");
2097 assert_eq!(parsed, json!({"error": "something"}));
2098
2099 // With details → additive `details` key.
2100 let with_details = ApiError::bad_request("run: flow eval: blocked".to_string())
2101 .with_details(json!({
2102 "failed_step": "gate",
2103 "verdict_value": {"verdict": "BLOCKED"},
2104 "partial_ctx": {"steps": {}},
2105 }));
2106 let resp = with_details.into_response();
2107 let (parts, body) = resp.into_parts();
2108 assert_eq!(parts.status, StatusCode::BAD_REQUEST);
2109 let bytes = to_bytes(body, 4096).await.expect("details body");
2110 let parsed: serde_json::Value = serde_json::from_slice(&bytes).expect("parse details");
2111 assert_eq!(parsed["error"], "run: flow eval: blocked");
2112 assert_eq!(parsed["details"]["failed_step"], "gate");
2113 assert_eq!(parsed["details"]["verdict_value"]["verdict"], "BLOCKED");
2114 assert!(parsed["details"]["partial_ctx"].is_object());
2115 }
2116
2117 /// The mapper lifts the structured `TaskLaunchError::FlowEval` fields
2118 /// into `ApiError.details` while preserving the pre-#76 message prefix
2119 /// (`"run: flow eval: <msg>"`). Regression: every other
2120 /// `TaskApplicationError` variant collapses to the pre-#76 shape (no
2121 /// `details`).
2122 #[test]
2123 fn flow_eval_error_to_api_error_lifts_structural_fields_into_details() {
2124 let err = TaskApplicationError::Launch(TaskLaunchError::FlowEval {
2125 message: "blocked: {\"verdict\":\"BLOCKED\"}".to_string(),
2126 failed_step: Some("gate".to_string()),
2127 verdict_value: Some(json!({"verdict": "BLOCKED"})),
2128 partial_ctx: Some(json!({"steps": {}})),
2129 });
2130 let api_err = flow_eval_error_to_api_error(err);
2131 assert_eq!(api_err.status, StatusCode::BAD_REQUEST);
2132 assert!(
2133 api_err.message.starts_with("run: flow eval: "),
2134 "message must preserve pre-#76 `run: flow eval: <msg>` prefix, got: {}",
2135 api_err.message
2136 );
2137 let details = api_err.details.expect("details must be Some for FlowEval");
2138 assert_eq!(details["failed_step"], "gate");
2139 assert_eq!(details["verdict_value"]["verdict"], "BLOCKED");
2140 assert!(details["partial_ctx"].is_object());
2141 }
2142
2143 /// Regression: a non-`FlowEval` error must still map to a `bad_request`
2144 /// without a `details` field — the pre-#76 shape for e.g.
2145 /// `TaskApplicationError::NoStore`.
2146 #[test]
2147 fn flow_eval_error_to_api_error_non_flow_eval_falls_back_to_message_only() {
2148 let api_err = flow_eval_error_to_api_error(TaskApplicationError::NoStore);
2149 assert_eq!(api_err.status, StatusCode::BAD_REQUEST);
2150 assert!(api_err.message.starts_with("run: "));
2151 assert!(
2152 api_err.details.is_none(),
2153 "non-FlowEval errors must not carry a details field (pre-#76 shape)"
2154 );
2155 }
2156
2157 /// A `FlowEval` with every optional field `None` (upstream flow-ir
2158 /// error path — no dispatcher breadcrumb, no run_ctx snapshot) still
2159 /// lifts into `details` — the shape is `null` per key, which serialize
2160 /// as JSON `null`. Consumers must treat missing / `null` as "not
2161 /// available", both are legitimate.
2162 #[test]
2163 fn flow_eval_error_to_api_error_with_all_none_still_populates_details_with_nulls() {
2164 let err = TaskApplicationError::Launch(TaskLaunchError::FlowEval {
2165 message: "unresolved extern".to_string(),
2166 failed_step: None,
2167 verdict_value: None,
2168 partial_ctx: None,
2169 });
2170 let api_err = flow_eval_error_to_api_error(err);
2171 let details = api_err
2172 .details
2173 .expect("details Some even when fields are None");
2174 assert_eq!(details["failed_step"], Value::Null);
2175 assert_eq!(details["verdict_value"], Value::Null);
2176 assert_eq!(details["partial_ctx"], Value::Null);
2177 }
2178
2179 // ─── GH #81 Layer 1: archived-BP guidance on run paths ──────────
2180
2181 #[test]
2182 fn from_task_resolve_archived_maps_to_409_with_unarchive_hint() {
2183 use mlua_swarm::blueprint::store::{BlueprintId, BlueprintStoreError};
2184 use mlua_swarm::TaskApplicationError;
2185 let bp_id = BlueprintId::new("greeter".to_string());
2186 let err = TaskApplicationError::Store(BlueprintStoreError::Archived(bp_id));
2187 let api = ApiError::from_task_resolve(&err, "bp resolve");
2188 assert_eq!(api.status, StatusCode::CONFLICT);
2189 // The wording is byte-identical to the register path
2190 // (`blueprints::seed_blueprint`) so downstream tooling can grep
2191 // either surface.
2192 assert_eq!(
2193 api.message,
2194 "bp resolve: blueprint greeter is archived; \
2195 POST /v1/blueprints/greeter/unarchive first"
2196 );
2197 }
2198
2199 #[test]
2200 fn from_task_resolve_archived_honours_the_caller_supplied_prefix() {
2201 use mlua_swarm::blueprint::store::{BlueprintId, BlueprintStoreError};
2202 use mlua_swarm::TaskApplicationError;
2203 // The rekick site prepends `task {task_id}: ` to distinguish
2204 // rekick failures from launch failures in logs.
2205 let bp_id = BlueprintId::new("scout".to_string());
2206 let err = TaskApplicationError::Store(BlueprintStoreError::Archived(bp_id));
2207 let api = ApiError::from_task_resolve(&err, "task T-abc: bp resolve");
2208 assert_eq!(api.status, StatusCode::CONFLICT);
2209 assert!(api.message.starts_with("task T-abc: bp resolve:"));
2210 assert!(api
2211 .message
2212 .contains("blueprint scout is archived; POST /v1/blueprints/scout/unarchive first"));
2213 }
2214
2215 #[test]
2216 fn from_task_resolve_non_archived_store_error_stays_400() {
2217 // A store IdNotFound → resolve() surfaces
2218 // `TaskApplicationError::Store(BlueprintStoreError::IdNotFound(...))`
2219 // which must remain the pre-#81 400 shape.
2220 use mlua_swarm::blueprint::store::{BlueprintId, BlueprintStoreError};
2221 use mlua_swarm::TaskApplicationError;
2222 let bp_id = BlueprintId::new("no-such".to_string());
2223 let err = TaskApplicationError::Store(BlueprintStoreError::IdNotFound(bp_id));
2224 let api = ApiError::from_task_resolve(&err, "bp resolve");
2225 assert_eq!(api.status, StatusCode::BAD_REQUEST);
2226 assert!(api.message.starts_with("bp resolve:"));
2227 }
2228
2229 #[test]
2230 fn from_task_resolve_no_store_error_stays_400() {
2231 // A non-Store variant (NoStore fires when `BlueprintRef::Id` is
2232 // used against an inline-only TaskApplication) also stays 400.
2233 use mlua_swarm::TaskApplicationError;
2234 let err = TaskApplicationError::NoStore;
2235 let api = ApiError::from_task_resolve(&err, "bp resolve");
2236 assert_eq!(api.status, StatusCode::BAD_REQUEST);
2237 }
2238}