1use crate::backlog;
9use crate::cli::{
10 Cli, Command, GrantCommand, PermissionCommand, QuestionCommand, RevisionCommand, TicketCommand,
11};
12use crate::output::{self, ansi};
13use crate::planning_tui::PlanningOutcome;
14use crate::tail::{self, EventRenderer};
15use anyhow::{anyhow, bail, Context, Result};
16use kranz_engine::backend::AgentBackend;
17use kranz_engine::backend_claude::ClaudeBackend;
18use kranz_engine::config;
19use kranz_engine::control;
20use kranz_engine::corpus_export;
21use kranz_engine::cost;
22use kranz_engine::event_log::{EventLog, LockForce};
23use kranz_engine::mission_catalog;
24use kranz_engine::orchestrator::{MissionEngine, PlanRequest};
25use kranz_engine::paths::MissionPaths;
26use kranz_engine::reducer;
27use kranz_engine::trace_export;
28use kranz_engine::types::{ControlCommand, MissionConfig, MissionState, MissionStatus};
29use std::io::{IsTerminal, Write};
30use std::path::{Path, PathBuf};
31use std::sync::atomic::{AtomicBool, Ordering};
32use std::sync::Arc;
33use std::time::{Duration, SystemTime};
34
35pub async fn run_cli(cli: Cli) -> Result<i32> {
38 let lock_force = cli.lock_force();
39 let repo = match cli.repo {
40 Some(repo) => repo,
41 None => std::env::current_dir().context("cannot determine the current directory")?,
42 };
43 if cli.dangerously_allow_all {
44 print_danger_banner();
45 }
46
47 match cli.command {
48 Command::Licenses => {
49 let mut out = std::io::stdout().lock();
50 out.write_all(include_bytes!("../LICENSE"))?;
51 out.write_all(b"\n\n")?;
52 out.write_all(include_bytes!("../assets/THIRD_PARTY_NOTICES.txt"))?;
53 out.write_all(b"\n\nEmbedded dashboard dependency notices\n")?;
54 out.write_all(include_bytes!(
55 "../assets/dashboard/dist/THIRD_PARTY_NOTICES.txt"
56 ))?;
57 Ok(0)
58 }
59 Command::Init {
60 gates,
61 register,
62 id,
63 display_name,
64 } => {
65 let options = crate::init::InitOptions {
66 gates,
67 registration: register.then_some(crate::init::Registration { id, display_name }),
68 global_config: kranz_engine::paths::global_config(),
69 };
70 let report = crate::init::initialize(&repo, &options)?;
71 print!("{}", crate::init::render(&report));
72 Ok(0)
73 }
74 Command::Plan { goal } => {
75 let cfg = load_config(&repo, cli.dangerously_allow_all)?;
76 cmd_plan(repo, goal, cfg, cli.mission.as_deref(), lock_force)
77 .await
78 .map_err(augment_limit_hint)
79 }
80 Command::Run => {
81 let mission = select_mission(&repo, cli.mission.as_deref())?;
82 cmd_run(repo, mission, lock_force, cli.dangerously_allow_all)
83 .await
84 .map_err(augment_limit_hint)
85 }
86 Command::Status { json } => {
87 let mission = select_mission(&repo, cli.mission.as_deref())?;
88 let state = load_state(&repo, &mission)?;
89 if json {
90 println!("{}", serde_json::to_string_pretty(&state)?);
91 } else {
92 print!("{}", output::render_status(&state));
93 }
94 Ok(0)
95 }
96 Command::SandboxProbe { json } => {
97 let report = kranz_engine::sandbox_windows::probe();
98 if json {
99 println!("{}", serde_json::to_string_pretty(&report)?);
100 } else {
101 print!("{}", report.render_text());
102 }
103 Ok(0)
104 }
105 Command::SandboxPrepare { targets } => {
106 for target in targets {
107 eprintln!(
108 "Preparing AppContainer host metadata: target={}",
109 target.display()
110 );
111 let changed = kranz_engine::sandbox_windows::prepare_appcontainer_host(&target)
112 .map_err(anyhow::Error::msg)?;
113 let result = if changed {
114 "applied"
115 } else {
116 "already prepared"
117 };
118 println!(
119 "AppContainer host preparation {result}: target={} mask=0x00120088 inheritance=none",
120 target.display()
121 );
122 }
123 eprintln!("Preparing AppContainer profile-parent metadata (derived from USERPROFILE)");
127 let (profile_parent, profile_changed) =
128 kranz_engine::sandbox_windows::prepare_appcontainer_profile_parent()
129 .map_err(anyhow::Error::msg)?;
130 let profile_result = if profile_changed {
131 "applied"
132 } else {
133 "already prepared"
134 };
135 println!(
136 "AppContainer profile-parent preparation {profile_result}: target={} mask=0x00120088 inheritance=none",
137 profile_parent.display()
138 );
139 eprintln!("Preparing AppContainer null-device metadata: target=\\Device\\Null");
140 kranz_engine::sandbox_windows::prepare_appcontainer_null_device()
141 .map_err(anyhow::Error::msg)?;
142 println!(
143 "AppContainer null-device preparation applied: target=\\Device\\Null inheritance=none"
144 );
145 Ok(0)
146 }
147 Command::Outcomes {
148 json,
149 all,
150 window_days,
151 } => {
152 if all {
153 let config = kranz_engine::paths::global_config()
156 .ok_or_else(|| anyhow::anyhow!("cannot locate the home directory"))?;
157 let report =
158 crate::merged_costs::assess_all(&config, window_days, chrono::Utc::now());
159 if json {
160 println!("{}", serde_json::to_string_pretty(&report)?);
161 } else {
162 print!("{}", crate::merged_costs::render_org(&report));
163 }
164 } else {
165 anyhow::ensure!(
166 window_days <= kranz_engine::outcomes::MAX_MERGED_CHANGE_WINDOW_DAYS,
167 "outcome reason window is too large"
168 );
169 let mut options = kranz_engine::outcomes::OutcomesOptions::resolve(&repo);
170 options.reason_window = Some((window_days, chrono::Utc::now()));
171 let outcomes =
172 kranz_engine::outcomes::compute_outcomes_with_options(&repo, &options)?;
173 if json {
174 println!("{}", output::render_outcomes_json(&outcomes)?);
175 } else {
176 print!("{}", output::render_outcomes(&outcomes));
177 }
178 }
179 Ok(0)
180 }
181 Command::EscalationMetrics { json } => {
182 let metrics = kranz_engine::escalation_metrics::compute_escalation_metrics(&repo)?;
183 if json {
184 println!("{}", output::render_escalation_metrics_json(&metrics)?);
185 } else {
186 print!("{}", output::render_escalation_metrics(&metrics));
187 }
188 Ok(0)
189 }
190 Command::Provenance { mission_id, json } => {
191 let mission = select_mission(&repo, mission_id.as_deref().or(cli.mission.as_deref()))?;
192 let chain = kranz_engine::provenance::compute_provenance(&repo, &mission)?;
193 if json {
194 println!("{}", output::render_provenance_json(&chain)?);
195 } else {
196 print!("{}", output::render_provenance(&chain));
197 }
198 Ok(0)
199 }
200 Command::ReviewPacket { mission_id, json } => {
201 let mission = select_mission(&repo, mission_id.as_deref().or(cli.mission.as_deref()))?;
202 let packet = kranz_engine::review_packet::compute_review_packet(&repo, &mission)?;
203 if json {
204 println!("{}", serde_json::to_string_pretty(&packet)?);
205 } else {
206 print!("{}", kranz_engine::review_packet::render_markdown(&packet));
207 }
208 Ok(0)
209 }
210 Command::GateScores { gate, json } => {
211 let series = kranz_engine::gate_scores::compute_gate_score_series(&repo, &gate)?;
212 if json {
213 println!("{}", output::render_gate_score_series_json(&series)?);
214 } else {
215 print!("{}", output::render_gate_score_series(&series));
216 }
217 Ok(0)
218 }
219 Command::EvidenceBundle { mission_id, out } => {
220 let mission = select_mission(&repo, mission_id.as_deref().or(cli.mission.as_deref()))?;
221 let out = out.unwrap_or_else(|| PathBuf::from(format!("evidence-bundle-{mission}")));
222 let outcome =
223 kranz_engine::evidence_bundle::export_evidence_bundle(&repo, &mission, &out)?;
224 println!(
225 "evidence bundle for mission {mission} written to {} ({} files; {} resolved artefacts, {} unresolved)",
226 outcome.out_dir.display(),
227 outcome.files_written,
228 outcome.resolved_artefacts,
229 outcome.unresolved_artefacts,
230 );
231 Ok(0)
232 }
233 Command::ExportTraces {
234 mission_id,
235 all,
236 out,
237 } => {
238 let jsonl = if all {
239 cmd_export_traces_all(&repo)
240 } else {
241 let mission =
242 select_mission(&repo, mission_id.as_deref().or(cli.mission.as_deref()))?;
243 cmd_export_traces(&repo, &mission)?
244 };
245 match out {
246 Some(path) => {
250 kranz_engine::trace_export::write_export_output(&path, jsonl.as_bytes())
251 .with_context(|| {
252 format!("writing export-traces output to {}", path.display())
253 })?
254 }
255 None => print!("{jsonl}"),
256 }
257 Ok(0)
258 }
259 Command::ExportCorpus {
260 mission_id,
261 all,
262 out,
263 } => {
264 let jsonl = if all {
265 cmd_export_corpus_all(&repo)
266 } else {
267 let mission =
268 select_mission(&repo, mission_id.as_deref().or(cli.mission.as_deref()))?;
269 cmd_export_corpus(&repo, &mission)?
270 };
271 match out {
272 Some(path) => {
274 kranz_engine::trace_export::write_export_output(&path, jsonl.as_bytes())
275 .with_context(|| {
276 format!("writing export-corpus output to {}", path.display())
277 })?
278 }
279 None => print!("{jsonl}"),
280 }
281 Ok(0)
282 }
283 Command::Pause => {
284 let mission = select_control_mission(&repo, cli.mission.as_deref())?;
285 cmd_pause(&repo, &mission)?;
286 println!("pause queued for mission {mission} (takes effect between worker runs)");
287 if let Some(hint) = control_queue_hint(&repo, &mission) {
288 println!("{hint}");
289 }
290 Ok(0)
291 }
292 Command::Resume => {
293 let mission = select_control_mission(&repo, cli.mission.as_deref())?;
294 cmd_resume(&repo, &mission)?;
295 println!("resume queued for mission {mission} (takes effect between worker runs)");
296 if let Some(hint) = control_queue_hint(&repo, &mission) {
297 println!("{hint}");
298 }
299 Ok(0)
300 }
301 Command::Msg { text, interrupt } => {
302 let mission = select_control_mission(&repo, cli.mission.as_deref())?;
303 cmd_msg(&repo, &mission, &text, interrupt)?;
304 if interrupt {
305 println!(
306 "message queued for mission {mission} with --interrupt: the current worker \
307 run will be aborted (recorded as partial) before the message is injected"
308 );
309 } else {
310 println!(
311 "message queued for mission {mission}; it is processed between worker runs"
312 );
313 }
314 if let Some(hint) = control_queue_hint(&repo, &mission) {
315 println!("{hint}");
316 }
317 Ok(0)
318 }
319 Command::Revise { id, instructions } => {
320 let instructions = instructions.join(" ");
321 cmd_request_revision(&repo, &id, &instructions)?;
322 println!("revision request queued for mission {id}");
323 if let Some(hint) = control_queue_hint(&repo, &id) {
324 println!("{hint}");
325 }
326 Ok(0)
327 }
328 Command::Revision { command } => {
329 match command {
330 RevisionCommand::Approve { id, revision } => {
331 cmd_approve_revision(&repo, &id, revision)?;
332 println!("revision {revision} approval queued for mission {id}");
333 if let Some(hint) = control_queue_hint(&repo, &id) {
334 println!("{hint}");
335 }
336 }
337 RevisionCommand::Reject { id, revision } => {
338 cmd_reject_revision(&repo, &id, revision)?;
339 println!("revision {revision} rejection queued for mission {id}");
340 if let Some(hint) = control_queue_hint(&repo, &id) {
341 println!("{hint}");
342 }
343 }
344 }
345 Ok(0)
346 }
347 Command::Permission { command } => {
348 match command {
349 PermissionCommand::List { id } => {
350 let state = load_state(&repo, &id)?;
351 let now = chrono::Utc::now();
352 let requests: Vec<_> = state
353 .permissions
354 .values()
355 .filter(|r| r.pending(now))
356 .collect();
357 println!(
358 "{}",
359 kranz_engine::presentation::visible(&serde_json::to_string_pretty(
360 &requests
361 )?)
362 );
363 }
364 PermissionCommand::Allow {
365 id,
366 request_id,
367 binding,
368 } => {
369 cmd_permission_answer(&repo, &id, &request_id, &binding, true)?;
370 println!("one-call permission answer queued; delivery is recorded separately");
371 }
372 PermissionCommand::Deny {
373 id,
374 request_id,
375 binding,
376 } => {
377 cmd_permission_answer(&repo, &id, &request_id, &binding, false)?;
378 println!("one-call permission refusal queued");
379 }
380 }
381 Ok(0)
382 }
383 Command::Grant { command } => {
384 match command {
385 GrantCommand::Approve { id, command } => {
386 cmd_approve_grant(&repo, &id, &command)?;
387 println!("grant approval for `{command}` queued for mission {id}");
388 if let Some(hint) = control_queue_hint(&repo, &id) {
389 println!("{hint}");
390 }
391 }
392 GrantCommand::Deny {
393 id,
394 command,
395 reason,
396 } => {
397 cmd_deny_grant(&repo, &id, &command, &reason)?;
398 println!("grant denial for `{command}` queued for mission {id}");
399 if let Some(hint) = control_queue_hint(&repo, &id) {
400 println!("{hint}");
401 }
402 }
403 }
404 Ok(0)
405 }
406 Command::Question { command } => {
407 match command {
408 QuestionCommand::List { id } => {
409 print!("{}", cmd_list_questions(&repo, &id)?);
410 }
411 QuestionCommand::Answer {
412 id,
413 question_id,
414 answer,
415 option,
416 } => {
417 cmd_answer_question(&repo, &id, &question_id, &answer, option)?;
418 println!("answer for question {question_id} queued for mission {id}");
419 if let Some(hint) = control_queue_hint(&repo, &id) {
420 println!("{hint}");
421 }
422 }
423 }
424 Ok(0)
425 }
426 Command::Missions => {
427 print!("{}", cmd_missions(&repo)?);
428 Ok(0)
429 }
430 Command::Abandon { id, reason } => {
431 let mission = select_mission(&repo, id.as_deref().or(cli.mission.as_deref()))?;
434 let reason = reason.as_deref().unwrap_or("abandoned by operator");
435 cmd_abandon(&repo, &mission, reason, lock_force)?;
436 println!("mission {mission} ABANDONED ({reason})");
437 Ok(0)
438 }
439 Command::Clean { yes, all } => cmd_clean(&repo, yes, all),
440 Command::Ticket { command } => dispatch_ticket(&repo, command, cli.mission.as_deref()),
441 Command::Draft {
442 slug,
443 yes,
444 from_mission,
445 } => backlog::cmd_draft(
446 repo,
447 &slug,
448 yes,
449 from_mission.as_deref(),
450 cli.dangerously_allow_all,
451 )
452 .await
453 .map_err(augment_limit_hint),
454 Command::Decompose { goal, yes } => {
455 backlog::cmd_decompose(repo, &goal, yes, cli.dangerously_allow_all)
456 .await
457 .map_err(augment_limit_hint)
458 }
459 Command::Exec {
460 file,
461 yes: _,
462 max_cycles,
463 enqueue,
464 enqueue_source,
465 enqueue_external_ref,
466 push,
467 allow_unvalidated,
468 } => crate::exec::cmd_exec(
469 repo,
470 file,
471 crate::exec::ExecOptions {
472 max_cycles,
473 enqueue,
474 enqueue_source: enqueue_source.zip(enqueue_external_ref).map(
475 |(producer, external_ref)| crate::exec::ExternalEnqueueSource {
476 producer,
477 external_ref,
478 },
479 ),
480 push,
481 dangerously_allow_all: cli.dangerously_allow_all,
482 allow_unvalidated,
483 },
484 )
485 .await
486 .map_err(augment_limit_hint),
487 Command::Queue { remove } => match remove {
488 Some(mission_id) => {
489 print!("{}", backlog::cmd_queue_remove(&repo, &mission_id)?);
490 Ok(0)
491 }
492 None => {
493 print!("{}", backlog::cmd_queue(&repo));
494 Ok(0)
495 }
496 },
497 Command::KnowledgeRefresh { json } => cmd_knowledge_refresh(&repo, json),
498 Command::Scan { staged, range } => cmd_scan(&repo, staged, range.as_deref()),
499 Command::DomainLint { seed_config, json } => {
500 cmd_domain_lint(&repo, seed_config.as_deref(), json)
501 }
502 Command::HookGuard { config } => {
503 let mut stdin = std::io::stdin();
504 Ok(crate::hook_guard::run_hook_guard(&config, &mut stdin))
505 }
506 Command::HookStatus { config } => {
507 let mut stdin = std::io::stdin();
508 Ok(crate::hook_status::run_hook_status(
509 &config,
510 &mut stdin,
511 &crate::hook_status::post_signal,
512 )
513 .await)
514 }
515 Command::Ready { json, all } => {
516 if all {
517 let config = kranz_engine::paths::global_config()
518 .ok_or_else(|| anyhow::anyhow!("cannot locate the home directory"))?;
519 let report = crate::ready::assess_all(&config);
520 if json {
521 println!("{}", serde_json::to_string_pretty(&report)?);
522 } else {
523 print!("{}", crate::ready::render_org(&report));
524 }
525 } else {
526 let report = crate::ready::assess(&repo);
527 if json {
528 println!("{}", serde_json::to_string_pretty(&report)?);
529 } else {
530 print!("{}", crate::ready::render(&report));
531 }
532 }
533 Ok(0)
534 }
535 Command::Work { once, expect } => backlog::cmd_work(repo, once, expect)
536 .await
537 .map_err(augment_limit_hint),
538 Command::Serve {
539 port,
540 host,
541 insecure_lan,
542 read_auth,
543 open,
544 dashboard,
545 token,
546 read_token,
547 slack,
548 } => {
549 cmd_serve(
550 repo,
551 host,
552 port,
553 insecure_lan,
554 read_auth,
555 open,
556 dashboard,
557 token,
558 read_token,
559 slack,
560 )
561 .await
562 }
563 Command::Release { url, token } => {
564 let mission = select_mission(&repo, cli.mission.as_deref())?;
565 cmd_release(&repo, &mission, &url, token).await
566 }
567 Command::Config { command } => {
568 crate::config_cmd::cmd_config(&repo, command, cli.mission.as_deref())
569 }
570 Command::Pack { command } => match command {
571 crate::cli::PackCommand::Lint { dir } => cmd_pack_lint(&repo, &dir),
572 },
573 Command::Standards { command } => match command {
574 crate::cli::StandardsCommand::Metrics { json } => {
575 let report = kranz_engine::standards_metrics::compute(&repo)?;
576 if json {
577 println!("{}", serde_json::to_string_pretty(&report)?);
578 } else {
579 print!("{}", crate::output::render_standards_metrics(&report));
580 }
581 Ok(0)
582 }
583 crate::cli::StandardsCommand::Lint { dir, against } => {
584 cmd_standards_lint(&repo, &dir, against.as_deref())
585 }
586 crate::cli::StandardsCommand::Waive {
587 rule,
588 revision,
589 finding,
590 reason,
591 expires,
592 } => {
593 let mission = select_mission(&repo, cli.mission.as_deref())?;
594 cmd_standards_waive(
595 &repo,
596 &mission,
597 &rule,
598 revision,
599 finding.as_deref(),
600 &reason,
601 &expires,
602 lock_force,
603 )
604 }
605 crate::cli::StandardsCommand::Attest { rule, reason } => {
606 let mission = select_mission(&repo, cli.mission.as_deref())?;
607 cmd_standards_attest(&repo, &mission, &rule, &reason, lock_force)
608 }
609 },
610 Command::Otel {
611 endpoint,
612 from_start,
613 } => crate::otel::run_otel(repo, cli.mission.clone(), endpoint, from_start).await,
614 }
615}
616
617fn dispatch_ticket(repo: &Path, command: TicketCommand, mission: Option<&str>) -> Result<i32> {
621 match command {
622 TicketCommand::List => {
623 print!("{}", backlog::cmd_ticket_list(repo));
624 Ok(0)
625 }
626 TicketCommand::Ready { include_deferred } => {
627 print!("{}", backlog::cmd_ticket_ready(repo, include_deferred));
628 Ok(0)
629 }
630 TicketCommand::Show { slug } => {
631 print!("{}", backlog::cmd_ticket_show(repo, &slug)?);
632 Ok(0)
633 }
634 TicketCommand::New { slug, title, goal } => {
635 let path = backlog::cmd_ticket_new(repo, &slug, &title, goal.as_deref())?;
636 println!("created ticket '{slug}' at {}", path.display());
637 Ok(0)
638 }
639 TicketCommand::ImportOpenspec { path, slug } => {
640 let written = crate::openspec::import_change(repo, &path, slug.as_deref())?;
641 println!("imported {} as {}", path.display(), written.display());
642 println!(
643 "acceptance criteria are still prose: replace the placeholder in \
644 '## Acceptance hints' with commands that can fail"
645 );
646 Ok(0)
647 }
648 TicketCommand::Note { slug, text } => {
649 print!(
650 "{}",
651 crate::ticket_notes::cmd_ticket_note(repo, &slug, &text.join(" "))?
652 );
653 Ok(0)
654 }
655 TicketCommand::Notes { slug } => {
656 print!("{}", crate::ticket_notes::cmd_ticket_notes(repo, &slug)?);
657 Ok(0)
658 }
659 TicketCommand::Queue {
660 slug,
661 mission: explicit,
662 force,
663 } => {
664 backlog::cmd_ticket_queue(repo, &slug, explicit.as_deref().or(mission), force)
666 }
667 TicketCommand::Approve {
668 slug,
669 mission: explicit,
670 force,
671 } => {
672 backlog::cmd_ticket_approve(repo, &slug, explicit.as_deref().or(mission), force)
675 }
676 TicketCommand::MigrateState { yes } => backlog::cmd_ticket_migrate_state(repo, yes),
677 }
678}
679
680pub fn load_config(repo: &Path, dangerously_allow_all: bool) -> Result<MissionConfig> {
687 let mut cfg = config::load(repo)?;
688 if dangerously_allow_all {
689 cfg.dangerously_allow_all = true;
690 }
691 config::validate(&cfg)?;
692 Ok(cfg)
693}
694
695pub(crate) fn build_backend(cfg: &MissionConfig) -> Result<Arc<dyn AgentBackend>> {
698 let backend = ClaudeBackend::discover(cfg.claude_binary.as_deref())?;
699 Ok(Arc::new(backend))
700}
701
702pub fn select_mission(repo: &Path, explicit: Option<&str>) -> Result<String> {
705 if let Some(id) = explicit {
706 require_mission(repo, id)?;
707 return Ok(id.to_string());
708 }
709 let ids = MissionPaths::list_missions(repo);
710 match ids.len() {
711 0 => bail!(
712 "no missions found under {}; create one with `kranz plan \"<goal>\"`",
713 repo.join(".kranz").join("missions").display()
714 ),
715 1 => Ok(ids.into_iter().next().expect("len checked")),
716 _ => {
717 let mut best: Option<(SystemTime, String)> = None;
718 for id in ids {
719 let events = MissionPaths::new(repo, &id).events_file();
720 let mtime = std::fs::metadata(&events)
721 .and_then(|m| m.modified())
722 .unwrap_or(SystemTime::UNIX_EPOCH);
723 let newer = match &best {
724 Some((t, _)) => mtime >= *t,
725 None => true,
726 };
727 if newer {
728 best = Some((mtime, id));
729 }
730 }
731 Ok(best.expect("non-empty list").1)
732 }
733 }
734}
735
736pub fn select_control_mission(repo: &Path, explicit: Option<&str>) -> Result<String> {
748 if explicit.is_some() {
749 return Ok(control::resolve_active_mission(repo, explicit)?);
750 }
751 let mission = select_mission(repo, None)?;
752 let status = load_state(repo, &mission)?.mission.status;
753 if mission_catalog::is_terminal_status(status) {
754 bail!(
755 "mission {mission} is {status:?}; control commands apply only to active \
756 missions (a terminal mission's inbox is never drained — see \
757 `kranz missions`)"
758 );
759 }
760 Ok(mission)
761}
762
763pub fn control_queue_hint(repo: &Path, mission_id: &str) -> Option<String> {
768 let paths = MissionPaths::new(repo, mission_id);
769 (!mission_catalog::mission_lock_is_live(&paths)).then(|| {
770 format!(
771 "note: mission {mission_id} is not currently running — the command is \
772 queued and applies when the mission next runs"
773 )
774 })
775}
776
777pub fn select_planning_mission(repo: &Path, explicit: Option<&str>) -> Result<String> {
781 if let Some(id) = explicit {
782 require_mission(repo, id)?;
783 return Ok(id.to_string());
784 }
785 let mut best: Option<(SystemTime, String)> = None;
786 for id in MissionPaths::list_missions(repo) {
787 let Ok(state) = load_state(repo, &id) else {
788 continue; };
790 if state.mission.status != MissionStatus::Planning {
791 continue;
792 }
793 let events = MissionPaths::new(repo, &id).events_file();
794 let mtime = std::fs::metadata(&events)
795 .and_then(|m| m.modified())
796 .unwrap_or(SystemTime::UNIX_EPOCH);
797 if best.as_ref().is_none_or(|(t, _)| mtime >= *t) {
798 best = Some((mtime, id));
799 }
800 }
801 best.map(|(_, id)| id).ok_or_else(|| {
802 anyhow!(
803 "no mission is currently in planning under {} — start one with \
804 `kranz plan \"<goal>\"`",
805 repo.join(".kranz").join("missions").display()
806 )
807 })
808}
809
810pub fn augment_limit_hint(e: anyhow::Error) -> anyhow::Error {
814 let msg = format!("{e:#}").to_ascii_lowercase();
815 if [
816 "session limit",
817 "usage limit",
818 "rate limit",
819 "hit your limit",
820 ]
821 .iter()
822 .any(|s| msg.contains(s))
823 {
824 e.context(
825 "this is your Claude subscription's usage window, not a Kranz failure. \
826 The mission and its conversation are saved: when the limit resets, \
827 `kranz plan` (no goal) resumes planning and `kranz run` resumes execution",
828 )
829 } else {
830 e
831 }
832}
833
834fn require_mission(repo: &Path, mission_id: &str) -> Result<MissionPaths> {
836 if !MissionPaths::is_safe_id(mission_id) {
837 bail!(
838 "invalid mission id '{mission_id}': ids cannot contain path separators, '..', or drive designators"
839 );
840 }
841 let paths = MissionPaths::new(repo, mission_id);
842 if !paths.events_file().is_file() {
843 bail!(
844 "mission '{mission_id}' not found under {} (see `kranz missions`)",
845 paths.missions_dir().display()
846 );
847 }
848 Ok(paths)
849}
850
851pub fn load_state(repo: &Path, mission_id: &str) -> Result<MissionState> {
854 let paths = require_mission(repo, mission_id)?;
855 let events = EventLog::read_events(&paths.events_file())
856 .with_context(|| format!("reading the event log of mission '{mission_id}'"))?;
857 let state = reducer::fold(&events)
858 .with_context(|| format!("folding the event log of mission '{mission_id}'"))?;
859 Ok(state)
860}
861
862pub fn cmd_export_traces(repo: &Path, mission_id: &str) -> Result<String> {
867 let paths = require_mission(repo, mission_id)?;
868 let events = EventLog::read_events(&paths.events_file())
869 .with_context(|| format!("reading the event log of mission '{mission_id}'"))?;
870 let state = reducer::fold(&events)
871 .with_context(|| format!("folding the event log of mission '{mission_id}'"))?;
872 let pairs = trace_export::export_validated_traces(&state, &events);
873 Ok(trace_export::to_jsonl(&pairs))
874}
875
876pub fn cmd_export_traces_all(repo: &Path) -> String {
881 let mut pairs = Vec::new();
882 for mission_id in MissionPaths::list_missions(repo) {
883 let paths = MissionPaths::new(repo, &mission_id);
884 let Ok(events) = EventLog::read_events(&paths.events_file()) else {
885 continue;
886 };
887 let Ok(state) = reducer::fold(&events) else {
888 continue;
889 };
890 pairs.extend(trace_export::export_validated_traces(&state, &events));
891 }
892 trace_export::to_jsonl(&pairs)
893}
894
895pub fn cmd_export_corpus(repo: &Path, mission_id: &str) -> Result<String> {
902 let paths = require_mission(repo, mission_id)?;
903 paths.require_no_follow()?;
904 let events = EventLog::read_events(&paths.events_file())
905 .with_context(|| format!("reading the event log of mission '{mission_id}'"))?;
906 let records = corpus_export::export_corpus(&paths.mission_dir(), mission_id, &events)
907 .with_context(|| format!("deriving the training corpus of mission '{mission_id}'"))?;
908 Ok(corpus_export::to_jsonl(&records))
909}
910
911pub fn cmd_export_corpus_all(repo: &Path) -> String {
917 let mut records = Vec::new();
918 for mission_id in MissionPaths::list_missions(repo) {
919 let paths = MissionPaths::new(repo, &mission_id);
920 if paths.require_no_follow().is_err() {
921 continue;
922 }
923 let Ok(events) = EventLog::read_events(&paths.events_file()) else {
924 continue;
925 };
926 let Ok(mission_records) =
927 corpus_export::export_corpus(&paths.mission_dir(), &mission_id, &events)
928 else {
929 continue;
930 };
931 records.extend(mission_records);
932 }
933 corpus_export::to_jsonl(&records)
934}
935
936pub fn print_danger_banner() {
938 eprintln!(
939 "\n\
940 ============================================================\n\
941 !! --dangerously-allow-all IS SET !!\n\
942 !! !!\n\
943 !! Permission gating is BYPASSED for every agent !!\n\
944 !! session (bypassPermissions). Workers can run ANY !!\n\
945 !! command: file writes, network access, git push, !!\n\
946 !! package publishes, sudo. !!\n\
947 !! !!\n\
948 !! Only use this on a sandboxed, disposable checkout. !!\n\
949 ============================================================\n"
950 );
951}
952
953struct StdinLines {
965 rx: Arc<tokio::sync::Mutex<tokio::sync::mpsc::UnboundedReceiver<String>>>,
966}
967
968impl StdinLines {
969 fn spawn() -> Self {
970 static CHANNEL: std::sync::OnceLock<
971 Arc<tokio::sync::Mutex<tokio::sync::mpsc::UnboundedReceiver<String>>>,
972 > = std::sync::OnceLock::new();
973 let rx = CHANNEL
974 .get_or_init(|| {
975 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
976 std::thread::spawn(move || {
977 let mut buf = String::new();
978 loop {
979 buf.clear();
980 match std::io::stdin().read_line(&mut buf) {
981 Ok(0) | Err(_) => break, Ok(_) => {
983 let line = buf.trim_end_matches(['\r', '\n']).to_string();
984 if tx.send(line).is_err() {
985 break;
986 }
987 }
988 }
989 }
990 });
991 Arc::new(tokio::sync::Mutex::new(rx))
992 })
993 .clone();
994 StdinLines { rx }
995 }
996
997 async fn next(&mut self) -> Option<String> {
999 self.rx.lock().await.recv().await
1000 }
1001
1002 fn drain(&mut self) -> usize {
1004 let Ok(mut rx) = self.rx.try_lock() else {
1007 return 0;
1008 };
1009 let mut n = 0;
1010 while rx.try_recv().is_ok() {
1011 n += 1;
1012 }
1013 n
1014 }
1015
1016 fn drain_noisily(&mut self, tty: bool) {
1020 if !std::io::stdin().is_terminal() {
1021 return;
1022 }
1023 let n = self.drain();
1024 if n > 0 {
1025 let (dim, reset) = if tty {
1026 (ansi::DIM, ansi::RESET)
1027 } else {
1028 ("", "")
1029 };
1030 eprintln!(
1031 "{dim}(ignored {n} line(s) typed while the orchestrator was working — \
1032 the prompt below wants fresh input){reset}"
1033 );
1034 }
1035 }
1036}
1037
1038async fn cmd_plan(
1046 repo: PathBuf,
1047 goal: Option<String>,
1048 cfg: MissionConfig,
1049 explicit_mission: Option<&str>,
1050 force_lock: LockForce,
1051) -> Result<i32> {
1052 let backend = build_backend(&cfg)?;
1053 let (mut engine, intro) = match goal {
1054 Some(goal) => {
1055 let engine = MissionEngine::create(backend, repo.clone(), &goal, cfg)?;
1056 let intro = format!("mission {} created (planning)", engine.mission_id());
1057 (engine, intro)
1058 }
1059 None => {
1060 let mission = select_planning_mission(&repo, explicit_mission)?;
1061 let engine = MissionEngine::resume(backend, repo.clone(), &mission, force_lock)?;
1062 if engine.state().mission.status != MissionStatus::Planning {
1063 return Err(anyhow!(
1064 "mission {mission} is {:?}, not in planning — use 'kranz run' \
1065 to execute it, or 'kranz plan \"<goal>\"' to start a new mission",
1066 engine.state().mission.status
1067 ));
1068 }
1069 let intro = format!(
1070 "resuming planning for mission {mission} — the conversation continues \
1071 where it left off"
1072 );
1073 (engine, intro)
1074 }
1075 };
1076
1077 if std::io::stdin().is_terminal() && std::io::stdout().is_terminal() {
1083 let mission_id = engine.mission_id().to_string();
1084 return match crate::planning_tui::run(engine, intro).await? {
1088 PlanningOutcome::ApprovedRun => start_run_after_plan(repo, mission_id).await,
1089 PlanningOutcome::ApprovedExit | PlanningOutcome::NotApproved => Ok(0),
1090 };
1091 }
1092 println!("{intro}");
1093 let tty = std::io::stdout().is_terminal();
1094
1095 let color = std::io::stderr().is_terminal();
1098 let stop = Arc::new(AtomicBool::new(false));
1099 let printer = tokio::spawn(tail::tail_events(
1100 engine.paths().events_file(),
1101 engine.state().last_seq,
1102 EventRenderer::planning(engine.state(), color),
1103 Arc::clone(&stop),
1104 ));
1105 let mut approved = false;
1106 let mut run_now = false;
1107 let mut stdin_lines = StdinLines::spawn();
1108
1109 println!("talk to the orchestrator to shape the plan:");
1110 println!(" /plan request the plan + cost estimate and review it for approval");
1111 println!(" /quit exit planning (Ctrl-D works too)");
1112
1113 loop {
1114 stdin_lines.drain_noisily(tty);
1115 if tty {
1116 print!("you> ");
1117 let _ = std::io::stdout().flush();
1118 }
1119 let Some(line) = stdin_lines.next().await else {
1120 break; };
1122 let line = line.trim().to_string();
1123 if line.is_empty() {
1124 continue;
1125 }
1126 match line.as_str() {
1127 "/quit" => break,
1128 "/plan" => {
1129 let request = engine.request_plan().await;
1130 if let Some(seed) = engine.take_seed_reply() {
1133 print_orchestrator_reply(&seed, tty);
1134 }
1135 let plan = match request {
1136 Ok(PlanRequest::Ready(plan)) => plan,
1137 Ok(PlanRequest::NotReady(text)) => {
1138 print_orchestrator_reply(&text, tty);
1141 println!(
1142 "the orchestrator isn't ready to emit the plan yet — answer it \
1143 above, then /plan again."
1144 );
1145 continue;
1146 }
1147 Ok(PlanRequest::WrongPlan { reason }) => {
1148 print_orchestrator_reply(&reason, tty);
1151 println!(
1152 "the orchestrator believes a plan here is likely WRONG — reframe \
1153 the goal or fix the premise above, then /plan again."
1154 );
1155 continue;
1156 }
1157 Err(e) => {
1158 eprintln!(
1159 "kranz: plan request failed: {:#}",
1160 augment_limit_hint(e.into())
1161 );
1162 continue;
1163 }
1164 };
1165 println!("{}", output::render_plan(&plan));
1166 let calibration = cost::calibrate(&repo);
1169 let estimate = cost::estimate(&plan, &engine.state().config, &calibration.params);
1170 let estimate = cost::apply_shape(estimate, &plan, &calibration);
1171 println!(
1172 "{}",
1173 output::render_cost_estimate(&estimate, calibration.missions_used)
1174 );
1175
1176 stdin_lines.drain_noisily(tty);
1177 print!("approve? [y/N] ");
1178 let _ = std::io::stdout().flush();
1179 let answer = stdin_lines.next().await.unwrap_or_default();
1180 if matches!(answer.trim().to_ascii_lowercase().as_str(), "y" | "yes") {
1181 match engine.approve_plan(plan) {
1182 Ok(()) => {
1183 let branch = engine.state().mission.mission_branch.clone();
1184 approved = true;
1185 if std::io::stdin().is_terminal() {
1190 println!("plan approved and committed on {branch}.");
1191 stdin_lines.drain_noisily(tty);
1192 print!("start execution now? [Y/n] ");
1193 let _ = std::io::stdout().flush();
1194 let reply = stdin_lines.next().await;
1195 if run_now_answer(reply.as_deref()) {
1196 run_now = true;
1197 } else {
1198 println!("run 'kranz run' to execute.");
1199 }
1200 } else {
1201 println!(
1202 "plan approved and committed on {branch}. \
1203 run 'kranz run' to execute."
1204 );
1205 }
1206 break;
1207 }
1208 Err(e) => eprintln!("kranz: plan approval failed: {e}"),
1209 }
1210 } else {
1211 println!("not approved — back to the conversation.");
1212 }
1213 }
1214 _ if line.starts_with('/') => {
1215 println!("unknown command {line}; use /plan or /quit");
1216 }
1217 _ => {
1218 let result = engine.planning_turn(&line).await;
1219 if let Some(seed) = engine.take_seed_reply() {
1222 print_orchestrator_reply(&seed, tty);
1223 }
1224 match result {
1225 Ok(reply) => print_orchestrator_reply(&reply, tty),
1226 Err(e) => eprintln!(
1227 "kranz: orchestrator turn failed: {:#}",
1228 augment_limit_hint(e.into())
1229 ),
1230 }
1231 }
1232 }
1233 }
1234 if !approved {
1235 println!(
1236 "leaving planning; mission {} was not approved. Resume anytime with `kranz plan`.",
1237 engine.mission_id()
1238 );
1239 }
1240 let mission_id = engine.mission_id().to_string();
1241 drop(engine);
1244 stop.store(true, Ordering::Relaxed);
1245 let _ = printer.await;
1246 if run_now {
1247 return start_run_after_plan(repo, mission_id).await;
1250 }
1251 Ok(0)
1252}
1253
1254pub fn run_now_answer(answer: Option<&str>) -> bool {
1259 match answer {
1260 None => false,
1261 Some(text) => !matches!(text.trim().to_ascii_lowercase().as_str(), "n" | "no"),
1262 }
1263}
1264
1265async fn start_run_after_plan(repo: PathBuf, mission_id: String) -> Result<i32> {
1269 println!(
1270 "starting mission {mission_id} — live event feed follows \
1271 (Ctrl-C safe; resume with 'kranz run')"
1272 );
1273 run_mission_loop(repo, mission_id, LockForce::No, true).await
1274}
1275
1276fn print_orchestrator_reply(text: &str, tty: bool) {
1278 let prefix = if tty {
1279 format!("{}orchestrator>{} ", ansi::DIM, ansi::RESET)
1280 } else {
1281 "orchestrator> ".to_string()
1282 };
1283 for line in text.lines() {
1284 println!("{prefix}{line}");
1285 }
1286}
1287
1288async fn cmd_run(
1295 repo: PathBuf,
1296 mission: String,
1297 force_lock: LockForce,
1298 dangerously_allow_all: bool,
1299) -> Result<i32> {
1300 if dangerously_allow_all {
1303 let paths = require_mission(&repo, &mission)?;
1304 control::enqueue(
1305 &paths,
1306 &ControlCommand::ConfigChange {
1307 patch: serde_json::json!({ "dangerouslyAllowAll": true }),
1308 },
1309 )?;
1310 }
1311 run_mission_loop(repo, mission, force_lock, true).await
1312}
1313
1314pub(crate) async fn run_mission_loop(
1318 repo: PathBuf,
1319 mission: String,
1320 force_lock: LockForce,
1321 interactive: bool,
1325) -> Result<i32> {
1326 let cfg = load_config(&repo, false)?;
1327 let backend = build_backend(&cfg)?;
1328 run_mission_loop_with_backend(repo, mission, force_lock, interactive, backend).await
1329}
1330
1331async fn run_mission_loop_with_backend(
1335 repo: PathBuf,
1336 mission: String,
1337 force_lock: LockForce,
1338 interactive: bool,
1339 backend: Arc<dyn AgentBackend>,
1340) -> Result<i32> {
1341 let paths = require_mission(&repo, &mission)?;
1342
1343 loop {
1344 let mut engine =
1345 MissionEngine::resume(Arc::clone(&backend), repo.clone(), &mission, force_lock)?;
1346
1347 let color = std::io::stderr().is_terminal();
1349 let renderer = EventRenderer::seeded(engine.state(), color);
1350 let stop = Arc::new(AtomicBool::new(false));
1351 let printer = tokio::spawn(tail::tail_events(
1352 engine.paths().events_file(),
1353 engine.state().last_seq,
1354 renderer,
1355 Arc::clone(&stop),
1356 ));
1357
1358 let run_result = engine.run().await;
1359 drop(engine);
1362 stop.store(true, Ordering::Relaxed);
1363 let _ = printer.await;
1364
1365 let status = run_result?;
1366 if let Err(e) = kranz_engine::work::reconcile_ticket_for_mission(&repo, &mission) {
1370 eprintln!("kranz run: warning: failed to reconcile linked ticket: {e}");
1371 }
1372
1373 match status {
1374 MissionStatus::Complete => {
1375 println!("mission {mission} COMPLETE");
1376 return Ok(0);
1377 }
1378 MissionStatus::Blocked => {
1379 eprintln!(
1380 "\n\
1381 ==================== MILESTONE BLOCKED ====================\n\
1382 A milestone is blocked (fix-cycle cap reached or blocked by\n\
1383 the orchestrator). Inspect it with `kranz status`.\n\
1384 ==========================================================="
1385 );
1386 if interactive && std::io::stdin().is_terminal() && std::io::stdout().is_terminal()
1390 {
1391 print!(
1392 "guidance for the orchestrator (what to do about the block; \
1393 empty line or Ctrl-D exits)\nguidance> "
1394 );
1395 let _ = std::io::stdout().flush();
1396 let mut lines = StdinLines::spawn();
1397 if let Some(text) = lines.next().await {
1398 let text = text.trim().to_string();
1399 if !text.is_empty() {
1400 control::enqueue(
1401 &paths,
1402 &ControlCommand::Msg {
1403 text,
1404 interrupt: false,
1405 },
1406 )?;
1407 println!("guidance queued — resuming the mission…");
1408 continue;
1409 }
1410 }
1411 }
1412 eprintln!(
1413 "unblock later by sending guidance via `kranz msg \"<text>\"` \
1414 and re-running `kranz run`."
1415 );
1416 println!("mission {mission} BLOCKED");
1417 return Ok(2);
1418 }
1419 MissionStatus::Failed => {
1420 println!("mission {mission} FAILED");
1421 return Ok(1);
1422 }
1423 other => {
1424 println!(
1425 "mission {mission} ended as {}",
1426 output::mission_status_label(other)
1427 );
1428 return Ok(1);
1429 }
1430 }
1431 }
1432}
1433
1434pub fn cmd_pause(repo: &Path, mission_id: &str) -> Result<PathBuf> {
1440 let paths = require_mission(repo, mission_id)?;
1441 Ok(control::enqueue(&paths, &ControlCommand::Pause)?)
1442}
1443
1444pub fn cmd_resume(repo: &Path, mission_id: &str) -> Result<PathBuf> {
1446 let paths = require_mission(repo, mission_id)?;
1447 Ok(control::enqueue(&paths, &ControlCommand::Resume)?)
1448}
1449
1450pub fn cmd_msg(repo: &Path, mission_id: &str, text: &str, interrupt: bool) -> Result<PathBuf> {
1452 let paths = require_mission(repo, mission_id)?;
1453 let cmd = ControlCommand::Msg {
1454 text: text.to_string(),
1455 interrupt,
1456 };
1457 Ok(control::enqueue(&paths, &cmd)?)
1458}
1459
1460pub fn cmd_request_revision(repo: &Path, mission_id: &str, instructions: &str) -> Result<PathBuf> {
1462 let instructions = instructions.trim();
1463 if instructions.is_empty() {
1464 bail!("revision instructions must not be empty");
1465 }
1466 let paths = require_revisable_mission(repo, mission_id)?;
1467 Ok(control::enqueue(
1468 &paths,
1469 &ControlCommand::RequestRevision {
1470 instructions: instructions.to_string(),
1471 },
1472 )?)
1473}
1474
1475pub fn cmd_approve_revision(repo: &Path, mission_id: &str, revision: u32) -> Result<PathBuf> {
1477 let paths = require_pending_revision(repo, mission_id, revision)?;
1478 Ok(control::enqueue(
1479 &paths,
1480 &ControlCommand::ApproveRevision { revision },
1481 )?)
1482}
1483
1484pub fn cmd_reject_revision(repo: &Path, mission_id: &str, revision: u32) -> Result<PathBuf> {
1486 let paths = require_pending_revision(repo, mission_id, revision)?;
1487 Ok(control::enqueue(
1488 &paths,
1489 &ControlCommand::RejectRevision { revision },
1490 )?)
1491}
1492
1493pub fn cmd_permission_answer(
1494 repo: &Path,
1495 mission_id: &str,
1496 request_id: &str,
1497 binding: &str,
1498 allow: bool,
1499) -> Result<PathBuf> {
1500 let paths = require_revisable_mission(repo, mission_id)?;
1501 let state = load_state(repo, mission_id)?;
1502 let record = state
1503 .permissions
1504 .get(request_id)
1505 .ok_or_else(|| anyhow!("unknown live permission"))?;
1506 record.validate_answer(binding, allow, chrono::Utc::now())?;
1507 Ok(control::enqueue(
1508 &paths,
1509 &ControlCommand::ResolvePermission {
1510 resolution: kranz_engine::live_permission::Resolution {
1511 request_id: request_id.into(),
1512 binding_digest: binding.into(),
1513 allow,
1514 actor: kranz_engine::live_permission::Actor::LocalRepositoryAuthority,
1515 reason: "operator answered through the local CLI".into(),
1516 },
1517 },
1518 )?)
1519}
1520
1521pub fn cmd_approve_grant(repo: &Path, mission_id: &str, command: &str) -> Result<PathBuf> {
1523 let paths = require_pending_grant(repo, mission_id, command)?;
1524 Ok(control::enqueue(
1525 &paths,
1526 &ControlCommand::ApproveGrant {
1527 command: command.to_string(),
1528 },
1529 )?)
1530}
1531
1532pub fn cmd_deny_grant(
1534 repo: &Path,
1535 mission_id: &str,
1536 command: &str,
1537 reason: &str,
1538) -> Result<PathBuf> {
1539 let paths = require_pending_grant(repo, mission_id, command)?;
1540 Ok(control::enqueue(
1541 &paths,
1542 &ControlCommand::DenyGrant {
1543 command: command.to_string(),
1544 reason: reason.to_string(),
1545 },
1546 )?)
1547}
1548
1549pub fn cmd_list_questions(repo: &Path, mission_id: &str) -> Result<String> {
1554 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1555 let state = load_state(repo, &mission_id)?;
1556 if state.pending_questions.is_empty() {
1557 return Ok(format!("mission {mission_id} has no open questions\n"));
1558 }
1559 let mut out = String::new();
1560 for q in &state.pending_questions {
1561 out.push_str(&format!(
1562 "{} ({}): {}\n",
1563 q.question_id,
1564 q.feature_id.as_deref().unwrap_or("mission"),
1565 q.text
1566 ));
1567 if q.options.is_empty() {
1568 out.push_str(" free-text answer expected\n");
1569 } else {
1570 for (index, option) in q.options.iter().enumerate() {
1571 out.push_str(&format!(" [{index}] {option}\n"));
1572 }
1573 }
1574 }
1575 Ok(out)
1576}
1577
1578pub fn cmd_answer_question(
1581 repo: &Path,
1582 mission_id: &str,
1583 question_id: &str,
1584 answer: &str,
1585 option: Option<u32>,
1586) -> Result<PathBuf> {
1587 let paths = require_pending_question(repo, mission_id, question_id, option, answer)?;
1588 Ok(control::enqueue(
1589 &paths,
1590 &ControlCommand::AnswerQuestion {
1591 question_id: question_id.to_string(),
1592 answer: answer.to_string(),
1593 option,
1594 },
1595 )?)
1596}
1597
1598fn require_pending_question(
1604 repo: &Path,
1605 mission_id: &str,
1606 question_id: &str,
1607 option: Option<u32>,
1608 answer: &str,
1609) -> Result<MissionPaths> {
1610 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1611 let state = load_state(repo, &mission_id)?;
1612 let Some(pending) = state
1613 .pending_questions
1614 .iter()
1615 .find(|q| q.question_id == question_id)
1616 else {
1617 bail!("mission {mission_id} has no open question '{question_id}'");
1618 };
1619 if let Some(index) = option {
1620 match pending.options.get(index as usize) {
1621 Some(expected) if expected == answer => {}
1622 Some(expected) => bail!(
1623 "answer `{answer}` does not match option {index} (`{expected}`) of question '{question_id}'"
1624 ),
1625 None => bail!(
1626 "question '{question_id}' has no option {index} (it offered {})",
1627 pending.options.len()
1628 ),
1629 }
1630 }
1631 Ok(MissionPaths::new(repo, &mission_id))
1632}
1633
1634fn require_revisable_mission(repo: &Path, mission_id: &str) -> Result<MissionPaths> {
1635 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1636 let state = load_state(repo, &mission_id)?;
1637 if state.mission.status == MissionStatus::Planning {
1638 bail!("mission {mission_id} has no approved plan to revise yet");
1639 }
1640 Ok(MissionPaths::new(repo, &mission_id))
1641}
1642
1643fn require_pending_revision(repo: &Path, mission_id: &str, revision: u32) -> Result<MissionPaths> {
1644 let paths = require_revisable_mission(repo, mission_id)?;
1645 let state = load_state(repo, mission_id)?;
1646 match state.pending_revision {
1647 Some(pending) if pending.revision == revision => Ok(paths),
1648 Some(pending) => bail!(
1649 "mission {mission_id} is awaiting revision {}, not {revision}",
1650 pending.revision
1651 ),
1652 None => bail!("mission {mission_id} has no pending revision"),
1653 }
1654}
1655
1656fn require_pending_grant(repo: &Path, mission_id: &str, command: &str) -> Result<MissionPaths> {
1660 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1661 let state = load_state(repo, &mission_id)?;
1662 match state.pending_grant_request {
1663 Some(pending) if pending.command == command => Ok(MissionPaths::new(repo, &mission_id)),
1664 Some(pending) => bail!(
1665 "mission {mission_id} is awaiting a grant for `{}`, not `{command}`",
1666 pending.command
1667 ),
1668 None => bail!("mission {mission_id} has no pending grant request"),
1669 }
1670}
1671
1672pub fn cmd_knowledge_refresh(repo: &Path, json: bool) -> Result<i32> {
1676 let report = kranz_engine::knowledge::refresh_knowledge(repo);
1677 if json {
1678 println!("{}", serde_json::to_string_pretty(&report)?);
1679 } else if report.findings.is_empty() {
1680 println!("knowledge refresh: no notes");
1681 } else {
1682 for finding in &report.findings {
1683 let labels: Vec<String> = finding
1684 .verdicts
1685 .iter()
1686 .map(|v| match v {
1687 kranz_engine::knowledge::RefreshVerdict::Ok => "ok".into(),
1688 kranz_engine::knowledge::RefreshVerdict::AlreadyStale => "already-stale".into(),
1689 kranz_engine::knowledge::RefreshVerdict::Unverified => "unverified".into(),
1690 kranz_engine::knowledge::RefreshVerdict::InvalidMetadata { field, value } => {
1691 format!(
1692 "invalid-metadata:{field}:{}",
1693 value.as_deref().unwrap_or("missing")
1694 )
1695 }
1696 kranz_engine::knowledge::RefreshVerdict::InvalidCitation { citation } => {
1697 format!("invalid-citation:{citation}")
1698 }
1699 kranz_engine::knowledge::RefreshVerdict::PathMissing { path } => {
1700 format!("path-missing:{path}")
1701 }
1702 kranz_engine::knowledge::RefreshVerdict::PathDrifted { path } => {
1703 format!("path-drifted:{path}")
1704 }
1705 kranz_engine::knowledge::RefreshVerdict::CommandSkipped { command } => {
1706 format!("command-skipped:{command}")
1707 }
1708 kranz_engine::knowledge::RefreshVerdict::ProbeFailed { target, error } => {
1709 format!("probe-failed:{target}:{error}")
1710 }
1711 })
1712 .collect();
1713 println!(
1714 "{} ({}) [{}] {}",
1715 finding.rel_path,
1716 finding.title,
1717 finding.freshness,
1718 labels.join(", ")
1719 );
1720 }
1721 if report.check_needed() {
1722 println!("knowledge refresh: check-needed");
1723 } else {
1724 println!("knowledge refresh: ok");
1725 }
1726 }
1727 Ok(if report.check_needed() { 1 } else { 0 })
1728}
1729
1730pub fn cmd_scan(repo: &Path, staged: bool, range: Option<&str>) -> Result<i32> {
1731 if staged && range.is_some() {
1732 bail!("choose either --staged or --range, not both");
1733 }
1734 let git = kranz_engine::git_ops::GitRepo::open(repo)?;
1735 let diff = if staged {
1736 git.diff_staged()?
1737 } else if let Some(range) = range {
1738 git.diff_range(range)?
1739 } else {
1740 git.diff_range("HEAD")?
1741 };
1742 let allowed = std::fs::read_to_string(repo.join(kranz_engine::scrub::SECRET_ALLOWLIST_PATH))
1743 .ok()
1744 .map(|text| kranz_engine::scrub::read_allowlist_text(&text))
1745 .unwrap_or_default();
1746 let findings = kranz_engine::scrub::filter_allowed(
1747 kranz_engine::scrub::scan_unified_diff(&diff),
1748 &allowed,
1749 );
1750 if findings.is_empty() {
1751 println!("secret scan passed");
1752 Ok(0)
1753 } else {
1754 println!(
1755 "secret scan failed; add a fingerprint to {} only for a reviewed false positive:\n{}",
1756 kranz_engine::scrub::SECRET_ALLOWLIST_PATH,
1757 kranz_engine::scrub::format_findings(&findings)
1758 );
1759 Ok(2)
1760 }
1761}
1762
1763pub fn cmd_domain_lint(repo: &Path, seed_config: Option<&Path>, json: bool) -> Result<i32> {
1774 use kranz_engine::domain_lint as dl;
1775 let config_path = repo.join(dl::DENYLIST_PATH);
1776
1777 if let Some(terms_file) = seed_config {
1778 let terms = std::fs::read_to_string(terms_file)
1779 .with_context(|| format!("read terms file {}", terms_file.display()))?;
1780 let existing = std::fs::read_to_string(&config_path).ok();
1781 let config = dl::seed_config(existing.as_deref(), &terms)?;
1782 if let Some(parent) = config_path.parent() {
1784 std::fs::create_dir_all(parent)
1785 .with_context(|| format!("create {}", parent.display()))?;
1786 }
1787 std::fs::write(&config_path, &config)
1788 .with_context(|| format!("write {}", config_path.display()))?;
1789 let denylist = dl::load_denylist(&config)?;
1790 println!(
1792 "seeded {} ({} hashed terms, {})",
1793 dl::DENYLIST_PATH,
1794 denylist.term_count(),
1795 if existing.is_some() {
1796 "salt preserved"
1797 } else {
1798 "fresh salt"
1799 }
1800 );
1801 warn_if_terms_file_unprotected(repo, terms_file);
1802 return Ok(0);
1803 }
1804
1805 let config_text = std::fs::read_to_string(&config_path).with_context(|| {
1806 format!(
1807 "read {} — seed it with `kranz domain-lint --seed-config <terms-file>` (docs/domain-lint.md)",
1808 dl::DENYLIST_PATH
1809 )
1810 })?;
1811 let denylist = dl::load_denylist(&config_text)?;
1812 let allowed = std::fs::read_to_string(repo.join(dl::ALLOWLIST_PATH))
1813 .ok()
1814 .map(|text| kranz_engine::scrub::read_allowlist_text(&text))
1815 .unwrap_or_default();
1816 let report = dl::lint_tree(repo, &denylist, &allowed)?;
1817
1818 if json {
1819 println!(
1820 "{}",
1821 serde_json::to_string_pretty(&serde_json::json!({
1822 "passed": report.is_clean(),
1823 "filesScanned": report.files_scanned,
1824 "filesSkipped": report.files_skipped,
1825 "findings": report.findings,
1826 }))?
1827 );
1828 }
1829 if report.is_clean() {
1830 if !json {
1831 println!(
1832 "domain lint passed ({} files scanned)",
1833 report.files_scanned
1834 );
1835 }
1836 Ok(0)
1837 } else {
1838 if !json {
1839 println!(
1840 "domain lint failed: {} unwaived hit(s); add a fingerprint to {} only for a reviewed false positive:",
1841 report.findings.len(),
1842 dl::ALLOWLIST_PATH
1843 );
1844 for finding in &report.findings {
1845 println!("{} {}:{}", finding.fingerprint, finding.path, finding.line);
1846 }
1847 }
1848 Ok(1)
1849 }
1850}
1851
1852fn warn_if_terms_file_unprotected(repo: &Path, terms_file: &Path) {
1857 let (Ok(repo), Ok(terms_file)) = (repo.canonicalize(), terms_file.canonicalize()) else {
1858 return;
1859 };
1860 let Ok(relative) = terms_file.strip_prefix(&repo) else {
1861 return; };
1863 let ignored = std::process::Command::new("git")
1864 .args(["check-ignore", "-q", "--"])
1865 .arg(relative)
1866 .current_dir(&repo)
1867 .status()
1868 .map(|status| status.success())
1869 .unwrap_or(true);
1871 if !ignored {
1872 eprintln!(
1873 "warning: {} is inside the repo and NOT gitignored — move it outside the repo or use {} (gitignored)",
1874 relative.display(),
1875 kranz_engine::domain_lint::TERMS_LOCAL_PATH
1876 );
1877 }
1878}
1879
1880fn cmd_pack_lint(repo: &Path, dir: &Path) -> Result<i32> {
1891 let trust = kranz_engine::pack::standards::trust_for_dir(repo, dir);
1892 match kranz_engine::pack::Pack::load_with_trust(dir, trust) {
1893 Ok(Some(pack)) => {
1894 print!("{}", kranz_engine::pack::render_lint(&pack));
1895 Ok(0)
1896 }
1897 Ok(None) => {
1898 println!(
1899 "no pack at {} (no {}) — nothing to lint",
1900 dir.display(),
1901 kranz_engine::pack::PACK_MANIFEST
1902 );
1903 Ok(0)
1904 }
1905 Err(err) => {
1906 eprintln!("invalid pack at {}: {err}", dir.display());
1907 Ok(1)
1908 }
1909 }
1910}
1911
1912fn cmd_standards_lint(repo: &Path, dir: &Path, against: Option<&str>) -> Result<i32> {
1919 let trust = kranz_engine::pack::standards::trust_for_dir(repo, dir);
1920 let pack = match kranz_engine::pack::Pack::load_with_trust(dir, trust) {
1921 Ok(Some(pack)) => pack,
1922 Ok(None) => {
1923 println!(
1924 "no pack at {} (no {}) — nothing to lint",
1925 dir.display(),
1926 kranz_engine::pack::PACK_MANIFEST
1927 );
1928 return Ok(0);
1929 }
1930 Err(err) => {
1931 eprintln!("invalid pack at {}: {err}", dir.display());
1932 return Ok(1);
1933 }
1934 };
1935 let Some(manifest) = &pack.standards else {
1936 println!(
1937 "pack `{}` (schema {}) at {} declares no [standards] root — nothing to lint",
1938 pack.name,
1939 pack.schema,
1940 dir.display()
1941 );
1942 return Ok(0);
1943 };
1944 print!(
1945 "{}",
1946 kranz_engine::pack::standards::render_manifest(manifest, trust)
1947 );
1948 if let Some(refname) = against {
1949 let Some(pack_rel) = kranz_engine::pack::standards::repo_relative_dir(repo, dir) else {
1953 eprintln!(
1954 "--against reads the base pack from tracked git blobs in {}; {} is outside \
1955 the repo — external packs have no base history to compare against (D-A)",
1956 repo.display(),
1957 dir.display()
1958 );
1959 return Ok(1);
1960 };
1961 let git = kranz_engine::git_ops::GitRepo::open(repo)?;
1962 let base = match kranz_engine::pack::standards::load_at_ref(&git, refname, &pack_rel) {
1963 Ok(base) => base,
1964 Err(err) => {
1965 eprintln!("cannot load the base standards at `{refname}`: {err}");
1966 return Ok(1);
1967 }
1968 };
1969 let errors = kranz_engine::pack::standards::check_transitions(base.as_ref(), manifest);
1970 print!(
1971 "{}",
1972 kranz_engine::pack::standards::render_transition_report(
1973 refname,
1974 base.as_ref(),
1975 &errors
1976 )
1977 );
1978 if !errors.is_empty() {
1979 return Ok(1);
1980 }
1981 }
1982 Ok(0)
1983}
1984
1985#[allow(clippy::too_many_arguments)]
1995fn cmd_standards_waive(
1996 repo: &Path,
1997 mission_id: &str,
1998 rule: &str,
1999 revision: Option<u64>,
2000 finding: Option<&str>,
2001 reason: &str,
2002 expires: &str,
2003 force_lock: LockForce,
2004) -> Result<i32> {
2005 let expires_at = match chrono::DateTime::parse_from_rfc3339(expires) {
2006 Ok(parsed) => parsed.with_timezone(&chrono::Utc),
2007 Err(err) => {
2008 eprintln!("waiver refused: --expires must be an RFC 3339 instant: {err}");
2009 return Ok(1);
2010 }
2011 };
2012 let request = kranz_engine::standards_waiver::WaiverRequest {
2013 rule_id: rule.to_string(),
2014 revision,
2015 finding_subject: finding.map(str::to_string),
2016 reason: reason.to_string(),
2017 expires_at,
2018 };
2019 let outcome = match kranz_engine::standards_waiver::approve_standards_waiver(
2020 repo, mission_id, &request, "cli", force_lock,
2021 ) {
2022 Ok(outcome) => outcome,
2023 Err(kranz_engine::error::EngineError::LockHeld(e)) => {
2024 eprintln!(
2025 "waiver refused: an engine still holds mission '{mission_id}'s lock — stop \
2026 the running mission first (a waiver against a live mission would race the \
2027 runner's own appends).\n (underlying: {e})"
2028 );
2029 return Ok(1);
2030 }
2031 Err(e) => {
2032 eprintln!("waiver refused: {e}");
2033 return Ok(1);
2034 }
2035 };
2036 let pinned = &outcome.rule;
2039 println!(
2040 "recorded standards.waiver.approved (seq {})",
2041 outcome.event.seq
2042 );
2043 println!("mission: {mission_id}");
2044 println!(
2045 "rule: {} r{} — {}, {}; checker {}; waivable: {}",
2046 pinned.id,
2047 pinned.revision,
2048 pinned.level,
2049 pinned.effective_status,
2050 pinned.checker.as_deref().unwrap_or("-"),
2051 pinned.waivable
2052 );
2053 println!(" statement: {}", pinned.statement);
2054 println!(
2055 "finding: {} (run {})\n evidence: {}",
2056 outcome.finding_subject, outcome.run_id, outcome.finding_evidence
2057 );
2058 println!(" fingerprint: sha256:{}", outcome.finding_fingerprint);
2059 if pinned.when_paths.is_empty() {
2060 println!(
2061 "affected paths: the whole mission diff ({} path(s)) — the rule is unscoped",
2062 outcome.affected_paths.len()
2063 );
2064 } else if outcome.affected_paths.is_empty() {
2065 println!(
2066 "affected paths: (none — the rule's when-paths match no changed path; the \
2067 waiver binds the empty scoped diff)"
2068 );
2069 } else {
2070 println!("affected paths: {}", outcome.affected_paths.join(", "));
2071 }
2072 println!(
2073 "diff digest: sha256:{} (covers the affected-path diff at the mission branch tip)",
2074 outcome.diff_digest
2075 );
2076 println!(
2077 "approver: {} via cli\nreason: {reason}\nexpires: {}",
2078 kranz_engine::standards_waiver::LOCAL_OPERATOR,
2079 expires_at.to_rfc3339()
2080 );
2081 Ok(0)
2082}
2083
2084fn cmd_standards_attest(
2089 repo: &Path,
2090 mission_id: &str,
2091 rule: &str,
2092 reason: &str,
2093 force_lock: LockForce,
2094) -> Result<i32> {
2095 let record = match kranz_engine::standards_attestation::approve_attestation(
2096 repo, mission_id, rule, reason, "cli", force_lock,
2097 ) {
2098 Ok(record) => record,
2099 Err(kranz_engine::error::EngineError::LockHeld(e)) => {
2100 eprintln!(
2101 "attestation refused: an engine still holds mission '{mission_id}'s lock — stop \
2102 the running mission first.\n (underlying: {e})"
2103 );
2104 return Ok(1);
2105 }
2106 Err(e) => {
2107 eprintln!("attestation refused: {e}");
2108 return Ok(1);
2109 }
2110 };
2111 println!(
2112 "recorded standards.attestation.approved (seq {})",
2113 record.seq
2114 );
2115 println!("mission: {mission_id}");
2116 println!("rule: {} r{}", record.rule_id, record.rule_revision);
2117 if record.paths.is_empty() {
2118 println!("affected paths: (none)");
2119 } else {
2120 println!("affected paths: {}", record.paths.join(", "));
2121 }
2122 println!("diff digest: sha256:{}", record.diff_digest);
2123 println!(
2124 "approver: {} via {}\nreason: {}",
2125 record.approver, record.surface, record.reason
2126 );
2127 Ok(0)
2128}
2129
2130pub fn cmd_missions(repo: &Path) -> Result<String> {
2137 let index_contents =
2138 std::fs::read_to_string(MissionPaths::new(repo, "_").missions_dir().join("index.md"))
2139 .unwrap_or_default();
2140 let mut ids = MissionPaths::list_missions(repo);
2141 for id in mission_catalog::mission_index_ids(&index_contents) {
2142 if !ids.contains(&id) {
2143 ids.push(id);
2144 }
2145 }
2146 ids.sort();
2147 if ids.is_empty() {
2148 return Ok("no missions\n".to_string());
2149 }
2150 let ticket_of: std::collections::HashMap<String, String> =
2153 kranz_engine::ticket::Ticket::list(repo)
2154 .into_iter()
2155 .filter_map(|t| {
2156 kranz_engine::ticket::Ticket::mission_for(repo, &t.slug).map(|m| (m, t.slug))
2157 })
2158 .collect();
2159 let mut out = String::new();
2160 for id in ids {
2161 let ticket = ticket_of
2162 .get(&id)
2163 .map(|s| format!(" [ticket: {s}]"))
2164 .unwrap_or_default();
2165 let paths = MissionPaths::new(repo, &id);
2166 if !paths.events_file().is_file() {
2167 out.push_str(&format!(
2168 "{id} {:<10} deleted mission (no data recorded)\n",
2169 "DELETED"
2170 ));
2171 continue;
2172 }
2173 if let Err(error) = paths.require_no_follow() {
2176 out.push_str(&format!("{id} {:<10} (unreadable: {error})\n", "FAILED"));
2177 continue;
2178 }
2179 match load_state(repo, &id) {
2180 Ok(state) => out.push_str(&format!(
2181 "{id} {:<10} {}{ticket}\n",
2182 output::mission_status_label(state.mission.status),
2183 state.mission.goal
2184 )),
2185 Err(e) => out.push_str(&format!("{id} {:<10} (unreadable: {e:#})\n", "FAILED")),
2186 }
2187 }
2188 Ok(out)
2189}
2190
2191pub fn cmd_abandon(
2199 repo: &Path,
2200 mission_id: &str,
2201 reason: &str,
2202 force_lock: LockForce,
2203) -> Result<()> {
2204 require_mission(repo, mission_id)?;
2205 mission_catalog::abandon_mission(repo, mission_id, reason, force_lock).map_err(|e| {
2206 if matches!(e, kranz_engine::error::EngineError::LockHeld(_)) {
2207 anyhow!(
2208 "cannot abandon mission '{mission_id}' — an engine still holds its lock. \
2209 Stop the running `kranz run` first. If the holder is a crashed leftover, \
2210 pass --force-lock (steals unless the holder is provably alive); a provably \
2211 LIVE holder that you have verified to be a zombie or foreign process \
2212 additionally requires --dangerously-steal-live-lock.\n (underlying: {e})"
2213 )
2214 } else {
2215 anyhow::Error::new(e).context(format!("abandoning mission '{mission_id}'"))
2216 }
2217 })
2218}
2219
2220#[derive(Debug, Clone, PartialEq, Eq)]
2222pub struct CleanEntry {
2223 pub id: String,
2224 pub status_label: String,
2225 pub goal: String,
2226}
2227
2228pub fn select_cleanable(repo: &Path, all: bool) -> Vec<CleanEntry> {
2236 let mut out = Vec::new();
2237 for id in MissionPaths::list_missions(repo) {
2238 let paths = MissionPaths::new(repo, &id);
2239 if mission_catalog::mission_lock_is_live(&paths) {
2241 continue;
2242 }
2243 let Ok(state) = load_state(repo, &id) else {
2244 continue; };
2246 let has_plan = paths.plan_file().is_file();
2247 if mission_catalog::cleanable_class(state.mission.status, has_plan).is_cleaned(all) {
2248 out.push(CleanEntry {
2249 id,
2250 status_label: output::mission_status_label(state.mission.status).to_string(),
2251 goal: state.mission.goal,
2252 });
2253 }
2254 }
2255 out
2256}
2257
2258pub fn render_clean_listing(entries: &[CleanEntry]) -> String {
2261 if entries.is_empty() {
2262 return "nothing to clean\n".to_string();
2263 }
2264 let mut out = String::new();
2265 for e in entries {
2266 out.push_str(&format!("{:<10} {} {}\n", e.status_label, e.id, e.goal));
2267 }
2268 out
2269}
2270
2271fn cmd_clean(repo: &Path, yes: bool, all: bool) -> Result<i32> {
2277 let entries = select_cleanable(repo, all);
2278 if entries.is_empty() {
2279 print!("{}", render_clean_listing(&entries));
2280 return Ok(0);
2281 }
2282
2283 print!("{}", render_clean_listing(&entries));
2284 println!(
2285 "\n{} mission director{} above would be removed.{}",
2286 entries.len(),
2287 if entries.len() == 1 { "y" } else { "ies" },
2288 if all {
2289 ""
2290 } else {
2291 " (Complete missions are kept; pass --all to include them.)"
2292 }
2293 );
2294
2295 if !yes && !confirm_clean()? {
2296 println!("clean aborted; nothing removed.");
2297 return Ok(0);
2298 }
2299
2300 let removed = remove_missions(repo, &entries, true);
2301 println!(
2302 "cleaned {} mission director{}",
2303 removed.len(),
2304 if removed.len() == 1 { "y" } else { "ies" }
2305 );
2306 Ok(0)
2307}
2308
2309pub fn remove_missions(repo: &Path, entries: &[CleanEntry], verbose: bool) -> Vec<String> {
2315 let mut removed = Vec::new();
2316 for e in entries {
2317 let paths = MissionPaths::new(repo, &e.id);
2318 if mission_catalog::mission_lock_is_live(&paths) {
2324 eprintln!("kranz: skipping {} — became live since listing", e.id);
2325 continue;
2326 }
2327 let dir = paths.mission_dir();
2328 match std::fs::remove_dir_all(&dir) {
2329 Ok(()) => {
2330 if verbose {
2331 println!("removed {}", dir.display());
2332 }
2333 mission_catalog::prune_mission_index_file(repo, &e.id);
2334 removed.push(e.id.clone());
2335 }
2336 Err(err) => eprintln!("kranz: could not remove {}: {err}", dir.display()),
2337 }
2338 }
2339 removed
2340}
2341
2342fn confirm_clean() -> Result<bool> {
2345 use std::io::BufRead;
2346 print!("proceed? [y/N] ");
2347 std::io::stdout().flush().ok();
2348 let mut line = String::new();
2349 let n = std::io::stdin().lock().read_line(&mut line)?;
2350 if n == 0 {
2351 return Ok(false); }
2353 Ok(matches!(
2354 line.trim().to_ascii_lowercase().as_str(),
2355 "y" | "yes"
2356 ))
2357}
2358
2359pub(crate) fn refuse_non_loopback_without_insecure_lan(
2369 bind: std::net::IpAddr,
2370 insecure_lan: bool,
2371) -> Result<()> {
2372 if !bind.is_loopback() && !insecure_lan {
2373 anyhow::bail!(
2374 "refusing to bind {bind}: non-loopback binds expose the API on the \
2375 network (reads require the read-only or mutation token; POSTs \
2376 require the mutation token). \
2377 Re-run with `--insecure-lan` if you intentionally trust this \
2378 network (LAN/tailnet), or keep the default `--host 127.0.0.1`."
2379 );
2380 }
2381 Ok(())
2382}
2383
2384pub(crate) fn effective_require_read_token(bind_is_loopback: bool, read_auth: bool) -> bool {
2390 !bind_is_loopback || read_auth
2391}
2392
2393#[allow(clippy::too_many_arguments)]
2403async fn cmd_serve(
2404 repo: PathBuf,
2405 host: String,
2406 port: u16,
2407 insecure_lan: bool,
2408 read_auth: bool,
2409 open: bool,
2410 dashboard: Option<PathBuf>,
2411 token: Option<String>,
2412 read_token: Option<String>,
2413 slack: bool,
2414) -> Result<i32> {
2415 let bind: std::net::IpAddr = host
2416 .parse()
2417 .map_err(|e| anyhow!("--host '{host}' is not an IP address: {e}"))?;
2418 refuse_non_loopback_without_insecure_lan(bind, insecure_lan)?;
2419 let token = token
2420 .or_else(|| std::env::var("KRANZ_TOKEN").ok())
2421 .unwrap_or_else(kranz_server::generate_token);
2422 let read_token = read_token
2423 .or_else(|| std::env::var("KRANZ_READ_TOKEN").ok())
2424 .unwrap_or_else(kranz_server::generate_token);
2425 kranz_server::MutationAuthority::new(token.clone())
2429 .context("--token / KRANZ_TOKEN is invalid")?;
2430 anyhow::ensure!(
2431 !read_token.is_empty() && read_token.bytes().all(|byte| byte.is_ascii_graphic()),
2432 "--read-token / KRANZ_READ_TOKEN must be non-empty visible ASCII without whitespace"
2433 );
2434 anyhow::ensure!(
2435 read_token != token,
2436 "--read-token / KRANZ_READ_TOKEN must differ from the mutation token"
2437 );
2438 if !bind.is_loopback() {
2439 eprintln!(
2440 "WARNING: binding {bind} with --insecure-lan — the API is reachable \
2441 beyond this machine. Reads require the read-only or mutation token; \
2442 POSTs require the mutation token. Use only on a network you trust."
2443 );
2444 }
2445 if read_auth && effective_require_read_token(bind.is_loopback(), read_auth) {
2446 eprintln!(
2447 "--read-auth: GETs and the WS upgrade now require the read-only or \
2448 mutation token, including on loopback; POSTs still require mutation \
2449 authority."
2450 );
2451 }
2452 let listener = kranz_server::bind_listener(bind, port)
2455 .await
2456 .map_err(|e| anyhow!("failed to bind {bind}:{port}: {e}"))?;
2457 let local_addr = listener
2458 .local_addr()
2459 .map_err(|e| anyhow!("failed to read bound address: {e}"))?;
2460 let multi_host = Arc::new(kranz_server::MultiRepoHost::from_global_config(
2464 kranz_engine::paths::global_config().as_deref(),
2465 repo.clone(),
2466 )?);
2467 multi_host.ensure_auto_work_started();
2470
2471 if slack {
2475 if multi_host.uses_operator_catalog() {
2476 let catalog = multi_repo_slack_catalog(&multi_host)?;
2477 tokio::spawn(async move {
2478 let never = std::future::pending::<()>();
2479 if let Err(error) = kranz_slack::serve_slack_catalog(catalog, never).await {
2480 tracing::error!(%error, "slack catalog bridge exited with an error");
2481 }
2482 });
2483 } else {
2484 let context = single_repo_slack_context(&multi_host)?;
2485 let repo_slack = context.root().to_path_buf();
2486 let hosted = context.host().ok_or_else(|| {
2487 anyhow!(
2488 "cannot start Slack bridge for unavailable repository '{}': {}",
2489 context.id(),
2490 context.unavailable_reason().unwrap_or("unavailable")
2491 )
2492 })?;
2493 let bridge_host: kranz_slack::SharedHost =
2494 Arc::new(crate::host_bridge::HostedPlanning(hosted.clone()));
2495 tokio::spawn(async move {
2496 let never = std::future::pending::<()>();
2500 if let Err(e) =
2501 kranz_slack::serve_slack(&repo_slack, Some(bridge_host), never).await
2502 {
2503 tracing::error!(error = %e, "slack bridge exited with an error");
2504 }
2505 });
2506 }
2507 }
2508
2509 let dashboard_assets = resolve_dashboard_assets(&repo, dashboard);
2510 let display_host = match local_addr.ip() {
2513 std::net::IpAddr::V6(v6) => format!("[{v6}]"),
2514 std::net::IpAddr::V4(v4) => v4.to_string(),
2515 };
2516 let url = format!("http://{display_host}:{}/", local_addr.port());
2517 println!("kranz server on {url}");
2518 println!("mutation token: {token}");
2519 println!("read token: {read_token} (GETs/WS only — safe for dashboards and agents)");
2520 if open {
2521 println!(
2522 "opening the bare URL; no token rides in the opener's argv. Paste the \
2523 mutation token above when the dashboard asks for one"
2524 );
2525 }
2526 match &dashboard_assets {
2527 Some(DashboardAssets::Embedded) => println!(
2528 "serving embedded dashboard ({})",
2529 crate::embedded_dashboard::EMBEDDED_DASHBOARD_SOURCE
2530 ),
2531 Some(DashboardAssets::Dir(dir)) => println!("serving dashboard from {}", dir.display()),
2532 None => println!(
2533 "no dashboard build found (--dashboard, $KRANZ_DASHBOARD_DIST, \
2534 <repo>/apps/dashboard/dist, installed asset dirs, or the kranz checkout) \
2535 and no embedded dashboard is available; serving API only"
2536 ),
2537 }
2538
2539 let static_assets = dashboard_assets.map(|assets| match assets {
2540 DashboardAssets::Dir(dir) => kranz_server::DashboardStatic::Dir(dir),
2541 DashboardAssets::Embedded => {
2542 kranz_server::DashboardStatic::Embedded(crate::embedded_dashboard::EMBEDDED_DASHBOARD)
2543 }
2544 });
2545
2546 if open {
2547 let url = url.clone();
2572 tokio::spawn(async move {
2573 tokio::time::sleep(Duration::from_millis(600)).await;
2574 open_browser(&url);
2575 });
2576 }
2577
2578 let shutdown = async {
2579 if let Err(e) = tokio::signal::ctrl_c().await {
2580 tracing::error!(error = %e, "failed to install ctrl-c handler");
2581 }
2582 };
2583 let result = serve_multi_with_token_cleanup(
2584 &repo,
2585 multi_host,
2586 listener,
2587 static_assets,
2588 token,
2589 read_token,
2590 read_auth,
2591 shutdown,
2592 )
2593 .await;
2594 match result {
2595 Ok(()) => Ok(0),
2596 Err(e) => Err(anyhow!("server failed: {e}")),
2597 }
2598}
2599
2600fn single_repo_slack_context(
2601 multi_host: &kranz_server::MultiRepoHost,
2602) -> Result<Arc<kranz_server::RepoContext>> {
2603 multi_host
2604 .compatibility_context()
2605 .ok_or_else(|| anyhow!("--slack requires one healthy configured repository"))
2606}
2607
2608fn multi_repo_slack_catalog(
2609 multi_host: &kranz_server::MultiRepoHost,
2610) -> Result<kranz_slack::SlackCatalog> {
2611 let global_config = kranz_engine::paths::global_config()
2612 .ok_or_else(|| anyhow!("cannot locate the operator config directory for Slack affinity"))?;
2613 let affinity_path = global_config
2614 .parent()
2615 .expect("global config has a parent")
2616 .join("slack")
2617 .join("thread-affinity.json");
2618 let repos = multi_host
2619 .contexts()
2620 .map(|context| {
2621 let host = context.host().map(|host| {
2622 Arc::new(crate::host_bridge::HostedPlanning(host.clone()))
2623 as kranz_slack::SharedHost
2624 });
2625 kranz_slack::SlackRepo {
2626 id: context.id().to_string(),
2627 root: context.root().to_path_buf(),
2628 display_name: context
2629 .config()
2630 .display_name
2631 .clone()
2632 .unwrap_or_else(|| context.id().to_string()),
2633 routes: context
2634 .config()
2635 .slack
2636 .channels
2637 .iter()
2638 .map(|route| kranz_slack::SlackRoute {
2639 team_id: route.team.clone(),
2640 channel_id: route.channel.clone(),
2641 })
2642 .collect(),
2643 allow_users: context.config().slack.allow_users.clone(),
2644 available: context.is_healthy(),
2645 host,
2646 unavailable_reason: context.unavailable_reason().map(str::to_string),
2647 default: multi_host.default_repo() == Some(context.id()),
2648 }
2649 })
2650 .collect();
2651 kranz_slack::SlackCatalog::new(repos, affinity_path)
2652}
2653
2654#[allow(clippy::too_many_arguments)]
2655async fn serve_multi_with_token_cleanup(
2656 repo: &Path,
2657 multi_host: Arc<kranz_server::MultiRepoHost>,
2658 listener: tokio::net::TcpListener,
2659 static_assets: Option<kranz_server::DashboardStatic>,
2660 token: String,
2661 read_token: String,
2662 read_auth: bool,
2663 shutdown: impl std::future::Future<Output = ()> + Send + 'static,
2664) -> anyhow::Result<()> {
2665 let (token_file, read_token_file) = if multi_host.uses_operator_catalog() {
2666 let address = listener
2667 .local_addr()
2668 .context("cannot resolve the bound address for operator token storage")?;
2669 let token_file = write_operator_serve_token(address, &token)
2670 .context("cannot securely store the operator serve token")?;
2671 let read_token_file = write_operator_serve_read_token(address, &read_token)
2672 .context("cannot securely store the operator serve read token")?;
2673 (token_file, read_token_file)
2674 } else {
2675 let token_file =
2676 write_serve_token(repo, &token).context("cannot securely store the serve token")?;
2677 let read_token_file = write_serve_read_token(repo, &read_token)
2678 .context("cannot securely store the serve read token")?;
2679 (token_file, read_token_file)
2680 };
2681 let result = kranz_server::serve_multi_on_listener(
2682 multi_host,
2683 listener,
2684 static_assets,
2685 kranz_server::MutationAuthority::new(token)?,
2686 Some(read_token),
2687 read_auth,
2688 shutdown,
2689 )
2690 .await;
2691 remove_token_file(&token_file);
2692 remove_token_file(&read_token_file);
2693 result
2694}
2695
2696#[cfg(test)]
2702async fn serve_with_token_cleanup(
2703 repo: &Path,
2704 host: Arc<kranz_server::MissionHost>,
2705 listener: tokio::net::TcpListener,
2706 static_assets: Option<kranz_server::DashboardStatic>,
2707 token: String,
2708 read_token: String,
2709 shutdown: impl std::future::Future<Output = ()> + Send + 'static,
2710) -> anyhow::Result<()> {
2711 let token_file = write_serve_token(repo, &token)?;
2715 let read_token_file = write_serve_read_token(repo, &read_token)?;
2716
2717 let result = kranz_server::serve_on_listener(
2718 host,
2719 listener,
2720 static_assets,
2721 kranz_server::MutationAuthority::new(token)?,
2722 shutdown,
2723 )
2724 .await;
2725
2726 remove_token_file(&token_file);
2727 remove_token_file(&read_token_file);
2728
2729 result
2730}
2731
2732fn write_serve_token(repo: &Path, token: &str) -> std::io::Result<PathBuf> {
2738 write_token_file(&repo.join(".kranz").join("serve.token"), token)
2739}
2740
2741fn write_serve_read_token(repo: &Path, read_token: &str) -> std::io::Result<PathBuf> {
2747 write_token_file(&repo.join(".kranz").join("serve.read.token"), read_token)
2748}
2749
2750fn write_operator_serve_token(
2755 address: std::net::SocketAddr,
2756 token: &str,
2757) -> std::io::Result<PathBuf> {
2758 let global = kranz_engine::paths::global_config().ok_or_else(|| {
2759 std::io::Error::new(
2760 std::io::ErrorKind::NotFound,
2761 "cannot resolve operator config directory",
2762 )
2763 })?;
2764 let path = operator_serve_token_path(&global, address);
2765 write_token_file(&path, token)
2766}
2767
2768fn write_operator_serve_read_token(
2770 address: std::net::SocketAddr,
2771 read_token: &str,
2772) -> std::io::Result<PathBuf> {
2773 let global = kranz_engine::paths::global_config().ok_or_else(|| {
2774 std::io::Error::new(
2775 std::io::ErrorKind::NotFound,
2776 "cannot resolve operator config directory",
2777 )
2778 })?;
2779 let path = operator_serve_read_token_path(&global, address);
2780 write_token_file(&path, read_token)
2781}
2782
2783fn operator_serve_read_token_path(global_config: &Path, address: std::net::SocketAddr) -> PathBuf {
2785 operator_serve_token_path(global_config, address).with_extension("read.token")
2786}
2787
2788fn operator_serve_token_path(global_config: &Path, address: std::net::SocketAddr) -> PathBuf {
2789 let endpoint = match address.ip() {
2790 std::net::IpAddr::V4(ip) => format!("v4-{:08x}-{}", u32::from(ip), address.port()),
2791 std::net::IpAddr::V6(ip) => format!("v6-{:032x}-{}", u128::from(ip), address.port()),
2792 };
2793 global_config
2794 .parent()
2795 .unwrap_or_else(|| Path::new("."))
2796 .join("serve")
2797 .join(format!("{endpoint}.token"))
2798}
2799
2800fn legacy_operator_serve_token_path(global_config: &Path, port: u16) -> PathBuf {
2804 global_config
2805 .parent()
2806 .unwrap_or_else(|| Path::new("."))
2807 .join("serve")
2808 .join(format!("{port}.token"))
2809}
2810
2811fn write_token_file(path: &Path, token: &str) -> std::io::Result<PathBuf> {
2812 let dir = path.parent().ok_or_else(|| {
2813 std::io::Error::new(std::io::ErrorKind::InvalidInput, "token path has no parent")
2814 })?;
2815 std::fs::create_dir_all(dir)?;
2816 let tmp = dir.join(format!(
2821 ".{}.tmp-{}",
2822 path.file_name()
2823 .and_then(|name| name.to_str())
2824 .unwrap_or("serve.token"),
2825 uuid::Uuid::new_v4().as_simple()
2826 ));
2827 let write_tmp = || -> std::io::Result<()> {
2828 let mut options = std::fs::OpenOptions::new();
2829 options.create_new(true).write(true);
2830 #[cfg(unix)]
2831 {
2832 use std::os::unix::fs::OpenOptionsExt;
2833 options.mode(0o600);
2834 }
2835 let mut file = options.open(&tmp)?;
2836 file.write_all(token.as_bytes())?;
2837 file.flush()?;
2838 Ok(())
2839 };
2840 if let Err(error) = write_tmp() {
2841 let _ = std::fs::remove_file(&tmp);
2842 return Err(error);
2843 }
2844 #[cfg(unix)]
2845 {
2846 use std::os::unix::fs::PermissionsExt;
2847 std::fs::set_permissions(&tmp, std::fs::Permissions::from_mode(0o600))?;
2848 }
2849 if let Err(error) = std::fs::rename(&tmp, path) {
2850 let _ = std::fs::remove_file(&tmp);
2851 return Err(error);
2852 }
2853 Ok(path.to_path_buf())
2854}
2855
2856fn remove_token_file(path: &Path) {
2857 if let Err(err) = std::fs::remove_file(path) {
2858 if err.kind() != std::io::ErrorKind::NotFound {
2859 eprintln!("kranz: could not remove {}: {err}", path.display());
2860 }
2861 }
2862}
2863
2864#[derive(Debug, Clone, PartialEq, Eq)]
2865enum DashboardAssets {
2866 Dir(PathBuf),
2867 Embedded,
2868}
2869
2870#[derive(Debug, Clone, Default)]
2871struct DashboardResolutionInputs {
2872 env_dist: Option<PathBuf>,
2873 home: Option<PathBuf>,
2874 exe: Option<PathBuf>,
2875 manifest_dir: Option<PathBuf>,
2876 embedded_available: bool,
2877}
2878
2879impl DashboardResolutionInputs {
2880 fn runtime() -> Self {
2881 Self {
2882 env_dist: std::env::var_os("KRANZ_DASHBOARD_DIST").map(PathBuf::from),
2883 home: std::env::var_os(if cfg!(windows) { "USERPROFILE" } else { "HOME" })
2884 .map(PathBuf::from),
2885 exe: std::env::current_exe().ok(),
2886 manifest_dir: Some(PathBuf::from(env!("CARGO_MANIFEST_DIR"))),
2887 embedded_available: !crate::embedded_dashboard::EMBEDDED_DASHBOARD.is_empty(),
2888 }
2889 }
2890}
2891
2892fn resolve_dashboard_assets(repo: &Path, explicit: Option<PathBuf>) -> Option<DashboardAssets> {
2893 resolve_dashboard_assets_from(repo, explicit, &DashboardResolutionInputs::runtime())
2894}
2895
2896fn resolve_dashboard_assets_from(
2904 repo: &Path,
2905 explicit: Option<PathBuf>,
2906 inputs: &DashboardResolutionInputs,
2907) -> Option<DashboardAssets> {
2908 if let Some(d) = explicit {
2909 return Some(DashboardAssets::Dir(d));
2912 }
2913
2914 if let Some(d) = first_dashboard_dir(dashboard_dir_candidates(repo, inputs)) {
2915 return Some(DashboardAssets::Dir(d));
2916 }
2917
2918 inputs
2919 .embedded_available
2920 .then_some(DashboardAssets::Embedded)
2921}
2922
2923pub fn resolve_dashboard_dist(repo: &Path, explicit: Option<PathBuf>) -> Option<PathBuf> {
2926 match resolve_dashboard_assets(repo, explicit) {
2927 Some(DashboardAssets::Dir(dir)) => Some(dir),
2928 Some(DashboardAssets::Embedded) | None => None,
2929 }
2930}
2931
2932fn dashboard_dir_candidates(repo: &Path, inputs: &DashboardResolutionInputs) -> Vec<PathBuf> {
2933 let mut candidates = Vec::new();
2934
2935 candidates.extend(inputs.env_dist.clone());
2936 candidates.push(repo.join("apps").join("dashboard").join("dist"));
2937
2938 if let Some(home) = &inputs.home {
2939 candidates.push(home.join(".kranz").join("dashboard").join("dist"));
2940 candidates.push(home.join(".kranz").join("dashboard"));
2941 }
2942
2943 if let Some(exe) = &inputs.exe {
2944 candidates.extend(installed_dashboard_dirs(exe));
2945 candidates.extend(source_checkout_dist_from_exe(exe));
2946 }
2947
2948 if let Some(manifest_dir) = &inputs.manifest_dir {
2949 candidates.extend(source_checkout_dist_from_manifest(manifest_dir));
2950 candidates.push(manifest_dir.join("assets").join("dashboard").join("dist"));
2951 }
2952
2953 candidates
2954}
2955
2956fn first_dashboard_dir(candidates: Vec<PathBuf>) -> Option<PathBuf> {
2957 candidates
2958 .into_iter()
2959 .find(|d| d.join("index.html").is_file())
2960}
2961
2962fn installed_dashboard_dirs(exe: &Path) -> Vec<PathBuf> {
2963 let Some(bin_dir) = exe.parent() else {
2964 return Vec::new();
2965 };
2966 let mut dirs = vec![
2967 bin_dir.join("dashboard").join("dist"),
2968 bin_dir.join("dashboard"),
2969 ];
2970 if let Some(prefix) = bin_dir.parent() {
2971 dirs.push(
2972 prefix
2973 .join("share")
2974 .join("kranz")
2975 .join("dashboard")
2976 .join("dist"),
2977 );
2978 dirs.push(prefix.join("share").join("kranz").join("dashboard"));
2979 }
2980 dirs
2981}
2982
2983fn source_checkout_dist_from_exe(exe: &Path) -> Option<PathBuf> {
2984 let profile_dir = exe.parent()?;
2985 let target_dir = profile_dir.parent()?;
2986 if target_dir.file_name()? != "target" {
2987 return None;
2988 }
2989 Some(
2990 target_dir
2991 .parent()?
2992 .join("apps")
2993 .join("dashboard")
2994 .join("dist"),
2995 )
2996}
2997
2998fn source_checkout_dist_from_manifest(manifest_dir: &Path) -> Option<PathBuf> {
2999 Some(
3000 manifest_dir
3001 .parent()?
3002 .parent()?
3003 .join("apps")
3004 .join("dashboard")
3005 .join("dist"),
3006 )
3007}
3008
3009fn open_browser(url: &str) {
3011 #[cfg(target_os = "macos")]
3012 let mut command = {
3013 let mut c = std::process::Command::new("open");
3014 c.arg(url);
3015 c
3016 };
3017 #[cfg(target_os = "windows")]
3018 let mut command = {
3019 let mut c = std::process::Command::new("cmd");
3020 c.args(["/C", "start", "", url]);
3022 c
3023 };
3024 #[cfg(not(any(target_os = "macos", target_os = "windows")))]
3025 let mut command = {
3026 let mut c = std::process::Command::new("xdg-open");
3027 c.arg(url);
3028 c
3029 };
3030
3031 match command.spawn() {
3032 Ok(mut child) => {
3033 std::thread::spawn(move || {
3035 let _ = child.wait();
3036 });
3037 }
3038 Err(e) => eprintln!("kranz: could not open the browser: {e}"),
3039 }
3040}
3041
3042fn resolve_release_token(
3054 repo: &Path,
3055 url: &str,
3056 flag: Option<String>,
3057) -> Result<String, ReleaseTokenError> {
3058 let parsed = reqwest::Url::parse(url).ok();
3059 let operator_lookup = parsed
3060 .as_ref()
3061 .and_then(|url| {
3062 kranz_engine::paths::global_config()
3063 .as_deref()
3064 .map(|global| operator_token_for_url(global, url))
3065 })
3066 .unwrap_or(OperatorTokenLookup::Absent);
3067 let allow_repo_fallback = parsed.as_ref().is_some_and(automatic_repo_token_allowed);
3068 resolve_release_token_from_sources(repo, operator_lookup, allow_repo_fallback, flag)
3069}
3070
3071#[derive(Debug, Clone, PartialEq, Eq)]
3076enum OperatorTokenLookup {
3077 Absent,
3078 Found(String),
3079 Legacy(String),
3080 Ambiguous,
3081}
3082
3083#[derive(Debug, Clone, PartialEq, Eq)]
3084enum ReleaseTokenError {
3085 Absent,
3086 Ambiguous,
3087}
3088
3089fn resolve_release_token_from_sources(
3090 repo: &Path,
3091 operator_lookup: OperatorTokenLookup,
3092 allow_repo_fallback: bool,
3093 flag: Option<String>,
3094) -> Result<String, ReleaseTokenError> {
3095 if let Some(flag) = flag {
3096 return Ok(flag);
3097 }
3098 if let Ok(env) = std::env::var("KRANZ_TOKEN") {
3099 return Ok(env);
3100 }
3101 match operator_lookup {
3102 OperatorTokenLookup::Found(authority) => Ok(authority),
3103 OperatorTokenLookup::Legacy(authority) => {
3108 let repo_authority = allow_repo_fallback
3109 .then(|| read_token_file(&repo.join(".kranz").join("serve.token")))
3110 .flatten();
3111 Ok(repo_authority.unwrap_or(authority))
3112 }
3113 OperatorTokenLookup::Ambiguous => Err(ReleaseTokenError::Ambiguous),
3114 OperatorTokenLookup::Absent => allow_repo_fallback
3115 .then(|| read_token_file(&repo.join(".kranz").join("serve.token")))
3116 .flatten()
3117 .ok_or(ReleaseTokenError::Absent),
3118 }
3119}
3120
3121fn scan_operator_token_addresses(
3122 global_config: &Path,
3123 addresses: &[std::net::SocketAddr],
3124) -> OperatorTokenLookup {
3125 let mut authority = None;
3126 for address in addresses {
3127 if let Some(found) = read_token_file(&operator_serve_token_path(global_config, *address)) {
3128 if authority.is_some() {
3129 return OperatorTokenLookup::Ambiguous;
3132 }
3133 authority = Some(found);
3134 }
3135 }
3136 match authority {
3137 Some(authority) => OperatorTokenLookup::Found(authority),
3138 None => OperatorTokenLookup::Absent,
3139 }
3140}
3141
3142fn operator_token_for_url(global_config: &Path, url: &reqwest::Url) -> OperatorTokenLookup {
3143 let Some(port) = url.port_or_known_default() else {
3144 return OperatorTokenLookup::Absent;
3145 };
3146 let Some(host) = normalized_url_host(url) else {
3147 return OperatorTokenLookup::Absent;
3148 };
3149 let endpoint_lookup = if host.eq_ignore_ascii_case("localhost") {
3153 let primary = [
3154 std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, port)),
3155 std::net::SocketAddr::from((std::net::Ipv6Addr::LOCALHOST, port)),
3156 ];
3157 match scan_operator_token_addresses(global_config, &primary) {
3158 OperatorTokenLookup::Absent => scan_operator_token_addresses(
3159 global_config,
3160 &[
3161 std::net::SocketAddr::from((std::net::Ipv4Addr::UNSPECIFIED, port)),
3162 std::net::SocketAddr::from((std::net::Ipv6Addr::UNSPECIFIED, port)),
3163 ],
3164 ),
3165 other => other,
3166 }
3167 } else {
3168 let Ok(ip) = host.parse::<std::net::IpAddr>() else {
3169 return OperatorTokenLookup::Absent;
3170 };
3171 if !ip.is_loopback() {
3174 return OperatorTokenLookup::Absent;
3175 }
3176 let primary = [std::net::SocketAddr::new(ip, port)];
3177 match scan_operator_token_addresses(global_config, &primary) {
3178 OperatorTokenLookup::Absent => {
3179 let unspecified = match ip {
3180 std::net::IpAddr::V4(_) => {
3181 std::net::IpAddr::V4(std::net::Ipv4Addr::UNSPECIFIED)
3182 }
3183 std::net::IpAddr::V6(_) => {
3184 std::net::IpAddr::V6(std::net::Ipv6Addr::UNSPECIFIED)
3185 }
3186 };
3187 scan_operator_token_addresses(
3188 global_config,
3189 &[std::net::SocketAddr::new(unspecified, port)],
3190 )
3191 }
3192 other => other,
3193 }
3194 };
3195 match endpoint_lookup {
3196 OperatorTokenLookup::Absent => {
3197 match read_token_file(&legacy_operator_serve_token_path(global_config, port)) {
3199 Some(token) => OperatorTokenLookup::Legacy(token),
3200 None => OperatorTokenLookup::Absent,
3201 }
3202 }
3203 other => other,
3204 }
3205}
3206
3207fn automatic_repo_token_allowed(url: &reqwest::Url) -> bool {
3208 normalized_url_host(url).is_some_and(|host| {
3209 host.eq_ignore_ascii_case("localhost")
3210 || host
3211 .parse::<std::net::IpAddr>()
3212 .is_ok_and(|ip| ip.is_loopback())
3213 })
3214}
3215
3216fn normalized_url_host(url: &reqwest::Url) -> Option<&str> {
3217 let host = url.host_str()?;
3218 Some(
3219 host.strip_prefix('[')
3220 .and_then(|host| host.strip_suffix(']'))
3221 .unwrap_or(host),
3222 )
3223}
3224
3225fn read_token_file(path: &Path) -> Option<String> {
3226 let contents = std::fs::read_to_string(path).ok()?;
3227 let authority = contents.trim_end().to_string();
3228 if authority.is_empty() {
3229 None
3230 } else {
3231 Some(authority)
3232 }
3233}
3234
3235fn release_http_client() -> Result<reqwest::Client> {
3236 reqwest::Client::builder()
3237 .redirect(reqwest::redirect::Policy::none())
3238 .connect_timeout(Duration::from_secs(5))
3239 .timeout(Duration::from_secs(30))
3240 .build()
3241 .context("building release HTTP client")
3242}
3243
3244fn release_repo_id_is_valid(id: &str) -> bool {
3245 let mut chars = id.chars();
3246 let valid_first = chars.next().is_some_and(|ch| ch.is_ascii_alphanumeric());
3247 let valid_rest = chars.all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_'));
3248 valid_first && valid_rest
3249}
3250
3251async fn cmd_release(
3252 repo: &Path,
3253 mission_id: &str,
3254 url: &str,
3255 token: Option<String>,
3256) -> Result<i32> {
3257 let token = match resolve_release_token(repo, url, token) {
3258 Ok(token) => token,
3259 Err(ReleaseTokenError::Ambiguous) => {
3260 bail!(
3261 "multiple ~/.kranz/serve/<endpoint>.token files match {url}; \
3262 pass --token for the intended serve instead of guessing"
3263 );
3264 }
3265 Err(ReleaseTokenError::Absent) => {
3266 bail!(
3267 "no mutation token available — pass --token, set $KRANZ_TOKEN, or use a local URL matching \
3268 a live ~/.kranz/serve/<endpoint>.token / single-repo .kranz/serve.token \
3269 (the token `kranz serve` prints on startup)"
3270 );
3271 }
3272 };
3273
3274 let client = release_http_client()?;
3275 let repo_id = resolve_release_repo_id(repo, url, &client, &token).await?;
3276 if let Some(id) = repo_id.as_deref() {
3277 if !release_repo_id_is_valid(id) {
3278 bail!("live serve returned an invalid repository id '{id}'; refusing release");
3279 }
3280 }
3281 let endpoint = release_endpoint(url, mission_id, repo_id.as_deref());
3282
3283 let response = client
3284 .post(&endpoint)
3285 .header("x-kranz-token", token)
3286 .json(&serde_json::json!({}))
3287 .send()
3288 .await
3289 .map_err(|e| {
3290 if e.is_connect() {
3291 anyhow!("no kranz serve reachable at {url} — is it running?")
3292 } else {
3293 anyhow::Error::new(e).context(format!("releasing mission '{mission_id}'"))
3294 }
3295 })?;
3296
3297 match response.status() {
3298 reqwest::StatusCode::OK => {
3299 let body: serde_json::Value = response.json().await.unwrap_or_default();
3300 let released = body
3301 .get("released")
3302 .and_then(serde_json::Value::as_bool)
3303 .unwrap_or(false);
3304 if released {
3305 println!("mission {mission_id} released — the lock is now free");
3306 } else {
3307 println!("mission {mission_id} was already free (no lock held)");
3308 }
3309 Ok(0)
3310 }
3311 reqwest::StatusCode::CONFLICT => {
3312 eprintln!("kranz: mission {mission_id} has a turn in flight — try again shortly");
3313 Ok(1)
3314 }
3315 reqwest::StatusCode::NOT_FOUND => {
3316 eprintln!("kranz: unknown mission '{mission_id}' at {url}");
3317 Ok(1)
3318 }
3319 reqwest::StatusCode::UNAUTHORIZED => {
3320 eprintln!(
3321 "kranz: the token was missing or invalid — check --token / $KRANZ_TOKEN \
3322 against the token `kranz serve` printed on startup"
3323 );
3324 Ok(1)
3325 }
3326 other => {
3327 let body = response.text().await.unwrap_or_default();
3328 eprintln!("kranz: release failed ({other}): {body}");
3329 Ok(1)
3330 }
3331 }
3332}
3333
3334async fn resolve_release_repo_id(
3339 repo: &Path,
3340 url: &str,
3341 client: &reqwest::Client,
3342 token: &str,
3343) -> Result<Option<String>> {
3344 live_release_repo_id(repo, url, client, token).await
3345}
3346
3347fn release_repo_id_from_summaries(
3348 repo: &Path,
3349 repos: &[kranz_server::RepoSummary],
3350) -> Result<Option<String>> {
3351 if repos.is_empty() {
3352 bail!("the live serve returned an empty repository catalog; refusing an unscoped release");
3353 }
3354 let root = std::fs::canonicalize(repo).unwrap_or_else(|_| repo.to_path_buf());
3355 let root_str = root.to_string_lossy();
3356 let matches: Vec<&kranz_server::RepoSummary> = repos
3357 .iter()
3358 .filter(|entry| {
3359 let entry_root = PathBuf::from(&entry.root);
3360 let entry_canon =
3361 std::fs::canonicalize(&entry_root).unwrap_or_else(|_| entry_root.clone());
3362 entry_canon == root || entry.root == root_str
3363 })
3364 .collect();
3365 match matches.as_slice() {
3366 [one] => Ok(Some(one.id.clone())),
3367 [] => Err(anyhow!(
3368 "selected repository '{}' is not present in the live serve catalog; refusing an unscoped release",
3369 repo.display()
3370 )),
3371 _ => Err(anyhow!(
3372 "selected repository '{}' matches multiple live serve catalog entries; refusing release",
3373 repo.display()
3374 )),
3375 }
3376}
3377
3378async fn live_release_repo_id(
3379 repo: &Path,
3380 url: &str,
3381 client: &reqwest::Client,
3382 token: &str,
3383) -> Result<Option<String>> {
3384 let base = url.trim_end_matches('/');
3385 let mut request = client.get(format!("{base}/api/repos"));
3389 let loopback_catalog = reqwest::Url::parse(url)
3390 .ok()
3391 .is_some_and(|parsed| automatic_repo_token_allowed(&parsed));
3392 if !loopback_catalog {
3393 request = request.header("x-kranz-token", token);
3394 }
3395 let send_error = |e: reqwest::Error| {
3396 if e.is_connect() {
3397 anyhow!("no kranz serve reachable at {url} — is it running?")
3398 } else {
3399 anyhow::Error::new(e).context("listing live serve repositories")
3400 }
3401 };
3402 let mut response = request.send().await.map_err(send_error)?;
3403 if loopback_catalog && response.status() == reqwest::StatusCode::UNAUTHORIZED {
3408 response = client
3409 .get(format!("{base}/api/repos"))
3410 .header("x-kranz-token", token)
3411 .send()
3412 .await
3413 .map_err(send_error)?;
3414 }
3415 if !response.status().is_success() {
3416 bail!(
3417 "cannot list repositories at {url}: HTTP {}",
3418 response.status()
3419 );
3420 }
3421 let repos: Vec<kranz_server::RepoSummary> = response
3422 .json()
3423 .await
3424 .context("parsing GET /api/repos response")?;
3425 release_repo_id_from_summaries(repo, &repos)
3426}
3427
3428fn release_endpoint(url: &str, mission_id: &str, repo_id: Option<&str>) -> String {
3429 let base = url.trim_end_matches('/');
3430 match repo_id {
3431 Some(repo_id) => {
3432 format!("{base}/api/repos/{repo_id}/missions/{mission_id}/release")
3433 }
3434 None => format!("{base}/api/missions/{mission_id}/release"),
3435 }
3436}
3437
3438#[cfg(test)]
3439mod tests {
3440 use super::*;
3441 use std::fs;
3442
3443 static KRANZ_TOKEN_ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
3448
3449 #[tokio::test]
3452 #[allow(clippy::await_holding_lock)]
3456 async fn release_without_a_token_errors_clearly() {
3457 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3458 unsafe {
3460 std::env::remove_var("KRANZ_TOKEN");
3461 }
3462 let tmp = tempfile::tempdir().unwrap();
3463 let repo = tmp.path().to_path_buf();
3464 let err = cmd_release(&repo, "m-1", "http://192.0.2.1:4560", None)
3468 .await
3469 .unwrap_err();
3470 let msg = err.to_string();
3471 assert!(msg.contains("--token"), "{msg}");
3472 assert!(msg.contains("KRANZ_TOKEN"), "{msg}");
3473 }
3474
3475 #[test]
3478 fn release_reads_token_file_when_flag_and_env_absent() {
3479 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3480 unsafe {
3482 std::env::remove_var("KRANZ_TOKEN");
3483 }
3484 let tmp = tempfile::tempdir().unwrap();
3485 let repo = tmp.path().to_path_buf();
3486 write_serve_token(&repo, "file-token").unwrap();
3487
3488 assert_eq!(
3489 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Absent, true, None),
3490 Ok("file-token".to_string())
3491 );
3492 }
3493
3494 #[test]
3495 fn release_reads_token_file_but_flag_overrides() {
3496 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3497 unsafe {
3499 std::env::remove_var("KRANZ_TOKEN");
3500 }
3501 let tmp = tempfile::tempdir().unwrap();
3502 let repo = tmp.path().to_path_buf();
3503 write_serve_token(&repo, "file-token").unwrap();
3504
3505 assert_eq!(
3506 resolve_release_token_from_sources(
3507 &repo,
3508 OperatorTokenLookup::Absent,
3509 true,
3510 Some("flag-token".to_string()),
3511 ),
3512 Ok("flag-token".to_string())
3513 );
3514 }
3515
3516 #[test]
3517 fn release_reads_token_file_but_env_overrides() {
3518 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3519 unsafe {
3521 std::env::set_var("KRANZ_TOKEN", "env-token");
3522 }
3523 let tmp = tempfile::tempdir().unwrap();
3524 let repo = tmp.path().to_path_buf();
3525 write_serve_token(&repo, "file-token").unwrap();
3526
3527 let result =
3528 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Absent, true, None);
3529 unsafe {
3530 std::env::remove_var("KRANZ_TOKEN");
3531 }
3532 assert_eq!(result, Ok("env-token".to_string()));
3533 }
3534
3535 #[test]
3538 fn release_reads_token_file_none_present_yields_none() {
3539 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3540 unsafe {
3542 std::env::remove_var("KRANZ_TOKEN");
3543 }
3544 let tmp = tempfile::tempdir().unwrap();
3545 let repo = tmp.path().to_path_buf();
3546
3547 assert_eq!(
3548 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Absent, true, None),
3549 Err(ReleaseTokenError::Absent)
3550 );
3551 }
3552
3553 #[test]
3554 fn release_prefers_operator_process_token_over_repo_compatibility_token() {
3555 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3556 unsafe {
3557 std::env::remove_var("KRANZ_TOKEN");
3558 }
3559 let tmp = tempfile::tempdir().unwrap();
3560 let repo = tmp.path().join("repo");
3561 write_serve_token(&repo, "repo-token").unwrap();
3562 let global = tmp.path().join("operator").join("config.json");
3563 let address = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3564 let operator = operator_serve_token_path(&global, address);
3565 write_token_file(&operator, "operator-token").unwrap();
3566 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3567
3568 assert_eq!(
3569 resolve_release_token_from_sources(
3570 &repo,
3571 operator_token_for_url(&global, &url),
3572 true,
3573 None,
3574 ),
3575 Ok("operator-token".to_string())
3576 );
3577 }
3578
3579 #[test]
3580 fn ambiguous_operator_token_does_not_fall_through_to_repo_file() {
3581 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3582 unsafe {
3583 std::env::remove_var("KRANZ_TOKEN");
3584 }
3585 let tmp = tempfile::tempdir().unwrap();
3586 let repo = tmp.path().join("repo");
3587 write_serve_token(&repo, "repo-token").unwrap();
3588
3589 assert_eq!(
3590 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Ambiguous, true, None,),
3591 Err(ReleaseTokenError::Ambiguous)
3592 );
3593 }
3594
3595 #[test]
3596 fn operator_token_discovery_is_ambiguous_for_dual_localhost_endpoint_files() {
3597 let tmp = tempfile::tempdir().unwrap();
3598 let global = tmp.path().join(".kranz").join("config.json");
3599 let v4 = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3600 let v6 = std::net::SocketAddr::from((std::net::Ipv6Addr::LOCALHOST, 4560));
3601 write_token_file(&operator_serve_token_path(&global, v4), "v4-token").unwrap();
3602 write_token_file(&operator_serve_token_path(&global, v6), "v6-token").unwrap();
3603 let url = reqwest::Url::parse("http://localhost:4560").unwrap();
3604
3605 assert_eq!(
3606 operator_token_for_url(&global, &url),
3607 OperatorTokenLookup::Ambiguous
3608 );
3609 }
3610
3611 #[test]
3612 fn operator_token_path_is_scoped_by_full_bound_endpoint() {
3613 let global = Path::new("/operator/.kranz/config.json");
3614 let a =
3615 operator_serve_token_path(global, std::net::SocketAddr::from(([127, 0, 0, 1], 4560)));
3616 let b =
3617 operator_serve_token_path(global, std::net::SocketAddr::from(([127, 0, 0, 2], 4560)));
3618 assert_ne!(a, b);
3619 assert_eq!(
3620 a,
3621 Path::new("/operator/.kranz/serve/v4-7f000001-4560.token")
3622 );
3623 }
3624
3625 #[test]
3626 fn operator_token_discovery_handles_ipv6_url_brackets() {
3627 let tmp = tempfile::tempdir().unwrap();
3628 let global = tmp.path().join(".kranz").join("config.json");
3629 let address = std::net::SocketAddr::from((std::net::Ipv6Addr::LOCALHOST, 4560));
3630 write_token_file(&operator_serve_token_path(&global, address), "ipv6-token").unwrap();
3631 let url = reqwest::Url::parse("http://[::1]:4560").unwrap();
3632
3633 assert_eq!(
3634 operator_token_for_url(&global, &url),
3635 OperatorTokenLookup::Found("ipv6-token".to_string())
3636 );
3637 }
3638
3639 #[test]
3640 fn automatic_token_discovery_refuses_remote_domain_urls() {
3641 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3642 unsafe {
3643 std::env::remove_var("KRANZ_TOKEN");
3644 }
3645 let tmp = tempfile::tempdir().unwrap();
3646 let repo = tmp.path().join("repo");
3647 write_serve_token(&repo, "repo-token").unwrap();
3648
3649 assert_eq!(
3650 resolve_release_token(&repo, "https://example.com:4560", None),
3651 Err(ReleaseTokenError::Absent)
3652 );
3653 }
3654
3655 #[test]
3656 fn operator_token_discovery_refuses_non_loopback_ip_literals() {
3657 let tmp = tempfile::tempdir().unwrap();
3658 let global = tmp.path().join(".kranz").join("config.json");
3659 let address = std::net::SocketAddr::from(([203, 0, 113, 10], 4560));
3660 write_token_file(&operator_serve_token_path(&global, address), "remote-token").unwrap();
3661 let url = reqwest::Url::parse("http://203.0.113.10:4560").unwrap();
3662
3663 assert_eq!(
3664 operator_token_for_url(&global, &url),
3665 OperatorTokenLookup::Absent
3666 );
3667 }
3668
3669 #[test]
3670 fn operator_token_discovery_prefers_exact_loopback_over_stale_unspecified() {
3671 let tmp = tempfile::tempdir().unwrap();
3672 let global = tmp.path().join(".kranz").join("config.json");
3673 let loopback = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3674 let unspecified = std::net::SocketAddr::from((std::net::Ipv4Addr::UNSPECIFIED, 4560));
3675 write_token_file(&operator_serve_token_path(&global, loopback), "live-token").unwrap();
3676 write_token_file(
3677 &operator_serve_token_path(&global, unspecified),
3678 "stale-token",
3679 )
3680 .unwrap();
3681 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3682
3683 assert_eq!(
3684 operator_token_for_url(&global, &url),
3685 OperatorTokenLookup::Found("live-token".to_string())
3686 );
3687 }
3688
3689 #[test]
3690 fn operator_token_discovery_falls_back_to_legacy_port_token() {
3691 let tmp = tempfile::tempdir().unwrap();
3692 let global = tmp.path().join(".kranz").join("config.json");
3693 write_token_file(
3694 &legacy_operator_serve_token_path(&global, 4560),
3695 "legacy-token",
3696 )
3697 .unwrap();
3698 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3699
3700 assert_eq!(
3701 operator_token_for_url(&global, &url),
3702 OperatorTokenLookup::Legacy("legacy-token".to_string())
3703 );
3704 }
3705
3706 #[test]
3707 fn stale_legacy_operator_token_does_not_mask_live_repo_token() {
3708 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3709 unsafe {
3710 std::env::remove_var("KRANZ_TOKEN");
3711 }
3712 let tmp = tempfile::tempdir().unwrap();
3713 let repo = tmp.path().join("repo");
3714 write_serve_token(&repo, "live-repo-token").unwrap();
3715 let global = tmp.path().join(".kranz").join("config.json");
3716 write_token_file(
3717 &legacy_operator_serve_token_path(&global, 4560),
3718 "stale-legacy-token",
3719 )
3720 .unwrap();
3721 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3722 let operator = operator_token_for_url(&global, &url);
3723
3724 assert_eq!(
3725 resolve_release_token_from_sources(&repo, operator, true, None),
3726 Ok("live-repo-token".to_string())
3727 );
3728 }
3729
3730 #[test]
3731 fn operator_token_discovery_prefers_endpoint_file_over_legacy_port_token() {
3732 let tmp = tempfile::tempdir().unwrap();
3733 let global = tmp.path().join(".kranz").join("config.json");
3734 let loopback = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3735 write_token_file(
3736 &operator_serve_token_path(&global, loopback),
3737 "endpoint-token",
3738 )
3739 .unwrap();
3740 write_token_file(
3741 &legacy_operator_serve_token_path(&global, 4560),
3742 "legacy-token",
3743 )
3744 .unwrap();
3745 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3746
3747 assert_eq!(
3748 operator_token_for_url(&global, &url),
3749 OperatorTokenLookup::Found("endpoint-token".to_string())
3750 );
3751 }
3752
3753 #[test]
3754 fn release_repo_id_from_summaries_refuses_duplicate_root() {
3755 let tmp = tempfile::tempdir().unwrap();
3756 let repo = tmp.path().join("repo");
3757 std::fs::create_dir_all(&repo).unwrap();
3758 let root = repo.to_string_lossy().into_owned();
3759 let summaries = vec![
3760 kranz_server::RepoSummary {
3761 id: "alpha".to_string(),
3762 root: root.clone(),
3763 display_name: "alpha".to_string(),
3764 group: None,
3765 pinned: false,
3766 is_default: true,
3767 status: "healthy".to_string(),
3768 error: None,
3769 activity: kranz_server::RepoActivity::default(),
3770 },
3771 kranz_server::RepoSummary {
3772 id: "beta".to_string(),
3773 root,
3774 display_name: "beta".to_string(),
3775 group: None,
3776 pinned: false,
3777 is_default: false,
3778 status: "healthy".to_string(),
3779 error: None,
3780 activity: kranz_server::RepoActivity::default(),
3781 },
3782 ];
3783 let error = release_repo_id_from_summaries(&repo, &summaries).unwrap_err();
3784 assert!(error.to_string().contains("matches multiple"));
3785 }
3786
3787 #[test]
3788 fn slack_operator_catalog_is_not_rejected_by_legacy_context_helper() {
3789 fn init_git(root: &Path) {
3790 std::fs::create_dir_all(root).unwrap();
3791 let status = std::process::Command::new("git")
3792 .args(["init", "-q"])
3793 .arg(root)
3794 .status()
3795 .unwrap();
3796 assert!(status.success());
3797 }
3798
3799 let tmp = tempfile::tempdir().unwrap();
3800 let a = tmp.path().join("a");
3801 init_git(&a);
3802 let multi = kranz_server::MultiRepoHost::from_config(kranz_server::HostConfig {
3803 default_repo: Some("a".to_string()),
3804 max_concurrent_repos: 1,
3805 repos: vec![kranz_server::RepoConfig {
3806 id: "a".to_string(),
3807 root: a,
3808 display_name: None,
3809 group: None,
3810 pinned: false,
3811 slack: kranz_server::RepoSlackConfig::default(),
3812 }],
3813 })
3814 .unwrap();
3815 let context = single_repo_slack_context(&multi).unwrap();
3816 assert_eq!(context.id(), "a");
3817 }
3818
3819 #[test]
3820 fn release_endpoint_is_repo_scoped_when_catalog_id_is_known() {
3821 assert_eq!(
3822 release_endpoint("http://127.0.0.1:4560/", "same-id", Some("repo-b")),
3823 "http://127.0.0.1:4560/api/repos/repo-b/missions/same-id/release"
3824 );
3825 }
3826
3827 #[test]
3828 fn release_repo_id_from_summaries_matches_live_catalog_root() {
3829 let tmp = tempfile::tempdir().unwrap();
3830 let repo = tmp.path().join("repo");
3831 std::fs::create_dir_all(&repo).unwrap();
3832 let summaries = vec![kranz_server::RepoSummary {
3833 id: "alpha".to_string(),
3834 root: repo.to_string_lossy().into_owned(),
3835 display_name: "alpha".to_string(),
3836 group: None,
3837 pinned: false,
3838 is_default: true,
3839 status: "healthy".to_string(),
3840 error: None,
3841 activity: kranz_server::RepoActivity::default(),
3842 }];
3843 assert_eq!(
3844 release_repo_id_from_summaries(&repo, &summaries)
3845 .unwrap()
3846 .as_deref(),
3847 Some("alpha")
3848 );
3849 }
3850
3851 #[test]
3852 fn release_repo_id_from_summaries_refuses_unmatched_root() {
3853 let tmp = tempfile::tempdir().unwrap();
3854 let repo = tmp.path().join("local");
3855 let other = tmp.path().join("other");
3856 std::fs::create_dir_all(&repo).unwrap();
3857 std::fs::create_dir_all(&other).unwrap();
3858 let summaries = vec![kranz_server::RepoSummary {
3859 id: "other".to_string(),
3860 root: other.to_string_lossy().into_owned(),
3861 display_name: "other".to_string(),
3862 group: None,
3863 pinned: false,
3864 is_default: false,
3865 status: "healthy".to_string(),
3866 error: None,
3867 activity: kranz_server::RepoActivity::default(),
3868 }];
3869 let error = release_repo_id_from_summaries(&repo, &summaries).unwrap_err();
3870 assert!(error
3871 .to_string()
3872 .contains("not present in the live serve catalog"));
3873 }
3874
3875 #[test]
3876 fn release_repo_id_from_summaries_refuses_empty_catalog() {
3877 let tmp = tempfile::tempdir().unwrap();
3878 let error = release_repo_id_from_summaries(tmp.path(), &[]).unwrap_err();
3879 assert!(error.to_string().contains("empty repository catalog"));
3880 }
3881
3882 #[tokio::test]
3883 async fn live_release_repo_lookup_authenticates_protected_catalog() {
3884 let tmp = tempfile::tempdir().unwrap();
3885 let repo = tmp.path().join("repo");
3886 std::fs::create_dir_all(&repo).unwrap();
3887 let status = std::process::Command::new("git")
3888 .args(["init", "-q"])
3889 .arg(&repo)
3890 .status()
3891 .unwrap();
3892 assert!(status.success());
3893
3894 let catalog = Arc::new(kranz_server::MultiRepoHost::single(repo.clone()).unwrap());
3895 let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
3896 .await
3897 .unwrap();
3898 let address = listener.local_addr().unwrap();
3899 let app = kranz_server::router_with_multi_repo_host_and_addr(
3901 catalog,
3902 None,
3903 kranz_server::MutationAuthority::new("catalog-token").unwrap(),
3904 Some(address),
3905 true,
3906 false,
3907 );
3908 let server = tokio::spawn(async move {
3909 axum::serve(listener, app).await.unwrap();
3910 });
3911 let client = release_http_client().unwrap();
3912 let url = format!("http://{address}");
3913
3914 let repo_id = live_release_repo_id(&repo, &url, &client, "catalog-token")
3915 .await
3916 .unwrap();
3917
3918 server.abort();
3919 let _ = server.await;
3920 assert_eq!(repo_id.as_deref(), Some("repo"));
3921 }
3922
3923 #[tokio::test]
3924 async fn live_release_repo_lookup_retries_with_token_when_read_gated() {
3925 let tmp = tempfile::tempdir().unwrap();
3926 let repo = tmp.path().join("repo");
3927 std::fs::create_dir_all(&repo).unwrap();
3928 let status = std::process::Command::new("git")
3929 .args(["init", "-q"])
3930 .arg(&repo)
3931 .status()
3932 .unwrap();
3933 assert!(status.success());
3934
3935 let catalog = Arc::new(kranz_server::MultiRepoHost::single(repo.clone()).unwrap());
3936 let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
3937 .await
3938 .unwrap();
3939 let address = listener.local_addr().unwrap();
3940 let app = kranz_server::router_with_multi_repo_host_and_addr(
3945 catalog,
3946 None,
3947 kranz_server::MutationAuthority::new("catalog-token").unwrap(),
3948 Some(address),
3949 false,
3950 true,
3951 );
3952 let server = tokio::spawn(async move {
3953 axum::serve(listener, app).await.unwrap();
3954 });
3955 let client = release_http_client().unwrap();
3956 let url = format!("http://{address}");
3957
3958 let repo_id = live_release_repo_id(&repo, &url, &client, "catalog-token")
3959 .await
3960 .unwrap();
3961
3962 server.abort();
3963 let _ = server.await;
3964 assert_eq!(repo_id.as_deref(), Some("repo"));
3965 }
3966
3967 #[test]
3968 fn empty_token_files_are_ignored() {
3969 let tmp = tempfile::tempdir().unwrap();
3970 let path = tmp.path().join("serve.token");
3971 std::fs::write(&path, "").unwrap();
3972 assert_eq!(read_token_file(&path), None);
3973 }
3974
3975 #[cfg(unix)]
3976 #[test]
3977 fn serve_token_file_is_written_with_owner_only_permissions() {
3978 use std::os::unix::fs::PermissionsExt;
3979 let tmp = tempfile::tempdir().unwrap();
3980 let repo = tmp.path().to_path_buf();
3981 let path = write_serve_token(&repo, "secret").unwrap();
3982 let mode = std::fs::metadata(&path).unwrap().permissions().mode();
3983 assert_eq!(mode & 0o777, 0o600);
3984 }
3985
3986 #[cfg(unix)]
3987 #[test]
3988 fn serve_read_token_file_is_written_with_owner_only_permissions() {
3989 use std::os::unix::fs::PermissionsExt;
3990 let tmp = tempfile::tempdir().unwrap();
3991 let repo = tmp.path().to_path_buf();
3992 let path = write_serve_read_token(&repo, "read-secret").unwrap();
3993 assert!(path.ends_with(".kranz/serve.read.token"));
3994 let mode = std::fs::metadata(&path).unwrap().permissions().mode();
3995 assert_eq!(mode & 0o777, 0o600);
3996 }
3997
3998 #[test]
3999 fn operator_read_token_path_sits_beside_the_operator_token() {
4000 let global = Path::new("/home/op/.kranz/config.json");
4001 let address = std::net::SocketAddr::from(([127, 0, 0, 1], 4560));
4002 let mutation = operator_serve_token_path(global, address);
4003 let read = operator_serve_read_token_path(global, address);
4004 assert_eq!(read, mutation.with_extension("read.token"));
4005 assert!(read.to_string_lossy().ends_with(".read.token"));
4006 }
4007
4008 #[cfg(unix)]
4009 #[test]
4010 fn existing_token_permissions_are_hardened_before_replacement() {
4011 use std::os::unix::fs::PermissionsExt;
4012 let tmp = tempfile::tempdir().unwrap();
4013 let path = tmp.path().join("serve.token");
4014 std::fs::write(&path, "old-token").unwrap();
4015 std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o644)).unwrap();
4016
4017 write_token_file(&path, "new-token").unwrap();
4018
4019 assert_eq!(std::fs::read_to_string(&path).unwrap(), "new-token");
4020 let mode = std::fs::metadata(&path).unwrap().permissions().mode();
4021 assert_eq!(mode & 0o777, 0o600);
4022 }
4023
4024 #[tokio::test]
4025 async fn serve_token_file_is_removed_after_graceful_shutdown() {
4026 let tmp = tempfile::tempdir().unwrap();
4027 let repo = tmp.path().to_path_buf();
4028 let path = repo.join(".kranz").join("serve.token");
4029 let read_path = repo.join(".kranz").join("serve.read.token");
4030
4031 let host = std::sync::Arc::new(kranz_server::MissionHost::new(repo.clone()));
4032 let listener =
4033 kranz_server::bind_listener(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST), 0)
4034 .await
4035 .unwrap();
4036 serve_with_token_cleanup(
4037 &repo,
4038 host,
4039 listener,
4040 None,
4041 "tok".to_string(),
4042 "read-tok".to_string(),
4043 std::future::ready(()),
4044 )
4045 .await
4046 .unwrap();
4047
4048 assert!(!path.exists());
4049 assert!(!read_path.exists());
4050 }
4051
4052 fn dashboard_at(path: PathBuf) -> PathBuf {
4053 fs::create_dir_all(&path).unwrap();
4054 fs::write(path.join("index.html"), "<!doctype html>").unwrap();
4055 path
4056 }
4057
4058 fn inputs() -> DashboardResolutionInputs {
4059 DashboardResolutionInputs {
4060 embedded_available: false,
4061 ..Default::default()
4062 }
4063 }
4064
4065 #[test]
4066 fn dashboard_resolution_honors_explicit_path_verbatim() {
4067 let tmp = tempfile::tempdir().unwrap();
4068 let repo = tmp.path().join("repo");
4069 let explicit = tmp.path().join("missing-dashboard");
4070
4071 assert_eq!(
4072 resolve_dashboard_assets_from(&repo, Some(explicit.clone()), &inputs()),
4073 Some(DashboardAssets::Dir(explicit))
4074 );
4075 }
4076
4077 #[test]
4078 fn dashboard_resolution_env_precedes_repo_and_invalid_env_is_skipped() {
4079 let tmp = tempfile::tempdir().unwrap();
4080 let repo = tmp.path().join("repo");
4081 let repo_dist = dashboard_at(repo.join("apps").join("dashboard").join("dist"));
4082 let env_dist = dashboard_at(tmp.path().join("env-dist"));
4083
4084 let mut with_env = inputs();
4085 with_env.env_dist = Some(env_dist.clone());
4086 assert_eq!(
4087 resolve_dashboard_assets_from(&repo, None, &with_env),
4088 Some(DashboardAssets::Dir(env_dist))
4089 );
4090
4091 let mut with_invalid_env = inputs();
4092 with_invalid_env.env_dist = Some(tmp.path().join("missing-env-dist"));
4093 assert_eq!(
4094 resolve_dashboard_assets_from(&repo, None, &with_invalid_env),
4095 Some(DashboardAssets::Dir(repo_dist))
4096 );
4097 }
4098
4099 #[test]
4100 fn dashboard_resolution_finds_installed_asset_dirs() {
4101 let tmp = tempfile::tempdir().unwrap();
4102 let repo = tmp.path().join("repo");
4103 let exe = tmp.path().join("prefix").join("bin").join("kranz");
4104 let installed = dashboard_at(
4105 tmp.path()
4106 .join("prefix")
4107 .join("share")
4108 .join("kranz")
4109 .join("dashboard")
4110 .join("dist"),
4111 );
4112
4113 let mut inputs = inputs();
4114 inputs.exe = Some(exe);
4115 assert_eq!(
4116 resolve_dashboard_assets_from(&repo, None, &inputs),
4117 Some(DashboardAssets::Dir(installed))
4118 );
4119 }
4120
4121 #[test]
4122 fn dashboard_resolution_finds_checkout_used_to_build_installed_binary() {
4123 let tmp = tempfile::tempdir().unwrap();
4124 let repo = tmp.path().join("mission-repo");
4125 let checkout = tmp.path().join("kranz");
4126 let manifest_dir = checkout.join("crates").join("cli");
4127 let checkout_dist = dashboard_at(checkout.join("apps").join("dashboard").join("dist"));
4128
4129 let mut inputs = inputs();
4130 inputs.manifest_dir = Some(manifest_dir);
4131 inputs.exe = Some(tmp.path().join("cargo-home").join("bin").join("kranz"));
4132 assert_eq!(
4133 resolve_dashboard_assets_from(&repo, None, &inputs),
4134 Some(DashboardAssets::Dir(checkout_dist))
4135 );
4136 }
4137
4138 #[test]
4139 fn dashboard_resolution_falls_back_to_embedded_assets() {
4140 let tmp = tempfile::tempdir().unwrap();
4141 let repo = tmp.path().join("repo");
4142 let mut inputs = inputs();
4143 inputs.embedded_available = true;
4144
4145 assert_eq!(
4146 resolve_dashboard_assets_from(&repo, None, &inputs),
4147 Some(DashboardAssets::Embedded)
4148 );
4149 }
4150
4151 #[test]
4152 fn embedded_dashboard_bundle_contains_index() {
4153 assert!(
4154 crate::embedded_dashboard::EMBEDDED_DASHBOARD
4155 .iter()
4156 .any(|file| file.path == "index.html"),
4157 "embedded dashboard source: {}",
4158 crate::embedded_dashboard::EMBEDDED_DASHBOARD_SOURCE
4159 );
4160 }
4161
4162 #[test]
4163 fn serve_refuses_non_loopback_without_insecure_lan() {
4164 let bind: std::net::IpAddr = "0.0.0.0".parse().unwrap();
4165 let err = refuse_non_loopback_without_insecure_lan(bind, false).unwrap_err();
4166 let msg = err.to_string();
4167 assert!(
4168 msg.contains("refusing to bind") && msg.contains("--insecure-lan"),
4169 "unexpected error: {msg}"
4170 );
4171 }
4172
4173 #[test]
4174 fn serve_allows_non_loopback_with_insecure_lan() {
4175 let bind: std::net::IpAddr = "0.0.0.0".parse().unwrap();
4176 refuse_non_loopback_without_insecure_lan(bind, true).unwrap();
4177 }
4178
4179 #[test]
4180 fn serve_allows_loopback_without_insecure_lan() {
4181 let bind: std::net::IpAddr = "127.0.0.1".parse().unwrap();
4182 refuse_non_loopback_without_insecure_lan(bind, false).unwrap();
4183 }
4184
4185 #[test]
4186 fn read_auth_on_loopback_requires_read_token() {
4187 assert!(effective_require_read_token(true, true));
4188 }
4189
4190 #[test]
4191 fn read_auth_off_loopback_bind_does_not_require_read_token() {
4192 assert!(!effective_require_read_token(true, false));
4193 }
4194
4195 #[test]
4196 fn read_auth_non_loopback_always_requires_read_token() {
4197 assert!(effective_require_read_token(false, true));
4198 assert!(effective_require_read_token(false, false));
4199 }
4200
4201 fn reconcile_turn(reply: &str) -> Vec<kranz_engine::backend::AgentEvent> {
4206 vec![
4207 kranz_engine::backend_mock::mock_text(reply),
4208 kranz_engine::backend_mock::mock_result_text(reply),
4209 ]
4210 }
4211
4212 fn reconcile_worker_pass() -> kranz_engine::backend_mock::MockScript {
4213 kranz_engine::backend_mock::MockScript::single_shot_json(&serde_json::json!({
4214 "result": "pass",
4215 "summary": "implemented and tested",
4216 "filesTouched": ["delivered.txt"],
4217 "testsAdded": [],
4218 "testEvidence": "all green",
4219 "commits": []
4220 }))
4221 .writes_file("delivered.txt", "delivered by the mock worker\n")
4222 }
4223
4224 fn reconcile_plan_json() -> serde_json::Value {
4225 serde_json::json!({
4226 "goal": "ship the demo",
4227 "validationContract": [],
4228 "milestones": [{
4229 "title": "M1",
4230 "features": [{
4231 "title": "F1",
4232 "spec": "build the thing",
4233 "validationCriteria": ["it works"]
4234 }]
4235 }]
4236 })
4237 }
4238
4239 #[tokio::test]
4248 async fn reconcile_on_terminal_after_cli_run_marks_ticket_done() {
4249 let tmp = tempfile::tempdir().unwrap();
4250 let repo = tmp.path().to_path_buf();
4251 let status = std::process::Command::new("git")
4252 .args(["init", "-b", "main"])
4253 .current_dir(&repo)
4254 .status()
4255 .unwrap();
4256 assert!(status.success());
4257 std::process::Command::new("git")
4258 .args(["config", "user.name", "test"])
4259 .current_dir(&repo)
4260 .status()
4261 .unwrap();
4262 std::process::Command::new("git")
4263 .args(["config", "user.email", "test@example.com"])
4264 .current_dir(&repo)
4265 .status()
4266 .unwrap();
4267 std::fs::write(repo.join("README.md"), "seed\n").unwrap();
4268 std::process::Command::new("git")
4269 .args(["add", "-A"])
4270 .current_dir(&repo)
4271 .status()
4272 .unwrap();
4273 std::process::Command::new("git")
4274 .args(["commit", "-m", "seed"])
4275 .current_dir(&repo)
4276 .status()
4277 .unwrap();
4278 let repo = std::fs::canonicalize(&repo).unwrap();
4279
4280 let judgement = serde_json::json!({
4281 "decision": "complete",
4282 "guidance": "",
4283 "summary": "worker did the job"
4284 });
4285 let orch_setup = kranz_engine::backend_mock::MockScript::streaming(vec![
4289 kranz_engine::backend_mock::mock_init("orch-session"),
4290 kranz_engine::backend_mock::mock_result_text("seed-hi"),
4291 ])
4292 .responding(vec![
4293 reconcile_turn("let's scope the demo"),
4294 reconcile_turn(&reconcile_plan_json().to_string()),
4295 ]);
4296 let orch_run = kranz_engine::backend_mock::MockScript::streaming(vec![
4297 kranz_engine::backend_mock::mock_init("orch-session-2"),
4298 kranz_engine::backend_mock::mock_result_text("ack"),
4299 ])
4300 .responding(vec![
4301 reconcile_turn("ack"),
4302 reconcile_turn(
4303 &serde_json::json!({"action": "commit-as-is", "note": "worker delivered files"})
4304 .to_string(),
4305 ),
4306 reconcile_turn(&judgement.to_string()),
4307 reconcile_turn("NONE"),
4308 reconcile_turn("NONE"),
4312 reconcile_turn("NONE"),
4313 reconcile_turn("NONE"),
4314 ]);
4315 let backend: Arc<dyn AgentBackend> =
4316 Arc::new(kranz_engine::backend_mock::MockBackend::with_scripts(vec![
4317 orch_setup,
4318 kranz_engine::backend_mock::MockScript::single_shot("ok"),
4323 reconcile_worker_pass(),
4326 orch_run,
4327 ]));
4328
4329 let cfg = MissionConfig {
4330 skip_scrutiny: true,
4331 skip_functional: true,
4332 ..Default::default()
4333 };
4334 let mut engine =
4335 MissionEngine::create(Arc::clone(&backend), repo.clone(), "ship the demo", cfg)
4336 .unwrap();
4337 let mission_id = engine.mission_id().to_string();
4338 engine.planning_turn("ship the demo").await.unwrap();
4339 let request = engine.request_plan().await.unwrap();
4340 let plan = match request {
4341 PlanRequest::Ready(plan) => plan,
4342 PlanRequest::NotReady(text) => panic!("expected a ready plan, got: {text}"),
4343 PlanRequest::WrongPlan { reason } => {
4344 panic!("expected a ready plan, got a wrong-plan escalation: {reason}")
4345 }
4346 };
4347 engine.approve_plan(plan).unwrap();
4348 drop(engine);
4349
4350 kranz_engine::ticket::Ticket::record_mission(&repo, "my-ticket", &mission_id).unwrap();
4353 kranz_engine::ticket::Ticket::write_state(
4354 &repo,
4355 "my-ticket",
4356 kranz_engine::ticket::TicketState::Failed,
4357 None,
4358 )
4359 .unwrap();
4360 assert_eq!(
4361 kranz_engine::ticket::Ticket::read_state(&repo, "my-ticket"),
4362 kranz_engine::ticket::TicketState::Failed
4363 );
4364
4365 let exit_code =
4366 run_mission_loop_with_backend(repo.clone(), mission_id, LockForce::No, false, backend)
4367 .await
4368 .unwrap();
4369 assert_eq!(exit_code, 0);
4370
4371 assert_eq!(
4372 kranz_engine::ticket::Ticket::read_state(&repo, "my-ticket"),
4373 kranz_engine::ticket::TicketState::Done,
4374 "run_mission_loop must reconcile the linked ticket to Done on Complete"
4375 );
4376 }
4377}