Skip to main content

mlua_swarm_server/
tasks.rs

1//! HTTP surface for the Task/Run persistence axis (issue #13 ID-hierarchy
2//! reconciliation: Blueprint → Task → Run → Step → Attempt).
3//!
4//! - `GET  /v1/tasks`          — list every persisted `TaskRecord`, newest first.
5//! - `GET  /v1/tasks/:id`      — a `TaskRecord` plus every `RunRecord` kicked from it.
6//! - `POST /v1/tasks/:id/runs` — re-kick an existing Task: mints a fresh `RunId`,
7//!   re-resolves the stored `blueprint_ref` (refreshing `Blueprint.default_init_ctx`
8//!   exactly like original launch time — issue #19 ST4), 3-layer-merges it with
9//!   `TaskRecord.input_ctx` and an **optional** [`RunKickRequest`] body's
10//!   `init_ctx_override` (see [`merge_init_ctx_3layer`]), dispatches through
11//!   `TaskApplication::handle_with_run`, and returns the new `{task_id, run_id}`
12//!   pair. A body-less request (or one that omits both fields) preserves the
13//!   pre-#19 rekick behavior byte-for-byte.
14//! - `GET  /v1/runs/:id`       — a single `RunRecord` (`step_entries` trace included).
15//!
16//! `POST /v1/tasks` itself (the flow-eval entry point, `tasks_start` /
17//! `run_flow_form`) stays in `crate::lib` — it is the pre-existing
18//! Operator-inject-aware dispatch path this module's handlers re-kick
19//! through, not a new one. This module owns the read/list/re-kick surface
20//! plus the [`finalize_run`] persistence helper both paths share.
21//!
22//! Authorization follows the same convention as the existing `POST /v1/tasks`
23//! entry: no `Authorization` header is required (the route is open), and the
24//! only Operator-session correlation available is the request-body-level
25//! `operator_sid` (see `crate::TaskLaunchRequest` doc) — this module invents no
26//! new auth mechanism.
27
28use axum::{
29    extract::{Path, Query, State},
30    http::StatusCode,
31    Json,
32};
33use mlua_swarm::application::{TaskApplicationError, TaskApplicationInput, TaskApplicationOutput};
34use mlua_swarm::service::merge_init_ctx_3layer;
35use mlua_swarm::store::run::{RunContext, RunRecord, RunStatus, RunStoreError};
36use mlua_swarm::store::task::{TaskRecord, TaskRecordStatus, TaskStoreError};
37use mlua_swarm::{Role, RunId, TaskId, TaskInputSpec};
38use serde::{Deserialize, Serialize};
39use serde_json::Value;
40use std::collections::HashMap;
41use std::time::Duration;
42
43use crate::{ApiError, AppState};
44
45/// Current Unix time in whole seconds. `TaskRecord` / `RunRecord` timestamps
46/// are `u64` seconds (not milliseconds) — see their field docs in
47/// `mlua_swarm::store::task` / `mlua_swarm::store::run`.
48pub(crate) fn now_secs() -> u64 {
49    std::time::SystemTime::now()
50        .duration_since(std::time::UNIX_EPOCH)
51        .map(|d| d.as_secs())
52        .unwrap_or(0)
53}
54
55/// Shared finalize step for a dispatched kick: updates the Run's
56/// `result_ref` + status and the owning Task's coarse status based on the
57/// `TaskApplication::handle_with_run` outcome, then returns that same
58/// outcome unchanged so callers keep shaping their own wire response /
59/// error.
60///
61/// Secondary persistence failures (the store call itself erroring) are
62/// logged via `tracing::warn!` and otherwise swallowed — they must not mask
63/// the primary dispatch outcome the caller already has in hand.
64pub(crate) async fn finalize_run(
65    state: &AppState,
66    task_id: &TaskId,
67    run_id: &RunId,
68    outcome: Result<TaskApplicationOutput, TaskApplicationError>,
69) -> Result<TaskApplicationOutput, TaskApplicationError> {
70    match &outcome {
71        Ok(out) => {
72            if let Err(e) = state
73                .run_store
74                .set_result(run_id, out.final_ctx.clone())
75                .await
76            {
77                tracing::warn!(%run_id, error = %e, "finalize_run: set_result failed");
78            }
79            if let Err(e) = state.run_store.update_status(run_id, RunStatus::Done).await {
80                tracing::warn!(%run_id, error = %e, "finalize_run: run update_status(Done) failed");
81            }
82            if let Err(e) = state
83                .task_store
84                .update_status(task_id, TaskRecordStatus::Done)
85                .await
86            {
87                tracing::warn!(%task_id, error = %e, "finalize_run: task update_status(Done) failed");
88            }
89        }
90        Err(e) => {
91            if let Err(store_err) = state
92                .run_store
93                .update_status(run_id, RunStatus::Failed)
94                .await
95            {
96                tracing::warn!(%run_id, error = %store_err, "finalize_run: run update_status(Failed) failed");
97            }
98            if let Err(store_err) = state
99                .task_store
100                .update_status(task_id, TaskRecordStatus::Failed)
101                .await
102            {
103                tracing::warn!(%task_id, error = %store_err, "finalize_run: task update_status(Failed) failed");
104            }
105            tracing::warn!(%task_id, %run_id, error = %e, "finalize_run: dispatch failed");
106        }
107    }
108    outcome
109}
110
111/// Query params for `GET /v1/tasks`.
112#[derive(Debug, Deserialize, Default)]
113pub struct TasksListQuery {
114    /// Caps the returned list to the first N entries (already newest-first
115    /// per `TaskStore::list`). Omitted = no cap.
116    #[serde(default)]
117    pub limit: Option<usize>,
118}
119
120/// `GET /v1/tasks?limit=N`. Lists every persisted `TaskRecord`, newest first.
121pub async fn tasks_list(
122    State(state): State<AppState>,
123    Query(q): Query<TasksListQuery>,
124) -> Result<Json<Vec<TaskRecord>>, ApiError> {
125    let mut records = state.task_store.list().await.map_err(ApiError::engine)?;
126    if let Some(limit) = q.limit {
127        records.truncate(limit);
128    }
129    Ok(Json(records))
130}
131
132/// Response body for `GET /v1/tasks/:id`.
133#[derive(Debug, Clone, Serialize, schemars::JsonSchema)]
134pub struct TaskDetailResponse {
135    /// The Task's own record.
136    pub task: TaskRecord,
137    /// Every Run kicked from this Task, oldest first (`RunStore::list_by_task` order).
138    pub runs: Vec<RunRecord>,
139}
140
141/// `GET /v1/tasks/:id`. Returns the `TaskRecord` plus every `RunRecord`
142/// kicked from it (`RunStore::list_by_task`, oldest kick first).
143pub async fn task_get(
144    State(state): State<AppState>,
145    Path(id): Path<String>,
146) -> Result<Json<TaskDetailResponse>, ApiError> {
147    let task_id =
148        TaskId::parse(id).map_err(|e| ApiError::bad_request(format!("invalid task id: {e}")))?;
149    let task = state
150        .task_store
151        .get(&task_id)
152        .await
153        .map_err(map_task_store_err)?;
154    let runs = state
155        .run_store
156        .list_by_task(&task_id)
157        .await
158        .map_err(ApiError::engine)?;
159    Ok(Json(TaskDetailResponse { task, runs }))
160}
161
162/// Request body for `POST /v1/tasks/:id/runs` (issue #19 ST4) — every
163/// field is optional, and the body itself is optional (see
164/// [`task_rekick`]'s `Option<Json<Self>>` parameter); a caller that sends
165/// no body, or `{}`, or omits a field gets exactly today's rekick
166/// behavior for that layer.
167#[derive(Debug, Deserialize, Default, schemars::JsonSchema)]
168pub struct RunKickRequest {
169    /// Per-Run override for the flow-ir initial ctx. Merged on top of
170    /// `TaskRecord.input_ctx` (itself already merged on top of
171    /// `Blueprint.default_init_ctx` at original launch time) via
172    /// [`merge_init_ctx_3layer`] — Run wins on key collision, same
173    /// shallow-merge / non-Object-fully-replaces rule as every other
174    /// layer in the cascade. `None` (absent field, or no body at all) is
175    /// a no-op: the BP+Task merge alone seeds this kick, identical to
176    /// pre-#19 rekick.
177    #[serde(default)]
178    #[schemars(with = "Option<Value>")]
179    pub init_ctx_override: Option<Value>,
180    /// Per-Run override for the Task-level canonical fields
181    /// (`project_root` / `work_dir` / `task_metadata`). `None` falls back
182    /// to `TaskRecord.task_input_spec` (the spec resolved and snapshotted
183    /// at original `POST /v1/tasks` time); `Some` replaces it wholesale
184    /// for this kick only — the stored `TaskRecord.task_input_spec` is
185    /// never mutated by a rekick.
186    #[serde(default)]
187    pub task_input_override: Option<TaskInputSpec>,
188    /// Per-Run ceiling (seconds) for this kick's synchronous dispatch
189    /// await (issue #35 ST3 — GH #33 Guard 2 parity). `Some(0)` is
190    /// rejected (400). `None` falls back to `AppState.sync_timeout_secs`
191    /// (the server-wide default), same cascade as
192    /// `TaskLaunchRequest.timeout_secs` (`lib.rs:818-826`).
193    #[serde(default)]
194    pub timeout_secs: Option<u64>,
195    /// GH #37: opt into the detached (asynchronous) rekick — same
196    /// semantics as `TaskLaunchRequest.detach`. `false` (default) keeps
197    /// the synchronous dispatch; `true` spawns the flow eval as a
198    /// detached background task bounded by the run TTL alone and returns
199    /// `202 Accepted` with `status: "running"` immediately. Mutually
200    /// exclusive with `timeout_secs` (`400` when combined).
201    #[serde(default)]
202    pub detach: bool,
203}
204
205/// Response body for `POST /v1/tasks/:id/runs`.
206#[derive(Debug, Clone, Serialize, schemars::JsonSchema)]
207pub struct RunKickResponse {
208    /// The re-kicked Task's id (echoes the path param).
209    #[schemars(with = "String")]
210    pub task_id: TaskId,
211    /// The freshly minted Run id for this kick.
212    #[schemars(with = "String")]
213    pub run_id: RunId,
214    /// Kick outcome at response time (GH #37). The synchronous path
215    /// reports the dispatched run's terminal-side status (`done`); a
216    /// detached kick reports `running` — poll `GET /v1/runs/:id` for the
217    /// terminal status and result.
218    pub status: RunStatus,
219}
220
221/// `POST /v1/tasks/:id/runs`. Re-kicks an existing Task: reads its stored
222/// `blueprint_ref`, re-resolves it through [`TaskApplication::resolve`]
223/// (issue #19 ST4 — refreshes `Blueprint.default_init_ctx` exactly like
224/// original launch time, rather than replaying a launch-time-only
225/// snapshot), 3-layer-merges `{bp default, TaskRecord.input_ctx, an
226/// optional per-Run override}` via [`merge_init_ctx_3layer`], resolves the
227/// Task-level canonical fields (`RunKickRequest.task_input_override`,
228/// falling back to `TaskRecord.task_input_spec`), mints a fresh `RunId`,
229/// dispatches through `TaskApplication::handle_with_run` (the unadorned
230/// Operator-default path — no per-request Operator override support here,
231/// unlike `POST /v1/tasks`; the stored Task carries no such preferences)
232/// plus a freshly-built `RunContext` (issue #13 run_id propagation, so
233/// this kick's steps get their own `step_entries` trace), and persists the
234/// outcome via [`finalize_run`].
235///
236/// The body is optional (`Option<Json<RunKickRequest>>`) — no body, or a
237/// body with both fields absent, preserves the pre-#19 rekick behavior
238/// byte-for-byte (`must_not_simplify #3`).
239///
240/// Issue #35 ST3 ports the GH #33 sync-hang guards from `run_flow_form` to
241/// this handler, both checked before any Task/Run store write: Guard 1
242/// (503) fails fast when the resolved Blueprint declares the
243/// `operator_delegate` spawner-hint layer and no operator is attached;
244/// Guard 2 (504) wraps the dispatch await in `RunKickRequest.timeout_secs`
245/// (falling back to the server-wide `sync_timeout_secs`), marking the
246/// Run/Task `Failed` rather than leaving them `Running` forever on expiry.
247pub async fn task_rekick(
248    State(state): State<AppState>,
249    Path(id): Path<String>,
250    body: Option<Json<RunKickRequest>>,
251) -> Result<(StatusCode, Json<RunKickResponse>), ApiError> {
252    let task_id =
253        TaskId::parse(id).map_err(|e| ApiError::bad_request(format!("invalid task id: {e}")))?;
254    let task = state
255        .task_store
256        .get(&task_id)
257        .await
258        .map_err(map_task_store_err)?;
259
260    let blueprint_ref: mlua_swarm::application::BlueprintRef =
261        serde_json::from_value(task.blueprint_ref.clone()).map_err(|e| {
262            ApiError::bad_request(format!(
263                "task {task_id}: stored blueprint_ref failed to decode: {e}"
264            ))
265        })?;
266
267    // issue #19 ST4 (must_not_simplify #5): re-resolve the Blueprint the
268    // same way `run_flow_form`'s TTL cascade does, so a store-backed
269    // `BlueprintRef::Id` gets its *current* `default_init_ctx` on every
270    // rekick rather than whatever was true at original launch time. The
271    // Inline path is a pure pass-through, so this is a no-op there.
272    let (resolved_bp, _bound_version) = state
273        .task_app
274        .resolve(&blueprint_ref)
275        .await
276        .map_err(|e| ApiError::bad_request(format!("task {task_id}: bp resolve: {e}")))?;
277
278    let req = body.map(|Json(r)| r).unwrap_or_default();
279
280    // GH #33 Guard 2 ceiling resolution (issue #35 ST3 — mirrors
281    // `run_flow_form`'s `lib.rs:813-826` cascade): request field > server
282    // config > built-in default. Validated up front, before Guard 1 and
283    // before any Task/Run store writes, so a caller-supplied `Some(0)`
284    // fails fast with `400` rather than minting records for a rekick that
285    // was never going to dispatch.
286    // GH #37: `detach: true` makes the sync ceiling meaningless (the
287    // detached kick is bounded by the run TTL alone) — combining the two
288    // is rejected here, same fail-fast-before-side-effects ordering.
289    let detach = req.detach;
290    let sync_timeout_secs = match (detach, req.timeout_secs) {
291        (true, Some(_)) => {
292            return Err(ApiError::bad_request(
293                "timeout_secs is the synchronous rekick ceiling and does not apply to a \
294                 detached rekick (detach: true), whose lifetime bound is the run TTL — omit \
295                 timeout_secs"
296                    .into(),
297            ));
298        }
299        (false, Some(0)) => {
300            return Err(ApiError::bad_request(
301                "timeout_secs: 0 is invalid; omit the field to use the server default".into(),
302            ));
303        }
304        (false, Some(v)) => v,
305        (_, None) => state.sync_timeout_secs,
306    };
307
308    // GH #33 Guard 1 (issue #35 ST3 — adapted signal): `RunKickRequest`
309    // carries no per-request Operator override field (unlike
310    // `run_flow_form`'s `op_req.operator_backend_id`, sourced from
311    // `TaskLaunchRequest.operator` — this module's doc, above, confirms
312    // that's by design). The adapted "operator backend referenced" signal
313    // is the Blueprint's own `spawner_hints.layers` instead: when the
314    // resolved Blueprint declares the `operator_delegate` layer and zero
315    // operators are attached at all, fail fast rather than dispatching
316    // into a session nothing can serve. Same ordering invariant
317    // `run_flow_form` observes: this check runs before any Task/Run row
318    // is touched (no side effects on the 503 path).
319    if resolved_bp
320        .spawner_hints
321        .layers
322        .iter()
323        .any(|l| l == "operator_delegate")
324    {
325        let attached = state.engine.list_operator_ids().await;
326        if attached.is_empty() {
327            return Err(ApiError::unavailable(format!(
328                "no operator attached to serve this rekick (task {task_id}'s \
329                 Blueprint declares the operator_delegate layer): attach an \
330                 operator via POST /v1/operators + WS, or use the poll-style \
331                 flow (GET /v1/worker/prompt + POST /v1/worker/submit)"
332            )));
333        }
334    }
335
336    let merged_init_ctx = merge_init_ctx_3layer(
337        resolved_bp.default_init_ctx.as_ref(),
338        &task.input_ctx,
339        req.init_ctx_override.as_ref(),
340    );
341
342    // must_not_simplify #4: `task_input_override` wins for this kick only;
343    // falling back to the Task-level snapshot never mutates
344    // `TaskRecord.task_input_spec` itself.
345    let task_input_spec: Option<TaskInputSpec> = match req.task_input_override {
346        Some(over) => Some(over),
347        None => task
348            .task_input_spec
349            .as_ref()
350            .map(|v| serde_json::from_value(v.clone()))
351            .transpose()
352            .map_err(|e| {
353                ApiError::bad_request(format!(
354                    "task {task_id}: stored task_input_spec failed to decode: {e}"
355                ))
356            })?,
357    };
358
359    let run_id = RunId::new();
360    let now = now_secs();
361    state
362        .task_store
363        .update_status(&task_id, TaskRecordStatus::Running)
364        .await
365        .map_err(ApiError::engine)?;
366    state
367        .run_store
368        .create(RunRecord {
369            id: run_id.clone(),
370            task_id: task_id.clone(),
371            status: RunStatus::Running,
372            step_entries: Vec::new(),
373            degradations: Vec::new(),
374            operator_sid: None,
375            result_ref: None,
376            created_at: now,
377            updated_at: now,
378        })
379        .await
380        .map_err(ApiError::engine)?;
381
382    let input = TaskApplicationInput {
383        blueprint: blueprint_ref,
384        operator_id: "http-run".to_string(),
385        role: Role::Operator,
386        ttl: Duration::from_secs(crate::default_run_ttl()),
387        init_ctx: merged_init_ctx,
388        operator_kind: None,
389        bridge_id: None,
390        hook_id: None,
391        operator_backend_id: None,
392        operator_kind_overrides: HashMap::new(),
393        task_input: task_input_spec,
394    };
395    let run_ctx = RunContext {
396        run_id: run_id.clone(),
397        run_store: state.run_store.clone(),
398    };
399
400    // GH #37 detached rekick: same driver-detach semantics as
401    // `run_flow_form` — the eval runs in its own spawned task bounded by
402    // the run TTL alone, `finalize_run` (or the ttl-expiry `Failed`
403    // marking) is owned by that task, and this handler returns `202
404    // Accepted` immediately.
405    if detach {
406        let ttl_secs = crate::default_run_ttl();
407        let bg_state = state.clone();
408        let bg_task_id = task_id.clone();
409        let bg_run_id = run_id.clone();
410        tokio::spawn(async move {
411            let outcome = match tokio::time::timeout(
412                Duration::from_secs(ttl_secs),
413                bg_state.task_app.handle_with_run(input, Some(run_ctx)),
414            )
415            .await
416            {
417                Ok(outcome) => outcome,
418                Err(_elapsed) => {
419                    let reason = serde_json::json!({
420                        "error": format!("detached rekick exceeded {ttl_secs}s ttl ceiling"),
421                    });
422                    if let Err(e) = bg_state.run_store.set_result(&bg_run_id, reason).await {
423                        tracing::warn!(%bg_run_id, error = %e, "task_rekick: detached ttl set_result failed");
424                    }
425                    if let Err(e) = bg_state
426                        .run_store
427                        .update_status(&bg_run_id, RunStatus::Failed)
428                        .await
429                    {
430                        tracing::warn!(%bg_run_id, error = %e, "task_rekick: detached ttl run update_status failed");
431                    }
432                    if let Err(e) = bg_state
433                        .task_store
434                        .update_status(&bg_task_id, TaskRecordStatus::Failed)
435                        .await
436                    {
437                        tracing::warn!(%bg_task_id, error = %e, "task_rekick: detached ttl task update_status failed");
438                    }
439                    return;
440                }
441            };
442            // `finalize_run` persists both the Ok and Err outcomes itself;
443            // the passthrough return value has no consumer here.
444            let _ = finalize_run(&bg_state, &bg_task_id, &bg_run_id, outcome).await;
445        });
446        return Ok((
447            StatusCode::ACCEPTED,
448            Json(RunKickResponse {
449                task_id,
450                run_id,
451                status: RunStatus::Running,
452            }),
453        ));
454    }
455
456    // GH #33 Guard 2 (issue #35 ST3 — mirrors `run_flow_form`'s
457    // `lib.rs:935-990` exactly): the single await point this handler
458    // blocks on. On expiry the timed-out future is dropped, cancelling the
459    // in-process flow eval — the flow is abandoned, not resumed. Best
460    // effort: mark the Run/Task so they do not stay `Running` forever.
461    let outcome = match tokio::time::timeout(
462        Duration::from_secs(sync_timeout_secs),
463        state.task_app.handle_with_run(input, Some(run_ctx)),
464    )
465    .await
466    {
467        Ok(outcome) => outcome,
468        Err(_elapsed) => {
469            let reason = serde_json::json!({
470                "error": format!("sync rekick exceeded {sync_timeout_secs}s timeout ceiling")
471            });
472            if let Err(e) = state.run_store.set_result(&run_id, reason).await {
473                tracing::warn!(%run_id, error = %e, "task_rekick: timeout set_result failed");
474            }
475            if let Err(e) = state
476                .run_store
477                .update_status(&run_id, RunStatus::Failed)
478                .await
479            {
480                tracing::warn!(%run_id, error = %e, "task_rekick: timeout run update_status failed");
481            }
482            if let Err(e) = state
483                .task_store
484                .update_status(&task_id, TaskRecordStatus::Failed)
485                .await
486            {
487                tracing::warn!(%task_id, error = %e, "task_rekick: timeout task update_status failed");
488            }
489            return Err(ApiError::timeout(format!(
490                "sync rekick exceeded {sync_timeout_secs}s timeout ceiling: task {task_id}, run {run_id}"
491            )));
492        }
493    };
494    finalize_run(&state, &task_id, &run_id, outcome)
495        .await
496        .map_err(|e| ApiError::bad_request(format!("run: {e}")))?;
497
498    Ok((
499        StatusCode::CREATED,
500        Json(RunKickResponse {
501            task_id,
502            run_id,
503            status: RunStatus::Done,
504        }),
505    ))
506}
507
508/// `GET /v1/runs/:id`. Returns a single `RunRecord` (its `step_entries`
509/// trace included).
510pub async fn run_get(
511    State(state): State<AppState>,
512    Path(id): Path<String>,
513) -> Result<Json<RunRecord>, ApiError> {
514    let run_id =
515        RunId::parse(id).map_err(|e| ApiError::bad_request(format!("invalid run id: {e}")))?;
516    let run = state
517        .run_store
518        .get(&run_id)
519        .await
520        .map_err(map_run_store_err)?;
521    Ok(Json(run))
522}
523
524/// `pub(crate)` so `crate::projection`'s `GET /v1/tasks/:id/ctx` handler can
525/// reuse this module's existing-Task-existence-check error mapping (same
526/// 404-vs-500 split `task_get` already applies).
527pub(crate) fn map_task_store_err(e: TaskStoreError) -> ApiError {
528    match e {
529        TaskStoreError::NotFound(id) => ApiError::not_found(format!("task not found: {id}")),
530        other => ApiError::engine(other),
531    }
532}
533
534fn map_run_store_err(e: RunStoreError) -> ApiError {
535    match e {
536        RunStoreError::NotFound(id) => ApiError::not_found(format!("run not found: {id}")),
537        other => ApiError::engine(other),
538    }
539}
540
541// ──────────────────────────────────────────────────────────────────────────
542// UT
543// ──────────────────────────────────────────────────────────────────────────
544
545#[cfg(test)]
546mod tests {
547    use super::*;
548    use mlua_swarm::application::BlueprintRef;
549    use mlua_swarm::blueprint::{
550        current_schema_version, AgentDef, AgentKind, Blueprint, BlueprintMetadata, CompilerHints,
551        CompilerStrategy,
552    };
553    use mlua_swarm::core::config::EngineCfg;
554    use mlua_swarm::core::engine::Engine;
555    use mlua_swarm::store::output::InMemoryOutputStore;
556    use mlua_swarm::store::run::InMemoryRunStore;
557    use mlua_swarm::store::task::InMemoryTaskStore;
558    use std::collections::HashMap;
559    use std::sync::Arc;
560    use tokio::sync::Mutex;
561
562    /// A single-step flow.ir Blueprint that always succeeds: `Step { ref:
563    /// "identity", in: lit("hello"), out: $.out }` against the baseline
564    /// `RustFn` identity worker (same shape as `seed_blueprint` in
565    /// `mlua-swarm-cli`'s `serve.rs`, self-contained here rather than
566    /// importing a binary crate).
567    fn identity_blueprint() -> Blueprint {
568        Blueprint {
569            schema_version: current_schema_version(),
570            id: "tasks-test-bp".into(),
571            flow: serde_json::from_value(serde_json::json!({
572                "kind": "step",
573                "ref": mlua_swarm::worker::baseline::AG_IDENTITY,
574                "in": {"op": "lit", "value": "hello"},
575                "out": {"op": "path", "at": "$.out"},
576            }))
577            .expect("flow parse"),
578            agents: vec![AgentDef {
579                name: mlua_swarm::worker::baseline::AG_IDENTITY.into(),
580                kind: AgentKind::RustFn,
581                spec: serde_json::json!({"fn_id": mlua_swarm::worker::baseline::AG_IDENTITY}),
582                profile: None,
583                meta: None,
584            }],
585            operators: vec![],
586            metas: vec![],
587            hints: CompilerHints::default(),
588            strategy: CompilerStrategy::default(),
589            metadata: BlueprintMetadata::default(),
590            spawner_hints: Default::default(),
591            default_agent_kind: AgentKind::Operator,
592            default_operator_kind: None,
593            default_init_ctx: None,
594            default_agent_ctx: None,
595            default_context_policy: None,
596            projection_placement: None,
597            audits: vec![],
598            degradation_policy: None,
599        }
600    }
601
602    /// Minimal `AppState` for handler-level tests — mirrors the construction
603    /// `build_router_full` does internally, but skips the `Router` wrapper so
604    /// tests can call handler functions directly (this crate's established
605    /// unit-test convention; see e.g. `operator_ws::login`'s tests).
606    fn test_state() -> AppState {
607        let engine = Engine::new_with_layers(EngineCfg::default(), crate::default_layer_registry());
608        let compiler = mlua_swarm::Compiler::new(crate::default_registry());
609        let launch = Arc::new(mlua_swarm::TaskLaunchService::new(engine.clone(), compiler));
610        AppState {
611            engine,
612            sessions: Arc::new(Mutex::new(crate::SessionStore::default())),
613            task_app: Arc::new(mlua_swarm::TaskApplication::new_inline_only(launch)),
614            ws_operator_factory: None,
615            data_store: Arc::new(InMemoryOutputStore::new()),
616            operator_sessions: Arc::new(Mutex::new(HashMap::new())),
617            roles_to_sid: Arc::new(Mutex::new(HashMap::new())),
618            task_store: Arc::new(InMemoryTaskStore::new()),
619            run_store: Arc::new(InMemoryRunStore::new()),
620            base_url: None,
621            sync_timeout_secs: 300,
622        }
623    }
624
625    fn post_tasks_req(goal: &str) -> crate::TaskLaunchRequest {
626        crate::TaskLaunchRequest {
627            blueprint: BlueprintRef::Inline {
628                value: Box::new(identity_blueprint()),
629            },
630            init_ctx: serde_json::json!({"in": "hello"}),
631            project_root: None,
632            work_dir: None,
633            task_metadata: None,
634            ttl_secs: None,
635            operator: None,
636            operator_sid: None,
637            timeout_secs: None,
638            goal: Some(goal.to_string()),
639            detach: false,
640        }
641    }
642
643    #[test]
644    fn task_id_serializes_as_bare_string() {
645        // Sanity check for the newtype-struct transparency relied on
646        // throughout this module's response shapes (`TaskId` / `RunId`
647        // serialize as plain JSON strings, not `{"0": "..."}`).
648        let v = serde_json::to_value(TaskId::parse("T-abc").unwrap()).expect("serialize");
649        assert_eq!(v, serde_json::json!("T-abc"));
650    }
651
652    #[tokio::test]
653    async fn post_then_get_drill_down() {
654        let state = test_state();
655
656        let posted = crate::tasks_start(State(state.clone()), Json(post_tasks_req("smoke goal")))
657            .await
658            .expect("tasks_start")
659            .0;
660        let task_id = posted.task_id.clone();
661        let run_id = posted.run_id.clone();
662
663        // GET /v1/tasks lists it.
664        let list = tasks_list(State(state.clone()), Query(TasksListQuery { limit: None }))
665            .await
666            .expect("tasks_list")
667            .0;
668        assert!(
669            list.iter().any(|t| t.id == task_id),
670            "task {task_id} missing from list of {} tasks",
671            list.len()
672        );
673
674        // GET /v1/tasks/:id drills down to the Task + its Run.
675        let detail = task_get(State(state.clone()), Path(task_id.to_string()))
676            .await
677            .expect("task_get")
678            .0;
679        assert_eq!(detail.task.id, task_id);
680        assert_eq!(detail.task.goal, "smoke goal");
681        assert_eq!(detail.task.status, TaskRecordStatus::Done);
682        assert_eq!(detail.runs.len(), 1);
683        assert_eq!(detail.runs[0].id, run_id);
684        assert_eq!(detail.runs[0].status, RunStatus::Done);
685
686        // GET /v1/runs/:id returns the same Run directly.
687        let run = run_get(State(state.clone()), Path(run_id.to_string()))
688            .await
689            .expect("run_get")
690            .0;
691        assert_eq!(run.id, run_id);
692        assert_eq!(run.task_id, task_id);
693        assert_eq!(run.result_ref, Some(posted.final_ctx));
694
695        // issue #13 run_id propagation: `POST /v1/tasks` (`run_flow_form`)
696        // wires a `RunContext` into `TaskApplication::handle_with_run`, so
697        // the single dispatched step must be traced into `step_entries`.
698        assert_eq!(
699            run.step_entries.len(),
700            1,
701            "expected one step_entry for the 1-step identity Blueprint, got {:?}",
702            run.step_entries
703        );
704        assert_eq!(
705            run.step_entries[0].step_ref,
706            Some(mlua_swarm::worker::baseline::AG_IDENTITY.to_string())
707        );
708        assert_eq!(run.step_entries[0].status, Some("passed".to_string()));
709    }
710
711    // ──────────────────────────────────────────────────────────────────
712    // GH #33 — sync-hang guards (readiness precheck / timeout ceiling)
713    // ──────────────────────────────────────────────────────────────────
714
715    /// Same 1-step identity flow as [`identity_blueprint`], but opts into
716    /// the Blueprint-global Operator delegate axis
717    /// (`spawner_hints.layers = ["operator_delegate"]`) so a registered
718    /// `Operator` backend can be exercised end-to-end through the real
719    /// `tasks_start` dispatch path (`OperatorDelegateMiddleware` bypasses
720    /// `inner.spawn` and calls `operator.execute` instead — see
721    /// `mlua_swarm::middleware::OperatorDelegateMiddleware` doc).
722    fn identity_blueprint_with_operator_delegate() -> Blueprint {
723        Blueprint {
724            spawner_hints: mlua_swarm::SpawnerHints {
725                layers: vec!["operator_delegate".to_string()],
726            },
727            ..identity_blueprint()
728        }
729    }
730
731    /// `Operator` stub whose `execute` never resolves — the GH #33 Guard 2
732    /// fixture ("a registered-but-never-acking operator").
733    struct StallingOperator;
734
735    #[async_trait::async_trait]
736    impl mlua_swarm::Operator for StallingOperator {
737        async fn execute(
738            &self,
739            _ctx: &mlua_swarm::Ctx,
740            _system: Option<String>,
741            _prompt: Value,
742            _worker: Option<mlua_swarm::WorkerBinding>,
743            _worker_token: mlua_swarm::CapToken,
744        ) -> Result<mlua_swarm::WorkerResult, mlua_swarm::WorkerError> {
745            std::future::pending::<()>().await;
746            unreachable!("StallingOperator.execute must never resolve")
747        }
748    }
749
750    /// A launch request that references an operator backend by id (via
751    /// `operator.operator_backend_id`, the coarse Guard 1 signal) against
752    /// [`identity_blueprint_with_operator_delegate`].
753    fn operator_launch_req(
754        backend_id: &str,
755        timeout_secs: Option<u64>,
756    ) -> crate::TaskLaunchRequest {
757        crate::TaskLaunchRequest {
758            blueprint: BlueprintRef::Inline {
759                value: Box::new(identity_blueprint_with_operator_delegate()),
760            },
761            init_ctx: serde_json::json!({"in": "hello"}),
762            project_root: None,
763            work_dir: None,
764            task_metadata: None,
765            ttl_secs: None,
766            operator: Some(crate::OperatorReq {
767                operator_backend_id: Some(backend_id.to_string()),
768                ..Default::default()
769            }),
770            operator_sid: None,
771            timeout_secs,
772            goal: Some("operator delegate test goal".to_string()),
773            detach: false,
774        }
775    }
776
777    /// Guard 1: an operator-requiring launch with zero attached operators
778    /// must fail immediately with a structured `503`, not hang waiting on
779    /// a session nothing can serve.
780    #[tokio::test]
781    async fn sync_launch_zero_operators_fails_fast() {
782        let state = test_state();
783        // No `state.engine.register_operator(...)` call — zero operators
784        // attached, matching `list_operator_ids()` being empty.
785        let req = operator_launch_req("nonexistent-op", None);
786
787        let started = std::time::Instant::now();
788        let result = crate::tasks_start(State(state), Json(req)).await;
789        let elapsed = started.elapsed();
790
791        let err = match result {
792            Err(e) => e,
793            Ok(_) => panic!("zero attached operators must fail the operator-delegate launch"),
794        };
795        assert_eq!(err.status, StatusCode::SERVICE_UNAVAILABLE);
796        assert!(
797            err.message.contains("no operator attached"),
798            "error message must mention the missing operator: {}",
799            err.message
800        );
801        assert!(
802            elapsed < Duration::from_secs(1),
803            "guard 1 must fail fast (no dispatch, no timeout wait): took {elapsed:?}"
804        );
805    }
806
807    /// Guard 2: a launch that resolves to a registered-but-stalled
808    /// operator session must return a structured `504` within the
809    /// requested `timeout_secs` ceiling, not hang the request forever.
810    #[tokio::test]
811    async fn sync_launch_stalled_times_out() {
812        let state = test_state();
813        state
814            .engine
815            .register_operator("stall-op", Arc::new(StallingOperator))
816            .await;
817        let req = operator_launch_req("stall-op", Some(1));
818
819        let started = std::time::Instant::now();
820        // Outer safety-net timeout: if guard 2 itself regressed into an
821        // infinite hang, fail this test loudly instead of stalling `cargo
822        // test` indefinitely.
823        let result = tokio::time::timeout(
824            Duration::from_secs(5),
825            crate::tasks_start(State(state), Json(req)),
826        )
827        .await
828        .expect("tasks_start must resolve well within 5s when guard 2's ceiling is 1s");
829        let elapsed = started.elapsed();
830
831        let err = match result {
832            Err(e) => e,
833            Ok(_) => panic!("a stalled operator session must time out, not succeed"),
834        };
835        assert_eq!(err.status, StatusCode::GATEWAY_TIMEOUT);
836        assert!(
837            err.message.contains('1'),
838            "error message must mention the configured 1s ceiling: {}",
839            err.message
840        );
841        assert!(
842            elapsed < Duration::from_secs(3),
843            "guard 2 must fire close to the requested 1s ceiling: took {elapsed:?}"
844        );
845    }
846
847    /// Invariant 2: a launch that never references an operator backend
848    /// must never be rejected by guard 1 — the simplest existing passing
849    /// fixture (`post_tasks_req`) still succeeds unaffected.
850    #[tokio::test]
851    async fn sync_launch_without_operator_path_unaffected() {
852        let state = test_state();
853        let result = crate::tasks_start(
854            State(state),
855            Json(post_tasks_req("non-operator launch goal")),
856        )
857        .await;
858        if let Err(e) = &result {
859            panic!(
860                "non-operator launch must succeed unaffected by guard 1: {}",
861                e.message
862            );
863        }
864    }
865
866    /// Guard 2 ceiling resolution: `timeout_secs: Some(0)` is invalid
867    /// (design doc: "0 = reject with 400 or treat as invalid — pick one
868    /// and test it") — rejected fast, before any Task/Run side effects.
869    #[tokio::test]
870    async fn sync_launch_zero_timeout_secs_rejected() {
871        let state = test_state();
872        let mut req = post_tasks_req("zero timeout goal");
873        req.timeout_secs = Some(0);
874
875        let result = crate::tasks_start(State(state), Json(req)).await;
876        let err = match result {
877            Err(e) => e,
878            Ok(_) => panic!("timeout_secs: Some(0) must be rejected, not treated as a no-op"),
879        };
880        assert_eq!(err.status, StatusCode::BAD_REQUEST);
881        assert!(
882            err.message.contains("timeout_secs"),
883            "error message must reference timeout_secs: {}",
884            err.message
885        );
886    }
887
888    // ──────────────────────────────────────────────────────────────────
889    // GH #37 — detached launch / rekick (driver decoupled from request)
890    // ──────────────────────────────────────────────────────────────────
891
892    /// Polls the run store until the given Run reaches a terminal status,
893    /// panicking after ~5s — the detached paths complete in the
894    /// background, so tests must wait on the store rather than the
895    /// response.
896    async fn wait_for_terminal_run(state: &AppState, run_id: &RunId) -> RunRecord {
897        for _ in 0..50 {
898            let rec = state.run_store.get(run_id).await.expect("run get");
899            if !matches!(rec.status, RunStatus::Pending | RunStatus::Running) {
900                return rec;
901            }
902            tokio::time::sleep(Duration::from_millis(100)).await;
903        }
904        panic!("run {run_id} did not reach a terminal status within ~5s");
905    }
906
907    /// GH #37: `detach: true` returns `202 Accepted` immediately with
908    /// `status: "running"` and a null `final_ctx`; the eval completes in
909    /// the background and the Run/Task reach `Done` with the result and
910    /// step trace persisted — the same terminal state the sync path
911    /// produces.
912    #[tokio::test]
913    async fn detached_launch_returns_202_and_completes_in_background() {
914        let state = test_state();
915        let mut req = post_tasks_req("detached goal");
916        req.detach = true;
917
918        let reply = crate::tasks_start(State(state.clone()), Json(req))
919            .await
920            .expect("tasks_start (detached)");
921        assert_eq!(reply.1, StatusCode::ACCEPTED);
922        let posted = reply.0;
923        assert_eq!(posted.status, RunStatus::Running);
924        assert_eq!(
925            posted.final_ctx,
926            serde_json::Value::Null,
927            "a detached launch has no final_ctx at response time"
928        );
929
930        let rec = wait_for_terminal_run(&state, &posted.run_id).await;
931        assert_eq!(rec.status, RunStatus::Done);
932        assert!(
933            rec.result_ref.is_some(),
934            "finalize_run must persist the background eval's final_ctx"
935        );
936        assert_eq!(
937            rec.step_entries.len(),
938            1,
939            "the background eval must trace its step_entries like the sync path: {:?}",
940            rec.step_entries
941        );
942        let task = state
943            .task_store
944            .get(&posted.task_id)
945            .await
946            .expect("task get");
947        assert_eq!(task.status, TaskRecordStatus::Done);
948    }
949
950    /// GH #37: `detach: true` + `timeout_secs` is contradictory (the sync
951    /// ceiling has no meaning for a detached run) — rejected with `400`
952    /// before any Task/Run side effects.
953    #[tokio::test]
954    async fn detached_launch_with_timeout_secs_rejected() {
955        let state = test_state();
956        let mut req = post_tasks_req("detached + ceiling goal");
957        req.detach = true;
958        req.timeout_secs = Some(60);
959
960        let err = match crate::tasks_start(State(state.clone()), Json(req)).await {
961            Err(e) => e,
962            Ok(_) => panic!("detach + timeout_secs must be rejected"),
963        };
964        assert_eq!(err.status, StatusCode::BAD_REQUEST);
965        assert!(
966            err.message.contains("detach"),
967            "error message must explain the detach/timeout_secs conflict: {}",
968            err.message
969        );
970        let tasks = state.task_store.list().await.expect("task list");
971        assert!(
972            tasks.is_empty(),
973            "the 400 must fire before any TaskRecord is minted"
974        );
975    }
976
977    /// GH #37: a detached rekick returns `202 Accepted` with `status:
978    /// "running"` immediately and completes in the background, adding a
979    /// second `Done` Run to the same Task.
980    #[tokio::test]
981    async fn rekick_detached_returns_202_and_completes_in_background() {
982        let state = test_state();
983        let posted = crate::tasks_start(
984            State(state.clone()),
985            Json(post_tasks_req("detached rekick goal")),
986        )
987        .await
988        .expect("tasks_start")
989        .0;
990
991        let (status, rekicked) = task_rekick(
992            State(state.clone()),
993            Path(posted.task_id.to_string()),
994            Some(Json(RunKickRequest {
995                init_ctx_override: None,
996                task_input_override: None,
997                timeout_secs: None,
998                detach: true,
999            })),
1000        )
1001        .await
1002        .expect("task_rekick (detached)");
1003        assert_eq!(status, StatusCode::ACCEPTED);
1004        assert_eq!(rekicked.0.status, RunStatus::Running);
1005        assert_ne!(rekicked.0.run_id, posted.run_id);
1006
1007        let rec = wait_for_terminal_run(&state, &rekicked.0.run_id).await;
1008        assert_eq!(rec.status, RunStatus::Done);
1009        assert!(
1010            rec.result_ref.is_some(),
1011            "finalize_run must persist the background rekick's final_ctx"
1012        );
1013    }
1014
1015    /// GH #37: `detach: true` + `timeout_secs` on the rekick path is the
1016    /// same contradiction as on the launch path — `400`, no new Run
1017    /// minted.
1018    #[tokio::test]
1019    async fn rekick_detached_with_timeout_secs_rejected() {
1020        let state = test_state();
1021        let posted = crate::tasks_start(
1022            State(state.clone()),
1023            Json(post_tasks_req("detached rekick ceiling goal")),
1024        )
1025        .await
1026        .expect("tasks_start")
1027        .0;
1028
1029        let err = match task_rekick(
1030            State(state.clone()),
1031            Path(posted.task_id.to_string()),
1032            Some(Json(RunKickRequest {
1033                init_ctx_override: None,
1034                task_input_override: None,
1035                timeout_secs: Some(60),
1036                detach: true,
1037            })),
1038        )
1039        .await
1040        {
1041            Err(e) => e,
1042            Ok(_) => panic!("detach + timeout_secs must be rejected on rekick"),
1043        };
1044        assert_eq!(err.status, StatusCode::BAD_REQUEST);
1045        assert!(
1046            err.message.contains("detach"),
1047            "error message must explain the detach/timeout_secs conflict: {}",
1048            err.message
1049        );
1050        let runs = state
1051            .run_store
1052            .list_by_task(&posted.task_id)
1053            .await
1054            .expect("runs list");
1055        assert_eq!(
1056            runs.len(),
1057            1,
1058            "the 400 must fire before a second Run is minted"
1059        );
1060    }
1061
1062    #[tokio::test]
1063    async fn rekick_adds_a_second_run_to_the_same_task() {
1064        let state = test_state();
1065        let posted = crate::tasks_start(State(state.clone()), Json(post_tasks_req("rekick goal")))
1066            .await
1067            .expect("tasks_start")
1068            .0;
1069        let task_id = posted.task_id.clone();
1070        let first_run_id = posted.run_id.clone();
1071
1072        let (status, rekicked) = task_rekick(State(state.clone()), Path(task_id.to_string()), None)
1073            .await
1074            .expect("task_rekick");
1075        assert_eq!(status, StatusCode::CREATED);
1076        let second_run_id = rekicked.0.run_id.clone();
1077        assert_ne!(first_run_id, second_run_id);
1078
1079        let detail = task_get(State(state.clone()), Path(task_id.to_string()))
1080            .await
1081            .expect("task_get")
1082            .0;
1083        assert_eq!(
1084            detail.runs.len(),
1085            2,
1086            "expected 2 runs, got {:?}",
1087            detail.runs
1088        );
1089        let ids: Vec<&RunId> = detail.runs.iter().map(|r| &r.id).collect();
1090        assert!(ids.contains(&&first_run_id));
1091        assert!(ids.contains(&&second_run_id));
1092
1093        // issue #13 run_id propagation: each kick's own `EngineDispatcher`
1094        // (built fresh per `TaskApplication::handle_with_run` call) must
1095        // trace its own dispatched step into its own `RunRecord` —
1096        // independent `step_entries`, not shared/accumulated across kicks.
1097        let first_run = detail
1098            .runs
1099            .iter()
1100            .find(|r| r.id == first_run_id)
1101            .expect("first run present in detail.runs");
1102        let second_run = detail
1103            .runs
1104            .iter()
1105            .find(|r| r.id == second_run_id)
1106            .expect("second run present in detail.runs");
1107        assert_eq!(
1108            first_run.step_entries.len(),
1109            1,
1110            "first run step_entries: {:?}",
1111            first_run.step_entries
1112        );
1113        assert_eq!(
1114            second_run.step_entries.len(),
1115            1,
1116            "second run step_entries: {:?}",
1117            second_run.step_entries
1118        );
1119        assert_eq!(
1120            first_run.step_entries[0].step_ref,
1121            Some(mlua_swarm::worker::baseline::AG_IDENTITY.to_string())
1122        );
1123        assert_eq!(
1124            second_run.step_entries[0].step_ref,
1125            Some(mlua_swarm::worker::baseline::AG_IDENTITY.to_string())
1126        );
1127        assert_eq!(first_run.step_entries[0].status, Some("passed".to_string()));
1128        assert_eq!(
1129            second_run.step_entries[0].status,
1130            Some("passed".to_string())
1131        );
1132        assert_ne!(
1133            first_run.step_entries[0].step_id, second_run.step_entries[0].step_id,
1134            "each kick dispatches its own StepId — runs must not share step_entries"
1135        );
1136    }
1137
1138    #[tokio::test]
1139    async fn rekick_unknown_task_returns_404() {
1140        let state = test_state();
1141        // `.expect_err()` needs the Ok variant to be `Debug`; `Json<T>`'s
1142        // `Debug` impl is not guaranteed for every `T` across axum versions,
1143        // so a plain match sidesteps that bound entirely.
1144        match task_rekick(State(state), Path("T-does-not-exist".to_string()), None).await {
1145            Ok(_) => panic!("expected 404 for an unknown task"),
1146            Err(e) => assert_eq!(e.status, StatusCode::NOT_FOUND),
1147        }
1148    }
1149
1150    // ──────────────────────────────────────────────────────────────────
1151    // issue #19 ST4: `RunKickRequest` (optional body / 3-layer merge)
1152    // ──────────────────────────────────────────────────────────────────
1153
1154    /// A single-step flow.ir Blueprint that echoes `$.greeting` into
1155    /// `$.out` — unlike [`identity_blueprint`] (a fixed `lit("hello")`
1156    /// input), this one reads its `Step.in` from `ctx`, so it observes
1157    /// whichever `init_ctx` layer actually won the merge.
1158    fn greeting_blueprint() -> Blueprint {
1159        Blueprint {
1160            schema_version: current_schema_version(),
1161            id: "tasks-test-greeting-bp".into(),
1162            flow: serde_json::from_value(serde_json::json!({
1163                "kind": "step",
1164                "ref": mlua_swarm::worker::baseline::AG_IDENTITY,
1165                "in": {"op": "path", "at": "$.greeting"},
1166                "out": {"op": "path", "at": "$.out"},
1167            }))
1168            .expect("flow parse"),
1169            agents: vec![AgentDef {
1170                name: mlua_swarm::worker::baseline::AG_IDENTITY.into(),
1171                kind: AgentKind::RustFn,
1172                spec: serde_json::json!({"fn_id": mlua_swarm::worker::baseline::AG_IDENTITY}),
1173                profile: None,
1174                meta: None,
1175            }],
1176            operators: vec![],
1177            metas: vec![],
1178            hints: CompilerHints::default(),
1179            strategy: CompilerStrategy::default(),
1180            metadata: BlueprintMetadata::default(),
1181            spawner_hints: Default::default(),
1182            default_agent_kind: AgentKind::Operator,
1183            default_operator_kind: None,
1184            default_init_ctx: None,
1185            default_agent_ctx: None,
1186            default_context_policy: None,
1187            projection_placement: None,
1188            audits: vec![],
1189            degradation_policy: None,
1190        }
1191    }
1192
1193    fn post_greeting_task_req(
1194        greeting: &str,
1195        project_root: Option<&str>,
1196    ) -> crate::TaskLaunchRequest {
1197        crate::TaskLaunchRequest {
1198            blueprint: BlueprintRef::Inline {
1199                value: Box::new(greeting_blueprint()),
1200            },
1201            init_ctx: serde_json::json!({ "greeting": greeting }),
1202            project_root: project_root.map(str::to_string),
1203            work_dir: None,
1204            task_metadata: None,
1205            ttl_secs: None,
1206            operator: None,
1207            operator_sid: None,
1208            timeout_secs: None,
1209            goal: Some("st4 rekick goal".to_string()),
1210            detach: false,
1211        }
1212    }
1213
1214    #[tokio::test]
1215    async fn rekick_no_body_preserves_stored_task_input_ctx_byte_for_byte() {
1216        // must_not_simplify #3: a body-less rekick must behave exactly
1217        // like pre-#19 — the Task's own `input_ctx` alone seeds the kick.
1218        let state = test_state();
1219        let posted = crate::tasks_start(
1220            State(state.clone()),
1221            Json(post_greeting_task_req("from-task", None)),
1222        )
1223        .await
1224        .expect("tasks_start")
1225        .0;
1226        assert_eq!(posted.final_ctx["out"]["echoed"], "from-task");
1227
1228        let (status, rekicked) =
1229            task_rekick(State(state.clone()), Path(posted.task_id.to_string()), None)
1230                .await
1231                .expect("task_rekick");
1232        assert_eq!(status, StatusCode::CREATED);
1233
1234        let run = run_get(State(state.clone()), Path(rekicked.0.run_id.to_string()))
1235            .await
1236            .expect("run_get")
1237            .0;
1238        assert_eq!(
1239            run.result_ref.expect("result_ref present")["out"]["echoed"],
1240            "from-task"
1241        );
1242    }
1243
1244    #[tokio::test]
1245    async fn rekick_with_init_ctx_override_wins_over_stored_task_input_ctx() {
1246        let state = test_state();
1247        let posted = crate::tasks_start(
1248            State(state.clone()),
1249            Json(post_greeting_task_req("from-task", None)),
1250        )
1251        .await
1252        .expect("tasks_start")
1253        .0;
1254        assert_eq!(posted.final_ctx["out"]["echoed"], "from-task");
1255
1256        let (status, rekicked) = task_rekick(
1257            State(state.clone()),
1258            Path(posted.task_id.to_string()),
1259            Some(Json(RunKickRequest {
1260                init_ctx_override: Some(serde_json::json!({ "greeting": "from-run" })),
1261                task_input_override: None,
1262                timeout_secs: None,
1263                detach: false,
1264            })),
1265        )
1266        .await
1267        .expect("task_rekick");
1268        assert_eq!(status, StatusCode::CREATED);
1269
1270        let run = run_get(State(state.clone()), Path(rekicked.0.run_id.to_string()))
1271            .await
1272            .expect("run_get")
1273            .0;
1274        assert_eq!(
1275            run.result_ref.expect("result_ref present")["out"]["echoed"],
1276            "from-run",
1277            "Run's init_ctx_override must win over the stored Task input_ctx"
1278        );
1279    }
1280
1281    #[tokio::test]
1282    async fn rekick_with_stored_task_input_spec_dispatches_and_leaves_it_unmutated() {
1283        // Done Criteria: "Task record が task-level canonical fields を
1284        // 保持している時の rekick test". A Task created with
1285        // `project_root` set gets a `task_input_spec` snapshot; a
1286        // body-less rekick must both dispatch successfully (the stored
1287        // spec decodes and resolves without erroring) and leave
1288        // `TaskRecord.task_input_spec` untouched (must_not_simplify #4 —
1289        // a rekick never mutates the stored Task-level snapshot).
1290        let state = test_state();
1291        let posted = crate::tasks_start(
1292            State(state.clone()),
1293            Json(post_greeting_task_req("from-task", Some("/repo"))),
1294        )
1295        .await
1296        .expect("tasks_start")
1297        .0;
1298
1299        let before = state
1300            .task_store
1301            .get(&posted.task_id)
1302            .await
1303            .expect("task fetch");
1304        let before_spec: Option<TaskInputSpec> = before
1305            .task_input_spec
1306            .as_ref()
1307            .map(|v| serde_json::from_value(v.clone()).expect("decode task_input_spec"));
1308        assert_eq!(
1309            before_spec,
1310            Some(TaskInputSpec {
1311                project_root: Some("/repo".to_string()),
1312                work_dir: None,
1313                task_metadata: None,
1314            })
1315        );
1316
1317        let (status, _rekicked) =
1318            task_rekick(State(state.clone()), Path(posted.task_id.to_string()), None)
1319                .await
1320                .expect("task_rekick");
1321        assert_eq!(status, StatusCode::CREATED);
1322
1323        let after = state
1324            .task_store
1325            .get(&posted.task_id)
1326            .await
1327            .expect("task fetch");
1328        assert_eq!(
1329            after.task_input_spec, before.task_input_spec,
1330            "rekick must not mutate the stored Task-level task_input_spec snapshot"
1331        );
1332    }
1333
1334    #[tokio::test]
1335    async fn rekick_with_task_input_override_does_not_mutate_stored_task_record() {
1336        // must_not_simplify #4: `task_input_override` wins for this kick
1337        // only — the stored `TaskRecord.task_input_spec` is untouched.
1338        let state = test_state();
1339        let posted = crate::tasks_start(
1340            State(state.clone()),
1341            Json(post_greeting_task_req("from-task", Some("/repo"))),
1342        )
1343        .await
1344        .expect("tasks_start")
1345        .0;
1346
1347        let (status, _rekicked) = task_rekick(
1348            State(state.clone()),
1349            Path(posted.task_id.to_string()),
1350            Some(Json(RunKickRequest {
1351                init_ctx_override: None,
1352                task_input_override: Some(TaskInputSpec {
1353                    project_root: Some("/override".to_string()),
1354                    work_dir: None,
1355                    task_metadata: None,
1356                }),
1357                timeout_secs: None,
1358                detach: false,
1359            })),
1360        )
1361        .await
1362        .expect("task_rekick");
1363        assert_eq!(status, StatusCode::CREATED);
1364
1365        let after = state
1366            .task_store
1367            .get(&posted.task_id)
1368            .await
1369            .expect("task fetch");
1370        let after_spec: Option<TaskInputSpec> = after
1371            .task_input_spec
1372            .as_ref()
1373            .map(|v| serde_json::from_value(v.clone()).expect("decode task_input_spec"));
1374        assert_eq!(
1375            after_spec,
1376            Some(TaskInputSpec {
1377                project_root: Some("/repo".to_string()),
1378                work_dir: None,
1379                task_metadata: None,
1380            }),
1381            "a per-Run task_input_override must not leak into the stored TaskRecord"
1382        );
1383    }
1384
1385    // ──────────────────────────────────────────────────────────────────
1386    // GH #33 → task_rekick — sync-hang guards (issue #35 ST3 parity)
1387    // ──────────────────────────────────────────────────────────────────
1388
1389    /// A launch request for [`identity_blueprint_with_operator_delegate`]
1390    /// that does **not** reference an operator backend (`operator: None`)
1391    /// — used to create a rekick-able Task without tripping
1392    /// `run_flow_form`'s own Guard 1 at initial-launch time (the launch
1393    /// itself dispatches through the plain baseline path since
1394    /// `ctx.operator.operator` stays unset either way; the BP's
1395    /// `operator_delegate` layer only matters to `task_rekick`'s Guard 1,
1396    /// which reads `resolved_bp.spawner_hints.layers` directly rather than
1397    /// a per-request field).
1398    fn delegate_launch_req(goal: &str) -> crate::TaskLaunchRequest {
1399        crate::TaskLaunchRequest {
1400            blueprint: BlueprintRef::Inline {
1401                value: Box::new(identity_blueprint_with_operator_delegate()),
1402            },
1403            init_ctx: serde_json::json!({"in": "hello"}),
1404            project_root: None,
1405            work_dir: None,
1406            task_metadata: None,
1407            ttl_secs: None,
1408            operator: None,
1409            operator_sid: None,
1410            timeout_secs: None,
1411            goal: Some(goal.to_string()),
1412            detach: false,
1413        }
1414    }
1415
1416    /// Guard 1 (adapted signal): a Task whose stored Blueprint declares
1417    /// the `operator_delegate` layer, rekicked with zero attached
1418    /// operators, must fail immediately with a structured `503` — not
1419    /// dispatch and not hang waiting on a session nothing can serve.
1420    #[tokio::test]
1421    async fn rekick_zero_operators_with_operator_delegate_blueprint_fails_fast() {
1422        let state = test_state();
1423        let posted = crate::tasks_start(
1424            State(state.clone()),
1425            Json(delegate_launch_req("operator delegate rekick goal")),
1426        )
1427        .await
1428        .expect("tasks_start (no operator referenced, dispatches through baseline)")
1429        .0;
1430        // No `state.engine.register_operator(...)` call — zero operators
1431        // attached, matching `list_operator_ids()` being empty.
1432
1433        let started = std::time::Instant::now();
1434        let result = task_rekick(State(state), Path(posted.task_id.to_string()), None).await;
1435        let elapsed = started.elapsed();
1436
1437        let err = match result {
1438            Err(e) => e,
1439            Ok(_) => panic!(
1440                "rekicking a Task whose Blueprint declares operator_delegate with zero \
1441                 attached operators must fail, not dispatch"
1442            ),
1443        };
1444        assert_eq!(err.status, StatusCode::SERVICE_UNAVAILABLE);
1445        assert!(
1446            err.message.contains("no operator attached"),
1447            "error message must mention the missing operator: {}",
1448            err.message
1449        );
1450        assert!(
1451            elapsed < Duration::from_secs(1),
1452            "guard 1 must fail fast (no dispatch, no timeout wait): took {elapsed:?}"
1453        );
1454    }
1455
1456    /// Guard 2: a rekick with a `timeout_secs` ceiling shorter than the
1457    /// dispatch takes must return a structured `504` within the outer
1458    /// safety-net timeout, not hang the request forever.
1459    #[tokio::test]
1460    async fn rekick_stalled_operator_times_out() {
1461        let state = test_state();
1462        state
1463            .engine
1464            .register_operator("stall-op", Arc::new(StallingOperator))
1465            .await;
1466        let posted = crate::tasks_start(
1467            State(state.clone()),
1468            Json(delegate_launch_req("stalled rekick goal")),
1469        )
1470        .await
1471        .expect("tasks_start")
1472        .0;
1473
1474        let started = std::time::Instant::now();
1475        // Outer safety-net timeout: if guard 2 itself regressed into an
1476        // infinite hang, fail this test loudly instead of stalling `cargo
1477        // test` indefinitely.
1478        let result = tokio::time::timeout(
1479            Duration::from_secs(5),
1480            task_rekick(
1481                State(state),
1482                Path(posted.task_id.to_string()),
1483                Some(Json(RunKickRequest {
1484                    init_ctx_override: None,
1485                    task_input_override: None,
1486                    timeout_secs: Some(1),
1487                    detach: false,
1488                })),
1489            ),
1490        )
1491        .await
1492        .expect("task_rekick must resolve well within 5s when guard 2's ceiling is 1s");
1493        let elapsed = started.elapsed();
1494
1495        match &result {
1496            Err(e) => {
1497                assert_eq!(e.status, StatusCode::GATEWAY_TIMEOUT);
1498                assert!(
1499                    e.message.contains('1'),
1500                    "error message must mention the configured 1s ceiling: {}",
1501                    e.message
1502                );
1503                assert!(
1504                    elapsed < Duration::from_secs(3),
1505                    "guard 2 must fire close to the requested 1s ceiling: took {elapsed:?}"
1506                );
1507            }
1508            Ok(_) => {
1509                // `task_rekick` hardcodes `operator_backend_id: None` for
1510                // every kick (module doc, above — "no per-request Operator
1511                // override support here"), so a registered-but-unattached
1512                // `StallingOperator` is never actually engaged by a
1513                // rekick's dispatch; the flow resolves through the plain
1514                // baseline path instead. Guard 2's `tokio::time::timeout`
1515                // wrap is exercised (and does not falsely fire) rather
1516                // than tripped — assert the fast-success shape so a
1517                // regression that makes rekick dispatch slow (or that
1518                // makes Guard 2 falsely trip on a fast dispatch) is still
1519                // caught by the elapsed-time assertion below.
1520                assert!(
1521                    elapsed < Duration::from_secs(1),
1522                    "a rekick that never engages an Operator (task_rekick has no \
1523                     per-request operator override) must resolve fast, not stall: took {elapsed:?}"
1524                );
1525            }
1526        }
1527    }
1528
1529    /// Guard 2 ceiling resolution: `timeout_secs: Some(0)` is invalid —
1530    /// rejected fast, before any Task/Run side effects (the pre-existing
1531    /// run count for the rekicked Task is unchanged).
1532    #[tokio::test]
1533    async fn rekick_timeout_secs_zero_rejected() {
1534        let state = test_state();
1535        let posted = crate::tasks_start(
1536            State(state.clone()),
1537            Json(post_tasks_req("zero timeout rekick goal")),
1538        )
1539        .await
1540        .expect("tasks_start")
1541        .0;
1542
1543        let before = task_get(State(state.clone()), Path(posted.task_id.to_string()))
1544            .await
1545            .expect("task_get")
1546            .0;
1547        let runs_before = before.runs.len();
1548
1549        let result = task_rekick(
1550            State(state.clone()),
1551            Path(posted.task_id.to_string()),
1552            Some(Json(RunKickRequest {
1553                init_ctx_override: None,
1554                task_input_override: None,
1555                timeout_secs: Some(0),
1556                detach: false,
1557            })),
1558        )
1559        .await;
1560        let err = match result {
1561            Err(e) => e,
1562            Ok(_) => panic!("timeout_secs: Some(0) must be rejected, not treated as a no-op"),
1563        };
1564        assert_eq!(err.status, StatusCode::BAD_REQUEST);
1565        assert!(
1566            err.message.contains("timeout_secs"),
1567            "error message must reference timeout_secs: {}",
1568            err.message
1569        );
1570
1571        let after = task_get(State(state), Path(posted.task_id.to_string()))
1572            .await
1573            .expect("task_get")
1574            .0;
1575        assert_eq!(
1576            after.runs.len(),
1577            runs_before,
1578            "a rejected timeout_secs: Some(0) rekick must not create a new Run"
1579        );
1580    }
1581
1582    /// Invariant: a plain (non-`operator_delegate`) Task rekick must
1583    /// never be rejected by Guard 1 — the simplest existing passing
1584    /// rekick fixture still succeeds unaffected.
1585    #[tokio::test]
1586    async fn rekick_non_operator_path_unaffected_by_guard_1() {
1587        let state = test_state();
1588        let posted = crate::tasks_start(
1589            State(state.clone()),
1590            Json(post_tasks_req("non-operator rekick goal")),
1591        )
1592        .await
1593        .expect("tasks_start")
1594        .0;
1595
1596        let result = task_rekick(State(state), Path(posted.task_id.to_string()), None).await;
1597        if let Err(e) = &result {
1598            panic!(
1599                "a plain (non-operator_delegate) Task rekick must succeed unaffected by \
1600                 guard 1: {}",
1601                e.message
1602            );
1603        }
1604    }
1605
1606    #[tokio::test]
1607    async fn run_get_unknown_id_returns_404() {
1608        let state = test_state();
1609        match run_get(State(state), Path("R-does-not-exist".to_string())).await {
1610            Ok(_) => panic!("expected 404 for an unknown run"),
1611            Err(e) => assert_eq!(e.status, StatusCode::NOT_FOUND),
1612        }
1613    }
1614
1615    #[tokio::test]
1616    async fn task_get_unknown_id_returns_404() {
1617        let state = test_state();
1618        match task_get(State(state), Path("T-does-not-exist".to_string())).await {
1619            Ok(_) => panic!("expected 404 for an unknown task"),
1620            Err(e) => assert_eq!(e.status, StatusCode::NOT_FOUND),
1621        }
1622    }
1623}