Skip to main content

kranz_server/
host.rs

1//! Hosted-engine registry (docs/protocol.md "Mission lifecycle
2//! (server-hosted engine; M2.5)").
3//!
4//! `kranz serve` can HOST missions: for missions created via
5//! `POST /api/missions` this server process IS the single-writer engine — it
6//! holds the mission lock, so a concurrent `kranz run` correctly refuses, and
7//! either side can resume what the other started (the event log is the source
8//! of truth).
9//!
10//! Concurrency model:
11//! - Each planning-phase mission sits behind an `Arc<tokio::sync::Mutex<..>>`
12//!   so planning turns serialize per mission; handlers `try_lock` and a
13//!   contended lock is a 409 ("a turn is in flight"), never a queue.
14//! - `start` consumes the engine out of the registry (`Arc::try_unwrap`
15//!   succeeds only when no turn holds a clone) and spawns `engine.run()` as a
16//!   background task. When the run ends — Complete, Blocked or Failed — the
17//!   task drops the engine (flushing the log and releasing the single-writer
18//!   lock) and removes its registry entry, so the mission is observable and
19//!   resumable from anywhere.
20//! - `start` on a mission NOT in the registry (blocked earlier, or the server
21//!   restarted) resumes it from the event log — the re-invocable semantics of
22//!   the protocol.
23//!
24//! The agent backend is constructed lazily on first use, so a read-only
25//! `kranz serve` never needs a `claude` binary installed.
26
27use crate::error::{ApiError, ApiErrorCode};
28use crate::ServerState;
29use axum::body::Bytes;
30use axum::extract::{Path as UrlPath, State};
31use axum::http::StatusCode;
32use axum::response::IntoResponse;
33use axum::Json;
34use kranz_engine::backend::{AgentBackend, AgentEvent, PromptMode, SessionExit, SessionSpec};
35use kranz_engine::backend_claude::ClaudeBackend;
36use kranz_engine::config;
37use kranz_engine::cost::{self, CostEstimate};
38use kranz_engine::deps;
39use kranz_engine::draft::{drive_draft, DraftOutcome};
40use kranz_engine::error::EngineError;
41use kranz_engine::event_log::{EventLog, LockForce};
42use kranz_engine::git_ops::GitRepo;
43use kranz_engine::git_ops::KranzCommitMetadata;
44use kranz_engine::merge::{
45    merge_mission_with_standards_evidence, MergeReport, StandardsMergeEvidence,
46};
47use kranz_engine::orchestrator::{MissionEngine, PlanRequest};
48use kranz_engine::paths::MissionPaths;
49use kranz_engine::planning::plan_identity;
50use kranz_engine::queue;
51use kranz_engine::ticket::Ticket;
52use kranz_engine::types::{MissionConfig, MissionStatus, Plan, TokenUsage};
53use serde_json::{json, Value};
54use std::collections::HashMap;
55use std::path::{Path, PathBuf};
56use std::sync::{Arc, Mutex};
57use std::time::{Duration, Instant};
58use tokio::sync::{OwnedSemaphorePermit, Semaphore};
59
60/// The engine cell of a planning-phase mission: turns lock it, `start`
61/// consumes it.
62type EngineCell = Arc<tokio::sync::Mutex<Box<MissionEngine>>>;
63
64/// One mission hosted by this server process.
65enum HostedMission {
66    /// In planning (or approved, awaiting start): the live engine, holding
67    /// the mission lock and the orchestrator conversation, plus when it was
68    /// last touched by a planning turn (for the idle sweeper).
69    Planning {
70        cell: EngineCell,
71        last_use: Arc<Mutex<Instant>>,
72        /// The last plan `request_plan` returned Ready, awaiting approval —
73        /// ONE cache for every surface's approve affordance (Slack buttons,
74        /// web, glasses ring). Consumed by [`MissionHost::approve_pending`];
75        /// volatile by design (a restart forfeits it — re-request the plan).
76        pending_plan: Arc<Mutex<Option<Plan>>>,
77    },
78    /// `engine.run()` owns the engine inside this background task; the task
79    /// removes this entry when the run ends. `_repo_busy` holds the
80    /// repo-wide busy lock for the lifetime of the hosted run so a sibling
81    /// queue drain / `kranz work` cannot claim the same repo.
82    Running {
83        handle: tokio::task::JoinHandle<()>,
84        _repo_busy: kranz_engine::queue::RepoBusyHold,
85    },
86}
87
88/// What [`MissionHost::try_approve_pending_matching`] did.
89///
90/// `Mismatch` is deliberately not an error: the caller renders an
91/// "awaiting X, not Y" refusal naming both plans, which tells the reviewer
92/// what happened rather than handing them a status code.
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub enum PendingApproval {
95    /// The parked plan was the expected one and is now committed on the
96    /// mission branch; carries the mission branch name.
97    Approved(String),
98    /// Nothing was parked: never requested, or forfeited by a serve restart
99    /// or an idle release. The caller's own state-aware routing takes over.
100    NothingParked,
101    /// A plan IS parked and it is not the one the caller reviewed. Carries
102    /// the PARKED plan's identity so the refusal can name both.
103    Mismatch { parked: String },
104}
105
106/// Registry of missions this server process hosts (see module docs).
107pub struct MissionHost {
108    repo_root: PathBuf,
109    /// Lazy real backend — discovered on first mutating use, so read-only
110    /// serving works without a `claude` binary. Tests inject a mock via
111    /// [`MissionHost::with_backend`].
112    backend: tokio::sync::OnceCell<Arc<dyn AgentBackend>>,
113    /// Shared with each run task so it can remove its own entry on exit.
114    missions: Arc<Mutex<HashMap<String, HostedMission>>>,
115    /// The lazily-spawned idle-release background task, started at most once
116    /// (see [`MissionHost::ensure_sweeper_started`]).
117    sweeper: Mutex<Option<tokio::task::JoinHandle<()>>>,
118    /// The single tracked background queue drain slot (see
119    /// [`MissionHost::drain`]).
120    drain: Mutex<DrainSlot>,
121    /// Shared by every repository in a [`crate::MultiRepoHost`]. A permit is
122    /// held for the complete background mission/drain lifetime, so the
123    /// operator's `host.maxConcurrentRepos` is a real spend/load bound.
124    global_run_permits: Option<Arc<Semaphore>>,
125    /// The gate-suite executor [`MissionHost::merge`] runs under
126    /// `spawn_blocking`; real shell commands by default, a scripted stub in
127    /// tests (see [`MissionHost::with_gate_executor`]).
128    gate_executor: GateExecutor,
129    /// Short-TTL cache for the queue-front readiness probe so a 3s dashboard
130    /// poll does not re-shell every backend CLI on every GET /api/queue.
131    readiness_front_cache: Mutex<Option<FrontReadinessCache>>,
132    /// Backend readiness follows the same dependency-injection boundary as
133    /// `backend`: real hosts probe configured CLIs, while hosts supplied an
134    /// already-constructed backend treat that backend as available.
135    readiness_probe: ReadinessProbe,
136}
137
138/// Cached `GET /api/queue` readiness for the current queue front only.
139struct FrontReadinessCache {
140    mission_id: String,
141    report: Value,
142    at: Instant,
143}
144
145const READINESS_FRONT_CACHE_TTL: Duration = Duration::from_secs(5);
146
147/// One background drain task's observable progress — shared between the task
148/// (which updates it as it goes) and [`MissionHost::drain`] /
149/// [`MissionHost::queue_state`] (which read it back as JSON).
150#[derive(Debug, Clone, Default)]
151struct DrainState {
152    live: bool,
153    current_mission_id: Option<String>,
154    ran: Vec<String>,
155    parked: Vec<String>,
156}
157
158/// One gate-suite command execution: `executor(command, cwd)` →
159/// `(success, combined_stdout_stderr)`. Boxed so [`MissionHost`] can hold a
160/// real shell-backed default and tests can inject a scripted stub — the same
161/// seam shape as [`MissionHost::with_backend`] for the agent backend.
162type GateExecutor = Arc<dyn Fn(&str, &Path) -> (bool, String) + Send + Sync>;
163
164type ReadinessProbe =
165    fn(
166        &Path,
167        &str,
168    ) -> kranz_engine::error::Result<kranz_engine::backend_readiness::ReadinessReport>;
169
170fn injected_backend_readiness(
171    _repo_root: &Path,
172    mission_id: &str,
173) -> kranz_engine::error::Result<kranz_engine::backend_readiness::ReadinessReport> {
174    Ok(kranz_engine::backend_readiness::ReadinessReport {
175        mission_id: mission_id.to_string(),
176        roles: Vec::new(),
177        overall: kranz_engine::backend_readiness::ReadinessStatus::Ok,
178        warnings: Vec::new(),
179    })
180}
181
182/// The real gate executor delegates to the engine's 600-second process-tree
183/// bounded shell runner with a sanitized environment. It runs only from
184/// inside `tokio::task::spawn_blocking` (see [`MissionHost::merge`]).
185fn real_gate_executor() -> GateExecutor {
186    Arc::new(|command, cwd| kranz_engine::command_exec::run_bounded_gate_command(cwd, command))
187}
188
189/// The autoWork watcher's decision function, factored out so it's testable
190/// without standing up a full mission: drain only when autoWork is enabled,
191/// the queue has something waiting, and no drain is already live.
192fn should_auto_drain(auto_work: bool, queue_non_empty: bool, drain_live: bool) -> bool {
193    auto_work && queue_non_empty && !drain_live
194}
195
196fn drain_state_json(state: &DrainState) -> Value {
197    json!({
198        "live": state.live,
199        "currentMissionId": state.current_mission_id,
200        "ran": state.ran,
201        "parked": state.parked,
202    })
203}
204
205/// A tracked background drain: the task handle plus the state it shares with
206/// this host. `join.is_finished()` is how [`MissionHost::drain`] decides
207/// whether a tracked drain is still live.
208struct DrainHandle {
209    join: tokio::task::JoinHandle<()>,
210    state: Arc<Mutex<DrainState>>,
211}
212
213/// The drain tracker's state machine. `Starting` is a reservation held while
214/// `config::load` + `self.backend(...)` run with NO lock held (both can
215/// `.await`); it closes the race where two concurrent [`MissionHost::drain`]
216/// calls both observe "nothing tracked yet" and both spawn a task. A racing
217/// caller that sees `Starting` returns its shared [`DrainState`] instead of
218/// starting a second drain; the caller that installed the reservation later
219/// upgrades it to `Running` (same `Arc<Mutex<DrainState>>`), or clears it back
220/// to `Idle` on failure so a later call can retry.
221enum DrainSlot {
222    Idle,
223    Starting(Arc<Mutex<DrainState>>),
224    Running(DrainHandle),
225}
226
227impl MissionHost {
228    /// Host for `repo_root`, discovering the Claude backend on first use.
229    pub fn new(repo_root: PathBuf) -> Self {
230        MissionHost {
231            repo_root,
232            backend: tokio::sync::OnceCell::new(),
233            missions: Arc::new(Mutex::new(HashMap::new())),
234            sweeper: Mutex::new(None),
235            drain: Mutex::new(DrainSlot::Idle),
236            global_run_permits: None,
237            gate_executor: real_gate_executor(),
238            readiness_front_cache: Mutex::new(None),
239            readiness_probe: kranz_engine::backend_readiness::probe_mission,
240        }
241    }
242
243    /// Host with an injected backend (tests drive the full lifecycle through
244    /// `kranz_engine::backend_mock` without a `claude` binary).
245    pub fn with_backend(repo_root: PathBuf, backend: Arc<dyn AgentBackend>) -> Self {
246        MissionHost {
247            repo_root,
248            backend: tokio::sync::OnceCell::new_with(Some(backend)),
249            missions: Arc::new(Mutex::new(HashMap::new())),
250            sweeper: Mutex::new(None),
251            drain: Mutex::new(DrainSlot::Idle),
252            global_run_permits: None,
253            gate_executor: real_gate_executor(),
254            readiness_front_cache: Mutex::new(None),
255            readiness_probe: injected_backend_readiness,
256        }
257    }
258
259    /// Host with an injected gate-suite executor (tests script CI gate
260    /// outcomes for [`MissionHost::merge`] hermetically, without ever
261    /// shelling out to `cargo`/`npm`). Mirrors [`MissionHost::with_backend`]'s
262    /// seam, for the gate suite instead of the agent backend.
263    pub fn with_gate_executor<F>(repo_root: PathBuf, gate_executor: F) -> Self
264    where
265        F: Fn(&str, &Path) -> (bool, String) + Send + Sync + 'static,
266    {
267        MissionHost {
268            repo_root,
269            backend: tokio::sync::OnceCell::new(),
270            missions: Arc::new(Mutex::new(HashMap::new())),
271            sweeper: Mutex::new(None),
272            drain: Mutex::new(DrainSlot::Idle),
273            global_run_permits: None,
274            gate_executor: Arc::new(gate_executor),
275            readiness_front_cache: Mutex::new(None),
276            readiness_probe: kranz_engine::backend_readiness::probe_mission,
277        }
278    }
279
280    /// Host participating in a process-wide multi-repository execution cap.
281    pub(crate) fn new_with_global_run_permits(
282        repo_root: PathBuf,
283        global_run_permits: Arc<Semaphore>,
284    ) -> Self {
285        MissionHost {
286            repo_root,
287            backend: tokio::sync::OnceCell::new(),
288            missions: Arc::new(Mutex::new(HashMap::new())),
289            sweeper: Mutex::new(None),
290            drain: Mutex::new(DrainSlot::Idle),
291            global_run_permits: Some(global_run_permits),
292            gate_executor: real_gate_executor(),
293            readiness_front_cache: Mutex::new(None),
294            readiness_probe: kranz_engine::backend_readiness::probe_mission,
295        }
296    }
297
298    /// Test helper: inject a backend into a host that already shares the
299    /// multi-repository run semaphore.
300    #[cfg(test)]
301    pub(crate) fn with_backend_and_global_run_permits(
302        repo_root: PathBuf,
303        backend: Arc<dyn AgentBackend>,
304        global_run_permits: Arc<Semaphore>,
305    ) -> Self {
306        MissionHost {
307            repo_root,
308            backend: tokio::sync::OnceCell::new_with(Some(backend)),
309            missions: Arc::new(Mutex::new(HashMap::new())),
310            sweeper: Mutex::new(None),
311            drain: Mutex::new(DrainSlot::Idle),
312            global_run_permits: Some(global_run_permits),
313            gate_executor: real_gate_executor(),
314            readiness_front_cache: Mutex::new(None),
315            readiness_probe: injected_backend_readiness,
316        }
317    }
318
319    /// The repository this host creates missions in.
320    pub fn repo_root(&self) -> &PathBuf {
321        &self.repo_root
322    }
323
324    pub(crate) fn try_global_run_permit(&self) -> Result<Option<OwnedSemaphorePermit>, ApiError> {
325        self.global_run_permits
326            .as_ref()
327            .map(|permits| {
328                Arc::clone(permits).try_acquire_owned().map_err(|_| {
329                    ApiError::conflict(
330                        "host.maxConcurrentRepos is saturated; retry when another repository finishes",
331                    )
332                    .with_code(ApiErrorCode::RepositoryBusy)
333                })
334            })
335            .transpose()
336    }
337
338    /// The backend, constructing [`ClaudeBackend`] on first use.
339    async fn backend(
340        &self,
341        claude_binary: Option<&str>,
342    ) -> Result<Arc<dyn AgentBackend>, ApiError> {
343        let configured = claude_binary.map(str::to_string);
344        self.backend
345            .get_or_try_init(|| async move {
346                let backend = ClaudeBackend::discover(configured.as_deref())?;
347                Ok::<Arc<dyn AgentBackend>, EngineError>(Arc::new(backend))
348            })
349            .await
350            .map(Arc::clone)
351            .map_err(ApiError::from)
352    }
353
354    // -----------------------------------------------------------------------
355    // Lifecycle operations (one per endpoint)
356    // -----------------------------------------------------------------------
357
358    /// `POST /api/missions`: layered config + optional request patch →
359    /// validate → create the mission → hold its engine in the registry.
360    /// Public: the Slack bridge drives the same lifecycle through this host
361    /// (wired by `kranz serve --slack`), so these five operations are the
362    /// shared client surface, not axum-private plumbing.
363    pub async fn create(
364        &self,
365        goal: &str,
366        config_patch: Option<&Value>,
367    ) -> Result<String, ApiError> {
368        let mut cfg = config::load(&self.repo_root)?;
369        if let Some(patch) = config_patch {
370            if !patch.is_object() {
371                return Err(ApiError::bad_request("'config' must be a JSON object"));
372            }
373            let mut merged = serde_json::to_value(&cfg)
374                .map_err(|e| ApiError::internal(format!("config does not serialize: {e}")))?;
375            config::deep_merge(&mut merged, patch);
376            cfg = serde_json::from_value(merged).map_err(|e| {
377                ApiError::bad_request(format!("'config' patch does not deserialize: {e}"))
378            })?;
379        }
380        config::validate(&cfg)?;
381
382        let backend = self.backend(cfg.claude_binary.as_deref()).await?;
383        let engine = MissionEngine::create(backend, self.repo_root.clone(), goal, cfg)?;
384        let id = engine.mission_id().to_string();
385        self.missions
386            .lock()
387            .expect("missions registry lock")
388            .insert(id.clone(), new_planning(new_cell(Box::new(engine))));
389        self.ensure_sweeper_started();
390        Ok(id)
391    }
392
393    /// Entry point any surface (REST, Slack, CLI-over-HTTP) can call to draft
394    /// a backlog ticket non-interactively: validate the slug, load the ticket,
395    /// create its planning mission through this host — so the create path
396    /// registers it in `missions` and its lifecycle events stream over
397    /// `GET /api/missions/:id/ws` exactly like `POST /api/missions` — then run
398    /// [`drive_draft`] (roadmap f-1-1) to completion against that hosted
399    /// engine. The engine is dropped and the registry entry removed once the
400    /// draft turn ends (mirroring [`run_to_end`]'s drop-then-remove ordering)
401    /// so the mission stays observable/resumable afterward; no operator
402    /// checkout restoration happens here — a headless server has no checkout
403    /// to restore.
404    pub async fn draft(&self, slug: &str, then_enqueue: bool) -> Result<DraftOutcome, ApiError> {
405        Ticket::ensure_valid_slug(slug)?;
406        let ticket_path = Ticket::tickets_dir(&self.repo_root).join(format!("{slug}.md"));
407        if !ticket_path.is_file() {
408            return Err(ApiError::not_found(format!("ticket '{slug}' not found")));
409        }
410        let ticket = Ticket::load(&ticket_path)?;
411
412        let cfg = config_for_ticket(config::load(&self.repo_root)?, &ticket);
413        let backend = self.backend(cfg.claude_binary.as_deref()).await?;
414        let engine =
415            MissionEngine::create(backend, self.repo_root.clone(), &ticket.mission_goal(), cfg)?;
416        let id = engine.mission_id().to_string();
417        let cell = new_cell(Box::new(engine));
418        self.missions
419            .lock()
420            .expect("missions registry lock")
421            .insert(id.clone(), new_planning(Arc::clone(&cell)));
422        self.ensure_sweeper_started();
423
424        let drive_result = {
425            let mut engine = cell.lock().await;
426            drive_draft(&mut engine, &self.repo_root, &ticket, then_enqueue).await
427        };
428
429        // Drop the engine (flushes the log, frees the single-writer lock),
430        // then remove the registry entry — from that moment the mission is
431        // observable and resumable anywhere, same as `run_to_end`.
432        self.missions
433            .lock()
434            .expect("missions registry lock")
435            .remove(&id);
436        drop(cell);
437
438        Ok(drive_result?.outcome)
439    }
440
441    /// `POST /api/tickets/:slug/draft`: fire-and-observe twin of
442    /// [`Self::draft`] for the REST surface — a draft can run for a while (a
443    /// planning conversation with the orchestrator), so this creates the
444    /// planning mission SYNCHRONOUSLY (registering it exactly like `create`,
445    /// so its lifecycle streams over `GET /api/missions/:id/ws` immediately),
446    /// then spawns [`drive_draft`] as a background task and returns the
447    /// mission id right away. The final outcome (Review vs NeedsContext) is
448    /// read back later via `GET /api/tickets/:slug`.
449    pub async fn draft_async(&self, slug: &str, then_enqueue: bool) -> Result<String, ApiError> {
450        Ticket::ensure_valid_slug(slug)?;
451        let ticket_path = Ticket::tickets_dir(&self.repo_root).join(format!("{slug}.md"));
452        if !ticket_path.is_file() {
453            return Err(ApiError::not_found(format!("ticket '{slug}' not found")));
454        }
455        let ticket = Ticket::load(&ticket_path)?;
456
457        let cfg = config_for_ticket(config::load(&self.repo_root)?, &ticket);
458        let backend = self.backend(cfg.claude_binary.as_deref()).await?;
459        let engine =
460            MissionEngine::create(backend, self.repo_root.clone(), &ticket.mission_goal(), cfg)?;
461        let id = engine.mission_id().to_string();
462        let cell = new_cell(Box::new(engine));
463        self.missions
464            .lock()
465            .expect("missions registry lock")
466            .insert(id.clone(), new_planning(Arc::clone(&cell)));
467        self.ensure_sweeper_started();
468
469        let repo_root = self.repo_root.clone();
470        let missions = Arc::clone(&self.missions);
471        let mission_id = id.clone();
472        tokio::spawn(async move {
473            let drive_result = {
474                let mut engine = cell.lock().await;
475                drive_draft(&mut engine, &repo_root, &ticket, then_enqueue).await
476            };
477            // Same drop-then-remove ordering as `draft`/`run_to_end`: the
478            // engine flushes its log and frees the single-writer lock before
479            // the mission stops being "hosted here".
480            missions
481                .lock()
482                .expect("missions registry lock")
483                .remove(&mission_id);
484            drop(cell);
485            if let Err(e) = drive_result {
486                tracing::error!(mission = %mission_id, error = %e, "hosted ticket draft errored");
487            }
488        });
489
490        Ok(id)
491    }
492
493    /// `POST /api/tickets/:slug/approve`: the shared `kranz_engine::deps`
494    /// gate (cycle detection, unsatisfied-blocker refusal) plus the
495    /// enqueue side effects — the exact same core `kranz_cli`'s `kranz
496    /// ticket approve` calls, so the CLI and REST surfaces can never drift.
497    pub fn approve_ticket(
498        &self,
499        slug: &str,
500        force: bool,
501    ) -> Result<deps::ApprovedTicket, ApiError> {
502        deps::approve_ticket(&self.repo_root, slug, None, force).map_err(ApiError::from)
503    }
504
505    /// `POST /api/missions/:id/planning/turn`: one conversational turn. A
506    /// captured seed reply (fresh session / re-seed) is prepended — it
507    /// happened first in the conversation.
508    pub async fn planning_turn(&self, id: &str, text: &str) -> Result<String, ApiError> {
509        let cell = self.planning_cell_or_attach(id).await?;
510        let mut engine = try_lock(&cell)?;
511        let reply = engine.planning_turn(text).await?;
512        Ok(prepend_seed(engine.take_seed_reply(), reply))
513    }
514
515    /// `POST /api/missions/:id/planning/request-plan`: demand the plan.
516    /// Ready → plan + cost estimate; NotReady → the orchestrator's prose
517    /// (back to the conversation).
518    pub async fn request_plan(&self, id: &str) -> Result<Value, ApiError> {
519        let cell = self.planning_cell_or_attach(id).await?;
520        let mut engine = try_lock(&cell)?;
521        let request = engine.request_plan().await?;
522        let seed = engine.take_seed_reply();
523        match request {
524            PlanRequest::Ready(plan) => {
525                // Estimate with params calibrated from this repo's completed
526                // missions (built-in defaults when there are none yet).
527                let calibration = cost::calibrate(&self.repo_root);
528                let estimate = cost::estimate(&plan, &engine.state().config, &calibration.params);
529                let estimate = cost::apply_shape(estimate, &plan, &calibration);
530                // Park the reviewed plan so ANY surface's approve affordance
531                // (Slack buttons, web, glasses ring) can commit it later.
532                self.set_pending_plan(id, Some(plan.clone()));
533                Ok(json!({
534                    "ready": true,
535                    "planIdentity": plan_identity(&plan),
536                    "plan": plan,
537                    "estimate": estimate_json(&estimate),
538                    "calibration": { "missionsUsed": calibration.missions_used },
539                }))
540            }
541            PlanRequest::NotReady(reply) | PlanRequest::WrongPlan { reason: reply } => {
542                // A wrong-plan escalation reaches this interactive surface as
543                // the planner's reason text, exactly like a not-ready reply —
544                // the ticket-parking side effect is the draft flow's job.
545                Ok(json!({ "ready": false, "reply": prepend_seed(seed, reply) }))
546            }
547        }
548    }
549
550    /// `POST /api/missions/:id/approve`: commit plan.json/plan.md/index.md on
551    /// the mission branch exactly like the CLI. Returns the mission branch.
552    pub async fn approve(&self, id: &str, plan: Plan) -> Result<String, ApiError> {
553        let cell = self.planning_cell_or_attach(id).await?;
554        let mut engine = try_lock(&cell)?;
555        engine.approve_plan(plan)?;
556        self.set_pending_plan(id, None);
557        Ok(engine.state().mission.mission_branch.clone())
558    }
559
560    /// `POST /api/missions/:id/start`: consume the hosted engine into a
561    /// background `engine.run()` task — or, for a mission not in the registry
562    /// (blocked earlier, server restarted, or CLI-created), resume it from
563    /// the event log and run that.
564    pub async fn start(&self, id: &str) -> Result<(), ApiError> {
565        // Try to consume a hosted planning-phase engine.
566        let taken: Option<Box<MissionEngine>> = {
567            let mut map = self.missions.lock().expect("missions registry lock");
568            match map.remove(id) {
569                None => None,
570                Some(HostedMission::Running { handle, _repo_busy }) => {
571                    if handle.is_finished() {
572                        // The task ended but its cleanup lost the race with
573                        // this request: treat as not hosted (resume below).
574                        // Drop the busy hold so a resume can re-acquire.
575                        drop(_repo_busy);
576                        None
577                    } else {
578                        map.insert(
579                            id.to_string(),
580                            HostedMission::Running { handle, _repo_busy },
581                        );
582                        return Err(ApiError::conflict(format!(
583                            "mission '{id}' is already running — observe it via GET \
584                             /api/missions/{id}/state or steer it via POST \
585                             /api/missions/{id}/control"
586                        )));
587                    }
588                }
589                Some(HostedMission::Planning {
590                    cell,
591                    last_use,
592                    pending_plan,
593                }) => match Arc::try_unwrap(cell) {
594                    Err(cell) => {
595                        // A handler holds a clone: a turn is (or is about to
596                        // be) in flight. Put the entry back untouched.
597                        map.insert(
598                            id.to_string(),
599                            HostedMission::Planning {
600                                cell,
601                                last_use,
602                                pending_plan,
603                            },
604                        );
605                        return Err(turn_in_flight());
606                    }
607                    Ok(mutex) => {
608                        let engine = mutex.into_inner();
609                        if engine.state().mission.status == MissionStatus::Planning {
610                            map.insert(id.to_string(), new_planning(new_cell(engine)));
611                            return Err(ApiError::conflict(format!(
612                                "mission '{id}' has no approved plan yet — approve one via \
613                                 POST /api/missions/{id}/approve first"
614                            )));
615                        }
616                        Some(engine)
617                    }
618                },
619            }
620        };
621
622        let (engine, from_registry) = match taken {
623            Some(engine) => (engine, true),
624            None => {
625                // Re-invocable path: resume from the log. A live engine
626                // elsewhere (CLI, or a hosted run racing this request) holds
627                // the single-writer lock → EngineError::LockHeld → 409.
628                if !MissionPaths::new(&self.repo_root, id)
629                    .events_file()
630                    .is_file()
631                {
632                    return Err(ApiError::not_found(format!("unknown mission '{id}'")));
633                }
634                let cfg = config::load(&self.repo_root)?;
635                let backend = self.backend(cfg.claude_binary.as_deref()).await?;
636                let engine = Box::new(MissionEngine::resume(
637                    backend,
638                    self.repo_root.clone(),
639                    id,
640                    LockForce::No,
641                )?);
642                match engine.state().mission.status {
643                    MissionStatus::Planning => {
644                        return Err(ApiError::conflict(format!(
645                            "mission '{id}' is still in planning — approve a plan first \
646                             (POST /api/missions/{id}/approve, or `kranz plan`)"
647                        )))
648                    }
649                    MissionStatus::Complete => {
650                        return Err(ApiError::conflict(format!(
651                            "mission '{id}' is already complete — nothing to run"
652                        )))
653                    }
654                    MissionStatus::Failed => {
655                        return Err(ApiError::conflict(format!(
656                            "mission '{id}' has failed — inspect its log; there is nothing \
657                             the engine can resume"
658                        )))
659                    }
660                    _ => {}
661                }
662                (engine, false)
663            }
664        };
665
666        let global_run_permit = match self.try_global_run_permit() {
667            Ok(permit) => permit,
668            Err(error) => {
669                if from_registry {
670                    self.missions
671                        .lock()
672                        .expect("missions registry lock")
673                        .insert(id.to_string(), new_planning(new_cell(engine)));
674                }
675                return Err(error);
676            }
677        };
678
679        // Acquire the repo-wide busy lock before spawning: a sibling
680        // `kranz work` / hosted drain must not run in parallel. Held for the
681        // lifetime of the Running entry (dropped when the run ends).
682        let repo_busy = match kranz_engine::queue::acquire_repo_busy(&self.repo_root, id) {
683            Ok(hold) => hold,
684            Err(e) => {
685                if from_registry {
686                    // Put the planning/approved engine back so the operator
687                    // can retry once the sibling run finishes.
688                    self.missions
689                        .lock()
690                        .expect("missions registry lock")
691                        .insert(id.to_string(), new_planning(new_cell(engine)));
692                }
693                // Resume path: dropping `engine` releases the mission lock.
694                return Err(match e {
695                    e @ EngineError::LockHeld(_) => {
696                        ApiError::from(e).with_code(ApiErrorCode::RepositoryBusy)
697                    }
698                    other => other.into(),
699                });
700            }
701        };
702
703        // Insert the Running entry while holding the map lock across the
704        // spawn: if the run ends instantly, its cleanup blocks on this lock
705        // until the entry exists, so it can never leave a stale entry behind.
706        {
707            let mut map = self.missions.lock().expect("missions registry lock");
708            let missions = Arc::clone(&self.missions);
709            let mission_id = id.to_string();
710            let handle = spawn_with_global_run_permit(
711                global_run_permit,
712                run_to_end(engine, mission_id, missions),
713            );
714            map.insert(
715                id.to_string(),
716                HostedMission::Running {
717                    handle,
718                    _repo_busy: repo_busy,
719                },
720            );
721        }
722        Ok(())
723    }
724
725    /// `POST /api/missions/:id/merge`: the human-triggered gated Merge
726    /// action (roadmap M6). Loads the mission's `base_branch`/`base_sha`/
727    /// `mission_branch` from its event log (no engine needs to be hosted —
728    /// merge is independent of the planning/run-loop registry) and runs
729    /// [`kranz_engine::merge::merge_mission`] under `spawn_blocking` (git and
730    /// the gate suite are both blocking work). Never pushes.
731    pub async fn merge(&self, id: &str) -> Result<Value, ApiError> {
732        if !MissionPaths::is_safe_id(id) {
733            return Err(ApiError::not_found(format!("unknown mission '{id}'")));
734        }
735        let paths = MissionPaths::new(&self.repo_root, id);
736        if !paths.events_file().is_file() {
737            return Err(ApiError::not_found(format!("unknown mission '{id}'")));
738        }
739        // Serialize the complete read/pin/integrate/gate/advance transaction
740        // against mission runs and other merges in this repo. The hold moves
741        // INTO the blocking task below: if the client disconnects mid-gate-
742        // suite this handler future is dropped, but the detached blocking
743        // merge keeps mutating the primary tree — a hold living here would
744        // be released early, letting a dispatcher claim the busy repo.
745        let repo_busy = kranz_engine::queue::acquire_repo_busy(&self.repo_root, id).map_err(
746            |error| match error {
747                error @ EngineError::LockHeld(_) => {
748                    ApiError::from(error).with_code(ApiErrorCode::RepositoryBusy)
749                }
750                other => other.into(),
751            },
752        )?;
753        let events = EventLog::read_events(&paths.events_file())?;
754        let state = kranz_engine::reducer::fold(&events).map_err(ApiError::from)?;
755        if state.mission.status != MissionStatus::Complete {
756            return Err(ApiError::conflict(format!(
757                "mission '{id}' is {:?}; only a complete mission can be merged",
758                state.mission.status
759            )));
760        }
761        let base_branch = state.mission.base_branch.clone();
762        let base_sha = state.mission.base_sha.clone().ok_or_else(|| {
763            ApiError::conflict(format!(
764                "mission '{id}' has no pinned base sha — approve a plan first"
765            ))
766        })?;
767        let mission_branch = state.mission.mission_branch.clone();
768        // The approved Flight Rules pin (KRZ-342, D-E) rides into the merge:
769        // a repo-tracked pin makes merge re-resolve the live base policy
770        // against the exact scratch integration diff and refuse on
771        // enforced-set drift; `None` keeps the merge byte-identical.
772        let standards_pin = state.mission.standards_manifest.clone();
773        let standards_coverage = kranz_engine::standards_coverage::standards_coverage(id, &events);
774        let standards_evidence = StandardsMergeEvidence::from_mission_events(
775            id,
776            standards_pin.as_ref(),
777            standards_coverage.as_ref(),
778            &events,
779            chrono::Utc::now(),
780        );
781        let metadata = KranzCommitMetadata {
782            mission_id: state.mission.id.clone(),
783            cost_usd: state.total_cost_usd,
784            tokens: state.totals.clone(),
785        };
786
787        let repo_root = self.repo_root.clone();
788        let gate_executor = Arc::clone(&self.gate_executor);
789        // engine-gates-sandbox-wrapped: the merged mission's own
790        // `worker.sandbox` posture decides whether the gate suite (which
791        // executes that mission's worker-authored test/build code) runs
792        // inside the resolved sandbox profile. `enforce == off` falls
793        // through to the injected `gate_executor` — byte-identical pre-wrap
794        // behavior, and the test seam (`with_gate_executor`) stays
795        // authoritative there. Every enforced posture routes INTO the
796        // sandboxed runner: the process provider wraps in the resolved
797        // profile; `provider: container` wraps the gates in the mission
798        // container when a runtime is detected (ticket
799        // container-gate-wrapper); and the fail-closed postures (an
800        // unsupported platform, linux without `bwrap`, container without a
801        // runtime) error loudly at resolve rather than running unsandboxed
802        // (13th-pass review, P1).
803        let gate_policy = kranz_engine::command_exec::MergeGatePolicy {
804            sandbox: state.config.worker.sandbox.clone(),
805            mission_dir: paths.mission_dir(),
806        };
807        // A container-provider mission whose host has NO container runtime
808        // cannot wrap its merge gates (ticket container-gate-wrapper): they
809        // fail closed at resolve instead of running unsandboxed. This merge
810        // path has no event log, so the SAME note the resolve error carries
811        // goes to the operator-visible server log first — the refusal then
812        // reads as the config problem it is, never a flaky gate.
813        if let Some(note) = gate_policy.degradation_note() {
814            tracing::warn!(mission = %id, note = %note, "merge gate sandbox cannot wrap; gates fail closed");
815        }
816        let report = tokio::task::spawn_blocking(move || {
817            let repo = GitRepo::open(&repo_root)?;
818            let report = merge_mission_with_standards_evidence(
819                &repo,
820                &base_branch,
821                &base_sha,
822                &mission_branch,
823                Some(metadata),
824                standards_pin.as_ref(),
825                &standards_evidence,
826                |cmd, cwd| {
827                    if gate_policy.enforces_on_this_host() {
828                        kranz_engine::command_exec::run_bounded_gate_command_sandboxed(
829                            cwd,
830                            cmd,
831                            &gate_policy,
832                        )
833                    } else {
834                        gate_executor(cmd, cwd)
835                    }
836                },
837            );
838            // Explicit: the repo-busy hold is released HERE, once the merge
839            // has fully finished — never earlier by a dropped handler future.
840            drop(repo_busy);
841            report
842        })
843        .await
844        .map_err(|e| ApiError::internal(format!("merge task panicked: {e}")))?
845        .map_err(ApiError::from)?;
846
847        match report {
848            MergeReport::Merged { commit, stale_base } => Ok(json!({
849                "merged": true,
850                "commit": commit,
851                "staleBase": stale_base.map(|warning| json!({
852                    "baseSha": warning.base_sha,
853                    "liveBase": warning.live_base,
854                    "mergeCommitsSinceBase": warning.merge_commits_since_base,
855                    "message": format!(
856                        "stale base: {} merge commit(s) landed on {} since the mission base; cross-branch semantic conflicts are more likely, and full gates have run",
857                        warning.merge_commits_since_base,
858                        warning.live_base,
859                    ),
860                })),
861            })),
862            MergeReport::RefusedDirtyTree => Err(ApiError::conflict(
863                "refusing to merge: tracked working tree is dirty",
864            )),
865            MergeReport::GateFailed { gate, output } => Err(ApiError::unprocessable(
866                kranz_engine::scrub::scrub(&format!("{gate} failed:\n{output}")),
867            )),
868            MergeReport::GateConfigInvalid { detail } => Err(ApiError::unprocessable(format!(
869                "refusing to merge without a valid repo gate suite: {detail}"
870            ))),
871            MergeReport::SecretScanFailed { findings } => Err(ApiError::unprocessable(format!(
872                "secret scan failed; add a fingerprint to {} only for a reviewed false positive:\n{}",
873                kranz_engine::scrub::SECRET_ALLOWLIST_PATH,
874                kranz_engine::scrub::format_findings(&findings)
875            ))),
876            MergeReport::Conflict { files } => Err(ApiError::conflict(format!(
877                "merge conflicted in: {}",
878                files.join(", ")
879            ))),
880            MergeReport::RefusedPreMerge { detail } => Err(ApiError::conflict(format!(
881                "merge refused before it started: {detail}"
882            ))),
883            MergeReport::StandardsDrifted {
884                approved_digest,
885                current_digest,
886                changed_rules,
887            } => {
888                // KRZ-342 (D-E/D-H): the refusal is the merge's answer; the
889                // `standards.drifted` event is its evidence. Append it to the
890                // mission log best-effort — the mission is Complete, so no
891                // engine should hold the log lock; a held lock downgrades to
892                // a server-log warning, never to a silent 4xx.
893                if let Err(error) = EventLog::acquire(
894                    &paths,
895                    id,
896                    std::time::Duration::ZERO,
897                    LockForce::No,
898                )
899                .and_then(|mut log| {
900                    log.append(kranz_engine::events::EventKind::StandardsDrifted {
901                        approved_digest: approved_digest.clone(),
902                        current_digest: current_digest.clone(),
903                        surface: "merge".to_string(),
904                        changed_rules: changed_rules.clone(),
905                    })
906                    .map(|_| ())
907                }) {
908                    tracing::warn!(mission = %id, %error, "standards.drifted event could not be appended; the merge refusal stands");
909                }
910                Err(ApiError::unprocessable(format!(
911                    "refusing to merge: the live base Flight Rules policy drifted from the \
912                     approved pin (approved sha256:{approved_digest}, current {}) — the \
913                     applicable enforced set changed; revalidate and re-approve the mission:\n{}",
914                    current_digest
915                        .as_deref()
916                        .map(|d| format!("sha256:{d}"))
917                        .unwrap_or_else(|| "<unreadable>".to_string()),
918                    changed_rules.join("\n")
919                )))
920            }
921            MergeReport::StandardsFailed {
922                rule_id,
923                checker,
924                output,
925            } => Err(ApiError::unprocessable(kranz_engine::scrub::scrub(
926                &format!(
927                    "Flight Rules merge checker refused {rule_id} ({checker}):\n{output}"
928                ),
929            ))),
930        }
931    }
932
933    /// Read-only, LLM-backed Q&A for `/kranz ask`: ground the model in current
934    /// mission/ticket state and return one answer plus usage. This deliberately
935    /// bypasses the hosted mission registry: it must never create a mission,
936    /// append mission events, enqueue work, approve, start, or merge.
937    pub async fn ask(&self, question: &str) -> Result<Value, ApiError> {
938        let question = question.trim();
939        if question.is_empty() {
940            return Err(ApiError::bad_request("ask requires a question"));
941        }
942        let cfg = config::load(&self.repo_root)?;
943        config::validate(&cfg)?;
944        let role = cfg.validator_scrutiny.clone();
945        let backend = self.backend(cfg.claude_binary.as_deref()).await?;
946        let prompt = ask_prompt(question, &ask_context(&self.repo_root));
947        let spec = SessionSpec {
948            cwd: self.repo_root.clone(),
949            prompt: PromptMode::SingleShot(prompt),
950            append_system_prompt: Some(
951                "You answer read-only questions about this Kranz repository. \
952                 Use only the supplied context; if it is insufficient, say what is missing. \
953                 Do not modify files, run commands, create missions, enqueue work, approve, \
954                 start, or merge anything."
955                    .to_string(),
956            ),
957            model: role.model,
958            effort: role.reasoning_effort,
959            session_id: format!("ask-{}", uuid::Uuid::new_v4()),
960            resume: None,
961            permission_mode: Some("plan".to_string()),
962            allowed_tools: vec![],
963            disallowed_tools: vec![
964                "Bash(*)".to_string(),
965                "Edit(*)".to_string(),
966                "Write(*)".to_string(),
967            ],
968            tools: vec![
969                "Read".to_string(),
970                "Grep".to_string(),
971                "Glob".to_string(),
972                "LS".to_string(),
973            ],
974            writable: false,
975            settings_json: None,
976            json_schema: None,
977            max_budget_usd: role.max_budget_usd,
978            max_turns: role.max_turns,
979            env: HashMap::new(),
980            sandbox: None,
981            hook_status: None,
982        };
983        let outcome = run_ask_session(backend, spec).await?;
984        Ok(json!({
985            "answer": outcome.answer,
986            "costUsd": outcome.cost_usd,
987            "tokens": outcome.tokens,
988        }))
989    }
990
991    /// Release a hosted idle engine: drop it from the registry (flushing its
992    /// log and freeing the single-writer lock) so an EXTERNAL runner — the
993    /// `kranz work` dispatcher, a terminal `kranz plan/run` — can take the
994    /// mission over. The approve-and-QUEUE path needs this: without it the
995    /// approved engine would sit attached here holding the lock, and the very
996    /// dispatcher the queue points at would be refused with `LockHeld`.
997    ///
998    /// Returns `true` when the mission is now free of THIS host (released, or
999    /// was never hosted), `false` when it is actively running here (never
1000    /// interrupted). A turn in flight is an error, mirroring the other
1001    /// planning operations.
1002    pub fn release(&self, id: &str) -> Result<bool, ApiError> {
1003        release_from(&self.missions, id)
1004    }
1005
1006    /// Release every `Planning` entry idle for at least `threshold` (a mission
1007    /// touched more recently than that is left alone). A mid-turn cell can
1008    /// never actually be released — [`release`](Self::release) refuses it via
1009    /// `turn_in_flight`, which this treats as "not idle yet" rather than an
1010    /// error. Returns the ids this call actually released.
1011    pub fn sweep_idle(&self, threshold: Duration) -> Vec<String> {
1012        sweep_idle_from(&self.missions, threshold)
1013    }
1014
1015    /// Spawn the idle-release sweeper at most once, the first time a mission
1016    /// is hosted. It loops for the lifetime of the host: sleep, read the
1017    /// configured window, sweep. `planningIdleReleaseMinutes == 0` means
1018    /// "never release" — checked fresh each tick so a live config edit takes
1019    /// effect without a restart.
1020    fn ensure_sweeper_started(&self) {
1021        let mut guard = self.sweeper.lock().expect("sweeper lock");
1022        if guard.is_some() {
1023            return;
1024        }
1025        let repo_root = self.repo_root.clone();
1026        let missions = Arc::clone(&self.missions);
1027        *guard = Some(tokio::spawn(async move {
1028            const SWEEP_INTERVAL: Duration = Duration::from_secs(60);
1029            loop {
1030                tokio::time::sleep(SWEEP_INTERVAL).await;
1031                let minutes = match config::load(&repo_root) {
1032                    Ok(cfg) => cfg.planning_idle_release_minutes,
1033                    Err(_) => continue,
1034                };
1035                if minutes == 0 {
1036                    continue;
1037                }
1038                let threshold = Duration::from_secs(minutes * 60);
1039                let released = sweep_idle_from(&missions, threshold);
1040                for id in released {
1041                    tracing::info!(mission = %id, "released idle planning engine");
1042                }
1043            }
1044        }));
1045    }
1046
1047    /// Whether a tracked drain is currently live (a `Starting` reservation or
1048    /// a `Running` handle that hasn't finished). Read-only: never installs a
1049    /// reservation, so it never races [`Self::drain`]'s own check.
1050    pub(crate) fn drain_is_live(&self) -> bool {
1051        match &*self.drain.lock().expect("drain tracker lock") {
1052            DrainSlot::Idle => false,
1053            DrainSlot::Starting(_) => true,
1054            DrainSlot::Running(handle) => !handle.join.is_finished(),
1055        }
1056    }
1057
1058    /// One autoWork check: re-read config fresh (so a live `autoWork` toggle
1059    /// takes effect without a restart, exactly like the idle sweeper reads
1060    /// `planningIdleReleaseMinutes`), and kick off a drain when
1061    /// [`should_auto_drain`] says to. Invoked by the process-wide
1062    /// [`crate::MultiRepoHost`] watcher (and by tests) — per-host watchers
1063    /// are not started.
1064    pub(crate) async fn auto_work_tick(&self) -> bool {
1065        let cfg = match config::load(&self.repo_root) {
1066            Ok(cfg) => cfg,
1067            Err(_) => return false,
1068        };
1069        let queue_non_empty = kranz_engine::queue::peek(&self.repo_root).is_some();
1070        if should_auto_drain(cfg.auto_work, queue_non_empty, self.drain_is_live()) {
1071            // A sibling dispatcher already owns this repository. Skip it
1072            // before taking a process-wide permit so the catalog scheduler
1073            // can try another ready root in this same pass. `drain_once`
1074            // also stops on busy if ownership races this check.
1075            if kranz_engine::queue::is_repo_busy(&self.repo_root).is_some() {
1076                return false;
1077            }
1078            match self.drain_once().await {
1079                Ok(_) => return true,
1080                Err(e) if e.code == Some(ApiErrorCode::RepositoryBusy) => {}
1081                Err(e) => tracing::error!(error = %e.message, "autoWork drain failed"),
1082            }
1083        }
1084        false
1085    }
1086
1087    /// `POST /api/missions/:id/abandon`: retire a mission through the
1088    /// engine's canonical abandon path (terminal-refusing, event-recorded).
1089    /// A mission hosted HERE is taken out of the registry first — an idle
1090    /// planning engine is dropped (freeing the lock), a running task is
1091    /// aborted and awaited (the engine's Drop flushes the log and kills its
1092    /// agent children) — so the abandon event lands on a quiet log. A lock
1093    /// held by a FOREIGN process (a terminal `kranz plan/run`) surfaces as
1094    /// the engine's LockHeld → 409; the web never force-steals.
1095    pub async fn abandon(&self, id: &str, reason: &str) -> Result<(), ApiError> {
1096        let taken = self
1097            .missions
1098            .lock()
1099            .expect("missions registry lock")
1100            .remove(id);
1101        match taken {
1102            None => {}
1103            Some(HostedMission::Planning {
1104                cell,
1105                last_use,
1106                pending_plan,
1107            }) => match Arc::try_unwrap(cell) {
1108                Ok(mutex) => drop(mutex.into_inner()),
1109                Err(cell) => {
1110                    self.missions
1111                        .lock()
1112                        .expect("missions registry lock")
1113                        .insert(
1114                            id.to_string(),
1115                            HostedMission::Planning {
1116                                cell,
1117                                last_use,
1118                                pending_plan,
1119                            },
1120                        );
1121                    return Err(turn_in_flight());
1122                }
1123            },
1124            Some(HostedMission::Running { handle, _repo_busy }) => {
1125                if !handle.is_finished() {
1126                    handle.abort();
1127                }
1128                // Cancelled or finished either way: await settles the task so
1129                // the engine is dropped (log flushed, lock freed) before we
1130                // append the abandon event. Dropping `_repo_busy` releases the
1131                // repo-wide busy lock.
1132                let _ = handle.await;
1133                drop(_repo_busy);
1134            }
1135        }
1136        kranz_engine::mission_catalog::abandon_mission(
1137            self.repo_root.clone(),
1138            id,
1139            reason,
1140            LockForce::No,
1141        )
1142        .map_err(ApiError::from)?;
1143        Ok(())
1144    }
1145
1146    /// `POST /api/missions/:id/delete`: remove a TERMINAL mission's directory,
1147    /// mirroring `kranz clean` exactly — `cleanable_class` decides, `all`
1148    /// opts in to deleting Complete missions (which otherwise stay: they feed
1149    /// the cost-calibration corpus), and a live lock is re-checked immediately
1150    /// before removal so nothing is ever deleted under a running engine.
1151    /// Only the mission directory and its own `missions/index.md` line go;
1152    /// branches, tags, and every other mission's index line are left intact
1153    /// (same contract as the CLI).
1154    pub fn clean(&self, id: &str, all: bool) -> Result<(), ApiError> {
1155        use kranz_engine::mission_catalog::{
1156            cleanable_class, mission_lock_is_live, prune_mission_index_file, CleanClass,
1157        };
1158        if self
1159            .missions
1160            .lock()
1161            .expect("missions registry lock")
1162            .contains_key(id)
1163        {
1164            return Err(ApiError::conflict(format!(
1165                "mission '{id}' is hosted by this server (attached or running) — abandon it \
1166                 first, or let its run finish"
1167            )));
1168        }
1169        let paths = MissionPaths::new(&self.repo_root, id);
1170        if !paths.events_file().is_file() {
1171            return Err(ApiError::not_found(format!("unknown mission '{id}'")));
1172        }
1173        let events = EventLog::read_events(&paths.events_file())?;
1174        let state = kranz_engine::reducer::fold(&events).map_err(ApiError::from)?;
1175        let has_plan = paths.plan_file().is_file();
1176        match cleanable_class(state.mission.status, has_plan) {
1177            CleanClass::Keep => {
1178                return Err(ApiError::conflict(format!(
1179                    "mission '{id}' is live ({:?}) — abandon it before deleting",
1180                    state.mission.status
1181                )))
1182            }
1183            CleanClass::CompleteKeepByDefault if !all => {
1184                return Err(ApiError::conflict(format!(
1185                    "mission '{id}' is Complete; completed missions feed the cost-calibration \
1186                     corpus — pass \"all\": true to delete it anyway"
1187                )))
1188            }
1189            CleanClass::Stale | CleanClass::CompleteKeepByDefault => {}
1190        }
1191        // Same last-instant liveness re-check as the CLI's remove_missions: a
1192        // husk can go live between the fold and the removal.
1193        if mission_lock_is_live(&paths) {
1194            return Err(ApiError::conflict(format!(
1195                "mission '{id}' became live — nothing was deleted"
1196            )));
1197        }
1198        kranz_engine::queue::remove(&self.repo_root, id);
1199        std::fs::remove_dir_all(paths.mission_dir())
1200            .map_err(|e| ApiError::internal(format!("removing mission '{id}': {e}")))?;
1201        prune_mission_index_file(&self.repo_root, id);
1202        Ok(())
1203    }
1204
1205    /// The reviewed plan awaiting approval, if any (clone). `GET
1206    /// /api/missions/:id/pending-plan` and the glasses PLAN page read this.
1207    pub fn pending_plan(&self, id: &str) -> Option<Plan> {
1208        let pending = {
1209            let map = self.missions.lock().expect("missions registry lock");
1210            match map.get(id) {
1211                Some(HostedMission::Planning { pending_plan, .. }) => Arc::clone(pending_plan),
1212                _ => return None,
1213            }
1214        };
1215        // Approval may hold this mission's plan lock through a commit. Do not
1216        // keep the registry locked while waiting and block unrelated missions.
1217        let plan = pending.lock().expect("pending plan lock").clone();
1218        plan
1219    }
1220
1221    fn set_pending_plan(&self, id: &str, plan: Option<Plan>) {
1222        let map = self.missions.lock().expect("missions registry lock");
1223        if let Some(HostedMission::Planning { pending_plan, .. }) = map.get(id) {
1224            *pending_plan.lock().expect("pending plan lock") = plan;
1225        }
1226    }
1227
1228    /// Approve the currently parked plan for an explicit untargeted command
1229    /// (`/kranz approve`). Preview-based clients must use the matching variant.
1230    pub async fn try_approve_pending(&self, id: &str) -> Result<Option<String>, ApiError> {
1231        match self.approve_parked(id, |_| true)? {
1232            PendingApproval::Approved(branch) => Ok(Some(branch)),
1233            PendingApproval::NothingParked => Ok(None),
1234            PendingApproval::Mismatch { .. } => unreachable!("unconditional approval"),
1235        }
1236    }
1237
1238    /// Approve only the plan the caller reviewed. The engine lock serializes
1239    /// replacement and approval; the pending-plan lock protects comparison and
1240    /// consumption. An approval failure leaves the original plan parked, with
1241    /// no restore that could overwrite a concurrent replacement.
1242    pub async fn try_approve_pending_matching(
1243        &self,
1244        id: &str,
1245        expected_identity: Option<&str>,
1246    ) -> Result<PendingApproval, ApiError> {
1247        self.approve_parked(id, |plan| {
1248            expected_identity == Some(plan_identity(plan).as_str())
1249        })
1250    }
1251
1252    fn approve_parked(
1253        &self,
1254        id: &str,
1255        matches: impl FnOnce(&Plan) -> bool,
1256    ) -> Result<PendingApproval, ApiError> {
1257        let (cell, pending) = {
1258            let map = self.missions.lock().expect("missions registry lock");
1259            let Some(HostedMission::Planning {
1260                cell,
1261                pending_plan,
1262                last_use,
1263            }) = map.get(id)
1264            else {
1265                return Ok(PendingApproval::NothingParked);
1266            };
1267            *last_use.lock().expect("last-use lock") = Instant::now();
1268            (Arc::clone(cell), Arc::clone(pending_plan))
1269        };
1270        // Use the same engine-before-pending order as request_plan/approve.
1271        // Keeping a cell clone also prevents release/start from removing it.
1272        let mut engine = try_lock(&cell)?;
1273        let mut parked = pending.lock().expect("pending plan lock");
1274        let Some(plan) = parked.as_ref() else {
1275            return Ok(PendingApproval::NothingParked);
1276        };
1277        if !matches(plan) {
1278            return Ok(PendingApproval::Mismatch {
1279                parked: plan_identity(plan),
1280            });
1281        }
1282        engine.approve_plan(plan.clone())?;
1283        parked.take();
1284        Ok(PendingApproval::Approved(
1285            engine.state().mission.mission_branch.clone(),
1286        ))
1287    }
1288
1289    /// REST approval requires the identity returned with the reviewed preview.
1290    pub async fn approve_pending(
1291        &self,
1292        id: &str,
1293        expected_identity: Option<&str>,
1294    ) -> Result<String, ApiError> {
1295        match self.try_approve_pending_matching(id, expected_identity).await? {
1296            PendingApproval::Approved(branch) => Ok(branch),
1297            PendingApproval::NothingParked => Err(ApiError::conflict(format!(
1298                "mission '{id}' has no reviewed plan pending — refresh the plan preview before approving"
1299            )).with_code(ApiErrorCode::StalePlan)),
1300            PendingApproval::Mismatch { .. } => Err(ApiError::conflict(
1301                "reviewed plan identity is missing or stale — refresh the plan preview before approving",
1302            ).with_code(ApiErrorCode::StalePlan)),
1303        }
1304    }
1305
1306    /// `POST /api/queue/drain`: run the queue drain/claim/skip loop
1307    /// ([`kranz_engine::work::drain_queue`]) as a background task on this
1308    /// serve process. This is just ANOTHER dispatcher: it does not register
1309    /// missions in the `missions` planning registry, and arbitrates against
1310    /// an external `kranz work` process exactly as today — through the queue
1311    /// claim files and the events.jsonl single-writer lock, no new locking.
1312    ///
1313    /// IDEMPOTENT while a drain is live: a second call while the tracked
1314    /// drain task has not finished returns THAT drain's current state
1315    /// instead of spawning a second one.
1316    pub async fn drain(&self) -> Result<Value, ApiError> {
1317        self.drain_with_mode(false).await
1318    }
1319
1320    /// Auto-work drains at most one queue front so the process-wide scheduler
1321    /// can rotate fairly to another ready repository after this mission.
1322    async fn drain_once(&self) -> Result<Value, ApiError> {
1323        self.drain_with_mode(true).await
1324    }
1325
1326    async fn drain_with_mode(&self, once: bool) -> Result<Value, ApiError> {
1327        // Fast path: a live drain already owns the slot.
1328        {
1329            let guard = self.drain.lock().expect("drain tracker lock");
1330            match &*guard {
1331                DrainSlot::Starting(state) => {
1332                    return Ok(drain_state_json(&state.lock().expect("drain state lock")));
1333                }
1334                DrainSlot::Running(handle) if !handle.join.is_finished() => {
1335                    return Ok(drain_state_json(
1336                        &handle.state.lock().expect("drain state lock"),
1337                    ));
1338                }
1339                DrainSlot::Idle | DrainSlot::Running(_) => {}
1340            }
1341        }
1342
1343        // Discover config/backend before reserving the drain slot or taking a
1344        // process-wide run permit. Holding either across `.await` would either
1345        // publish a false-live Starting reservation (on later saturation) or
1346        // starve sibling repositories during Claude discovery.
1347        let cfg = config::load(&self.repo_root)?;
1348        let backend = self.backend(cfg.claude_binary.as_deref()).await?;
1349        let repo_root = self.repo_root.clone();
1350
1351        // Re-check under the drain lock, then acquire the global permit and
1352        // install Starting in one critical section so saturation never leaves
1353        // a rolled-back live reservation for concurrent callers to observe.
1354        let (state, global_run_permit) = {
1355            let mut guard = self.drain.lock().expect("drain tracker lock");
1356            match &*guard {
1357                DrainSlot::Starting(state) => {
1358                    return Ok(drain_state_json(&state.lock().expect("drain state lock")));
1359                }
1360                DrainSlot::Running(handle) if !handle.join.is_finished() => {
1361                    return Ok(drain_state_json(
1362                        &handle.state.lock().expect("drain state lock"),
1363                    ));
1364                }
1365                DrainSlot::Idle | DrainSlot::Running(_) => {}
1366            }
1367            let global_run_permit = self.try_global_run_permit()?;
1368            let state = Arc::new(Mutex::new(DrainState {
1369                live: true,
1370                current_mission_id: None,
1371                ran: Vec::new(),
1372                parked: Vec::new(),
1373            }));
1374            *guard = DrainSlot::Starting(Arc::clone(&state));
1375            (state, global_run_permit)
1376        };
1377
1378        // Cold spawn path only (never the early-return branches above): the
1379        // dispatch branch is still whatever the operator's checkout was, so
1380        // capture it now, before the spawned task (or any concurrent racer)
1381        // can ever land on a mission branch. See `drain_task` for the
1382        // restore-on-exit half of this contract.
1383        let task_state = Arc::clone(&state);
1384        let readiness_probe = self.readiness_probe;
1385        let join = tokio::spawn(async move {
1386            let _global_run_permit = global_run_permit;
1387            drain_task(
1388                repo_root.clone(),
1389                task_state,
1390                once,
1391                move |mission_id| {
1392                    let backend = Arc::clone(&backend);
1393                    let repo_root = repo_root.clone();
1394                    async move { run_mission_headless(backend, repo_root, mission_id).await }
1395                },
1396                readiness_probe,
1397            )
1398            .await;
1399        });
1400
1401        let initial = drain_state_json(&state.lock().expect("drain state lock"));
1402        *self.drain.lock().expect("drain tracker lock") =
1403            DrainSlot::Running(DrainHandle { join, state });
1404        Ok(initial)
1405    }
1406
1407    /// `GET /api/queue`: the queue front-to-back, who (if anyone) currently
1408    /// holds the busy lock, and this host's own drain tracker.
1409    ///
1410    /// Readiness is probed for the **front entry only** (with a short TTL
1411    /// cache). Deeper entries omit `readiness` so a long queue cannot turn
1412    /// every dashboard poll into N CLI shells.
1413    pub fn queue_state(&self) -> Value {
1414        let entries = kranz_engine::queue::list(&self.repo_root);
1415        let busy_with = kranz_engine::queue::is_repo_busy(&self.repo_root);
1416        let drain = match &*self.drain.lock().expect("drain tracker lock") {
1417            DrainSlot::Running(handle) => {
1418                drain_state_json(&handle.state.lock().expect("drain state lock"))
1419            }
1420            DrainSlot::Starting(state) => {
1421                drain_state_json(&state.lock().expect("drain state lock"))
1422            }
1423            DrainSlot::Idle => drain_state_json(&DrainState::default()),
1424        };
1425
1426        let front_readiness = entries.first().map(|e| {
1427            let mid = e.mission_id.as_str();
1428            {
1429                let cache = self
1430                    .readiness_front_cache
1431                    .lock()
1432                    .expect("readiness front cache lock");
1433                if let Some(cached) = cache.as_ref() {
1434                    if cached.mission_id == mid && cached.at.elapsed() < READINESS_FRONT_CACHE_TTL {
1435                        return (mid.to_string(), cached.report.clone());
1436                    }
1437                }
1438            }
1439            let report = (self.readiness_probe)(&self.repo_root, mid)
1440                .ok()
1441                .and_then(|r| serde_json::to_value(r).ok())
1442                .unwrap_or(Value::Null);
1443            *self
1444                .readiness_front_cache
1445                .lock()
1446                .expect("readiness front cache lock") = Some(FrontReadinessCache {
1447                mission_id: mid.to_string(),
1448                report: report.clone(),
1449                at: Instant::now(),
1450            });
1451            (mid.to_string(), report)
1452        });
1453
1454        let entries_json: Vec<Value> = entries
1455            .into_iter()
1456            .map(|e| {
1457                let readiness = front_readiness.as_ref().and_then(|(id, report)| {
1458                    if id == &e.mission_id && !report.is_null() {
1459                        Some(report.clone())
1460                    } else {
1461                        None
1462                    }
1463                });
1464                json!({
1465                    "missionId": e.mission_id,
1466                    "ticketSlug": e.ticket_slug,
1467                    "priority": e.priority,
1468                    "seq": e.seq,
1469                    "readiness": readiness,
1470                })
1471            })
1472            .collect();
1473        let mut state = json!({
1474            "entries": entries_json,
1475            "busyWith": busy_with,
1476            "drain": drain,
1477        });
1478        // Additive observation for automation: when this host participates in
1479        // host.maxConcurrentRepos, surface whether the process-wide budget is
1480        // currently exhausted so agents can distinguish "no work" from "capped".
1481        if let Some(permits) = &self.global_run_permits {
1482            let available = permits.available_permits();
1483            state["maxConcurrentReposAvailable"] = json!(available);
1484            state["maxConcurrentReposSaturated"] = json!(available == 0);
1485        }
1486        state
1487    }
1488
1489    // -----------------------------------------------------------------------
1490    // Registry plumbing
1491    // -----------------------------------------------------------------------
1492
1493    /// The engine cell of a planning-phase hosted mission, with helpful 409s
1494    /// for every other state.
1495    fn planning_cell(&self, id: &str) -> Result<EngineCell, ApiError> {
1496        let map = self.missions.lock().expect("missions registry lock");
1497        match map.get(id) {
1498            Some(HostedMission::Planning { cell, last_use, .. }) => {
1499                *last_use.lock().expect("last-use lock") = Instant::now();
1500                Ok(Arc::clone(cell))
1501            }
1502            Some(HostedMission::Running { .. }) => Err(ApiError::conflict(format!(
1503                "mission '{id}' is running — steer it via POST /api/missions/{id}/control"
1504            ))),
1505            None => Err(self.not_hosted(id)),
1506        }
1507    }
1508
1509    /// [`planning_cell`], attaching an un-hosted in-planning mission from disk
1510    /// first when needed. This is what lets a mission whose engine was released
1511    /// (CLI-created, bridge seed turn, server restart) continue planning through
1512    /// this host: resume it under the single-writer lock, adopt it into the
1513    /// registry, and hand back its cell. A mission held live elsewhere surfaces
1514    /// as the engine's `LockHeld` (409) — unless the holder is this registry
1515    /// itself racing us, in which case the second lookup finds the winner.
1516    async fn planning_cell_or_attach(&self, id: &str) -> Result<EngineCell, ApiError> {
1517        let miss = match self.planning_cell(id) {
1518            Ok(cell) => return Ok(cell),
1519            Err(miss) => miss,
1520        };
1521        // Only "exists on disk but not hosted" is attachable; Running entries
1522        // and unknown missions keep their original error.
1523        if !MissionPaths::new(&self.repo_root, id)
1524            .events_file()
1525            .is_file()
1526            || self
1527                .missions
1528                .lock()
1529                .expect("missions registry lock")
1530                .contains_key(id)
1531        {
1532            return Err(miss);
1533        }
1534        let cfg = config::load(&self.repo_root)?;
1535        let backend = self.backend(cfg.claude_binary.as_deref()).await?;
1536        let engine = match MissionEngine::resume(backend, self.repo_root.clone(), id, LockForce::No)
1537        {
1538            Ok(engine) => Box::new(engine),
1539            // LockHeld can mean a concurrent request won the attach race and
1540            // the winner's engine now sits in the registry: prefer that cell.
1541            Err(EngineError::LockHeld(holder)) => {
1542                return self
1543                    .planning_cell(id)
1544                    .map_err(|_| ApiError::from(EngineError::LockHeld(holder)))
1545            }
1546            Err(e) => return Err(e.into()),
1547        };
1548        if engine.state().mission.status != MissionStatus::Planning {
1549            // Dropping the engine releases the just-taken lock.
1550            return Err(ApiError::conflict(format!(
1551                "mission '{id}' is not in planning (status {:?}) — planning turns only \
1552                 apply before a plan is approved",
1553                engine.state().mission.status
1554            )));
1555        }
1556        let cell = new_cell(engine);
1557        let mut map = self.missions.lock().expect("missions registry lock");
1558        // We hold the mission's file lock, so nobody else can have inserted a
1559        // LIVE engine meanwhile; insert unconditionally.
1560        map.insert(id.to_string(), new_planning(Arc::clone(&cell)));
1561        drop(map);
1562        self.ensure_sweeper_started();
1563        Ok(cell)
1564    }
1565
1566    /// A mission that exists on disk but has no engine in this registry:
1567    /// its engine lives elsewhere (CLI) or was released (run ended, server
1568    /// restarted). Unknown missions are a plain 404.
1569    fn not_hosted(&self, id: &str) -> ApiError {
1570        let paths = MissionPaths::new(&self.repo_root, id);
1571        if paths.events_file().is_file() {
1572            ApiError::conflict(format!(
1573                "mission '{id}' is not hosted by this server — resume planning with \
1574                 `kranz plan --mission {id}`, or start execution via POST \
1575                 /api/missions/{id}/start"
1576            ))
1577            .with_code(ApiErrorCode::MissionNotHosted)
1578        } else {
1579            ApiError::not_found(format!("unknown mission '{id}'"))
1580        }
1581    }
1582}
1583
1584/// Drive one hosted mission to a terminal state, then release everything:
1585/// drop the engine FIRST (flushes the log, releases the single-writer lock),
1586/// THEN remove the registry entry — from that moment the mission is
1587/// observable and resumable anywhere (server or CLI).
1588fn spawn_with_global_run_permit<F>(
1589    global_run_permit: Option<OwnedSemaphorePermit>,
1590    task: F,
1591) -> tokio::task::JoinHandle<()>
1592where
1593    F: std::future::Future<Output = ()> + Send + 'static,
1594{
1595    tokio::spawn(async move {
1596        // Task ownership is deliberate: abort and panic both drop the permit
1597        // even when registry cleanup inside the future never runs.
1598        let _global_run_permit = global_run_permit;
1599        task.await;
1600    })
1601}
1602
1603async fn run_to_end(
1604    mut engine: Box<MissionEngine>,
1605    mission_id: String,
1606    missions: Arc<Mutex<HashMap<String, HostedMission>>>,
1607) {
1608    let repo_root = engine.paths().repo_root.clone();
1609    let result = engine.run().await;
1610    match &result {
1611        Ok(status) => {
1612            tracing::info!(mission = %mission_id, status = ?status, "hosted mission run ended")
1613        }
1614        Err(e) => {
1615            tracing::error!(mission = %mission_id, error = %e, "hosted mission run errored")
1616        }
1617    }
1618    drop(engine);
1619    // Reconcile the linked ticket's .status sidecar to match the mission's
1620    // terminal/blocked status. Non-fatal: a reconcile failure must never
1621    // affect the registry cleanup below.
1622    if let Err(e) = kranz_engine::work::reconcile_ticket_for_mission(&repo_root, &mission_id) {
1623        tracing::warn!(mission = %mission_id, error = %e, "failed to reconcile linked ticket");
1624    }
1625    missions
1626        .lock()
1627        .expect("missions registry lock")
1628        .remove(&mission_id);
1629}
1630
1631struct AskRunOutcome {
1632    answer: String,
1633    cost_usd: f64,
1634    tokens: TokenUsage,
1635}
1636
1637async fn run_ask_session(
1638    backend: Arc<dyn AgentBackend>,
1639    spec: SessionSpec,
1640) -> Result<AskRunOutcome, ApiError> {
1641    let mut session = backend.start(spec).await.map_err(ApiError::from)?;
1642    let mut streamed_text = String::new();
1643    let mut result_text = None;
1644    let mut tokens = TokenUsage::default();
1645    let mut cost_usd = 0.0;
1646    let mut result_error = false;
1647    while let Some(event) = session.next_event().await.map_err(ApiError::from)? {
1648        match event {
1649            AgentEvent::Text { text, .. } => streamed_text.push_str(&text),
1650            AgentEvent::Result {
1651                text,
1652                is_error,
1653                usage,
1654                cost_usd: cost,
1655                ..
1656            } => {
1657                result_error |= is_error;
1658                tokens.add(&usage);
1659                cost_usd += cost.unwrap_or(0.0);
1660                if !text.trim().is_empty() {
1661                    result_text = Some(text);
1662                }
1663            }
1664            _ => {}
1665        }
1666    }
1667    match session.exit_status() {
1668        Some(SessionExit::Completed) if !result_error => {
1669            let answer = result_text.unwrap_or(streamed_text).trim().to_string();
1670            if answer.is_empty() {
1671                return Err(ApiError::internal("ask turn produced an empty answer"));
1672            }
1673            Ok(AskRunOutcome {
1674                answer,
1675                cost_usd,
1676                tokens,
1677            })
1678        }
1679        Some(SessionExit::Completed) => Err(ApiError::internal("ask turn failed")),
1680        Some(SessionExit::Failed(reason)) => {
1681            Err(ApiError::internal(format!("ask turn failed: {reason}")))
1682        }
1683        Some(SessionExit::Aborted) => Err(ApiError::internal("ask turn aborted")),
1684        None => Err(ApiError::internal("ask turn ended without an exit status")),
1685    }
1686}
1687
1688fn ask_prompt(question: &str, context: &str) -> String {
1689    format!(
1690        "Answer this operator question about the Kranz repository.\n\n\
1691         Rules:\n\
1692         - Ground the answer only in the context below.\n\
1693         - If the context is insufficient, say what is missing.\n\
1694         - Keep the answer concise but specific, citing mission ids or ticket slugs when relevant.\n\
1695         - This is read-only: do not propose that you have changed state.\n\n\
1696         Question:\n{question}\n\nContext:\n{context}"
1697    )
1698}
1699
1700fn ask_context(repo_root: &Path) -> String {
1701    let mut out = String::new();
1702    out.push_str("## Missions\n");
1703    let mut ids = MissionPaths::list_missions(repo_root);
1704    ids.sort();
1705    ids.reverse();
1706    if ids.is_empty() {
1707        out.push_str("(none)\n");
1708    }
1709    for id in ids.into_iter().take(20) {
1710        let paths = MissionPaths::new(repo_root, &id);
1711        let Ok(events) = EventLog::read_events(&paths.events_file()) else {
1712            continue;
1713        };
1714        let Ok(state) = kranz_engine::reducer::fold(&events) else {
1715            continue;
1716        };
1717        out.push_str(&format!(
1718            "- {}: {:?}; goal: {}; branch: {}; cost: ${:.4}; tokens in/out/cacheRead/cacheWrite: {}/{}/{}/{}\n",
1719            state.mission.id,
1720            state.mission.status,
1721            one_line(&state.mission.goal),
1722            state.mission.mission_branch,
1723            state.total_cost_usd,
1724            state.totals.input,
1725            state.totals.output,
1726            state.totals.cache_read,
1727            state.totals.cache_write,
1728        ));
1729        for decision in state.recent_decisions.iter().rev().take(3) {
1730            out.push_str(&format!("  decision: {}\n", one_line(decision)));
1731        }
1732        // Threat (follow-up review H-4): a worker under checkout isolation
1733        // can replace this leaf with a symlink, and 500 chars of whatever it
1734        // points at (`serve.token`, `~/.ssh/id_ed25519`) would ride into the
1735        // ask-session LLM context. The read pins the `.kranz` chain AND the
1736        // leaf; a refused read simply contributes no excerpt.
1737        let report = paths.report_file();
1738        let mut text = String::new();
1739        let read = kranz_engine::paths::open_read_nofollow(&report)
1740            .and_then(|mut file| {
1741                use std::io::Read as _;
1742                file.read_to_string(&mut text)?;
1743                Ok(())
1744            })
1745            .is_ok();
1746        if read {
1747            out.push_str(&format!(
1748                "  report excerpt: {}\n",
1749                truncate(&one_line(&text), 500)
1750            ));
1751        }
1752    }
1753
1754    out.push_str("\n## Tickets\n");
1755    let tickets = Ticket::list(repo_root);
1756    if tickets.is_empty() {
1757        out.push_str("(none)\n");
1758    }
1759    for ticket in tickets.iter().take(40) {
1760        let state = Ticket::read_state(repo_root, &ticket.slug);
1761        out.push_str(&format!(
1762            "- {} [{:?}, p{}]: {}; blocked-by: {}\n",
1763            ticket.slug,
1764            state,
1765            ticket.priority,
1766            one_line(&ticket.title),
1767            if ticket.blocked_by.is_empty() {
1768                "none".to_string()
1769            } else {
1770                ticket.blocked_by.join(", ")
1771            }
1772        ));
1773    }
1774
1775    out.push_str("\n## Queue\n");
1776    let entries = queue::list(repo_root);
1777    if entries.is_empty() {
1778        out.push_str("(empty)\n");
1779    }
1780    for entry in entries.iter().take(20) {
1781        out.push_str(&format!(
1782            "- {} priority={} ticket={}\n",
1783            entry.mission_id,
1784            entry.priority,
1785            entry.ticket_slug.as_deref().unwrap_or("-")
1786        ));
1787    }
1788    out
1789}
1790
1791fn one_line(text: &str) -> String {
1792    text.split_whitespace().collect::<Vec<_>>().join(" ")
1793}
1794
1795fn truncate(text: &str, max: usize) -> String {
1796    if text.chars().count() <= max {
1797        return text.to_string();
1798    }
1799    let mut out: String = text.chars().take(max.saturating_sub(1)).collect();
1800    out.push('…');
1801    out
1802}
1803
1804/// The spawned-task body behind [`MissionHost::drain`], factored out so a
1805/// host-level test can drive it directly (no `tokio::spawn`, so it stays
1806/// deterministic) with a fake `run_mission`. Captures the operator's dispatch
1807/// checkout, runs [`kranz_engine::work::drain_queue`] to completion, then
1808/// restores that checkout — honoring the SAME contract as the CLI
1809/// dispatcher's `restore_work_checkout` (`crates/cli/src/backlog.rs`), so a
1810/// hosted drain can never leave the repo stranded on a
1811/// `kranz/mission-*` branch.
1812async fn drain_task<R, Fut>(
1813    repo_root: PathBuf,
1814    state: Arc<Mutex<DrainState>>,
1815    once: bool,
1816    run_mission: R,
1817    readiness_probe: ReadinessProbe,
1818) where
1819    R: Fn(String) -> Fut,
1820    Fut: std::future::Future<Output = anyhow::Result<i32>>,
1821{
1822    drain_task_with_probe(repo_root, state, once, run_mission, readiness_probe).await;
1823}
1824
1825/// [`drain_task`] with an injectable readiness probe so checkout-restoration
1826/// tests remain hermetic on clean CI runners that intentionally have no agent
1827/// CLI installed.
1828async fn drain_task_with_probe<R, Fut, P>(
1829    repo_root: PathBuf,
1830    state: Arc<Mutex<DrainState>>,
1831    once: bool,
1832    run_mission: R,
1833    readiness_probe: P,
1834) where
1835    R: Fn(String) -> Fut,
1836    Fut: std::future::Future<Output = anyhow::Result<i32>>,
1837    P: Fn(
1838        &Path,
1839        &str,
1840    ) -> kranz_engine::error::Result<kranz_engine::backend_readiness::ReadinessReport>,
1841{
1842    // Capture BEFORE `drain_queue` runs anything — nothing has touched the
1843    // checkout yet, so this is genuinely the operator's dispatch-time branch.
1844    let dispatch_branch = GitRepo::open(&repo_root)
1845        .ok()
1846        .and_then(|g| g.current_branch().ok());
1847
1848    let result = kranz_engine::work::drain_queue_with_probe(
1849        &repo_root,
1850        once,
1851        |mission_id| {
1852            let state = Arc::clone(&state);
1853            let fut = run_mission(mission_id.clone());
1854            async move {
1855                state.lock().expect("drain state lock").current_mission_id =
1856                    Some(mission_id.clone());
1857                let outcome = fut.await;
1858                let mut guard = state.lock().expect("drain state lock");
1859                guard.current_mission_id = None;
1860                if outcome.is_ok() {
1861                    guard.ran.push(mission_id);
1862                }
1863                outcome
1864            }
1865        },
1866        readiness_probe,
1867    )
1868    .await;
1869
1870    match &result {
1871        Ok(report) if !report.stopped_busy => {
1872            {
1873                let mut guard = state.lock().expect("drain state lock");
1874                for id in &report.parked {
1875                    if !guard.parked.contains(id) {
1876                        guard.parked.push(id.clone());
1877                    }
1878                }
1879            }
1880            restore_drain_checkout(&repo_root, dispatch_branch.as_deref());
1881        }
1882        Ok(report) => {
1883            let mut guard = state.lock().expect("drain state lock");
1884            for id in &report.parked {
1885                if !guard.parked.contains(id) {
1886                    guard.parked.push(id.clone());
1887                }
1888            }
1889            // `stopped_busy` (only possible with `once`, which the hosted
1890            // drain never sets — wired for parity with `cmd_work` anyway):
1891            // a sibling dispatcher may still be mid-mission, so leave the
1892            // checkout exactly where it is.
1893        }
1894        Err(e) => {
1895            tracing::error!(error = %e, "hosted queue drain errored");
1896            // Same restore as the success path: an errored drain must not
1897            // leave the operator stranded on a mission branch.
1898            restore_drain_checkout(&repo_root, dispatch_branch.as_deref());
1899        }
1900    }
1901    state.lock().expect("drain state lock").live = false;
1902}
1903
1904/// Dispatcher-exit checkout restore for the hosted queue drain: mirrors
1905/// `restore_work_checkout` in `crates/cli/src/backlog.rs` verbatim. `None`
1906/// (capture failed, or nothing to restore) is a no-op; restoring TO a
1907/// `kranz/mission-*` branch is refused (that would recreate the very
1908/// stranding this exists to end); already back on the captured branch is a
1909/// no-op; a dirty TRACKED working tree aborts the restore (never carry
1910/// uncommitted operator edits across a branch switch) and leaves the
1911/// checkout on the mission branch with a warning logged.
1912fn restore_drain_checkout(repo_root: &Path, original: Option<&str>) {
1913    let Some(original) = original else { return };
1914    if original.starts_with("kranz/mission-") {
1915        return;
1916    }
1917    let Ok(git) = GitRepo::open(repo_root) else {
1918        return;
1919    };
1920    if git.current_branch().ok().as_deref() == Some(original) {
1921        return;
1922    }
1923    match git.is_clean_tracked() {
1924        Ok(true) => match git.checkout(original) {
1925            Ok(()) => tracing::info!(branch = %original, "hosted drain restored operator checkout"),
1926            Err(e) => {
1927                tracing::warn!(branch = %original, error = %e, "hosted drain could not restore checkout")
1928            }
1929        },
1930        Ok(false) => tracing::warn!(
1931            "hosted drain leaving checkout in place: tracked files have uncommitted changes"
1932        ),
1933        Err(e) => {
1934            tracing::warn!(error = %e, "hosted drain could not probe the working tree; checkout left in place")
1935        }
1936    }
1937}
1938
1939/// Headless `run_mission` injected into [`kranz_engine::work::drain_queue`]
1940/// by [`MissionHost::drain`]: resume the mission under the single-writer
1941/// lock and run it to a terminal state, with no live tail/printer attached
1942/// (unlike the CLI's `kranz work`) since no terminal is attached to a serve
1943/// process.
1944async fn run_mission_headless(
1945    backend: Arc<dyn AgentBackend>,
1946    repo_root: PathBuf,
1947    mission_id: String,
1948) -> anyhow::Result<i32> {
1949    let mut engine = MissionEngine::resume(backend, repo_root, &mission_id, LockForce::No)?;
1950    let status = engine.run().await?;
1951    Ok(exit_code_for(status))
1952}
1953
1954/// Map a terminal [`MissionStatus`] to the exit code the CLI's
1955/// `kranz work`/`kranz exec` report, matching `kranz_cli::exec::exit_code_for`.
1956fn exit_code_for(status: MissionStatus) -> i32 {
1957    match status {
1958        MissionStatus::Complete => 0,
1959        MissionStatus::Blocked => 2,
1960        _ => 1,
1961    }
1962}
1963
1964/// Apply a ticket's per-ticket budget override to the orchestrator role
1965/// (mirrors `kranz_cli::backlog::config_for_ticket`), so draft spend is
1966/// bounded by the ticket's `maxBudgetUsd` when it sets one.
1967fn config_for_ticket(mut cfg: MissionConfig, ticket: &Ticket) -> MissionConfig {
1968    if let Some(budget) = ticket.max_budget_usd {
1969        cfg.orchestrator.max_budget_usd = Some(budget);
1970    }
1971    cfg
1972}
1973
1974/// [`DraftOutcome`] as protocol camelCase JSON (the engine type is a plain
1975/// contract enum without serde derives) — the shape a REST `draft` handler
1976/// hands back. Not yet wired to a route (that's a later feature); kept here
1977/// so [`MissionHost::draft`]'s result has a ready serialization.
1978#[allow(dead_code)]
1979fn draft_outcome_json(outcome: &DraftOutcome) -> Value {
1980    match outcome {
1981        DraftOutcome::ParkedForReview {
1982            mission_id,
1983            mission_branch,
1984        } => json!({
1985            "outcome": "parkedForReview",
1986            "missionId": mission_id,
1987            "missionBranch": mission_branch,
1988        }),
1989        DraftOutcome::Enqueued { mission_id } => json!({
1990            "outcome": "enqueued",
1991            "missionId": mission_id,
1992        }),
1993        DraftOutcome::PlanAsProse { mission_id } => json!({
1994            "outcome": "planAsProse",
1995            "missionId": mission_id,
1996            "message": "the orchestrator produced a plan but emitted it as prose instead of \
1997                        through the plan channel, so nothing was queued; re-run draft for \
1998                        this ticket",
1999        }),
2000        DraftOutcome::NeedsContext {
2001            mission_id,
2002            questions,
2003        } => json!({
2004            "outcome": "needsContext",
2005            "missionId": mission_id,
2006            "questions": questions,
2007        }),
2008        DraftOutcome::WrongPlan { mission_id, reason } => json!({
2009            "outcome": "wrongPlan",
2010            "missionId": mission_id,
2011            "reason": reason,
2012        }),
2013    }
2014}
2015
2016fn new_cell(engine: Box<MissionEngine>) -> EngineCell {
2017    Arc::new(tokio::sync::Mutex::new(engine))
2018}
2019
2020/// A fresh `Planning` entry, last-used now.
2021fn new_planning(cell: EngineCell) -> HostedMission {
2022    HostedMission::Planning {
2023        cell,
2024        last_use: Arc::new(Mutex::new(Instant::now())),
2025        pending_plan: Arc::new(Mutex::new(None)),
2026    }
2027}
2028
2029/// Shared by [`MissionHost::release`] and the sweeper: drop an idle planning
2030/// engine from the registry (flushing its log and freeing the single-writer
2031/// lock), refuse a mid-turn one, and leave running/absent entries be.
2032fn release_from(
2033    missions: &Mutex<HashMap<String, HostedMission>>,
2034    id: &str,
2035) -> Result<bool, ApiError> {
2036    let mut map = missions.lock().expect("missions registry lock");
2037    match map.remove(id) {
2038        None => Ok(true),
2039        Some(HostedMission::Running { handle, _repo_busy }) => {
2040            let finished = handle.is_finished();
2041            if !finished {
2042                map.insert(
2043                    id.to_string(),
2044                    HostedMission::Running { handle, _repo_busy },
2045                );
2046            }
2047            Ok(finished)
2048        }
2049        Some(HostedMission::Planning {
2050            cell,
2051            last_use,
2052            pending_plan,
2053        }) => match Arc::try_unwrap(cell) {
2054            Ok(mutex) => {
2055                drop(mutex.into_inner()); // flushes the log, frees the lock
2056                Ok(true)
2057            }
2058            Err(cell) => {
2059                map.insert(
2060                    id.to_string(),
2061                    HostedMission::Planning {
2062                        cell,
2063                        last_use,
2064                        pending_plan,
2065                    },
2066                );
2067                Err(turn_in_flight())
2068            }
2069        },
2070    }
2071}
2072
2073/// Collect ids of `Planning` entries idle for at least `threshold`, release
2074/// each via [`release_from`], and return the ids actually freed. A mid-turn
2075/// cell (its `try_unwrap` fails inside `release_from`) is skipped, not an
2076/// error — it simply isn't idle yet from the sweeper's point of view.
2077/// `Running` entries are never candidates.
2078fn sweep_idle_from(
2079    missions: &Mutex<HashMap<String, HostedMission>>,
2080    threshold: Duration,
2081) -> Vec<String> {
2082    let idle_ids: Vec<String> = {
2083        let map = missions.lock().expect("missions registry lock");
2084        map.iter()
2085            .filter_map(|(id, mission)| match mission {
2086                HostedMission::Planning { last_use, .. } => {
2087                    let elapsed = last_use.lock().expect("last-use lock").elapsed();
2088                    (elapsed >= threshold).then(|| id.clone())
2089                }
2090                HostedMission::Running { .. } => None,
2091            })
2092            .collect()
2093    };
2094    idle_ids
2095        .into_iter()
2096        .filter(|id| matches!(release_from(missions, id), Ok(true)))
2097        .collect()
2098}
2099
2100/// Planning endpoints never queue behind each other: contended = 409.
2101fn try_lock(
2102    cell: &EngineCell,
2103) -> Result<tokio::sync::MutexGuard<'_, Box<MissionEngine>>, ApiError> {
2104    cell.try_lock().map_err(|_| turn_in_flight())
2105}
2106
2107fn turn_in_flight() -> ApiError {
2108    ApiError::conflict("a turn is in flight for this mission — wait for it to finish")
2109        .with_code(ApiErrorCode::TurnInFlight)
2110}
2111
2112/// Seed replies (fresh session / resume-ack / re-seed) happened first in the
2113/// conversation, so they go first in the combined reply.
2114fn prepend_seed(seed: Option<String>, reply: String) -> String {
2115    match seed {
2116        Some(seed) => format!("{seed}\n\n{reply}"),
2117        None => reply,
2118    }
2119}
2120
2121/// [`CostEstimate`] as protocol camelCase JSON (the engine type is a plain
2122/// contract struct without serde derives).
2123fn estimate_json(estimate: &CostEstimate) -> Value {
2124    let confidence = match estimate.confidence {
2125        kranz_engine::cost::Confidence::High => "high",
2126        kranz_engine::cost::Confidence::Low => "low",
2127    };
2128    json!({
2129        "workerRuns": estimate.worker_runs,
2130        "validatorRuns": estimate.validator_runs,
2131        "lowUsd": estimate.low_usd,
2132        "expectedUsd": estimate.expected_usd,
2133        "highUsd": estimate.high_usd,
2134        "confidence": confidence,
2135    })
2136}
2137
2138// ---------------------------------------------------------------------------
2139// Axum handlers (docs/protocol.md "Mission lifecycle" table)
2140// ---------------------------------------------------------------------------
2141
2142/// `POST /api/missions` — body `{"goal":"...", "config":{...}}` →
2143/// `201 {"id":"m-…"}`.
2144pub(crate) async fn create_mission(
2145    State(server): State<Arc<ServerState>>,
2146    body: Bytes,
2147) -> Result<impl IntoResponse, ApiError> {
2148    let value = parse_body(&body)?;
2149    let goal = value
2150        .get("goal")
2151        .and_then(Value::as_str)
2152        .map(str::trim)
2153        .filter(|goal| !goal.is_empty())
2154        .ok_or_else(|| {
2155            ApiError::bad_request(r#"body must be {"goal":"..."} with a non-empty goal"#)
2156        })?;
2157    let id = server.host.create(goal, value.get("config")).await?;
2158    Ok((StatusCode::CREATED, Json(json!({ "id": id }))))
2159}
2160
2161/// `POST /api/missions/:id/planning/turn` — body `{"text":"..."}` →
2162/// `200 {"reply":"..."}`.
2163pub(crate) async fn planning_turn(
2164    State(server): State<Arc<ServerState>>,
2165    UrlPath(id): UrlPath<String>,
2166    body: Bytes,
2167) -> Result<Json<Value>, ApiError> {
2168    let id = valid_id(&server, &id)?;
2169    let value = parse_body(&body)?;
2170    let text = value
2171        .get("text")
2172        .and_then(Value::as_str)
2173        .map(str::trim)
2174        .filter(|text| !text.is_empty())
2175        .ok_or_else(|| {
2176            ApiError::bad_request(r#"body must be {"text":"..."} with non-empty text"#)
2177        })?;
2178    let reply = server.host.planning_turn(&id, text).await?;
2179    Ok(Json(json!({ "reply": reply })))
2180}
2181
2182/// `POST /api/missions/:id/planning/request-plan` →
2183/// `200 {"ready":true,"plan":{...},"estimate":{...}}` or
2184/// `200 {"ready":false,"reply":"..."}`.
2185pub(crate) async fn request_plan(
2186    State(server): State<Arc<ServerState>>,
2187    UrlPath(id): UrlPath<String>,
2188) -> Result<Json<Value>, ApiError> {
2189    let id = valid_id(&server, &id)?;
2190    Ok(Json(server.host.request_plan(&id).await?))
2191}
2192
2193/// `POST /api/missions/:id/approve` — body `{"plan":{...}}` →
2194/// `200 {"branch":"kranz/mission-…"}`.
2195pub(crate) async fn approve_mission(
2196    State(server): State<Arc<ServerState>>,
2197    UrlPath(id): UrlPath<String>,
2198    body: Bytes,
2199) -> Result<Json<Value>, ApiError> {
2200    let id = valid_id(&server, &id)?;
2201    let value = parse_body(&body)?;
2202    let plan = value
2203        .get("plan")
2204        .cloned()
2205        .ok_or_else(|| ApiError::bad_request(r#"body must be {"plan":{...}}"#))?;
2206    let plan: Plan = serde_json::from_value(plan)
2207        .map_err(|e| ApiError::bad_request(format!("'plan' is not a valid Plan: {e}")))?;
2208    let branch = server.host.approve(&id, plan).await?;
2209    Ok(Json(json!({ "branch": branch })))
2210}
2211
2212/// `GET /api/missions/:id/pending-plan` → `200 {"pending":true,"plan":{…}}`
2213/// or `200 {"pending":false}`. The parked plan from the last Ready
2214/// request-plan — what the approve affordances (buttons, ring) will commit.
2215pub(crate) async fn pending_plan_route(
2216    axum::Extension(reads): axum::Extension<crate::read_work::ReadWork>,
2217    State(server): State<Arc<ServerState>>,
2218    UrlPath(id): UrlPath<String>,
2219) -> Result<Json<Value>, ApiError> {
2220    reads
2221        .run(move || {
2222            let id = valid_id(&server, &id)?;
2223            Ok(Json(match server.host.pending_plan(&id) {
2224                Some(plan) => {
2225                    json!({ "pending": true, "planIdentity": plan_identity(&plan), "plan": plan })
2226                }
2227                None => json!({ "pending": false }),
2228            }))
2229        })
2230        .await
2231}
2232
2233/// `POST /api/missions/:id/approve-pending` — body
2234/// `{"planIdentity": "…", "start": true}` → approve the matching parked plan
2235/// (409 when missing or stale), then optionally start.
2236pub(crate) async fn approve_pending_route(
2237    State(server): State<Arc<ServerState>>,
2238    UrlPath(id): UrlPath<String>,
2239    body: Bytes,
2240) -> Result<Json<Value>, ApiError> {
2241    let id = valid_id(&server, &id)?;
2242    let value = parse_body(&body)?;
2243    let start = value.get("start").and_then(Value::as_bool).unwrap_or(false);
2244    let expected_identity = value.get("planIdentity").and_then(Value::as_str);
2245    let branch = server.host.approve_pending(&id, expected_identity).await?;
2246    if start {
2247        server.host.start(&id).await?;
2248    }
2249    Ok(Json(json!({ "branch": branch, "started": start })))
2250}
2251
2252/// `POST /api/missions/:id/abandon` — optional body `{"reason":"..."}` →
2253/// `200 {"abandoned": true}`. See [`MissionHost::abandon`].
2254pub(crate) async fn abandon_mission_route(
2255    State(server): State<Arc<ServerState>>,
2256    UrlPath(id): UrlPath<String>,
2257    body: Bytes,
2258) -> Result<Json<Value>, ApiError> {
2259    let id = valid_id(&server, &id)?;
2260    let value = parse_body(&body)?;
2261    let reason = value
2262        .get("reason")
2263        .and_then(Value::as_str)
2264        .map(str::trim)
2265        .filter(|r| !r.is_empty())
2266        .unwrap_or("abandoned by operator");
2267    server.host.abandon(&id, reason).await?;
2268    Ok(Json(json!({ "abandoned": true })))
2269}
2270
2271/// `POST /api/missions/:id/release` — no body → `200 {"released": bool}`. See
2272/// [`MissionHost::release`]. A mission absent from disk is 404; a not-hosted
2273/// but on-disk mission is an idempotent 200 (already free). POST (not a
2274/// dedicated verb) so the mutation-token gate applies by construction.
2275pub(crate) async fn release_mission_route(
2276    State(server): State<Arc<ServerState>>,
2277    UrlPath(id): UrlPath<String>,
2278    body: Bytes,
2279) -> Result<Json<Value>, ApiError> {
2280    let id = valid_id(&server, &id)?;
2281    let _ = parse_body(&body)?;
2282    if !MissionPaths::new(server.host.repo_root(), &id)
2283        .events_file()
2284        .is_file()
2285    {
2286        return Err(ApiError::not_found(format!("mission '{id}' not found")));
2287    }
2288    let released = server.host.release(&id)?;
2289    Ok(Json(json!({ "released": released })))
2290}
2291
2292/// `POST /api/missions/:id/delete` — optional body `{"all": true}` (opt in to
2293/// deleting a Complete mission) → `200 {"deleted": true}`. See
2294/// [`MissionHost::clean`]. POST (not the DELETE verb) so the mutation-token
2295/// gate — which covers `POST /api/...` — applies by construction.
2296pub(crate) async fn delete_mission_route(
2297    State(server): State<Arc<ServerState>>,
2298    UrlPath(id): UrlPath<String>,
2299    body: Bytes,
2300) -> Result<Json<Value>, ApiError> {
2301    let id = valid_id(&server, &id)?;
2302    let value = parse_body(&body)?;
2303    let all = value.get("all").and_then(Value::as_bool).unwrap_or(false);
2304    server.host.clean(&id, all)?;
2305    Ok(Json(json!({ "deleted": true })))
2306}
2307
2308/// `POST /api/missions/:id/start` → `202 {"running":true}`.
2309pub(crate) async fn start_mission(
2310    State(server): State<Arc<ServerState>>,
2311    UrlPath(id): UrlPath<String>,
2312) -> Result<impl IntoResponse, ApiError> {
2313    let id = valid_id(&server, &id)?;
2314    server.host.start(&id).await?;
2315    Ok((StatusCode::ACCEPTED, Json(json!({ "running": true }))))
2316}
2317
2318/// `POST /api/missions/:id/merge` → `200 {"merged":true,"commit":"..."}` on
2319/// success. See [`MissionHost::merge`] for the non-2xx shapes (dirty tree /
2320/// gate failure / conflict).
2321pub(crate) async fn merge_mission_route(
2322    State(server): State<Arc<ServerState>>,
2323    UrlPath(id): UrlPath<String>,
2324) -> Result<Json<Value>, ApiError> {
2325    let id = valid_id(&server, &id)?;
2326    Ok(Json(server.host.merge(&id).await?))
2327}
2328
2329/// `POST /api/queue/drain` — no required body → `200 <drain-state JSON>`.
2330/// See [`MissionHost::drain`]; idempotent while a drain is already live.
2331pub(crate) async fn drain_queue_route(
2332    State(server): State<Arc<ServerState>>,
2333    body: Bytes,
2334) -> Result<Json<Value>, ApiError> {
2335    let _ = parse_body(&body)?;
2336    Ok(Json(server.host.drain().await?))
2337}
2338
2339/// `GET /api/queue` → `200 {"entries":[...], "busyWith": <id|null>, "drain": {...}}`.
2340/// See [`MissionHost::queue_state`]. Tokenless: read-only.
2341pub(crate) async fn queue_state_route(
2342    axum::Extension(reads): axum::Extension<crate::read_work::ReadWork>,
2343    State(server): State<Arc<ServerState>>,
2344) -> Result<Json<Value>, ApiError> {
2345    reads.run(move || Ok(Json(server.host.queue_state()))).await
2346}
2347
2348/// Validate the URL id with the same traversal rules as the read endpoints.
2349fn valid_id(server: &ServerState, id: &str) -> Result<String, ApiError> {
2350    crate::rest::mission_paths(server, id)?;
2351    Ok(id.to_string())
2352}
2353
2354pub(crate) fn parse_body(body: &Bytes) -> Result<Value, ApiError> {
2355    if body.is_empty() {
2356        return Ok(json!({}));
2357    }
2358    serde_json::from_slice(body)
2359        .map_err(|e| ApiError::bad_request(format!("invalid JSON body: {e}")))
2360}
2361
2362// ---------------------------------------------------------------------------
2363// Unit tests: try_lock contention (deterministic — the HTTP-level race of
2364// two concurrent in-flight turns is covered here instead, by holding the
2365// per-mission mutex exactly like an in-flight turn does)
2366// ---------------------------------------------------------------------------
2367
2368#[cfg(test)]
2369mod tests {
2370    use super::*;
2371    use axum::http::StatusCode;
2372    use kranz_engine::backend_mock::{mock_init, mock_result_text, MockBackend, MockScript};
2373    use std::process::Command;
2374    use std::sync::Once;
2375
2376    static ENV_ISOLATION: Once = Once::new();
2377
2378    /// Mask the host's global/system git config (same discipline as the
2379    /// engine's mission tests) AND the home directory: `MissionHost::create`
2380    /// goes through `config::load`, which would otherwise read the
2381    /// developer's real ~/.kranz/config.json.
2382    fn isolate_git_env() {
2383        ENV_ISOLATION.call_once(|| {
2384            let missing = std::env::temp_dir()
2385                .join(format!("kranz-host-test-no-config-{}", std::process::id()));
2386            std::env::set_var("GIT_CONFIG_GLOBAL", &missing);
2387            std::env::set_var("GIT_CONFIG_SYSTEM", &missing);
2388            if let Ok(ceiling) = std::fs::canonicalize(std::env::temp_dir()) {
2389                std::env::set_var("GIT_CEILING_DIRECTORIES", ceiling);
2390            }
2391            let home =
2392                std::env::temp_dir().join(format!("kranz-host-test-home-{}", std::process::id()));
2393            let _ = std::fs::create_dir_all(&home);
2394            std::env::set_var(if cfg!(windows) { "USERPROFILE" } else { "HOME" }, &home);
2395        });
2396    }
2397
2398    fn git(dir: &std::path::Path, args: &[&str]) {
2399        let out = Command::new("git")
2400            .args(args)
2401            .current_dir(dir)
2402            .output()
2403            .expect("spawn git");
2404        assert!(
2405            out.status.success(),
2406            "git {args:?}: {}",
2407            String::from_utf8_lossy(&out.stderr)
2408        );
2409    }
2410
2411    /// Throwaway repo with one commit; `None` (skip) when git is missing.
2412    fn init_repo() -> Option<(tempfile::TempDir, PathBuf)> {
2413        isolate_git_env();
2414        let git_works = Command::new("git")
2415            .arg("--version")
2416            .output()
2417            .map(|o| o.status.success())
2418            .unwrap_or(false);
2419        if !git_works {
2420            kranz_engine::test_capability::skip(
2421                kranz_engine::test_capability::capability::GIT,
2422                "git is not on PATH",
2423            );
2424            return None;
2425        }
2426        let dir = tempfile::tempdir().expect("tempdir");
2427        let init = Command::new("git")
2428            .args(["init", "-b", "main"])
2429            .current_dir(dir.path())
2430            .output()
2431            .expect("spawn git init");
2432        if !init.status.success() {
2433            git(dir.path(), &["init"]);
2434            git(dir.path(), &["symbolic-ref", "HEAD", "refs/heads/main"]);
2435        }
2436        git(dir.path(), &["config", "user.name", "test"]);
2437        git(dir.path(), &["config", "user.email", "test@example.com"]);
2438        std::fs::write(dir.path().join("README.md"), "seed\n").unwrap();
2439        git(dir.path(), &["add", "-A"]);
2440        git(dir.path(), &["commit", "-m", "seed"]);
2441        let root = std::fs::canonicalize(dir.path()).expect("canonicalize");
2442        Some((dir, root))
2443    }
2444
2445    #[tokio::test]
2446    async fn ask_runs_read_only_one_shot_without_creating_mission_state() {
2447        let Some((_dir, root)) = init_repo() else {
2448            return;
2449        };
2450        let backend = Arc::new(MockBackend::with_scripts(vec![MockScript::single_shot(
2451            "Nothing is currently blocked.",
2452        )]));
2453        let host = MissionHost::with_backend(root.clone(), backend.clone());
2454
2455        let before = MissionPaths::list_missions(&root);
2456        let value = host.ask("what is blocked?").await.unwrap();
2457
2458        assert_eq!(value["answer"], "Nothing is currently blocked.");
2459        assert_eq!(
2460            MissionPaths::list_missions(&root),
2461            before,
2462            "ask must not create or mutate mission directories"
2463        );
2464        let specs = backend.started_specs();
2465        assert_eq!(specs.len(), 1);
2466        assert!(!specs[0].writable, "ask session is read-only");
2467        assert_eq!(specs[0].permission_mode.as_deref(), Some("plan"));
2468        let prompt = match &specs[0].prompt {
2469            PromptMode::SingleShot(prompt) => prompt,
2470            other => panic!("ask must be one-shot, got {other:?}"),
2471        };
2472        assert!(prompt.contains("what is blocked?"));
2473        assert!(prompt.contains("## Missions"));
2474    }
2475
2476    #[tokio::test]
2477    async fn http_api_error_codes_match_dashboard_wire_fixtures() {
2478        use http_body_util::BodyExt as _;
2479
2480        let dir = tempfile::tempdir().unwrap();
2481        let paths = MissionPaths::new(dir.path(), "m-fixture");
2482        std::fs::create_dir_all(paths.mission_dir()).unwrap();
2483        std::fs::write(paths.events_file(), "").unwrap();
2484        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2485        let mut host = MissionHost::with_backend(dir.path().to_path_buf(), backend);
2486        host.global_run_permits = Some(Arc::new(Semaphore::new(0)));
2487
2488        // Exercise real host refusal paths, then the same IntoResponse used by
2489        // Axum. Both languages consume this fixture; client-only mocks cannot
2490        // silently invent a wire shape the server never emits.
2491        let errors = [
2492            ("mission_not_hosted", host.not_hosted("m-fixture")),
2493            ("turn_in_flight", turn_in_flight()),
2494            ("repository_busy", host.try_global_run_permit().unwrap_err()),
2495            (
2496                "stale_plan",
2497                host.approve_pending("m-fixture", None).await.unwrap_err(),
2498            ),
2499            ("legacy", ApiError::conflict("mission is not hosted")),
2500        ];
2501        let mut actual = Vec::new();
2502        for (name, error) in errors {
2503            let response = error.into_response();
2504            let status = response.status().as_u16();
2505            assert_eq!(response.headers()["content-type"], "application/json");
2506            let bytes = response.into_body().collect().await.unwrap().to_bytes();
2507            let body: Value = serde_json::from_slice(&bytes).unwrap();
2508            actual.push(json!({ "name": name, "status": status, "body": body }));
2509        }
2510        let fixture: Value = serde_json::from_str(include_str!(
2511            "../../../apps/dashboard/src/lib/fixtures/api-errors.json"
2512        ))
2513        .unwrap();
2514        assert_eq!(json!(actual), fixture);
2515    }
2516
2517    #[tokio::test]
2518    async fn contended_planning_mutex_is_409_for_turns_and_start() {
2519        let Some((_dir, root)) = init_repo() else {
2520            return;
2521        };
2522        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2523        let host = MissionHost::with_backend(root, backend);
2524        let id = host.create("ship it", None).await.expect("create mission");
2525
2526        // Hold the per-mission engine mutex exactly like an in-flight turn.
2527        let cell = host.planning_cell(&id).expect("hosted planning cell");
2528        let _guard = cell.try_lock().expect("uncontended lock");
2529
2530        let err = host
2531            .planning_turn(&id, "hello")
2532            .await
2533            .expect_err("turn must 409");
2534        assert_eq!(err.status, StatusCode::CONFLICT);
2535        assert!(err.message.contains("turn is in flight"), "{}", err.message);
2536        assert_eq!(err.code, Some(ApiErrorCode::TurnInFlight));
2537
2538        let err = host
2539            .request_plan(&id)
2540            .await
2541            .expect_err("request-plan must 409");
2542        assert_eq!(err.status, StatusCode::CONFLICT);
2543
2544        // `start` also refuses while a turn holds the engine (the Arc clone
2545        // keeps try_unwrap failing) — and the entry survives the attempt.
2546        let err = host.start(&id).await.expect_err("start must 409");
2547        assert_eq!(err.status, StatusCode::CONFLICT);
2548        assert!(
2549            host.planning_cell(&id).is_ok(),
2550            "registry entry must survive"
2551        );
2552    }
2553
2554    #[tokio::test]
2555    async fn start_without_an_approved_plan_is_409() {
2556        let Some((_dir, root)) = init_repo() else {
2557            return;
2558        };
2559        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2560        let host = MissionHost::with_backend(root, backend);
2561        let id = host.create("ship it", None).await.expect("create mission");
2562
2563        let err = host
2564            .start(&id)
2565            .await
2566            .expect_err("start must 409 in planning");
2567        assert_eq!(err.status, StatusCode::CONFLICT);
2568        assert!(err.message.contains("no approved plan"), "{}", err.message);
2569        // The engine went back into the registry: planning can continue.
2570        assert!(host.planning_cell(&id).is_ok());
2571    }
2572
2573    /// M-13 (follow-up review): the Slack bridge used to read the parked
2574    /// plan's identity and then call `try_approve_pending`: two independent
2575    /// lock takes, with concurrent spawned tasks free to park a different
2576    /// plan in between. This variant compares and takes under one
2577    /// acquisition, and a refusal must leave the plan parked so the right
2578    /// card can still approve it.
2579    #[tokio::test]
2580    async fn approve_pending_matching_refuses_a_different_plan_without_consuming_it() {
2581        let Some((_dir, root)) = init_repo() else {
2582            return;
2583        };
2584        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2585        let host = MissionHost::with_backend(root.clone(), backend);
2586        let id = host.create("ship it", None).await.expect("create mission");
2587        let plan: Plan = serde_json::from_value(plan_json()).expect("plan");
2588        let identity = plan_identity(&plan);
2589
2590        assert_eq!(
2591            host.try_approve_pending_matching(&id, Some(&identity))
2592                .await
2593                .unwrap(),
2594            PendingApproval::NothingParked,
2595            "nothing parked is not an approval"
2596        );
2597
2598        host.set_pending_plan(&id, Some(plan.clone()));
2599
2600        assert_eq!(
2601            host.try_approve_pending_matching(&id, Some("an older plan"))
2602                .await
2603                .unwrap(),
2604            PendingApproval::Mismatch {
2605                parked: identity.clone()
2606            },
2607            "a card naming a different plan must be refused, naming the parked one"
2608        );
2609        assert!(
2610            host.pending_plan(&id).is_some(),
2611            "a refused approve must not consume the parked plan"
2612        );
2613
2614        assert!(
2615            matches!(
2616                host.try_approve_pending_matching(&id, None).await.unwrap(),
2617                PendingApproval::Mismatch { .. }
2618            ),
2619            "a card that names no plan cannot match one"
2620        );
2621        assert!(host.pending_plan(&id).is_some());
2622
2623        assert_eq!(
2624            host.try_approve_pending_matching(&id, Some(&identity))
2625                .await
2626                .unwrap(),
2627            PendingApproval::Approved(format!("kranz/mission-{id}"))
2628        );
2629        assert!(
2630            host.pending_plan(&id).is_none(),
2631            "an approve consumes the parked plan"
2632        );
2633        assert_eq!(
2634            host.try_approve_pending_matching(&id, Some(&identity))
2635                .await
2636                .unwrap(),
2637            PendingApproval::NothingParked,
2638            "a second click has nothing left to commit"
2639        );
2640    }
2641
2642    #[tokio::test]
2643    async fn approve_pending_matching_leaves_pending_untouched_on_busy_or_failure() {
2644        let Some((_dir, root)) = init_repo() else {
2645            return;
2646        };
2647        let host = MissionHost::with_backend(root, Arc::new(MockBackend::new()));
2648        let id = host.create("ship it", None).await.unwrap();
2649        let mut plan: Plan = serde_json::from_value(plan_json()).unwrap();
2650        host.set_pending_plan(&id, Some(plan.clone()));
2651        let cell = host.planning_cell(&id).unwrap();
2652        let guard = cell.try_lock().unwrap();
2653        let identity = plan_identity(&plan);
2654        let err = host
2655            .try_approve_pending_matching(&id, Some(&identity))
2656            .await
2657            .unwrap_err();
2658        assert_eq!(err.status, StatusCode::CONFLICT);
2659        assert_eq!(plan_identity(&host.pending_plan(&id).unwrap()), identity);
2660        drop(guard);
2661
2662        plan.milestones.clear();
2663        let invalid_identity = plan_identity(&plan);
2664        host.set_pending_plan(&id, Some(plan));
2665        let err = host
2666            .try_approve_pending_matching(&id, Some(&invalid_identity))
2667            .await
2668            .unwrap_err();
2669        assert!(err.message.contains("no milestones"), "{}", err.message);
2670        assert_eq!(
2671            plan_identity(&host.pending_plan(&id).unwrap()),
2672            invalid_identity
2673        );
2674
2675        // A later replacement survives an old retry after the failed approval.
2676        let replacement: Plan = serde_json::from_value(plan_json()).unwrap();
2677        host.set_pending_plan(&id, Some(replacement));
2678        assert_eq!(
2679            host.try_approve_pending_matching(&id, Some(&invalid_identity))
2680                .await
2681                .unwrap(),
2682            PendingApproval::Mismatch {
2683                parked: identity.clone()
2684            },
2685        );
2686        assert_eq!(plan_identity(&host.pending_plan(&id).unwrap()), identity);
2687    }
2688
2689    #[tokio::test]
2690    async fn start_is_409_when_repo_busy() {
2691        let Some((_dir, root)) = init_repo() else {
2692            return;
2693        };
2694        // Hold the repo busy lock BEFORE hosting a mission — a live
2695        // events.jsonl.lock for the mission under test would block a
2696        // sibling acquire (legacy probe), so take the hold first.
2697        let _hold =
2698            kranz_engine::queue::acquire_repo_busy(&root, "m-sibling").expect("sibling busy hold");
2699        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2700        let host = MissionHost::with_backend(root.clone(), backend);
2701        let id = host.create("ship it", None).await.expect("create mission");
2702        let plan: Plan = serde_json::from_value(plan_json()).expect("plan");
2703        host.approve(&id, plan).await.expect("approve");
2704
2705        let err = host.start(&id).await.expect_err("start must 409 when busy");
2706        assert_eq!(err.status, StatusCode::CONFLICT);
2707        assert_eq!(err.code, Some(ApiErrorCode::RepositoryBusy));
2708        assert!(
2709            err.message.contains("busy"),
2710            "expected busy conflict, got: {}",
2711            err.message
2712        );
2713        // Engine restored to the registry so the operator can retry.
2714        assert!(host.planning_cell(&id).is_ok());
2715    }
2716
2717    #[tokio::test]
2718    async fn start_is_409_when_global_repository_limit_is_saturated() {
2719        let Some((_dir, root)) = init_repo() else {
2720            return;
2721        };
2722        let permits = Arc::new(Semaphore::new(1));
2723        let _other_repo = Arc::clone(&permits).try_acquire_owned().unwrap();
2724        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2725        let mut host = MissionHost::with_backend(root, backend);
2726        host.global_run_permits = Some(permits);
2727        let id = host.create("ship it", None).await.expect("create mission");
2728        let plan: Plan = serde_json::from_value(plan_json()).expect("plan");
2729        host.approve(&id, plan).await.expect("approve");
2730
2731        let error = host.start(&id).await.expect_err("global cap must refuse");
2732
2733        assert_eq!(error.status, StatusCode::CONFLICT);
2734        assert!(error.message.contains("maxConcurrentRepos"));
2735        assert_eq!(error.code, Some(ApiErrorCode::RepositoryBusy));
2736        assert!(
2737            host.planning_cell(&id).is_ok(),
2738            "refused start must restore the hosted engine"
2739        );
2740    }
2741
2742    #[tokio::test]
2743    async fn global_run_permit_is_released_when_hosted_task_panics() {
2744        let permits = Arc::new(Semaphore::new(1));
2745        let permit = Arc::clone(&permits).try_acquire_owned().unwrap();
2746        assert_eq!(permits.available_permits(), 0);
2747
2748        let handle = spawn_with_global_run_permit(Some(permit), async {
2749            panic!("simulated hosted-run panic");
2750        });
2751        assert!(handle.await.unwrap_err().is_panic());
2752
2753        assert_eq!(permits.available_permits(), 1);
2754    }
2755
2756    #[tokio::test]
2757    async fn sweep_idle_leaves_a_mid_turn_mission_hosted() {
2758        let Some((_dir, root)) = init_repo() else {
2759            return;
2760        };
2761        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2762        let host = MissionHost::with_backend(root, backend);
2763        let id = host.create("ship it", None).await.expect("create mission");
2764
2765        // Hold the per-mission engine mutex exactly like an in-flight turn.
2766        let cell = host.planning_cell(&id).expect("hosted planning cell");
2767        let _guard = cell.try_lock().expect("uncontended lock");
2768
2769        let released = host.sweep_idle(std::time::Duration::ZERO);
2770        assert!(!released.contains(&id), "{released:?}");
2771        assert!(
2772            host.planning_cell(&id).is_ok(),
2773            "mission must remain hosted"
2774        );
2775    }
2776
2777    #[tokio::test]
2778    async fn release_route_is_409_mid_turn() {
2779        use axum::body::Body;
2780        use axum::http::Request;
2781        use tower::ServiceExt;
2782
2783        let Some((_dir, root)) = init_repo() else {
2784            return;
2785        };
2786        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2787        let host = MissionHost::with_backend(root, backend);
2788        let id = host.create("ship it", None).await.expect("create mission");
2789
2790        // Hold the per-mission engine mutex exactly like an in-flight turn.
2791        let cell = host.planning_cell(&id).expect("hosted planning cell");
2792        let _guard = cell.try_lock().expect("uncontended lock");
2793
2794        let app =
2795            crate::router_with_host(host, None, crate::MutationAuthority::new("tok").unwrap());
2796        let response = app
2797            .oneshot(
2798                Request::builder()
2799                    .method("POST")
2800                    .uri(format!("/api/missions/{id}/release"))
2801                    .header("content-type", "application/json")
2802                    .header("x-kranz-token", "tok")
2803                    .body(Body::from("{}"))
2804                    .unwrap(),
2805            )
2806            .await
2807            .unwrap();
2808        assert_eq!(response.status(), StatusCode::CONFLICT);
2809    }
2810
2811    #[tokio::test]
2812    async fn bodyless_post_with_valid_token_is_not_rejected_as_unsupported_media_type() {
2813        use axum::body::Body;
2814        use axum::http::Request;
2815        use tower::ServiceExt;
2816
2817        let Some((_dir, root)) = init_repo() else {
2818            return;
2819        };
2820        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2821        let host = MissionHost::with_backend(root, backend);
2822        let id = host.create("ship it", None).await.expect("create mission");
2823
2824        let app =
2825            crate::router_with_host(host, None, crate::MutationAuthority::new("tok").unwrap());
2826        let response = app
2827            .oneshot(
2828                Request::builder()
2829                    .method("POST")
2830                    .uri(format!("/api/missions/{id}/start"))
2831                    // No content-type and no content-length: an empty
2832                    // bodyless POST, the case `curl -X POST .../start`
2833                    // (no `-d`) sends.
2834                    .header("x-kranz-token", "tok")
2835                    .body(Body::empty())
2836                    .unwrap(),
2837            )
2838            .await
2839            .unwrap();
2840        // A freshly created mission has no approved plan, so `start` 409s —
2841        // the point of this test is that it is NOT 415, i.e. the missing
2842        // content-type on an empty body no longer trips the media-type gate.
2843        assert_ne!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE);
2844        assert_eq!(response.status(), StatusCode::CONFLICT);
2845    }
2846
2847    #[tokio::test]
2848    async fn bodyless_post_gate_still_rejects_non_empty_non_json_bodies() {
2849        use axum::body::Body;
2850        use axum::http::Request;
2851        use tower::ServiceExt;
2852
2853        let Some((_dir, root)) = init_repo() else {
2854            return;
2855        };
2856        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2857        let host = MissionHost::with_backend(root, backend);
2858        let id = host.create("ship it", None).await.expect("create mission");
2859
2860        let app =
2861            crate::router_with_host(host, None, crate::MutationAuthority::new("tok").unwrap());
2862        let payload = "not json";
2863        let response = app
2864            .oneshot(
2865                Request::builder()
2866                    .method("POST")
2867                    .uri(format!("/api/missions/{id}/release"))
2868                    .header("content-type", "text/plain")
2869                    .header("content-length", payload.len().to_string())
2870                    .header("x-kranz-token", "tok")
2871                    .body(Body::from(payload))
2872                    .unwrap(),
2873            )
2874            .await
2875            .unwrap();
2876        assert_eq!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE);
2877    }
2878
2879    #[tokio::test]
2880    async fn create_rejects_an_invalid_config_patch() {
2881        let Some((_dir, root)) = init_repo() else {
2882            return;
2883        };
2884        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2885        let host = MissionHost::with_backend(root, backend);
2886
2887        // 9 is out of the 1..=8 range config::validate allows (M3), so the
2888        // create must be rejected as a bad request. (2..=8 is now valid — it
2889        // opts into parallel workers — so an out-of-range value is used here.)
2890        let patch = json!({ "maxParallelWorkers": 9 });
2891        let err = host
2892            .create("ship it", Some(&patch))
2893            .await
2894            .expect_err("must reject");
2895        assert_eq!(err.status, StatusCode::BAD_REQUEST);
2896    }
2897
2898    // -----------------------------------------------------------------------
2899    // Queue drain (roadmap f-1-2)
2900    // -----------------------------------------------------------------------
2901
2902    #[tokio::test]
2903    async fn empty_queue_drain_returns_ok_and_settles_idle() {
2904        let Some((_dir, root)) = init_repo() else {
2905            return;
2906        };
2907        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2908        let host = MissionHost::with_backend(root, backend);
2909
2910        let body = host
2911            .drain()
2912            .await
2913            .expect("drain must not error on an empty queue");
2914        assert!(body.get("live").is_some(), "{body}");
2915
2916        // The background task finds nothing queued and settles quickly.
2917        let deadline = std::time::Instant::now() + Duration::from_secs(5);
2918        loop {
2919            let state = host.queue_state();
2920            if state["drain"]["live"] == false {
2921                break;
2922            }
2923            assert!(
2924                std::time::Instant::now() < deadline,
2925                "drain never settled idle: {state}"
2926            );
2927            tokio::time::sleep(Duration::from_millis(20)).await;
2928        }
2929    }
2930
2931    #[tokio::test]
2932    async fn drain_is_409_when_global_repository_limit_is_saturated() {
2933        let Some((_dir, root)) = init_repo() else {
2934            return;
2935        };
2936        let permits = Arc::new(Semaphore::new(1));
2937        let _other_repo = Arc::clone(&permits).try_acquire_owned().unwrap();
2938        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2939        let mut host = MissionHost::with_backend(root, backend);
2940        host.global_run_permits = Some(permits);
2941
2942        let error = host.drain().await.expect_err("global cap must refuse");
2943
2944        assert_eq!(error.status, StatusCode::CONFLICT);
2945        assert!(error.message.contains("maxConcurrentRepos"));
2946        assert!(matches!(
2947            &*host.drain.lock().expect("drain tracker lock"),
2948            DrainSlot::Idle
2949        ));
2950    }
2951
2952    #[tokio::test]
2953    async fn queue_state_reports_global_concurrency_saturation() {
2954        let Some((_dir, root)) = init_repo() else {
2955            return;
2956        };
2957        let permits = Arc::new(Semaphore::new(1));
2958        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2959        let mut host = MissionHost::with_backend(root, backend);
2960        host.global_run_permits = Some(Arc::clone(&permits));
2961
2962        let open = host.queue_state();
2963        assert_eq!(open["maxConcurrentReposAvailable"], 1);
2964        assert_eq!(open["maxConcurrentReposSaturated"], false);
2965
2966        let _hold = permits.try_acquire_owned().unwrap();
2967        let saturated = host.queue_state();
2968        assert_eq!(saturated["maxConcurrentReposAvailable"], 0);
2969        assert_eq!(saturated["maxConcurrentReposSaturated"], true);
2970    }
2971
2972    #[tokio::test]
2973    async fn second_drain_while_live_returns_tracked_state_without_spawning_second() {
2974        let Some((_dir, root)) = init_repo() else {
2975            return;
2976        };
2977        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
2978        let host = MissionHost::with_backend(root, backend);
2979
2980        // Fabricate a live drain tracker directly — deterministic, instead
2981        // of racing a real queue against a fast mock backend.
2982        let state = Arc::new(Mutex::new(DrainState {
2983            live: true,
2984            current_mission_id: Some("m-fake".to_string()),
2985            ran: vec!["m-earlier".to_string()],
2986            parked: Vec::new(),
2987        }));
2988        let never_finishes = tokio::spawn(async {
2989            std::future::pending::<()>().await;
2990        });
2991        *host.drain.lock().expect("drain tracker lock") = DrainSlot::Running(DrainHandle {
2992            join: never_finishes,
2993            state: Arc::clone(&state),
2994        });
2995        let before = Arc::as_ptr(&state);
2996
2997        let first = host.drain().await.expect("drain must not error");
2998        let second = host.drain().await.expect("drain must not error");
2999        assert_eq!(first, second);
3000        assert_eq!(first["live"], true);
3001        assert_eq!(first["currentMissionId"], "m-fake");
3002        assert_eq!(first["ran"], json!(["m-earlier"]));
3003
3004        // The tracker still points at the SAME state Arc: no second task
3005        // was spawned to replace it.
3006        let after = {
3007            let guard = host.drain.lock().expect("drain tracker lock");
3008            match &*guard {
3009                DrainSlot::Running(handle) => Arc::as_ptr(&handle.state),
3010                _ => panic!("expected the tracker to still be Running"),
3011            }
3012        };
3013        assert_eq!(before, after, "a second drain must not replace the tracker");
3014    }
3015
3016    #[tokio::test]
3017    async fn two_concurrent_cold_drains_spawn_exactly_one() {
3018        let Some((_dir, root)) = init_repo() else {
3019            return;
3020        };
3021        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
3022        let host = MissionHost::with_backend(root, backend);
3023
3024        // Fire two drains concurrently from a cold (Idle) tracker. Neither
3025        // `config::load` nor `self.backend(...)` yields here (the backend
3026        // is pre-populated via `with_backend`, and this test runs on the
3027        // default current-thread flavor), so the first call's poll runs
3028        // synchronously all the way through installing the `Starting`
3029        // reservation, spawning the task, and upgrading to `Running` before
3030        // the second call is ever polled. The second call therefore always
3031        // observes an in-progress drain (`Starting` or `Running`, task not
3032        // yet scheduled) and returns its tracked state instead of spawning
3033        // a second drain task.
3034        //
3035        // NOTE: because nothing yields here, this test alone cannot catch a
3036        // regression that deletes the `DrainSlot::Starting` deflection arm —
3037        // see `starting_reservation_is_not_overwritten_or_double_spawned`
3038        // below for the deterministic test that actually guards that arm.
3039        let (first, second) = tokio::join!(host.drain(), host.drain());
3040        let first = first.expect("first drain must not error");
3041        let second = second.expect("second drain must not error");
3042        assert_eq!(first["live"], true, "{first}");
3043        assert_eq!(second["live"], true, "{second}");
3044
3045        // Exactly one drain is tracked: a single Starting-or-Running slot,
3046        // never two independently spawned tasks.
3047        match &*host.drain.lock().expect("drain tracker lock") {
3048            DrainSlot::Running(_) | DrainSlot::Starting(_) => {}
3049            DrainSlot::Idle => {
3050                panic!("expected a live drain to be tracked after two concurrent calls")
3051            }
3052        }
3053
3054        // The single tracked drain settles idle on its own — nothing is
3055        // left running forever, which would indicate a leaked second task.
3056        let deadline = std::time::Instant::now() + Duration::from_secs(5);
3057        loop {
3058            let state = host.queue_state();
3059            if state["drain"]["live"] == false {
3060                break;
3061            }
3062            assert!(
3063                std::time::Instant::now() < deadline,
3064                "drain never settled idle: {state}"
3065            );
3066            tokio::time::sleep(Duration::from_millis(20)).await;
3067        }
3068    }
3069
3070    // -----------------------------------------------------------------------
3071    // Hosted-drain checkout capture/restore (mirrors `restore_work_checkout`
3072    // in `crates/cli/src/backlog.rs` — see `drain_task`/`restore_drain_checkout`)
3073    // -----------------------------------------------------------------------
3074
3075    /// A fresh `DrainState` and one queued entry for `mission_id`, ready to
3076    /// feed [`drain_task`] directly (bypassing `tokio::spawn` for a
3077    /// deterministic test).
3078    fn seed_one_queued(root: &Path, mission_id: &str) -> Arc<Mutex<DrainState>> {
3079        kranz_engine::queue::enqueue(
3080            root,
3081            kranz_engine::queue::QueueEntry {
3082                mission_id: mission_id.to_string(),
3083                ticket_slug: None,
3084                priority: 5,
3085                seq: 0,
3086            },
3087        )
3088        .expect("enqueue");
3089        Arc::new(Mutex::new(DrainState::default()))
3090    }
3091
3092    fn proceed_readiness(
3093        _repo_root: &Path,
3094        mission_id: &str,
3095    ) -> kranz_engine::error::Result<kranz_engine::backend_readiness::ReadinessReport> {
3096        Ok(kranz_engine::backend_readiness::ReadinessReport {
3097            mission_id: mission_id.to_string(),
3098            roles: Vec::new(),
3099            overall: kranz_engine::backend_readiness::ReadinessStatus::Ok,
3100            warnings: Vec::new(),
3101        })
3102    }
3103
3104    #[tokio::test]
3105    async fn auto_work_drain_mode_processes_only_one_queue_front() {
3106        let Some((_dir, root)) = init_repo() else {
3107            return;
3108        };
3109        let state = seed_one_queued(&root, "m-first");
3110        kranz_engine::queue::enqueue(
3111            &root,
3112            kranz_engine::queue::QueueEntry {
3113                mission_id: "m-second".to_string(),
3114                ticket_slug: None,
3115                priority: 5,
3116                seq: 0,
3117            },
3118        )
3119        .expect("enqueue second mission");
3120
3121        drain_task_with_probe(
3122            root.clone(),
3123            Arc::clone(&state),
3124            true,
3125            |_mission_id| async { Ok(0) },
3126            proceed_readiness,
3127        )
3128        .await;
3129
3130        assert_eq!(state.lock().expect("drain state lock").ran, ["m-first"]);
3131        let remaining = kranz_engine::queue::list(&root);
3132        assert_eq!(remaining.len(), 1);
3133        assert_eq!(remaining[0].mission_id, "m-second");
3134    }
3135
3136    #[tokio::test]
3137    async fn hosted_drain_restores_dispatch_checkout() {
3138        let Some((_dir, root)) = init_repo() else {
3139            return;
3140        };
3141        let state = seed_one_queued(&root, "m-restore");
3142
3143        let run_root = root.clone();
3144        drain_task_with_probe(
3145            root.clone(),
3146            Arc::clone(&state),
3147            false,
3148            move |mission_id| {
3149                let root = run_root.clone();
3150                async move {
3151                    let git = GitRepo::open(&root)?;
3152                    let branch = format!("kranz/mission-{mission_id}");
3153                    git.create_branch(&branch, None)?;
3154                    git.checkout(&branch)?;
3155                    Ok(0)
3156                }
3157            },
3158            proceed_readiness,
3159        )
3160        .await;
3161
3162        assert_eq!(
3163            state.lock().expect("drain state lock").ran,
3164            ["m-restore"],
3165            "the injected mission runner must execute"
3166        );
3167
3168        let git = GitRepo::open(&root).expect("open repo");
3169        assert_eq!(
3170            git.current_branch().expect("current branch"),
3171            "main",
3172            "the operator's dispatch-time checkout must be restored on drain exit"
3173        );
3174    }
3175
3176    #[tokio::test]
3177    async fn hosted_drain_restores_dispatch_checkout_on_err() {
3178        let Some((_dir, root)) = init_repo() else {
3179            return;
3180        };
3181        let state = seed_one_queued(&root, "m-err-restore");
3182
3183        let run_root = root.clone();
3184        let runner_called = Arc::new(std::sync::atomic::AtomicBool::new(false));
3185        let called = Arc::clone(&runner_called);
3186        drain_task_with_probe(
3187            root.clone(),
3188            state,
3189            false,
3190            move |mission_id| {
3191                let root = run_root.clone();
3192                let called = Arc::clone(&called);
3193                async move {
3194                    called.store(true, std::sync::atomic::Ordering::SeqCst);
3195                    let git = GitRepo::open(&root)?;
3196                    let branch = format!("kranz/mission-{mission_id}");
3197                    git.create_branch(&branch, None)?;
3198                    git.checkout(&branch)?;
3199                    Err(anyhow::anyhow!("simulated drain runner failure"))
3200                }
3201            },
3202            proceed_readiness,
3203        )
3204        .await;
3205
3206        assert!(
3207            runner_called.load(std::sync::atomic::Ordering::SeqCst),
3208            "the injected mission runner must execute"
3209        );
3210
3211        let git = GitRepo::open(&root).expect("open repo");
3212        assert_eq!(
3213            git.current_branch().expect("current branch"),
3214            "main",
3215            "an errored drain must still restore the operator's dispatch-time checkout"
3216        );
3217    }
3218
3219    #[tokio::test]
3220    async fn hosted_drain_skips_restore_when_started_on_mission_branch() {
3221        let Some((_dir, root)) = init_repo() else {
3222            return;
3223        };
3224        {
3225            let git = GitRepo::open(&root).expect("open repo");
3226            git.create_branch("kranz/mission-existing", None)
3227                .expect("create existing mission branch");
3228            git.checkout("kranz/mission-existing")
3229                .expect("checkout existing mission branch");
3230        }
3231        let state = seed_one_queued(&root, "m-skip");
3232
3233        drain_task_with_probe(
3234            root.clone(),
3235            Arc::clone(&state),
3236            false,
3237            |_mission_id| async { Ok(0) },
3238            proceed_readiness,
3239        )
3240        .await;
3241
3242        assert_eq!(
3243            state.lock().expect("drain state lock").ran,
3244            ["m-skip"],
3245            "the injected mission runner must execute"
3246        );
3247
3248        let git = GitRepo::open(&root).expect("open repo");
3249        assert_eq!(
3250            git.current_branch().expect("current branch"),
3251            "kranz/mission-existing",
3252            "started on a mission branch: no restore must be attempted"
3253        );
3254    }
3255
3256    #[tokio::test]
3257    async fn hosted_drain_leaves_checkout_when_tracked_tree_dirty() {
3258        let Some((_dir, root)) = init_repo() else {
3259            return;
3260        };
3261        let state = seed_one_queued(&root, "m-dirty");
3262
3263        let run_root = root.clone();
3264        drain_task_with_probe(
3265            root.clone(),
3266            Arc::clone(&state),
3267            false,
3268            move |mission_id| {
3269                let root = run_root.clone();
3270                async move {
3271                    let git = GitRepo::open(&root)?;
3272                    let branch = format!("kranz/mission-{mission_id}");
3273                    git.create_branch(&branch, None)?;
3274                    git.checkout(&branch)?;
3275                    std::fs::write(root.join("README.md"), "dirty tracked edit\n")?;
3276                    Ok(0)
3277                }
3278            },
3279            proceed_readiness,
3280        )
3281        .await;
3282
3283        assert_eq!(
3284            state.lock().expect("drain state lock").ran,
3285            ["m-dirty"],
3286            "the injected mission runner must execute"
3287        );
3288
3289        let git = GitRepo::open(&root).expect("open repo");
3290        assert_eq!(
3291            git.current_branch().expect("current branch"),
3292            "kranz/mission-m-dirty",
3293            "a dirty tracked tree must abort the restore, leaving the checkout on the mission \
3294             branch"
3295        );
3296    }
3297
3298    #[tokio::test]
3299    async fn hosted_drain_second_call_does_not_capture_or_restore() {
3300        let Some((_dir, root)) = init_repo() else {
3301            return;
3302        };
3303        {
3304            let git = GitRepo::open(&root).expect("open repo");
3305            git.create_branch("feature-branch", None)
3306                .expect("create feature branch");
3307            git.checkout("feature-branch")
3308                .expect("checkout feature branch");
3309        }
3310        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
3311        let host = MissionHost::with_backend(root.clone(), backend);
3312
3313        // Fabricate a live drain tracker (same pattern as
3314        // `second_drain_while_live_returns_tracked_state_without_spawning_second`)
3315        // so the idempotent early-return path is exercised without racing a
3316        // real spawn.
3317        let tracked_state = Arc::new(Mutex::new(DrainState {
3318            live: true,
3319            current_mission_id: Some("m-inflight".to_string()),
3320            ran: Vec::new(),
3321            parked: Vec::new(),
3322        }));
3323        let never_finishes = tokio::spawn(async {
3324            std::future::pending::<()>().await;
3325        });
3326        *host.drain.lock().expect("drain tracker lock") = DrainSlot::Running(DrainHandle {
3327            join: never_finishes,
3328            state: Arc::clone(&tracked_state),
3329        });
3330
3331        let result = host
3332            .drain()
3333            .await
3334            .expect("second drain call must not error");
3335        assert_eq!(result["live"], true, "{result}");
3336
3337        // No capture/restore happened: the checkout this test set up before
3338        // the second call is untouched.
3339        let git = GitRepo::open(&root).expect("open repo");
3340        assert_eq!(
3341            git.current_branch().expect("current branch"),
3342            "feature-branch",
3343            "the idempotent second drain() must not mutate the checkout"
3344        );
3345
3346        // No second task was spawned: the tracker still points at the same
3347        // state Arc installed above.
3348        match &*host.drain.lock().expect("drain tracker lock") {
3349            DrainSlot::Running(handle) => {
3350                assert_eq!(
3351                    Arc::as_ptr(&handle.state),
3352                    Arc::as_ptr(&tracked_state),
3353                    "a second drain must not replace the tracker or spawn a second task"
3354                );
3355            }
3356            _ => panic!("expected the tracker to still be Running"),
3357        };
3358    }
3359
3360    /// Deterministically guards the `DrainSlot::Starting(state) => return
3361    /// ...` deflection arm in [`MissionHost::drain`]: manually install a
3362    /// `Starting` reservation, call `drain()`, and assert it returns the
3363    /// tracked live state WITHOUT overwriting the slot or spawning a task.
3364    /// If that match arm is deleted (falling through to the Idle/Running
3365    /// catch-all), this test fails because the slot gets overwritten with a
3366    /// fresh `Starting`/`Running` reservation (different `Arc::as_ptr`) and a
3367    /// real drain task gets spawned against this test's (git-less) repo.
3368    #[tokio::test]
3369    async fn starting_reservation_is_not_overwritten_or_double_spawned() {
3370        let Some((_dir, root)) = init_repo() else {
3371            return;
3372        };
3373        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
3374        let host = MissionHost::with_backend(root, backend);
3375
3376        let state = Arc::new(Mutex::new(DrainState {
3377            live: true,
3378            current_mission_id: Some("m-reserved".to_string()),
3379            ran: Vec::new(),
3380            parked: Vec::new(),
3381        }));
3382        *host.drain.lock().expect("drain tracker lock") = DrainSlot::Starting(Arc::clone(&state));
3383        let before = Arc::as_ptr(&state);
3384
3385        let result = host.drain().await.expect("drain must not error");
3386        assert_eq!(result["live"], true, "{result}");
3387        assert_eq!(result["currentMissionId"], "m-reserved");
3388
3389        // The slot must STILL be the same Starting reservation: not
3390        // overwritten to a new Starting/Running, and no task spawned.
3391        let after = match &*host.drain.lock().expect("drain tracker lock") {
3392            DrainSlot::Starting(tracked) => Arc::as_ptr(tracked),
3393            DrainSlot::Running(_) => panic!(
3394                "the Starting reservation was upgraded/replaced by this call — the deflection \
3395                 arm was bypassed and a second drain was spawned"
3396            ),
3397            DrainSlot::Idle => panic!("the Starting reservation was cleared by this call"),
3398        };
3399        assert_eq!(
3400            before, after,
3401            "drain() must return the SAME tracked reservation, not install a new one"
3402        );
3403    }
3404
3405    /// A failed drain construction (here: an unparseable `.kranz/config.json`)
3406    /// must clear the reservation back to `Idle` so a later call can retry —
3407    /// otherwise every future drain would deflect forever onto a dead
3408    /// reservation that no task will ever settle.
3409    #[tokio::test]
3410    async fn failed_drain_construction_clears_the_reservation_to_idle() {
3411        let Some((_dir, root)) = init_repo() else {
3412            return;
3413        };
3414        std::fs::create_dir_all(root.join(".kranz")).expect("mkdir .kranz");
3415        std::fs::write(root.join(".kranz").join("config.json"), "not json")
3416            .expect("write malformed config");
3417
3418        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
3419        let host = MissionHost::with_backend(root, backend);
3420
3421        host.drain()
3422            .await
3423            .expect_err("malformed config must fail drain construction");
3424
3425        let is_idle = matches!(
3426            &*host.drain.lock().expect("drain tracker lock"),
3427            DrainSlot::Idle
3428        );
3429        assert!(
3430            is_idle,
3431            "a failed drain construction must reset the tracker to Idle"
3432        );
3433    }
3434
3435    /// `queue_state()` must report the transient `Starting` reservation
3436    /// window as a live drain — a caller polling `GET /api/queue` right after
3437    /// `POST /api/queue/drain` must not observe a false "not live" gap.
3438    #[tokio::test]
3439    async fn queue_state_reports_a_starting_reservation_as_live() {
3440        let Some((_dir, root)) = init_repo() else {
3441            return;
3442        };
3443        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
3444        let host = MissionHost::with_backend(root, backend);
3445
3446        let state = Arc::new(Mutex::new(DrainState {
3447            live: true,
3448            current_mission_id: Some("m-starting".to_string()),
3449            ran: Vec::new(),
3450            parked: Vec::new(),
3451        }));
3452        *host.drain.lock().expect("drain tracker lock") = DrainSlot::Starting(state);
3453
3454        let queue_state = host.queue_state();
3455        assert_eq!(queue_state["drain"]["live"], true, "{queue_state}");
3456        assert_eq!(queue_state["drain"]["currentMissionId"], "m-starting");
3457    }
3458
3459    #[tokio::test]
3460    async fn queue_drain_route_requires_token_but_queue_route_does_not() {
3461        use axum::body::Body;
3462        use axum::http::Request;
3463        use tower::ServiceExt;
3464
3465        let Some((_dir, root)) = init_repo() else {
3466            return;
3467        };
3468        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
3469        let host = MissionHost::with_backend(root, backend);
3470        let app =
3471            crate::router_with_host(host, None, crate::MutationAuthority::new("tok").unwrap());
3472
3473        let response = app
3474            .clone()
3475            .oneshot(
3476                Request::builder()
3477                    .method("POST")
3478                    .uri("/api/queue/drain")
3479                    .body(Body::empty())
3480                    .unwrap(),
3481            )
3482            .await
3483            .unwrap();
3484        assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
3485
3486        let response = app
3487            .clone()
3488            .oneshot(
3489                Request::builder()
3490                    .method("POST")
3491                    .uri("/api/queue/drain")
3492                    .header("x-kranz-token", "tok")
3493                    .body(Body::empty())
3494                    .unwrap(),
3495            )
3496            .await
3497            .unwrap();
3498        assert_ne!(response.status(), StatusCode::UNAUTHORIZED);
3499        assert_eq!(response.status(), StatusCode::OK);
3500
3501        let response = app
3502            .oneshot(
3503                Request::builder()
3504                    .uri("/api/queue")
3505                    .body(Body::empty())
3506                    .unwrap(),
3507            )
3508            .await
3509            .unwrap();
3510        assert_eq!(response.status(), StatusCode::OK);
3511    }
3512
3513    // -----------------------------------------------------------------------
3514    // autoWork watcher (roadmap f-2-3)
3515    // -----------------------------------------------------------------------
3516
3517    #[test]
3518    fn should_auto_drain_truth_table() {
3519        // Only true when all three conditions line up.
3520        assert!(should_auto_drain(true, true, false));
3521        // autoWork off: never drain, regardless of the queue or live state.
3522        assert!(!should_auto_drain(false, true, false));
3523        assert!(!should_auto_drain(false, false, false));
3524        // Queue empty: nothing to drain even with autoWork on.
3525        assert!(!should_auto_drain(true, false, false));
3526        // A drain is already live: never start a second one.
3527        assert!(!should_auto_drain(true, true, true));
3528        assert!(!should_auto_drain(false, false, true));
3529    }
3530
3531    /// Writes `{"autoWork": enabled}` to the repo's `.kranz/config.json`
3532    /// (the project config layer `config::load` reads on every call,
3533    /// including the watcher's per-tick reload).
3534    fn write_auto_work_config(root: &std::path::Path, enabled: bool) {
3535        let dir = root.join(".kranz");
3536        std::fs::create_dir_all(&dir).expect("create .kranz dir");
3537        std::fs::write(
3538            dir.join("config.json"),
3539            json!({ "autoWork": enabled }).to_string(),
3540        )
3541        .expect("write config.json");
3542    }
3543
3544    /// One orchestrator turn batch: text + matching Result (mirrors the
3545    /// identical helper in `tests/host_test.rs`).
3546    fn turn(reply: &str) -> Vec<kranz_engine::backend::AgentEvent> {
3547        vec![
3548            kranz_engine::backend_mock::mock_text(reply),
3549            mock_result_text(reply),
3550        ]
3551    }
3552
3553    /// A completed single-shot preflight probe session whose reply
3554    /// authenticates, consumed once by `MissionEngine::worker_auth_verdict`
3555    /// before the first worker/validator session of the mission spawns.
3556    fn preflight_authenticated_script() -> MockScript {
3557        MockScript::single_shot("ack")
3558    }
3559
3560    /// Worker script: completed single-shot run with a passing WorkerReport.
3561    fn worker_pass() -> MockScript {
3562        MockScript::single_shot_json(&json!({
3563            "result": "pass",
3564            "summary": "implemented and tested",
3565            "filesTouched": [],
3566            "testsAdded": [],
3567            "testEvidence": "all green",
3568            "commits": []
3569        }))
3570    }
3571
3572    /// A minimal one-milestone/one-feature plan in wire (camelCase) shape.
3573    fn plan_json() -> Value {
3574        json!({
3575            "goal": "ship the demo",
3576            "validationContract": [],
3577            "milestones": [{
3578                "title": "M1",
3579                "features": [{
3580                    "title": "F1",
3581                    "spec": "build the thing",
3582                    "validationCriteria": ["it works"]
3583                }]
3584            }]
3585        })
3586    }
3587
3588    #[tokio::test(flavor = "multi_thread")]
3589    async fn auto_work_tick_drains_a_queued_mission_when_enabled() {
3590        let Some((_dir, root)) = init_repo() else {
3591            return;
3592        };
3593        write_auto_work_config(&root, true);
3594
3595        let judgement =
3596            json!({ "decision": "complete", "guidance": "", "summary": "worker did the job" });
3597        let orch = MockScript::streaming(vec![mock_init("orch-auto"), mock_result_text("seed-hi")])
3598            .responding(vec![
3599                turn("scoping the demo"),
3600                turn(&plan_json().to_string()),
3601            ]);
3602        let orch_run = MockScript::streaming(vec![
3603            mock_init("orch-auto-run"),
3604            mock_result_text("resumed"),
3605        ])
3606        .responding(vec![turn(&judgement.to_string()), turn("NONE")]);
3607        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![
3608            orch,
3609            preflight_authenticated_script(),
3610            worker_pass(),
3611            orch_run,
3612        ]));
3613        let host = MissionHost::with_backend(root.clone(), backend);
3614
3615        let id = host
3616            .create(
3617                "drain me via autoWork",
3618                Some(&json!({ "skipScrutiny": true, "skipFunctional": true })),
3619            )
3620            .await
3621            .expect("create mission");
3622        host.planning_turn(&id, "go").await.expect("planning turn");
3623        let plan_body = host.request_plan(&id).await.expect("request plan");
3624        assert_eq!(plan_body["ready"], true, "{plan_body}");
3625        let plan: Plan =
3626            serde_json::from_value(plan_body["plan"].clone()).expect("plan deserializes");
3627        host.approve(&id, plan).await.expect("approve");
3628        host.release(&id).expect("release");
3629
3630        kranz_engine::queue::enqueue(
3631            &root,
3632            kranz_engine::queue::QueueEntry {
3633                mission_id: id.clone(),
3634                ticket_slug: None,
3635                priority: 2,
3636                seq: 0,
3637            },
3638        )
3639        .expect("enqueue");
3640
3641        // No explicit drain()/POST call — the watcher's tick alone must
3642        // notice the queued entry and kick a drain off.
3643        host.auto_work_tick().await;
3644        assert!(
3645            host.drain_is_live(),
3646            "autoWork tick with autoWork=true and a non-empty queue must start a drain"
3647        );
3648
3649        let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
3650        loop {
3651            let state = host.queue_state();
3652            if state["entries"]
3653                .as_array()
3654                .map(|a| a.is_empty())
3655                .unwrap_or(false)
3656                && state["drain"]["live"] == false
3657            {
3658                break;
3659            }
3660            assert!(
3661                tokio::time::Instant::now() < deadline,
3662                "autoWork drain never completed: {state}"
3663            );
3664            tokio::time::sleep(Duration::from_millis(50)).await;
3665        }
3666    }
3667
3668    #[tokio::test]
3669    async fn auto_work_tick_leaves_the_queue_untouched_when_disabled() {
3670        let Some((_dir, root)) = init_repo() else {
3671            return;
3672        };
3673        // Absent key: default is false, exercised the same as an explicit
3674        // `{"autoWork": false}` layer. No mission needs to actually be
3675        // runnable here — the watcher must never even attempt a drain, so a
3676        // bare queue entry is enough to prove it's left alone.
3677        let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
3678        let host = MissionHost::with_backend(root.clone(), backend);
3679
3680        kranz_engine::queue::enqueue(
3681            &root,
3682            kranz_engine::queue::QueueEntry {
3683                mission_id: "m-untouched".to_string(),
3684                ticket_slug: None,
3685                priority: 2,
3686                seq: 0,
3687            },
3688        )
3689        .expect("enqueue");
3690
3691        host.auto_work_tick().await;
3692
3693        assert!(
3694            !host.drain_is_live(),
3695            "autoWork=false must never start a drain"
3696        );
3697        let entries = kranz_engine::queue::list(&root);
3698        assert_eq!(
3699            entries.len(),
3700            1,
3701            "queue entry must be left untouched when autoWork is disabled: {entries:?}"
3702        );
3703        assert_eq!(entries[0].mission_id, "m-untouched");
3704    }
3705}