Skip to main content

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//! - `POST /v1/runs/:id/resume` — resume an `Interrupted` Run under the same `run_id`.
30//! - `POST /v1/runs/:id/rerun-from` — GH #71 Layer A. Rerun a terminal Run from
31//!   a caller-specified step under the same `run_id` (physically truncates the
32//!   replay log at the cut point). See `tasks::run_rerun_from`.
33//! - `POST /v1/operators` / `GET /v1/operators/:sid` / `DELETE /v1/operators/:sid` /
34//!   `GET /v1/operators/:sid/ws` (WS upgrade) — REST-like Operator login flow,
35//!   Bearer-mandatory; the sole WS Operator session route. See `operator_ws::login`
36//!   module doc.
37//!
38//! The Enhance issue axis (`/issues`) lives in the `issues` module; callers merge
39//! `build_issues_router` to integrate it into the same server.
40//!
41//! # The 3 faces of the Operator role (= registered directly on the engine SoT)
42//!
43//! The engine stateless-executor refactor removed the three
44//! `AppState` registries (former `HookRegistry` / `BridgeRegistry` / `OperatorRegistry`);
45//! all registration now goes directly to the engine SoT via
46//! `engine.register_spawn_hook` / `register_senior_bridge` / `register_operator`.
47//! `WSOperatorSession` (in the `operator_ws` module) registers all three traits
48//! simultaneously under a single sid — one WS connection covers all 3 faces of
49//! the Operator role, the canonical pattern.
50//!
51//! # `build_*` family
52//!
53//! - [`build_router`] — minimal entry (= `default_registry()`)
54//! - [`build_router_with`] — caller provides a `SpawnerRegistry` and optional `BlueprintStore`
55//!
56//! The engine should be started with [`default_layer_registry`] (= `Engine::new_with_layers`);
57//! otherwise `Blueprint.spawner_hints` is ignored.
58
59#![warn(missing_docs)]
60
61/// HTTP surface for inspecting/registering Blueprint state (`/v1/blueprints/*`).
62pub mod blueprints;
63/// Server config file support (`~/.mse/config.toml`, CLI > file > default merge).
64pub mod config;
65/// `/v1/data/*` endpoints (v9 Big Response handling, Store-owner direct path).
66pub mod data;
67/// `GET /v1/doctor` — read-only startup config / Store snapshot.
68pub mod doctor;
69/// HTTP surface for the `/v1/enhance/log` axis.
70pub mod enhance_log;
71/// `EnhanceSetting` HTTP CRUD (`/v1/enhance-settings*`).
72pub mod enhance_settings;
73/// HTTP surface for the Enhance issue axis (`/v1/issues*`).
74pub mod issues;
75/// WebSocket Operator Callback IF (`/v1/operators*`).
76pub mod operator_ws;
77/// `GET /v1/tasks/:id/runs/:run/steps*` (the metadata + content debug
78/// plane over a Run's step OUTPUT — `McpQueryAdapter`, a server-side
79/// `mlua_swarm::core::projection::ProjectionAdapter` impl reading through
80/// the Data-plane `OutputStore` with a persisted `RunRecord.result_ref`
81/// fallback). See the module doc for how this relates to
82/// `operator_ws::session`'s in-flight `FileProjectionAdapter` hook and
83/// `worker`'s Worker-axis `context.steps` pointer assembly.
84pub mod projection;
85/// HTTP surface for the Task/Run persistence axis (issue #13 ID hierarchy;
86/// `GET /v1/tasks`, `GET /v1/tasks/:id`, `POST /v1/tasks/:id/runs`,
87/// `GET /v1/runs/:id`). `POST /v1/tasks` itself stays in this module (it is
88/// the entry point `tasks_start` shares with the flow-eval path) — see the
89/// `tasks` module doc for the split rationale.
90pub mod tasks;
91/// `/v1/worker/*` endpoints (SubAgent self-fetch path).
92pub mod worker;
93pub use blueprints::{build_blueprints_router, build_blueprints_router_with_refs};
94pub use enhance_log::build_enhance_log_router;
95pub use enhance_settings::build_enhance_settings_router;
96pub use issues::{build_issues_router, GetIssueResponse, PostIssueRequest, PostIssueResponse};
97pub use operator_ws::{
98    operators_create, operators_delete, operators_info, operators_ws_connect, ClientMsg,
99    OperatorSessionEntry, ServerMsg, WSOperatorSession,
100};
101pub use projection::{McpQueryAdapter, ProjectionSource, StepList, StepPathQuery, StepSummary};
102pub use tasks::{RunKickRequest, RunKickResponse, RunResumeResponse, TaskDetailResponse};
103pub use worker::{
104    worker_artifact, worker_prompt, worker_result, ArtifactQuery, PromptQuery, WorkerResultReq,
105};
106
107use axum::{
108    extract::{DefaultBodyLimit, State},
109    http::{header::AUTHORIZATION, HeaderMap, StatusCode},
110    response::{IntoResponse, Response},
111    routing::{get, post},
112    Json, Router,
113};
114use mlua_swarm::application::{BlueprintRef, TaskApplication};
115use mlua_swarm::blueprint::store::BlueprintStore;
116use mlua_swarm::core::config::CheckPolicy;
117use mlua_swarm::service::TaskLaunchService;
118use mlua_swarm::store::replay::{InMemoryReplayStore, ReplayStore};
119use mlua_swarm::store::run::{RunContext, RunRecord, RunStatus, RunStore};
120use mlua_swarm::store::task::{TaskRecord, TaskRecordStatus, TaskStore};
121use mlua_swarm::{
122    CapToken, Compiler, Engine, LayerRegistry, LuaInProcessSpawnerFactory, MainAIMiddleware,
123    OperatorDelegateMiddleware, OperatorSpawnerFactory, Role, RunId, RustFnInProcessSpawnerFactory,
124    SeniorEscalationMiddleware, SessionId, SpawnerRegistry, SubprocessProcessSpawnerFactory,
125    TaskId,
126};
127use serde::{Deserialize, Serialize};
128use serde_json::{json, Value};
129use std::collections::HashMap;
130use std::sync::Arc;
131use std::time::Duration;
132use tokio::sync::Mutex;
133
134/// In-memory session map backing `/v1/sessions` attach/detach.
135///
136/// The `sid` handed to the client on this REST path is the token nonce
137/// itself (a bearer secret), so the server never uses it as a map key —
138/// entries are keyed by its fingerprint
139/// (`mlua_swarm::types::token_fingerprint`; issue #14).
140#[derive(Default)]
141pub struct SessionStore {
142    /// Live session tokens keyed by the sid's fingerprint.
143    pub map: HashMap<String, CapToken>,
144}
145
146/// Shared axum handler state for the whole router. Cloned per-request (all
147/// fields are `Arc`/cheap-clone), constructed once in [`build_router_with_ws_factory`].
148#[derive(Clone)]
149pub struct AppState {
150    /// The engine SoT (attach/detach, dispatch, registries).
151    pub engine: Engine,
152    /// Live `/v1/sessions` attach records (Operator/Worker/etc session tokens).
153    pub sessions: Arc<Mutex<SessionStore>>,
154    /// Application used at the task entry to resolve `BlueprintRef`. Without a Store, runs in Inline-only mode.
155    pub task_app: Arc<TaskApplication>,
156    /// When `Some`, on WS connect a new `WSOperatorSession` is automatically registered
157    /// with this factory under the sid name (= a `kind=operator` + `operator_ref=<sid>` AgentDef
158    /// binds to the `WSOperatorSession` backend).
159    /// When `None`, no auto-registration happens; the session is only registered on
160    /// `engine.OperatorRegistry` (= only the `OperatorDelegateMiddleware` path is effective;
161    /// the `OperatorSpawnerFactory` path is dead).
162    pub ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
163    /// Owner of the Store on the Data path (Big Response handling). Added in v9.
164    /// Independent layer — the Engine core and the Domain path (`/v1/worker/result`)
165    /// are not involved.
166    /// Default = `InMemoryOutputStore` (constructed inside `build_router_with_ws_factory`);
167    /// callers can swap in an sqlite/fs backend later (future carry).
168    pub data_store: Arc<dyn mlua_swarm::store::output::OutputStore>,
169    /// Login-flow session store (`POST /v1/operators` mint records). `sid` →
170    /// `OperatorSessionEntry`. This is the sole session store for the WS
171    /// Operator role. See `operator_ws::login` module doc.
172    pub operator_sessions:
173        Arc<Mutex<HashMap<SessionId, Arc<crate::operator_ws::login::OperatorSessionEntry>>>>,
174    /// S1 login-flow roles-exclusivity map. Role name → owning `sid`. Checked
175    /// (and updated) atomically under a single lock in
176    /// `operator_ws::login::operators_create` — a role already present here
177    /// causes `POST /v1/operators` to return `409 CONFLICT`. Entries are
178    /// released on `DELETE /v1/operators/:sid`.
179    pub roles_to_sid: Arc<Mutex<HashMap<String, SessionId>>>,
180    /// Persistence for `Task` records (issue #13 ID-hierarchy work-item
181    /// identity; see `mlua_swarm::store::task` module doc). Default =
182    /// `InMemoryTaskStore` (constructed inside `build_router_full`); callers
183    /// can swap in a `SqliteTaskStore` via the `task_store` argument.
184    pub task_store: Arc<dyn TaskStore>,
185    /// Persistence for `Run` records (one kick of a Task; see
186    /// `mlua_swarm::store::run` module doc). Default = `InMemoryRunStore`;
187    /// callers can swap in a `SqliteRunStore` via the `run_store` argument.
188    pub run_store: Arc<dyn RunStore>,
189    /// Per-run replay log — the Ctx-snapshot + step-output store the engine
190    /// appends to after every completed step (see `mlua_swarm::store::replay`
191    /// module doc). Threaded into `RunContext` at every dispatch site so a
192    /// later restart-equivalent recovery can reconstruct the run. Default =
193    /// `InMemoryReplayStore` (process-volatile); callers can swap in a
194    /// `SqliteReplayStore` via the `replay_store` argument.
195    pub replay_store: Arc<dyn ReplayStore>,
196    /// Public HTTP base URL the server is reachable at (e.g.
197    /// `"http://127.0.0.1:7777"`), sourced from the binary at boot time.
198    /// When `Some`, `WSOperatorSession` renders it literally into the
199    /// Spawn `directive`'s `base_url` line so the receiving operator can
200    /// paste the frame into a SubAgent prompt without a `mse_doctor`
201    /// detour (issue #8). `None` preserves the historical fallback
202    /// (a placeholder that points at `mse_doctor`).
203    pub base_url: Option<Arc<str>>,
204    /// Server-wide fallback ceiling (seconds) for the `POST /v1/tasks`
205    /// synchronous launch await (GH #33 Guard 2; see `run_flow_form`'s doc
206    /// comment). Sourced from `config::ResolvedConfig::sync_timeout_secs`.
207    /// A per-request `TaskLaunchRequest.timeout_secs` override, when
208    /// present, takes priority over this value.
209    pub sync_timeout_secs: u64,
210}
211
212/// Minimal entry point: builds a router with [`default_registry`] and no
213/// `BlueprintStore` (Inline-only mode) or `ws_operator_factory`.
214pub fn build_router(engine: Engine) -> Router {
215    build_router_with(engine, default_registry(), None)
216}
217
218/// Default `LayerRegistry` for the server. Hint keys:
219/// - `"main_ai"` → `MainAIMiddleware` (= fires SpawnHook before/after)
220/// - `"senior_escalation"` → `SeniorEscalationMiddleware` (= on `ok=false`, escalates via `SeniorBridge.ask`)
221/// - `"operator_delegate"` → `OperatorDelegateMiddleware` (= when an operator backend is registered, delegates the entire spawn)
222///
223/// Including any of these keys in `Blueprint.spawner_hints.layers` causes them to
224/// be wrapped into a `SpawnerStack` at `service::linker::link` time (= per-launch;
225/// the old `engine.bind` global-state path is retired).
226/// Callers (the engine builder side) receive it via
227/// `Engine::new_with_layers(cfg, mse_server::default_layer_registry())`.
228pub fn default_layer_registry() -> LayerRegistry {
229    LayerRegistry::new()
230        .with_hint("main_ai", |_engine| Arc::new(MainAIMiddleware::new()))
231        .with_hint("senior_escalation", |_engine| {
232            Arc::new(SeniorEscalationMiddleware::new())
233        })
234        .with_hint("operator_delegate", |_engine| {
235            Arc::new(OperatorDelegateMiddleware::new())
236        })
237}
238
239/// Build form where the caller supplies a registry and an optional `BlueprintStore`.
240/// The Operator callback path (= external HTTP / WS callers acting as an Operator)
241/// must be pre-registered via `engine.register_*` (= the engine is the SoT).
242/// See the `operator_ws` module doc and `OperatorInfo` (engine-side `ctx.rs`) for details.
243pub fn build_router_with(
244    engine: Engine,
245    registry: SpawnerRegistry,
246    store: Option<Arc<dyn BlueprintStore>>,
247) -> Router {
248    build_router_with_ws_factory(engine, registry, store, None)
249}
250
251/// 4-argument variant of `build_router_with`. Passing `ws_operator_factory = Some(arc)`
252/// causes each WS connect to auto-register a new `WSOperatorSession` under its sid
253/// name with the factory (= a `kind=operator` AgentDef with `operator_ref: <sid>`
254/// can then bind to the WS client backend). Callers are expected to also install
255/// the same `Arc` into the `SpawnerRegistry` via
256/// `reg.register::<OperatorSpawnerFactory>(arc.clone())`.
257pub fn build_router_with_ws_factory(
258    engine: Engine,
259    registry: SpawnerRegistry,
260    store: Option<Arc<dyn BlueprintStore>>,
261    ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
262) -> Router {
263    build_router_with_ws_factory_and_output(engine, registry, store, ws_operator_factory, None)
264}
265
266/// 5-argument variant of [`build_router_with_ws_factory`]. Passing
267/// `output_store = Some(arc)` swaps the default `InMemoryOutputStore` for a
268/// caller-supplied backend (a `SqliteOutputStore`, for instance). `None`
269/// preserves the historical behaviour (fresh in-memory store per call).
270pub fn build_router_with_ws_factory_and_output(
271    engine: Engine,
272    registry: SpawnerRegistry,
273    store: Option<Arc<dyn BlueprintStore>>,
274    ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
275    output_store: Option<Arc<dyn mlua_swarm::store::output::OutputStore>>,
276) -> Router {
277    build_router_full(
278        engine,
279        registry,
280        store,
281        ws_operator_factory,
282        output_store,
283        None,
284        None,
285        None,
286        None,
287        crate::config::default_sync_timeout_secs(),
288    )
289}
290
291/// 8-argument variant of [`build_router_with_ws_factory_and_output`].
292/// Passing `base_url = Some(...)` (e.g. `"http://127.0.0.1:7777"`) makes
293/// `WSOperatorSession` render the actual server bind into the Spawn
294/// directive's `base_url` line, so the receiving operator can copy the
295/// frame straight into a SubAgent prompt (issue #8). `None` preserves
296/// the historical fallback (`<check with mse_doctor>` placeholder).
297/// `task_store` / `run_store` swap the default `InMemoryTaskStore` /
298/// `InMemoryRunStore` (issue #13 ID-hierarchy persistence) for a
299/// caller-supplied backend (`SqliteTaskStore` / `SqliteRunStore`, for
300/// instance); `None` preserves the process-volatile default.
301/// `sync_timeout_secs` is the server-wide fallback ceiling for the `POST
302/// /v1/tasks` synchronous launch await (GH #33 Guard 2) — see
303/// `AppState::sync_timeout_secs` / `run_flow_form`'s doc comment.
304// This is the terminal builder in the `build_router*` delegation chain
305// (each variant adds one more caller-overridable store/factory); the
306// argument count grows with the number of pluggable backends, not with
307// unrelated responsibilities, so a plain allow is preferable to bundling
308// them into a config struct only this one function would consume.
309#[allow(clippy::too_many_arguments)]
310pub fn build_router_full(
311    engine: Engine,
312    registry: SpawnerRegistry,
313    store: Option<Arc<dyn BlueprintStore>>,
314    ws_operator_factory: Option<Arc<OperatorSpawnerFactory>>,
315    output_store: Option<Arc<dyn mlua_swarm::store::output::OutputStore>>,
316    base_url: Option<Arc<str>>,
317    task_store: Option<Arc<dyn TaskStore>>,
318    run_store: Option<Arc<dyn RunStore>>,
319    replay_store: Option<Arc<dyn ReplayStore>>,
320    sync_timeout_secs: u64,
321) -> Router {
322    let compiler = Compiler::new(registry);
323    let launch = Arc::new(TaskLaunchService::new(engine.clone(), compiler));
324    let task_app = Arc::new(match store {
325        Some(s) => TaskApplication::new(launch, s),
326        None => TaskApplication::new_inline_only(launch),
327    });
328    let data_store: Arc<dyn mlua_swarm::store::output::OutputStore> = match output_store {
329        Some(s) => s,
330        None => Arc::new(mlua_swarm::store::output::InMemoryOutputStore::new()),
331    };
332    // subtask-4 / ST2 rework: wire the SAME `data_store` instance into the
333    // engine's submit-time projection sink (`Engine::submit_output` /
334    // `submit_worker_result_trusted`), so an ordinary worker
335    // `/v1/worker/submit` — not just the explicit `POST /v1/data/emit` —
336    // lands in this store too. `projection::McpQueryAdapter` (`GET
337    // /v1/tasks/:id/runs/:run/steps*`) reads through this same `Arc`,
338    // which is what makes an in-flight run's already-submitted step
339    // OUTPUT queryable.
340    engine.set_output_store(data_store.clone());
341    let task_store: Arc<dyn TaskStore> = match task_store {
342        Some(s) => s,
343        None => Arc::new(mlua_swarm::store::task::InMemoryTaskStore::new()),
344    };
345    let run_store: Arc<dyn RunStore> = match run_store {
346        Some(s) => s,
347        None => Arc::new(mlua_swarm::store::run::InMemoryRunStore::new()),
348    };
349    let replay_store: Arc<dyn ReplayStore> = match replay_store {
350        Some(s) => s,
351        None => Arc::new(InMemoryReplayStore::new()),
352    };
353    let state = AppState {
354        engine,
355        sessions: Arc::new(Mutex::new(SessionStore::default())),
356        task_app,
357        ws_operator_factory,
358        data_store,
359        operator_sessions: Arc::new(Mutex::new(HashMap::new())),
360        roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
361        task_store,
362        run_store,
363        replay_store,
364        base_url,
365        sync_timeout_secs,
366    };
367    Router::new()
368        .route("/v1/healthz", get(healthz))
369        .route("/v1/status", get(status_get))
370        // session = collection (POST = attach, DELETE = detach, sid via Authorization)
371        .route(
372            "/v1/sessions",
373            post(sessions_attach).delete(sessions_detach),
374        )
375        // task = flat, single level; authz resolved via Authorization: Bearer <sid>
376        .route("/v1/tasks", post(tasks_start).get(tasks::tasks_list))
377        .route("/v1/tasks/:id", get(tasks::task_get))
378        .route("/v1/tasks/:id/runs", post(tasks::task_rekick))
379        .route("/v1/tasks/:id/runs/:run/steps", get(projection::steps_list))
380        .route(
381            "/v1/tasks/:id/runs/:run/steps/:step",
382            get(projection::step_get),
383        )
384        .route(
385            "/v1/tasks/:id/runs/:run/steps/:step/content",
386            get(projection::step_content),
387        )
388        .route("/v1/runs/:id", get(tasks::run_get))
389        // Resume an Interrupted Run under the SAME run_id (replay cursor +
390        // stored launch-input snapshot); see `tasks::run_resume`.
391        .route("/v1/runs/:id/resume", post(tasks::run_resume))
392        // Rerun-from-step on a terminal Run under the SAME run_id (physically
393        // truncates the replay log at the cut point); see
394        // `tasks::run_rerun_from` for the full contract (GH #71 Layer A).
395        .route("/v1/runs/:id/rerun-from", post(tasks::run_rerun_from))
396        // REST-like Operator login flow (Bearer-mandatory, roles exclusivity).
397        // Sole WS Operator session route; see `operator_ws::login` module doc.
398        .route("/v1/operators", post(operators_create))
399        .route("/v1/operators/:sid/ws", get(operators_ws_connect))
400        .route(
401            "/v1/operators/:sid",
402            get(operators_info).delete(operators_delete),
403        )
404        // SubAgent self-fetch path (the SubAgent self-fetch design). The SubAgent puts the
405        // CapToken handed over via WS Spawn into Bearer and hits the prompt / result
406        // endpoints directly over HTTP. See the `worker` module doc for details.
407        .route("/v1/worker/prompt", get(worker::worker_prompt))
408        .route("/v1/worker/result", post(worker::worker_result))
409        // Simplified endpoint (= worker POSTs with just token + raw body; task_id is auto-looked-up).
410        // `DefaultBodyLimit::max` is applied explicitly here (and on the sibling
411        // `/v1/worker/artifact` below) — same 2MB axum ships as its implicit
412        // global default, made visible rather than relied on silently.
413        .route(
414            "/v1/worker/submit",
415            post(worker::worker_submit).layer(DefaultBodyLimit::max(2 * 1024 * 1024)),
416        )
417        // GH #36 ST1: named multi-part worker output. A worker stages one
418        // named part per POST here, then completes the attempt with the
419        // ordinary `/v1/worker/submit` above — see the `worker` module doc.
420        .route(
421            "/v1/worker/artifact",
422            post(worker::worker_artifact).layer(DefaultBodyLimit::max(2 * 1024 * 1024)),
423        )
424        // GH #31: `Http`-mode fetch target for `system_ref.uri` (raw baked system
425        // bytes, same Bearer flow as `/v1/worker/prompt`) + live per-agent render-size
426        // lookup for `bp_doctor` (no Bearer, same trust tier as blueprints `get_head`).
427        .route(
428            "/v1/worker/prompt/system",
429            get(worker::worker_prompt_system),
430        )
431        .route(
432            "/v1/agents/:name/render-size",
433            get(worker::agent_render_size),
434        )
435        // GH #32: structured worker degradation reporting — independent channel,
436        // never touches OutputStore / the fold path. See the `worker` module doc.
437        .route("/v1/worker/degradation", post(worker::worker_degradation))
438        // Data path (v9 Big Response handling, independent from Domain / verdict flow)
439        .route("/v1/data/emit", post(data::data_emit))
440        .route(
441            "/v1/data/:key",
442            get(data::data_get).post(data::data_emit_named),
443        )
444        .with_state(state)
445}
446
447/// Default registry = Subprocess + RustFn (baseline `identity` worker pre-baked) + empty Operator factory.
448///
449/// `RustFnInProcessSpawnerFactory` gets one baseline entry (`fn_id = "identity"`)
450/// baked in via [`mlua_swarm::worker::baseline::extend_with_baseline`]. This
451/// is the shared bootstrap / smoke worker SoT across each binary (the server / MCP adapter /
452/// one-shot runner) — it structurally replaces the old per-binary inline echo injection.
453///
454/// Usage: default Task path at server startup. If production needs additional
455/// backends, callers bring in a different registry via
456/// `build_router_with(engine, custom_registry)`. The enhance flow
457/// (= patch-spawner / patch-applier / verifier-router / committer axes) uses
458/// [`default_registry_with_enhance_flow`].
459///
460/// The Operator factory is an empty shell with zero registrations (= sids are
461/// dynamically registered per WS connect; see the `operator_ws` module).
462pub fn default_registry() -> SpawnerRegistry {
463    let rustfn_factory =
464        mlua_swarm::worker::baseline::extend_with_baseline(RustFnInProcessSpawnerFactory::new());
465
466    let mut reg = SpawnerRegistry::new();
467    reg.register::<SubprocessProcessSpawnerFactory>(Arc::new(SubprocessProcessSpawnerFactory));
468    reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(rustfn_factory));
469    // Empty `LuaInProcessSpawnerFactory`: no `fn_id` is pre-registered here,
470    // but BP agents can still declare `kind: lua` by carrying an inline
471    // `spec.source` (or a `$file`-expanded Lua chunk). This lets a BP ship
472    // deterministic Lua gates on the vanilla registry, without opting into
473    // the enhance flow. See `LuaInProcessSpawnerFactory` docs for the spec
474    // shape.
475    reg.register::<LuaInProcessSpawnerFactory>(Arc::new(LuaInProcessSpawnerFactory::new()));
476    reg.register::<OperatorSpawnerFactory>(Arc::new(OperatorSpawnerFactory::new()));
477    reg
478}
479
480/// Opt-in registry that merges [`default_registry`] with the enhance flow
481/// (Lua factory + AgentBlock factory).
482///
483/// Selected via the `the server` CLI flag `--enable-enhance-flow`. The enhance
484/// flow is a separate-axis wrapper: the Lua factory (= 3 Lua workers + 3 primitive
485/// bridges) and the AgentBlock factory (= patch-spawner path, expects
486/// `assets/operator_scripts/blueprint_patch_spawner.lua` + `ANTHROPIC_API_KEY`)
487/// are baked in as pipeline defaults. The baseline RustFn (`identity`) is pre-baked
488/// the same way as in `default_registry`.
489pub fn default_registry_with_enhance_flow() -> SpawnerRegistry {
490    let lua_factory =
491        mlua_swarm::enhance::blueprint::extend_factory(LuaInProcessSpawnerFactory::new());
492    // The Factory is stateless (= 1 process → 1 factory shared by all AgentDefs).
493    // Per-agent specialization (script_path / project_root, etc.) goes through AgentDef.spec.
494    // The enhance-flow patch-spawner is declared literally in agents[].spec of `default_blueprint.yaml`.
495    let agent_block_factory =
496        mlua_swarm::worker::agent_block::AgentBlockInProcessSpawnerFactory::new();
497    let rustfn_factory =
498        mlua_swarm::worker::baseline::extend_with_baseline(RustFnInProcessSpawnerFactory::new());
499
500    let mut reg = SpawnerRegistry::new();
501    reg.register::<SubprocessProcessSpawnerFactory>(Arc::new(SubprocessProcessSpawnerFactory));
502    reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(rustfn_factory));
503    reg.register::<LuaInProcessSpawnerFactory>(Arc::new(lua_factory));
504    reg.register::<mlua_swarm::worker::agent_block::AgentBlockInProcessSpawnerFactory>(Arc::new(
505        agent_block_factory,
506    ));
507    reg.register::<OperatorSpawnerFactory>(Arc::new(OperatorSpawnerFactory::new()));
508    reg
509}
510
511// ─── handlers ────────────────────────────────────────────────────────────
512
513async fn healthz() -> &'static str {
514    "ok"
515}
516
517/// Response body for `GET /v1/status` (issue #35 ST4 — lifecycle
518/// occupancy guard). Cheap-to-poll summary of "is it safe to kill this
519/// server right now".
520#[derive(Debug, Clone, Serialize, schemars::JsonSchema)]
521pub struct StatusResponse {
522    /// Count of `Run`s currently `Running` (`RunStore::list_running`).
523    /// Degrades to `0` on a store error rather than 500ing — see
524    /// module doc rationale.
525    pub running_runs: usize,
526    /// Count of attached Operator ids (`engine.list_operator_ids()`,
527    /// same idiom as `run_flow_form`'s Guard 1).
528    pub attached_operators: usize,
529}
530
531/// `GET /v1/status`. Infallible summary for the ST4 occupancy guard —
532/// store/engine query failures degrade the corresponding count to `0`
533/// (logged via `tracing::warn!`) rather than 500ing, since this
534/// endpoint may be polled frequently by a lifecycle-check caller that
535/// should not itself become a hang/error surface.
536async fn status_get(State(state): State<AppState>) -> Json<StatusResponse> {
537    let running_runs = state
538        .run_store
539        .list_running()
540        .await
541        .map(|v| v.len())
542        .unwrap_or_else(|e| {
543            tracing::warn!(error = %e, "status_get: list_running failed");
544            0
545        });
546    let attached_operators = state.engine.list_operator_ids().await.len();
547    Json(StatusResponse {
548        running_runs,
549        attached_operators,
550    })
551}
552
553#[derive(Deserialize)]
554struct AttachReq {
555    agent_id: String,
556    role: String,
557    ttl_secs: u64,
558}
559
560#[derive(Serialize)]
561struct AttachResp {
562    session_id: String,
563    role: String,
564}
565
566async fn sessions_attach(
567    State(state): State<AppState>,
568    Json(req): Json<AttachReq>,
569) -> Result<Json<AttachResp>, ApiError> {
570    let role = parse_role(&req.role)?;
571    let token = state
572        .engine
573        .attach(req.agent_id, role, Duration::from_secs(req.ttl_secs))
574        .await
575        .map_err(ApiError::engine)?;
576    // The wire `session_id` stays the nonce (Bearer credential contract);
577    // the server-side map key is its fingerprint (issue #14).
578    let sid = token.nonce.clone();
579    let key = token.fingerprint();
580    state.sessions.lock().await.map.insert(key, token);
581    Ok(Json(AttachResp {
582        session_id: sid,
583        role: req.role,
584    }))
585}
586
587async fn sessions_detach(
588    State(state): State<AppState>,
589    headers: HeaderMap,
590) -> Result<StatusCode, ApiError> {
591    let sid = extract_bearer(&headers)?;
592    let token = take_session_token(&state, &sid).await?;
593    state
594        .engine
595        .detach(&token)
596        .await
597        .map_err(ApiError::engine)?;
598    Ok(StatusCode::NO_CONTENT)
599}
600
601// ─── Unified /v1/tasks schema (= flow-eval path, Operator inject supported) ───────
602
603/// `/v1/tasks` POST schema. Uses the flow-eval path and supports Operator inject
604/// (kind / spawn_hook / senior_bridge). Expressing a one-shot task as a 1-Step
605/// Blueprint is the only correct model.
606///
607/// `pub` (issue #19 ST5) so its `schemars`-derived JSON Schema can be
608/// generated cross-crate by `mlua-swarm-cli`'s `mse://api/http-endpoints`
609/// MCP resource; fields stay module-private (no public field-level API
610/// surface is intended).
611#[derive(Deserialize, schemars::JsonSchema)]
612pub struct TaskLaunchRequest {
613    /// `BlueprintRef` selects Inline (a full Blueprint value) or Id (a
614    /// store lookup). Left opaque here — its own schema nests the full
615    /// `Blueprint` schema (owned by `mse://api/blueprint-schema`), and
616    /// mixing the two into this HTTP-endpoint resource would violate
617    /// their separation of concerns (see the resource's module doc).
618    #[schemars(with = "Value")]
619    blueprint: BlueprintRef,
620    /// flow.ir's initial `ctx` — every `Step.in` `$.<path>` reads from
621    /// here. This field's role is limited to the flow-ir eval seed
622    /// (issue #19); the Task-level execution context lives in the
623    /// sibling top-level fields below (`project_root` / `work_dir` /
624    /// `task_metadata`), promoted out of `init_ctx` to remove the
625    /// prior "free bag nested in free JSON" duplication.
626    ///
627    /// Backward compat: the pre-#19 shape — the same three keys nested
628    /// directly inside this object — is still honored as a fallback
629    /// when the sibling field is absent; see `run_flow_form`'s 2-stage
630    /// resolution and `TaskInputMiddleware::from_init_ctx`.
631    #[schemars(with = "Value")]
632    init_ctx: Value,
633    /// Task-level project root (issue #19 canonical Task IF field —
634    /// promoted out of `init_ctx`). Takes priority over a same-named
635    /// key nested inside `init_ctx` (backward-compat fallback).
636    #[serde(default)]
637    project_root: Option<String>,
638    /// Task-level working directory (issue #19), same priority rule as
639    /// `project_root`.
640    #[serde(default)]
641    work_dir: Option<String>,
642    /// Task-level arbitrary metadata bag (issue #19), same priority
643    /// rule as `project_root`.
644    #[serde(default)]
645    #[schemars(with = "Option<Value>")]
646    task_metadata: Option<Value>,
647    /// TTL in seconds. When unspecified (`None`), falls back in this order:
648    /// (1) `metadata.default_run_ttl_secs` from the resolved BP,
649    /// (2) if absent, the server global `default_run_ttl()` (1800s).
650    #[serde(default)]
651    ttl_secs: Option<u64>,
652    #[serde(default)]
653    operator: Option<OperatorReq>,
654    /// Explicit Operator session sid (or role alias) this task's entire Spawn
655    /// stream should be routed to (runtime Operator match stage 1).
656    ///
657    /// When `Some`, it is validated at request time against
658    /// `state.engine.list_operator_ids()` (the live `engine.operators`
659    /// registry key set): an unknown/never-registered id returns `400`
660    /// immediately — this is a deliberate hard-fail, in contrast to
661    /// `OperatorDelegateWrapped::spawn`, which silently falls through to
662    /// `inner.spawn` on a registry miss. A sid that *was* registered but has
663    /// since disconnected (WS `tx` cleared, session entry retained for
664    /// reconnect) passes this check and surfaces as an explicit dispatch-time
665    /// error instead (`WSOperatorSession::send_and_await` returns `Err` when
666    /// `tx` is `None`), which also propagates as a request failure rather
667    /// than a silent fallback.
668    ///
669    /// On success this value **overrides** `operator.operator_backend_id`
670    /// (last-write-wins, `operator_sid` takes priority) before the flow is
671    /// dispatched — see `run_flow_form`. Dispatch still only delegates if the
672    /// Blueprint opts into `spawner_hints.layers = ["operator_delegate"]`
673    /// (unchanged precondition, same as the existing `operator_backend_id`
674    /// field).
675    ///
676    /// When unset, behavior is unchanged: whatever
677    /// `operator.operator_backend_id` / BP-level `operator_ref` alias
678    /// resolution already does still applies.
679    #[serde(default)]
680    operator_sid: Option<String>,
681    /// Per-request override for the sync launch's timeout ceiling (GH #33
682    /// Guard 2, see `run_flow_form`'s doc comment). `None` (the default;
683    /// existing clients are unaffected) falls back to
684    /// `AppState::sync_timeout_secs` (server config), then the built-in
685    /// default (300s). `Some(0)` is rejected with `400` — omit the field
686    /// to defer to the server default rather than sending an explicit
687    /// zero.
688    #[serde(default)]
689    timeout_secs: Option<u64>,
690    /// Human-facing description of the work item (e.g. "resolve issue #10"),
691    /// stashed verbatim into the minted `TaskRecord.goal`. Omitted / `None`
692    /// stores an empty string — the flow-eval path itself never reads it.
693    #[serde(default)]
694    goal: Option<String>,
695    /// The "launch request" tier (tier 1, highest
696    /// priority) of the `check_policy` cascade
697    /// (`launch request > blueprint > server config`). `None` (the default;
698    /// existing clients are unaffected) leaves the tier unspecified so the
699    /// Blueprint-declared `check_policy` and, failing that, the server-wide
700    /// `EngineCfg.check_policy` default decide. Wire form is snake_case
701    /// (`"silent"` / `"warn"` / `"strict"`). Threaded verbatim into
702    /// `TaskApplicationInput.check_policy`.
703    #[serde(default)]
704    check_policy: Option<CheckPolicy>,
705    /// GH #37: opt into the detached (asynchronous) launch. `false` (the
706    /// default; existing clients are unaffected) keeps the synchronous
707    /// launch: the handler drives the flow eval inline and returns the
708    /// `final_ctx` on completion. `true` spawns the flow eval as a
709    /// detached background task and returns `202 Accepted` immediately
710    /// with `{task_id, run_id, status: "running"}` (`final_ctx` is
711    /// `null`) — the run's only lifetime bound is `ttl_secs`, and its
712    /// outcome is observed via `GET /v1/runs/:id` (or the `swarm_status`
713    /// MCP tool). Mutually exclusive with `timeout_secs` (the sync-launch
714    /// ceiling has no meaning for a detached run; combining them is a
715    /// `400`).
716    #[serde(default)]
717    detach: bool,
718}
719
720/// Operator inject sub-schema of [`TaskLaunchRequest`] (`kind` / `id` /
721/// `spawn_hook_id` / `senior_bridge_id` / `operator_backend_id` /
722/// `per_agent_kinds`). `pub` for the same cross-crate schema-generation
723/// reason as `TaskLaunchRequest`.
724#[derive(Deserialize, Default, schemars::JsonSchema)]
725pub struct OperatorReq {
726    /// `main_ai` / `automate` / `composite`. This is the "Runtime Global"
727    /// tier of the 4-tier `OperatorKind` cascade (see `mlua_swarm
728    /// ::ctx::collapse_operator_kind`); when unspecified, falls through to
729    /// the BP-level tiers (`OperatorDef.kind` / `Blueprint
730    /// .default_operator_kind`) instead of eagerly defaulting to `automate`.
731    #[serde(default)]
732    kind: Option<String>,
733    /// Operator id at attach time (= sessions tracking key in the EventLog); unspecified defaults to `"http-run"`.
734    #[serde(default)]
735    id: Option<String>,
736    /// Name of a hook pre-registered via `engine.register_spawn_hook`; `None` if unspecified.
737    #[serde(default)]
738    spawn_hook_id: Option<String>,
739    /// Name of a bridge pre-registered via `engine.register_senior_bridge`; `None` if unspecified.
740    #[serde(default)]
741    senior_bridge_id: Option<String>,
742    /// Name of an Operator backend pre-registered via `engine.register_operator`
743    /// (= the path that delegates the entire spawn to an external Operator);
744    /// `None` if unspecified. When `kind == MainAi/Composite` and this id is `Some`,
745    /// `OperatorDelegateMiddleware` bypasses `inner.spawn` and calls `operator.execute` instead.
746    /// This is a different axis from `operator.id` (= session tracking label);
747    /// `operator_backend_id` is the registry lookup key.
748    #[serde(default)]
749    operator_backend_id: Option<String>,
750    /// "Runtime Agent-level" tier (highest priority) of the `OperatorKind`
751    /// cascade — per-agent override, keyed by `AgentDef.name`, value is
752    /// `main_ai` / `automate` / `composite` (same parsing as `kind`).
753    /// `None` / absent means no per-agent override.
754    #[serde(default)]
755    per_agent_kinds: Option<HashMap<String, String>>,
756}
757
758/// Parse a wire-level kind string (`"main_ai"` / `"automate"` / `"composite"`)
759/// into `OperatorKind`. Shared by `OperatorReq.kind` and
760/// `OperatorReq.per_agent_kinds` values.
761fn parse_operator_kind_str(s: &str) -> Result<mlua_swarm::OperatorKind, ApiError> {
762    use mlua_swarm::OperatorKind;
763    match s {
764        "main_ai" => Ok(OperatorKind::MainAi),
765        "composite" => Ok(OperatorKind::Composite),
766        "automate" => Ok(OperatorKind::Automate),
767        other => Err(ApiError::bad_request(format!(
768            "operator kind: unknown value '{other}' (expected main_ai|automate|composite)"
769        ))),
770    }
771}
772
773/// `/v1/tasks` POST response body. `pub` for the same cross-crate
774/// schema-generation reason as [`TaskLaunchRequest`].
775#[derive(Serialize, schemars::JsonSchema)]
776pub struct TaskLaunchResponse {
777    /// The final flow.ir `ctx` after every `Step.out` has been written.
778    #[schemars(with = "Value")]
779    final_ctx: Value,
780    /// Debug-formatted `BlueprintVersion` the run resolved against, when
781    /// the Blueprint came from a store lookup (`None` for `Inline` refs).
782    bound_version: Option<String>,
783    /// Resolved TTL (seconds) actually applied to the run. Exposes the
784    /// 3-layer cascade (request body → BP metadata → server default) so
785    /// clients can verify which value took effect without re-deriving it.
786    effective_ttl_secs: u64,
787    /// Which layer of the TTL cascade won.
788    ttl_source: TtlSource,
789    /// The `TaskRecord` minted for this request (issue #13 ID-hierarchy
790    /// persistence). `GET /v1/tasks/:id` re-fetches it; `POST
791    /// /v1/tasks/:id/runs` re-kicks it under a fresh `RunId`.
792    #[schemars(with = "String")]
793    task_id: TaskId,
794    /// The `RunRecord` minted for this specific kick. `GET /v1/runs/:id`
795    /// re-fetches it (`step_entries` included).
796    #[schemars(with = "String")]
797    run_id: RunId,
798    /// Launch outcome at response time (GH #37). The synchronous path
799    /// (default) reports `done` — the flow eval completed before this
800    /// response was built. A detached launch (`detach: true`) reports
801    /// `running` — the eval continues in the background; poll `GET
802    /// /v1/runs/:id` for the terminal status and result.
803    status: RunStatus,
804}
805
806/// `tasks_start`'s reply — a [`TaskLaunchResponse`] plus the HTTP status
807/// it rides out on (`200 OK` for the synchronous path, `202 Accepted` for
808/// a detached launch, GH #37). A tuple struct with the body first so
809/// handler-level tests keep their established `.0` access to the response
810/// body regardless of which path produced it.
811pub struct TaskLaunchReply(pub TaskLaunchResponse, pub StatusCode);
812
813impl IntoResponse for TaskLaunchReply {
814    fn into_response(self) -> Response {
815        (self.1, Json(self.0)).into_response()
816    }
817}
818
819/// Which layer of the TTL cascade (request body → BP metadata → server
820/// default) resolved [`TaskLaunchResponse::effective_ttl_secs`]. `pub` for
821/// the same cross-crate schema-generation reason as `TaskLaunchRequest`.
822#[derive(Serialize, Clone, Copy, Debug, PartialEq, Eq, schemars::JsonSchema)]
823#[serde(rename_all = "snake_case")]
824pub enum TtlSource {
825    /// The request body's `ttl_secs` was set explicitly.
826    RequestBody,
827    /// The request body omitted `ttl_secs`; the resolved Blueprint's
828    /// `metadata.default_run_ttl_secs` was set.
829    BpMetadata,
830    /// Both the request body and the Blueprint metadata omitted a TTL;
831    /// the server-global `default_run_ttl()` (1800s) applied.
832    ServerDefault,
833}
834
835/// Unified `/v1/tasks` POST entry (= Flow form only).
836/// Runs `Blueprint.flow` to completion via flow eval in a single round-trip.
837/// One-shot tasks are also expressed as a 1-Step Blueprint. Operator
838/// (kind / spawn_hook / senior_bridge) can be injected per request body.
839/// `operator_sid` (S2, runtime Operator match stage 1) additionally
840/// lets the caller pin the task to a specific already-registered Operator
841/// session sid, bypassing BP-level alias lookup — see `TaskLaunchRequest` doc.
842async fn tasks_start(
843    State(state): State<AppState>,
844    Json(req): Json<TaskLaunchRequest>,
845) -> Result<TaskLaunchReply, ApiError> {
846    run_flow_form(&state, req).await
847}
848
849/// Flow-form path (= via `TaskApplication::handle_with_run`).
850/// Core handler behind the `/v1/tasks` entry (`tasks_start`).
851///
852/// Engine stateless-executor refactor: the per-request
853/// sub_engine + 3-registry propagate loop is retired; the startup-built
854/// `state.task_app` (= a `TaskLaunchService` wrap around `state.engine`) is
855/// used directly. The Operator callback IF (`spawn_hook_id` /
856/// `senior_bridge_id` / `operator_backend_id`) is registered on
857/// `state.engine.register_*` at WS connect time — the engine is the SoT.
858/// See the `operator_ws` module doc for details.
859///
860/// # GH #33 — sync-hang guards
861///
862/// This handler is always synchronous end-to-end (no sync/async branch);
863/// two fail-loud guards keep a bad launch from hanging the HTTP request
864/// forever:
865///
866/// - **Guard 1 (readiness precheck, `503`)**: when the request/BP
867///   references an operator backend (`operator.operator_backend_id`, set
868///   directly or via `operator_sid`) and `state.engine.list_operator_ids()`
869///   is empty, the request fails immediately rather than dispatching into
870///   a session with nothing attached to serve it. Coarse by design — a
871///   launch that cannot be cheaply determined to route through an operator
872///   is never rejected here (Guard 2 still covers the hang in that case).
873/// - **Guard 2 (sync timeout, `504`)**: the single
874///   `state.task_app.handle_with_run` await is wrapped in
875///   `tokio::time::timeout`. Ceiling cascade, highest priority first:
876///   request `timeout_secs` (rejecting `Some(0)` with `400`), then
877///   `AppState::sync_timeout_secs` (server config), then the built-in
878///   default (300s). On expiry the timed-out future is dropped — this
879///   cancels the in-process flow eval (the flow is abandoned, not
880///   resumed; intended v1 semantics) — and the Task/Run records are
881///   best-effort marked `Failed` so they do not stay `Running` forever.
882///
883/// # GH #37 — detached launch (`detach: true`)
884///
885/// The sync semantics above tie the flow-eval driver's lifetime to this
886/// request's future — a long-running detached worker that outlives the
887/// ceiling gets its (individually successful) `/v1/worker/*` submits
888/// orphaned when the driver is cancelled. `detach: true` decouples them:
889/// the eval (plus `finalize_run`) runs in a `tokio::spawn`ed background
890/// task whose only lifetime bound is the resolved `ttl_secs` (marked
891/// `Failed` on expiry, same best-effort persistence as Guard 2), and the
892/// handler returns `202 Accepted` with `status: "running"` immediately.
893/// Guard 1 still applies (checked before any store write); Guard 2's
894/// ceiling does not (`timeout_secs` + `detach` together is a `400`).
895/// Client disconnect after the `202` cannot cancel the run.
896async fn run_flow_form(
897    state: &AppState,
898    req: TaskLaunchRequest,
899) -> Result<TaskLaunchReply, ApiError> {
900    use mlua_swarm::application::{BlueprintRef as AppBlueprintRef, TaskApplicationInput};
901    use mlua_swarm::OperatorKind;
902
903    // Snapshot everything the TaskRecord needs before `req.blueprint` /
904    // `req.init_ctx` are moved into the dispatch path below.
905    let blueprint_ref_json = serde_json::to_value(&req.blueprint)
906        .map_err(|e| ApiError::bad_request(format!("blueprint snapshot: {e}")))?;
907    let input_ctx_snapshot = req.init_ctx.clone();
908    let goal = req.goal.clone().unwrap_or_default();
909
910    // issue #19 ST2: resolve the Task-level canonical fields
911    // (`project_root` / `work_dir` / `task_metadata`) once, at the wire
912    // boundary. Sibling top-level fields on the request body take
913    // priority; the pre-#19 shape (same key nested inside `init_ctx`) is
914    // only a fallback for legacy callers. The result is threaded straight
915    // through as `TaskApplicationInput.task_input` — `init_ctx` itself is
916    // NOT mutated, so it stays a pure flow-ir eval seed identical to
917    // whatever the caller sent.
918    let task_input_spec = build_task_input_spec_from_request(&req);
919    // Issue #19 ST4: snapshot the resolved spec into the `TaskRecord` (JSON,
920    // same "bare `Value`" rationale as `blueprint_ref_json` /
921    // `input_ctx_snapshot` above) so `POST /v1/tasks/:id/runs` can resolve
922    // it back out on rekick without re-deriving it from a since-stale
923    // request body. Cloned rather than computed from `task_input_spec`
924    // after the fact — the original is still moved into
925    // `TaskApplicationInput.task_input` below.
926    let task_input_spec_snapshot = task_input_spec
927        .clone()
928        .map(|spec| serde_json::to_value(&spec))
929        .transpose()
930        .map_err(|e| ApiError::bad_request(format!("task_input_spec snapshot: {e}")))?;
931    let init_ctx = req.init_ctx.clone();
932
933    let mut op_req = req.operator.unwrap_or_default();
934
935    // S2: explicit `operator_sid` override (runtime Operator match stage 1).
936    // Resolved *before* building `operator_kind` / dispatching so an
937    // unknown sid fails fast with a 400, never silently falling back to the
938    // BP-level alias lookup. See `TaskLaunchRequest::operator_sid` doc for the
939    // disconnected-vs-unknown distinction.
940    if let Some(sid) = &req.operator_sid {
941        let known_ids = state.engine.list_operator_ids().await;
942        if !known_ids.iter().any(|id| id == sid) {
943            return Err(ApiError::bad_request(format!(
944                "operator_sid: no such registered operator session '{sid}'"
945            )));
946        }
947        op_req.operator_backend_id = Some(sid.clone());
948    }
949
950    // GH #33 Guard 2 ceiling resolution: request field > server config >
951    // built-in default (300s, `config::default_sync_timeout_secs`).
952    // Validated up front — before any TaskRecord/RunRecord side effects —
953    // so a caller-supplied `Some(0)` fails fast with `400` rather than
954    // minting records for a launch that was never going to dispatch.
955    // GH #37: `detach: true` makes the sync ceiling meaningless (the
956    // detached run is bounded by `ttl_secs` alone) — combining the two
957    // is rejected here, same fail-fast-before-side-effects ordering.
958    let detach = req.detach;
959    let sync_timeout_secs = match (detach, req.timeout_secs) {
960        (true, Some(_)) => {
961            return Err(ApiError::bad_request(
962                "timeout_secs is the synchronous launch ceiling and does not apply to a \
963                 detached launch (detach: true), whose lifetime bound is ttl_secs — omit \
964                 timeout_secs"
965                    .into(),
966            ));
967        }
968        (false, Some(0)) => {
969            return Err(ApiError::bad_request(
970                "timeout_secs: 0 is invalid; omit the field to use the server default".into(),
971            ));
972        }
973        (false, Some(v)) => v,
974        (_, None) => state.sync_timeout_secs,
975    };
976
977    // GH #33 Guard 1: operator readiness precheck. Coarse signal — this
978    // handler can cheaply see whether the request/BP references an
979    // operator backend (`operator.operator_backend_id`, set directly or
980    // resolved above from `operator_sid`), but not the full
981    // `OperatorDelegateMiddleware` routing decision (that also considers
982    // BP-level `kind` tiers, resolved only at dispatch time). When a
983    // backend is referenced and *zero* operators are attached at all,
984    // fail fast rather than dispatching into a session nothing can serve.
985    // A launch this coarse check cannot positively identify as
986    // operator-delegate is never rejected here — Guard 2 (the timeout
987    // wrap below) still covers the hang in that case.
988    if let Some(backend_id) = op_req.operator_backend_id.as_deref() {
989        let attached = state.engine.list_operator_ids().await;
990        if attached.is_empty() {
991            return Err(ApiError::unavailable(format!(
992                "no operator attached to serve this launch (operator backend '{backend_id}' \
993                 requested): attach an operator via POST /v1/operators + WS, or use the \
994                 poll-style flow (GET /v1/worker/prompt + POST /v1/worker/submit)"
995            )));
996        }
997    }
998
999    // "Runtime Global" tier: `Some(_)` — including `Some(Automate)` — is
1000    // always an explicit request that outranks the BP-level tiers; an
1001    // absent/unset `kind` in the request body stays `None`, leaving the
1002    // BP-level tiers (`OperatorDef.kind` / `Blueprint.default_operator_kind`)
1003    // to decide instead of eagerly defaulting to `Automate`.
1004    let operator_kind = op_req
1005        .kind
1006        .as_deref()
1007        .map(parse_operator_kind_str)
1008        .transpose()?;
1009    let operator_id = op_req.id.unwrap_or_else(|| "http-run".to_string());
1010    // "Runtime Agent-level" tier: per-agent overrides. Absent/empty = no
1011    // override for any agent, letting the BP-level tiers decide per agent.
1012    let mut operator_kind_overrides: HashMap<String, OperatorKind> = HashMap::new();
1013    for (agent, kind_str) in op_req.per_agent_kinds.take().unwrap_or_default() {
1014        operator_kind_overrides.insert(agent, parse_operator_kind_str(&kind_str)?);
1015    }
1016
1017    let blueprint: AppBlueprintRef = match req.blueprint {
1018        AppBlueprintRef::Inline { value } => AppBlueprintRef::Inline { value },
1019        AppBlueprintRef::Id { id, version } => AppBlueprintRef::Id { id, version },
1020    };
1021
1022    // TTL resolution cascade: (1) request body value, (2) BP metadata `default_run_ttl_secs`,
1023    // (3) server global default (`default_run_ttl()`, 1800s).
1024    let (ttl_secs, ttl_source) = match req.ttl_secs {
1025        Some(v) => (v, TtlSource::RequestBody),
1026        None => {
1027            let (resolved_bp, _ver) = state
1028                .task_app
1029                .resolve(&blueprint)
1030                .await
1031                .map_err(|e| ApiError::bad_request(format!("bp resolve: {e}")))?;
1032            match resolved_bp.metadata.default_run_ttl_secs {
1033                Some(v) => (v, TtlSource::BpMetadata),
1034                None => (default_run_ttl(), TtlSource::ServerDefault),
1035            }
1036        }
1037    };
1038
1039    // Build the launch input up front so a snapshot of it can be persisted
1040    // into the RunRecord below — an Interrupted Run is resumed from that
1041    // snapshot (`POST /v1/runs/:id/resume`) under the same run_id.
1042    let input = TaskApplicationInput {
1043        blueprint,
1044        operator_id: operator_id.clone(),
1045        role: Role::Operator,
1046        ttl: Duration::from_secs(ttl_secs),
1047        init_ctx,
1048        operator_kind,
1049        bridge_id: op_req.senior_bridge_id,
1050        hook_id: op_req.spawn_hook_id,
1051        operator_backend_id: op_req.operator_backend_id,
1052        operator_kind_overrides,
1053        task_input: task_input_spec,
1054        // The request-body top-level `check_policy` (tier 1)
1055        // flows straight into the cascade resolved once in
1056        // `TaskLaunchService::launch`.
1057        check_policy: req.check_policy,
1058    };
1059    let input_json = Some(tasks::snapshot_launch_input(&input)?);
1060
1061    // issue #13 ID-hierarchy persistence: mint the work-item identity (Task)
1062    // and this kick's identity (Run) *before* dispatching, so a Task/Run
1063    // pair always exists even if the flow itself fails mid-way (the
1064    // Failed-status paths below still have a row to update).
1065    let task_id = TaskId::new();
1066    let run_id = RunId::new();
1067    let now = tasks::now_secs();
1068    state
1069        .task_store
1070        .create(TaskRecord {
1071            id: task_id.clone(),
1072            goal,
1073            blueprint_ref: blueprint_ref_json,
1074            input_ctx: input_ctx_snapshot,
1075            task_input_spec: task_input_spec_snapshot,
1076            status: TaskRecordStatus::Running,
1077            created_at: now,
1078            updated_at: now,
1079        })
1080        .await
1081        .map_err(ApiError::engine)?;
1082    state
1083        .run_store
1084        .create(RunRecord {
1085            id: run_id.clone(),
1086            task_id: task_id.clone(),
1087            status: RunStatus::Running,
1088            step_entries: Vec::new(),
1089            degradations: Vec::new(),
1090            operator_sid: req.operator_sid.clone(),
1091            result_ref: None,
1092            input_json,
1093            created_at: now,
1094            updated_at: now,
1095        })
1096        .await
1097        .map_err(ApiError::engine)?;
1098
1099    let run_ctx = RunContext::new(run_id.clone(), state.run_store.clone())
1100        .with_replay_store(state.replay_store.clone());
1101
1102    // GH #37 detached launch: the eval driver runs in its own spawned
1103    // task — its lifetime is bound to `ttl_secs`, not to this request's
1104    // future (client disconnect / handler completion cannot cancel it).
1105    // The spawned task owns the run to its terminal status: `finalize_run`
1106    // on completion, or the same best-effort `Failed` marking as Guard 2
1107    // if the ttl ceiling expires first.
1108    if detach {
1109        let bg_state = state.clone();
1110        let bg_task_id = task_id.clone();
1111        let bg_run_id = run_id.clone();
1112        tokio::spawn(async move {
1113            let outcome = match tokio::time::timeout(
1114                Duration::from_secs(ttl_secs),
1115                bg_state.task_app.handle_with_run(input, Some(run_ctx)),
1116            )
1117            .await
1118            {
1119                Ok(outcome) => outcome,
1120                Err(_elapsed) => {
1121                    let reason = json!({
1122                        "error": format!("detached run exceeded {ttl_secs}s ttl ceiling"),
1123                    });
1124                    if let Err(e) = bg_state.run_store.set_result(&bg_run_id, reason).await {
1125                        tracing::warn!(%bg_run_id, error = %e, "run_flow_form: detached ttl set_result failed");
1126                    }
1127                    if let Err(e) = bg_state
1128                        .run_store
1129                        .update_status(&bg_run_id, RunStatus::Failed)
1130                        .await
1131                    {
1132                        tracing::warn!(%bg_run_id, error = %e, "run_flow_form: detached ttl run update_status(Failed) failed");
1133                    }
1134                    if let Err(e) = bg_state
1135                        .task_store
1136                        .update_status(&bg_task_id, TaskRecordStatus::Failed)
1137                        .await
1138                    {
1139                        tracing::warn!(%bg_task_id, error = %e, "run_flow_form: detached ttl task update_status(Failed) failed");
1140                    }
1141                    return;
1142                }
1143            };
1144            // `finalize_run` persists both the Ok and Err outcomes itself;
1145            // the passthrough return value has no consumer here.
1146            let _ = tasks::finalize_run(&bg_state, &bg_task_id, &bg_run_id, outcome).await;
1147        });
1148        return Ok(TaskLaunchReply(
1149            TaskLaunchResponse {
1150                final_ctx: Value::Null,
1151                bound_version: None,
1152                effective_ttl_secs: ttl_secs,
1153                ttl_source,
1154                task_id,
1155                run_id,
1156                status: RunStatus::Running,
1157            },
1158            StatusCode::ACCEPTED,
1159        ));
1160    }
1161
1162    // GH #33 Guard 2: the single await point this handler blocks on. On
1163    // expiry the timed-out future is dropped, cancelling the in-process
1164    // flow eval — the flow is abandoned, not resumed (intended v1
1165    // semantics; stage-granularity resume is a coarser guarantee than
1166    // this handler makes, out of scope here).
1167    let outcome = match tokio::time::timeout(
1168        Duration::from_secs(sync_timeout_secs),
1169        state.task_app.handle_with_run(input, Some(run_ctx)),
1170    )
1171    .await
1172    {
1173        Ok(outcome) => outcome,
1174        Err(_elapsed) => {
1175            // Best effort: mark the Task/Run so they do not stay `Running`
1176            // forever. Reuses the existing `Failed` variant (no new
1177            // schema-crate enum additions) and stashes a reason string
1178            // into `RunRecord.result_ref` — the only free-form field the
1179            // Run schema carries; secondary persistence failures here are
1180            // logged and swallowed, mirroring `tasks::finalize_run`'s
1181            // error-path convention.
1182            let reason = json!({
1183                "error": format!("sync launch exceeded {sync_timeout_secs}s timeout ceiling"),
1184            });
1185            if let Err(e) = state.run_store.set_result(&run_id, reason).await {
1186                tracing::warn!(%run_id, error = %e, "run_flow_form: timeout run set_result failed");
1187            }
1188            if let Err(e) = state
1189                .run_store
1190                .update_status(&run_id, RunStatus::Failed)
1191                .await
1192            {
1193                tracing::warn!(%run_id, error = %e, "run_flow_form: timeout run update_status(Failed) failed");
1194            }
1195            if let Err(e) = state
1196                .task_store
1197                .update_status(&task_id, TaskRecordStatus::Failed)
1198                .await
1199            {
1200                tracing::warn!(%task_id, error = %e, "run_flow_form: timeout task update_status(Failed) failed");
1201            }
1202            return Err(ApiError::timeout(format!(
1203                "sync launch exceeded {sync_timeout_secs}s timeout ceiling: the in-process flow \
1204                 eval was abandoned (dropping the future cancels it); attach an operator that \
1205                 acks promptly (POST /v1/operators + WS), or raise timeout_secs / sync_timeout_secs"
1206            )));
1207        }
1208    };
1209
1210    let out = tasks::finalize_run(state, &task_id, &run_id, outcome)
1211        .await
1212        .map_err(|e| ApiError::bad_request(format!("run: {e}")))?;
1213
1214    Ok(TaskLaunchReply(
1215        TaskLaunchResponse {
1216            final_ctx: out.final_ctx,
1217            bound_version: out.bound_version.map(|v| format!("{:?}", v)),
1218            effective_ttl_secs: ttl_secs,
1219            ttl_source,
1220            task_id,
1221            run_id,
1222            status: RunStatus::Done,
1223        },
1224        StatusCode::OK,
1225    ))
1226}
1227
1228/// issue #19 ST2 direct sibling-field resolver — extracts the three
1229/// Task-level canonical fields (`project_root` / `work_dir` /
1230/// `task_metadata`) once at the wire boundary. Sibling top-level body
1231/// fields take priority; the pre-#19 shape (same key nested inside
1232/// `init_ctx`) is only a fallback for legacy callers. Unlike the ST1
1233/// `resolve_task_level_init_ctx` bridge this replaced, `init_ctx` is
1234/// NOT mutated — the resolved values are handed straight to
1235/// [`mlua_swarm::service::TaskLaunchInput::task_input`], keeping
1236/// `init_ctx` a pure flow-ir eval seed.
1237///
1238/// Returns `None` when all three fields resolve to `None` (no
1239/// middleware is layered onto the spawner stack downstream — the
1240/// [`mlua_swarm::middleware::task_input::TaskInputMiddleware::new_from_fields`]
1241/// contract).
1242fn build_task_input_spec_from_request(
1243    req: &TaskLaunchRequest,
1244) -> Option<mlua_swarm::service::TaskInputSpec> {
1245    let project_root = req.project_root.clone().or_else(|| {
1246        req.init_ctx
1247            .get("project_root")
1248            .and_then(Value::as_str)
1249            .map(String::from)
1250    });
1251    let work_dir = req.work_dir.clone().or_else(|| {
1252        req.init_ctx
1253            .get("work_dir")
1254            .and_then(Value::as_str)
1255            .map(String::from)
1256    });
1257    let task_metadata = req.task_metadata.clone().or_else(|| {
1258        req.init_ctx
1259            .get("task_metadata")
1260            .filter(|v| v.is_object())
1261            .cloned()
1262    });
1263
1264    if project_root.is_none() && work_dir.is_none() && task_metadata.is_none() {
1265        None
1266    } else {
1267        Some(mlua_swarm::service::TaskInputSpec {
1268            project_root,
1269            work_dir,
1270            task_metadata,
1271        })
1272    }
1273}
1274
1275// ─── helpers ─────────────────────────────────────────────────────────────
1276
1277async fn take_session_token(state: &AppState, sid: &str) -> Result<CapToken, ApiError> {
1278    // `sid` on this path is the token nonce itself (a bearer secret), so
1279    // both the map key and the not-found diagnostic use its fingerprint
1280    // (issue #14 — never echo the nonce back in an error body).
1281    let key = mlua_swarm::types::token_fingerprint(sid);
1282    state
1283        .sessions
1284        .lock()
1285        .await
1286        .map
1287        .remove(&key)
1288        .ok_or_else(|| ApiError::not_found(format!("session: fp={key}")))
1289}
1290
1291/// Extracts sid from `Authorization: Bearer <sid>`. Strict — does not accept any other scheme prefix.
1292fn extract_bearer(headers: &HeaderMap) -> Result<String, ApiError> {
1293    let v = headers
1294        .get(AUTHORIZATION)
1295        .ok_or_else(|| ApiError::bad_request("missing Authorization header".into()))?
1296        .to_str()
1297        .map_err(|_| ApiError::bad_request("invalid Authorization header encoding".into()))?;
1298    let sid = v
1299        .strip_prefix("Bearer ")
1300        .ok_or_else(|| ApiError::bad_request("Authorization must be 'Bearer <sid>'".into()))?
1301        .trim();
1302    if sid.is_empty() {
1303        return Err(ApiError::bad_request("Bearer sid is empty".into()));
1304    }
1305    Ok(sid.to_string())
1306}
1307
1308fn parse_role(s: &str) -> Result<Role, ApiError> {
1309    match s.to_ascii_lowercase().as_str() {
1310        "operator" => Ok(Role::Operator),
1311        "worker" => Ok(Role::Worker),
1312        "observer" => Ok(Role::Observer),
1313        "senior" => Ok(Role::Senior),
1314        other => Err(ApiError::bad_request(format!("unknown role: {other}"))),
1315    }
1316}
1317
1318// ─── error type ──────────────────────────────────────────────────────────
1319
1320/// Uniform error response type for the handlers in this module. Converts to
1321/// a JSON `{"error": message}` body with the given status via [`IntoResponse`].
1322#[derive(Debug)]
1323pub struct ApiError {
1324    status: StatusCode,
1325    message: String,
1326}
1327
1328impl ApiError {
1329    /// Wraps an engine-side error as `500 Internal Server Error`.
1330    pub fn engine(e: impl std::fmt::Display) -> Self {
1331        Self {
1332            status: StatusCode::INTERNAL_SERVER_ERROR,
1333            message: format!("engine: {e}"),
1334        }
1335    }
1336    /// Builds a `404 Not Found` with the given message.
1337    pub fn not_found(m: String) -> Self {
1338        Self {
1339            status: StatusCode::NOT_FOUND,
1340            message: m,
1341        }
1342    }
1343    /// Builds a `400 Bad Request` with the given message.
1344    pub fn bad_request(m: String) -> Self {
1345        Self {
1346            status: StatusCode::BAD_REQUEST,
1347            message: m,
1348        }
1349    }
1350    /// Builds a `409 Conflict` with the given message (`POST
1351    /// /v1/runs/:id/resume` — the Run is not `Interrupted`, or a concurrent
1352    /// resume already won the `Interrupted -> Running` compare-and-set).
1353    pub fn conflict(m: String) -> Self {
1354        Self {
1355            status: StatusCode::CONFLICT,
1356            message: m,
1357        }
1358    }
1359    /// Builds a `503 Service Unavailable` with the given message (GH #33
1360    /// Guard 1 — operator readiness precheck).
1361    pub fn unavailable(m: String) -> Self {
1362        Self {
1363            status: StatusCode::SERVICE_UNAVAILABLE,
1364            message: m,
1365        }
1366    }
1367    /// Builds a `504 Gateway Timeout` with the given message (GH #33
1368    /// Guard 2 — sync launch timeout ceiling).
1369    pub fn timeout(m: String) -> Self {
1370        Self {
1371            status: StatusCode::GATEWAY_TIMEOUT,
1372            message: m,
1373        }
1374    }
1375    /// Builds a `410 Gone` with the given message (GH #37 — worker
1376    /// submit/artifact addressed at a Run that already reached a terminal
1377    /// status; the silent-`204`-then-orphan alternative is the failure
1378    /// shape this replaces).
1379    pub fn gone(m: String) -> Self {
1380        Self {
1381            status: StatusCode::GONE,
1382            message: m,
1383        }
1384    }
1385    /// Builds a `413 Payload Too Large` with the given message (GH #42 —
1386    /// `@file:` sentinel resolves to a file larger than the shared
1387    /// `DefaultBodyLimit`; same size ceiling as the inline body path).
1388    pub fn payload_too_large(m: String) -> Self {
1389        Self {
1390            status: StatusCode::PAYLOAD_TOO_LARGE,
1391            message: m,
1392        }
1393    }
1394    /// Builds a `422 Unprocessable Entity` with the given message (GH #50
1395    /// — a `worker_submit` / `worker_artifact` value violates the
1396    /// dispatching agent's declared `VerdictContract`: rejected before it
1397    /// reaches `submit_worker_result_trusted` / `stage_worker_artifact_trusted`,
1398    /// i.e. before it can land in the flow ctx).
1399    pub fn unprocessable(m: impl Into<String>) -> Self {
1400        Self {
1401            status: StatusCode::UNPROCESSABLE_ENTITY,
1402            message: m.into(),
1403        }
1404    }
1405}
1406
1407impl IntoResponse for ApiError {
1408    fn into_response(self) -> Response {
1409        (self.status, Json(json!({"error": self.message}))).into_response()
1410    }
1411}
1412
1413fn default_run_ttl() -> u64 {
1414    // 1800s (= 30 min). Prevents op_token expiry across a flow.ir multi-step chain
1415    // (= 5+ SubAgent dispatches at 30–60s each). Origin: the observed fvloop smoke
1416    // where a post-gate mock-commit dispatch blew past 300s and expired — sibling of worker_token TTL.
1417    1800
1418}
1419
1420/// TTL cascade resolve helper (Blueprint metadata → server default fallback).
1421/// Second-stage fallback, called when the POST `/v1/tasks` body does not set `ttl_secs`.
1422/// (1) If BP metadata `default_run_ttl_secs` is `Some`, use it.
1423/// (2) If `None`, fall back to the server global `default_run_ttl()` (1800s).
1424///
1425/// # Full cascade (combined in `run_flow_form`)
1426///
1427/// - request body `ttl_secs=Some(v)` → v (this helper is not called)
1428/// - request body `None` + metadata `Some(v)` → v
1429/// - request body `None` + metadata `None` → `default_run_ttl()` = 1800s
1430#[cfg(test)]
1431fn resolve_ttl_from_metadata(metadata_ttl: Option<u64>) -> u64 {
1432    metadata_ttl.unwrap_or_else(default_run_ttl)
1433}
1434
1435#[cfg(test)]
1436mod tests {
1437    use super::*;
1438
1439    /// TTL cascade case 1: when the request body sets it, that value is used as-is
1440    /// (upper branch that does not go through the helper; semantic verify of the
1441    /// `Some(v) => v` direct-return path in `run_flow_form`).
1442    #[test]
1443    fn ttl_cascade_request_body_wins_over_metadata() {
1444        let req_ttl: Option<u64> = Some(100);
1445        let metadata_ttl: Option<u64> = Some(3600);
1446        let effective = match req_ttl {
1447            Some(v) => v,
1448            None => resolve_ttl_from_metadata(metadata_ttl),
1449        };
1450        assert_eq!(
1451            effective, 100,
1452            "request body ttl_secs=100 must win over metadata=3600 (cascade priority (1) > (2))"
1453        );
1454    }
1455
1456    /// TTL cascade case 2: request body omitted + BP metadata `Some(N)` → `N` is effective.
1457    #[test]
1458    fn ttl_cascade_metadata_used_when_body_missing() {
1459        let req_ttl: Option<u64> = None;
1460        let metadata_ttl: Option<u64> = Some(3600);
1461        let effective = match req_ttl {
1462            Some(v) => v,
1463            None => resolve_ttl_from_metadata(metadata_ttl),
1464        };
1465        assert_eq!(
1466            effective, 3600,
1467            "body None + metadata=3600 must resolve to 3600 (cascade (2))"
1468        );
1469    }
1470
1471    /// TTL cascade case 3: request body omitted + BP metadata `None` → server default (1800s).
1472    #[test]
1473    fn ttl_cascade_server_default_when_both_missing() {
1474        let req_ttl: Option<u64> = None;
1475        let metadata_ttl: Option<u64> = None;
1476        let effective = match req_ttl {
1477            Some(v) => v,
1478            None => resolve_ttl_from_metadata(metadata_ttl),
1479        };
1480        assert_eq!(
1481            effective,
1482            default_run_ttl(),
1483            "body None + metadata None must fall back to default_run_ttl() = 1800s"
1484        );
1485        assert_eq!(effective, 1800, "default_run_ttl() literal = 1800s");
1486    }
1487
1488    /// Helper unit: metadata `None` → 1800 (server default expansion).
1489    #[test]
1490    fn resolve_ttl_from_metadata_none_returns_server_default() {
1491        assert_eq!(resolve_ttl_from_metadata(None), 1800);
1492    }
1493
1494    /// Helper unit: metadata `Some(N)` → `N` (server default ignored).
1495    #[test]
1496    fn resolve_ttl_from_metadata_some_returns_value() {
1497        assert_eq!(resolve_ttl_from_metadata(Some(7200)), 7200);
1498        assert_eq!(resolve_ttl_from_metadata(Some(60)), 60);
1499    }
1500
1501    // ──────────────────────────────────────────────────────────────────
1502    // `TaskLaunchRequest.check_policy` wire field (T5)
1503    // ──────────────────────────────────────────────────────────────────
1504
1505    /// T5: a `POST /v1/tasks` body carrying a top-level `check_policy`
1506    /// deserializes into `TaskLaunchRequest.check_policy` using the
1507    /// snake_case wire form.
1508    #[test]
1509    fn task_launch_request_parses_check_policy_wire_field() {
1510        let body = json!({
1511            "blueprint": { "kind": "id", "id": "some-bp" },
1512            "init_ctx": {},
1513            "check_policy": "silent",
1514        });
1515        let req: TaskLaunchRequest =
1516            serde_json::from_value(body).expect("request must deserialize");
1517        assert_eq!(req.check_policy, Some(CheckPolicy::Silent));
1518    }
1519
1520    /// A body that omits `check_policy` leaves the field `None` (existing
1521    /// clients are unaffected — `#[serde(default)]`).
1522    #[test]
1523    fn task_launch_request_check_policy_defaults_to_none_when_omitted() {
1524        let body = json!({
1525            "blueprint": { "kind": "id", "id": "some-bp" },
1526            "init_ctx": {},
1527        });
1528        let req: TaskLaunchRequest =
1529            serde_json::from_value(body).expect("request must deserialize");
1530        assert_eq!(req.check_policy, None);
1531    }
1532
1533    // ──────────────────────────────────────────────────────────────────
1534    // issue #19 ST2: `build_task_input_spec_from_request` direct resolver
1535    // ──────────────────────────────────────────────────────────────────
1536
1537    fn task_req(
1538        init_ctx: Value,
1539        project_root: Option<&str>,
1540        work_dir: Option<&str>,
1541        task_metadata: Option<Value>,
1542    ) -> TaskLaunchRequest {
1543        TaskLaunchRequest {
1544            blueprint: BlueprintRef::Id {
1545                id: mlua_swarm::blueprint::store::BlueprintId::new("ut"),
1546                version: Default::default(),
1547            },
1548            init_ctx,
1549            project_root: project_root.map(String::from),
1550            work_dir: work_dir.map(String::from),
1551            task_metadata,
1552            ttl_secs: None,
1553            operator: None,
1554            operator_sid: None,
1555            timeout_secs: None,
1556            goal: None,
1557            detach: false,
1558            check_policy: None,
1559        }
1560    }
1561
1562    /// (a) Sibling fields only — no legacy keys in `init_ctx` — are
1563    /// returned in the `TaskInputSpec` unchanged. `init_ctx` itself is
1564    /// untouched by this resolver (checked separately at the call site).
1565    #[test]
1566    fn build_task_input_spec_from_request_returns_sibling_fields_when_present() {
1567        let req = task_req(
1568            json!({"free": "form"}),
1569            Some("/repo/sibling"),
1570            Some("/repo/sibling/work"),
1571            Some(json!({"issue": 19})),
1572        );
1573        let spec = build_task_input_spec_from_request(&req).expect("spec must be Some");
1574        assert_eq!(spec.project_root.as_deref(), Some("/repo/sibling"));
1575        assert_eq!(spec.work_dir.as_deref(), Some("/repo/sibling/work"));
1576        assert_eq!(spec.task_metadata, Some(json!({"issue": 19})));
1577    }
1578
1579    /// (b) No sibling fields — the pre-#19 shape (same 3 keys nested
1580    /// inside `init_ctx`) is used as the fallback source.
1581    #[test]
1582    fn build_task_input_spec_from_request_falls_back_to_legacy_init_ctx_shape() {
1583        let req = task_req(
1584            json!({
1585                "project_root": "/repo/legacy",
1586                "work_dir": "/repo/legacy/work",
1587                "task_metadata": {"issue": 17},
1588            }),
1589            None,
1590            None,
1591            None,
1592        );
1593        let spec = build_task_input_spec_from_request(&req).expect("spec must be Some");
1594        assert_eq!(spec.project_root.as_deref(), Some("/repo/legacy"));
1595        assert_eq!(spec.work_dir.as_deref(), Some("/repo/legacy/work"));
1596        assert_eq!(spec.task_metadata, Some(json!({"issue": 17})));
1597    }
1598
1599    /// (c) Both present — the sibling field must win over the legacy
1600    /// `init_ctx`-nested value.
1601    #[test]
1602    fn build_task_input_spec_from_request_sibling_wins_over_legacy_shape() {
1603        let req = task_req(
1604            json!({
1605                "project_root": "/repo/legacy",
1606                "work_dir": "/repo/legacy/work",
1607                "task_metadata": {"issue": 17},
1608            }),
1609            Some("/repo/sibling"),
1610            Some("/repo/sibling/work"),
1611            Some(json!({"issue": 19})),
1612        );
1613        let spec = build_task_input_spec_from_request(&req).expect("spec must be Some");
1614        assert_eq!(
1615            spec.project_root.as_deref(),
1616            Some("/repo/sibling"),
1617            "sibling field must win over the legacy init_ctx-nested value"
1618        );
1619        assert_eq!(spec.work_dir.as_deref(), Some("/repo/sibling/work"));
1620        assert_eq!(spec.task_metadata, Some(json!({"issue": 19})));
1621    }
1622
1623    /// (d) All three fields absent from both sibling and legacy shapes —
1624    /// resolver returns `None`, and no middleware is layered downstream.
1625    #[test]
1626    fn build_task_input_spec_from_request_returns_none_when_no_fields_present() {
1627        let req = task_req(json!({"unrelated": "value"}), None, None, None);
1628        assert!(build_task_input_spec_from_request(&req).is_none());
1629    }
1630
1631    /// Minimal `AppState` for the `status_get` handler-fn-direct-call test
1632    /// below — same construction shape as `tasks.rs::test_state()`
1633    /// (mirrors what `build_router_full` does internally, skipping the
1634    /// `Router` wrapper).
1635    fn status_test_state() -> AppState {
1636        let engine = Engine::new(mlua_swarm::EngineCfg::default());
1637        let compiler = mlua_swarm::Compiler::new(default_registry());
1638        let launch = Arc::new(mlua_swarm::TaskLaunchService::new(engine.clone(), compiler));
1639        AppState {
1640            engine,
1641            sessions: Arc::new(Mutex::new(SessionStore::default())),
1642            task_app: Arc::new(mlua_swarm::TaskApplication::new_inline_only(launch)),
1643            ws_operator_factory: None,
1644            data_store: Arc::new(mlua_swarm::store::output::InMemoryOutputStore::new()),
1645            operator_sessions: Arc::new(Mutex::new(HashMap::new())),
1646            roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
1647            task_store: Arc::new(mlua_swarm::store::task::InMemoryTaskStore::new()),
1648            run_store: Arc::new(mlua_swarm::store::run::InMemoryRunStore::new()),
1649            replay_store: Arc::new(mlua_swarm::store::replay::InMemoryReplayStore::new()),
1650            base_url: None,
1651            sync_timeout_secs: 300,
1652        }
1653    }
1654
1655    /// issue #35 ST4 Acceptance Criteria: `GET /v1/status` reports the
1656    /// count of `Running` `Run`s (`RunStore::list_running`) and attached
1657    /// Operator ids (`engine.list_operator_ids()`), called directly as a
1658    /// handler fn (no `Router` wrapper — this crate's established
1659    /// unit-test convention).
1660    #[tokio::test]
1661    async fn status_get_reports_running_runs_and_operators() {
1662        let state = status_test_state();
1663
1664        let now = std::time::SystemTime::now()
1665            .duration_since(std::time::UNIX_EPOCH)
1666            .map(|d| d.as_secs())
1667            .unwrap_or(0);
1668        state
1669            .run_store
1670            .create(RunRecord {
1671                id: RunId::new(),
1672                task_id: TaskId::new(),
1673                status: RunStatus::Running,
1674                step_entries: Vec::new(),
1675                degradations: Vec::new(),
1676                operator_sid: None,
1677                result_ref: None,
1678                input_json: None,
1679                created_at: now,
1680                updated_at: now,
1681            })
1682            .await
1683            .expect("seed running RunRecord");
1684
1685        // Throwaway `Operator` impl — only registration/list-count matters
1686        // for this test, `execute` is never dispatched (same idiom as
1687        // `tasks.rs::StallingOperator`).
1688        struct NoopOperator;
1689        #[async_trait::async_trait]
1690        impl mlua_swarm::Operator for NoopOperator {
1691            async fn execute(
1692                &self,
1693                _ctx: &mlua_swarm::Ctx,
1694                _system: Option<String>,
1695                _prompt: Value,
1696                _worker: Option<mlua_swarm::WorkerBinding>,
1697                _worker_token: mlua_swarm::CapToken,
1698            ) -> Result<mlua_swarm::WorkerResult, mlua_swarm::WorkerError> {
1699                unimplemented!("not exercised by this test — only registration/list matters")
1700            }
1701        }
1702        state
1703            .engine
1704            .register_operator("test-op", Arc::new(NoopOperator))
1705            .await;
1706
1707        let Json(resp) = status_get(State(state)).await;
1708        assert_eq!(resp.running_runs, 1);
1709        assert_eq!(resp.attached_operators, 1);
1710    }
1711}