1use 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
60type EngineCell = Arc<tokio::sync::Mutex<Box<MissionEngine>>>;
63
64enum HostedMission {
66 Planning {
70 cell: EngineCell,
71 last_use: Arc<Mutex<Instant>>,
72 pending_plan: Arc<Mutex<Option<Plan>>>,
77 },
78 Running {
83 handle: tokio::task::JoinHandle<()>,
84 _repo_busy: kranz_engine::queue::RepoBusyHold,
85 },
86}
87
88#[derive(Debug, Clone, PartialEq, Eq)]
94pub enum PendingApproval {
95 Approved(String),
98 NothingParked,
101 Mismatch { parked: String },
104}
105
106pub struct MissionHost {
108 repo_root: PathBuf,
109 backend: tokio::sync::OnceCell<Arc<dyn AgentBackend>>,
113 missions: Arc<Mutex<HashMap<String, HostedMission>>>,
115 sweeper: Mutex<Option<tokio::task::JoinHandle<()>>>,
118 drain: Mutex<DrainSlot>,
121 global_run_permits: Option<Arc<Semaphore>>,
125 gate_executor: GateExecutor,
129 readiness_front_cache: Mutex<Option<FrontReadinessCache>>,
132 readiness_probe: ReadinessProbe,
136}
137
138struct 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#[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
158type 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
182fn real_gate_executor() -> GateExecutor {
186 Arc::new(|command, cwd| kranz_engine::command_exec::run_bounded_gate_command(cwd, command))
187}
188
189fn 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
205struct DrainHandle {
209 join: tokio::task::JoinHandle<()>,
210 state: Arc<Mutex<DrainState>>,
211}
212
213enum DrainSlot {
222 Idle,
223 Starting(Arc<Mutex<DrainState>>),
224 Running(DrainHandle),
225}
226
227impl MissionHost {
228 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 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 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 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 #[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 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 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 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 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 self.missions
433 .lock()
434 .expect("missions registry lock")
435 .remove(&id);
436 drop(cell);
437
438 Ok(drive_result?.outcome)
439 }
440
441 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 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 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 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 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 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 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 Ok(json!({ "ready": false, "reply": prepend_seed(seed, reply) }))
546 }
547 }
548 }
549
550 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 pub async fn start(&self, id: &str) -> Result<(), ApiError> {
565 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 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 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 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 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 self.missions
689 .lock()
690 .expect("missions registry lock")
691 .insert(id.to_string(), new_planning(new_cell(engine)));
692 }
693 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 {
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 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 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 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 let gate_policy = kranz_engine::command_exec::MergeGatePolicy {
804 sandbox: state.config.worker.sandbox.clone(),
805 mission_dir: paths.mission_dir(),
806 };
807 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 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 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 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 pub fn release(&self, id: &str) -> Result<bool, ApiError> {
1003 release_from(&self.missions, id)
1004 }
1005
1006 pub fn sweep_idle(&self, threshold: Duration) -> Vec<String> {
1012 sweep_idle_from(&self.missions, threshold)
1013 }
1014
1015 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 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 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 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 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 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 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 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 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 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 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 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 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 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 pub async fn drain(&self) -> Result<Value, ApiError> {
1317 self.drain_with_mode(false).await
1318 }
1319
1320 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 {
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 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 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 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 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 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 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 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 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 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 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 map.insert(id.to_string(), new_planning(Arc::clone(&cell)));
1561 drop(map);
1562 self.ensure_sweeper_started();
1563 Ok(cell)
1564 }
1565
1566 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
1584fn 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 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 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 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
1804async 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
1825async 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 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 }
1894 Err(e) => {
1895 tracing::error!(error = %e, "hosted queue drain errored");
1896 restore_drain_checkout(&repo_root, dispatch_branch.as_deref());
1899 }
1900 }
1901 state.lock().expect("drain state lock").live = false;
1902}
1903
1904fn 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
1939async 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
1954fn exit_code_for(status: MissionStatus) -> i32 {
1957 match status {
1958 MissionStatus::Complete => 0,
1959 MissionStatus::Blocked => 2,
1960 _ => 1,
1961 }
1962}
1963
1964fn 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#[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
2020fn 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
2029fn 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()); 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
2073fn 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
2100fn 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
2112fn 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
2121fn 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
2138pub(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
2161pub(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
2182pub(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
2193pub(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
2212pub(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
2233pub(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
2252pub(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
2271pub(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
2292pub(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
2308pub(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
2318pub(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
2329pub(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
2339pub(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
2348fn 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#[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 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 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 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 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 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 assert!(host.planning_cell(&id).is_ok());
2571 }
2572
2573 #[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 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 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 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 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 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 .header("x-kranz-token", "tok")
2835 .body(Body::empty())
2836 .unwrap(),
2837 )
2838 .await
2839 .unwrap();
2840 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 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 #[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 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 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 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 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 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 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 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 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 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 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 #[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 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 #[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 #[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 #[test]
3518 fn should_auto_drain_truth_table() {
3519 assert!(should_auto_drain(true, true, false));
3521 assert!(!should_auto_drain(false, true, false));
3523 assert!(!should_auto_drain(false, false, false));
3524 assert!(!should_auto_drain(true, false, false));
3526 assert!(!should_auto_drain(true, true, true));
3528 assert!(!should_auto_drain(false, false, true));
3529 }
3530
3531 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 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 fn preflight_authenticated_script() -> MockScript {
3557 MockScript::single_shot("ack")
3558 }
3559
3560 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 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 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 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}