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!("{}", serde_json::to_string_pretty(&requests)?);
358 }
359 PermissionCommand::Allow {
360 id,
361 request_id,
362 binding,
363 } => {
364 cmd_permission_answer(&repo, &id, &request_id, &binding, true)?;
365 println!("one-call permission answer queued; delivery is recorded separately");
366 }
367 PermissionCommand::Deny {
368 id,
369 request_id,
370 binding,
371 } => {
372 cmd_permission_answer(&repo, &id, &request_id, &binding, false)?;
373 println!("one-call permission refusal queued");
374 }
375 }
376 Ok(0)
377 }
378 Command::Grant { command } => {
379 match command {
380 GrantCommand::Approve { id, command } => {
381 cmd_approve_grant(&repo, &id, &command)?;
382 println!("grant approval for `{command}` queued for mission {id}");
383 if let Some(hint) = control_queue_hint(&repo, &id) {
384 println!("{hint}");
385 }
386 }
387 GrantCommand::Deny {
388 id,
389 command,
390 reason,
391 } => {
392 cmd_deny_grant(&repo, &id, &command, &reason)?;
393 println!("grant denial for `{command}` queued for mission {id}");
394 if let Some(hint) = control_queue_hint(&repo, &id) {
395 println!("{hint}");
396 }
397 }
398 }
399 Ok(0)
400 }
401 Command::Question { command } => {
402 match command {
403 QuestionCommand::List { id } => {
404 print!("{}", cmd_list_questions(&repo, &id)?);
405 }
406 QuestionCommand::Answer {
407 id,
408 question_id,
409 answer,
410 option,
411 } => {
412 cmd_answer_question(&repo, &id, &question_id, &answer, option)?;
413 println!("answer for question {question_id} queued for mission {id}");
414 if let Some(hint) = control_queue_hint(&repo, &id) {
415 println!("{hint}");
416 }
417 }
418 }
419 Ok(0)
420 }
421 Command::Missions => {
422 print!("{}", cmd_missions(&repo)?);
423 Ok(0)
424 }
425 Command::Abandon { id, reason } => {
426 let mission = select_mission(&repo, id.as_deref().or(cli.mission.as_deref()))?;
429 let reason = reason.as_deref().unwrap_or("abandoned by operator");
430 cmd_abandon(&repo, &mission, reason, lock_force)?;
431 println!("mission {mission} ABANDONED ({reason})");
432 Ok(0)
433 }
434 Command::Clean { yes, all } => cmd_clean(&repo, yes, all),
435 Command::Ticket { command } => dispatch_ticket(&repo, command, cli.mission.as_deref()),
436 Command::Draft {
437 slug,
438 yes,
439 from_mission,
440 } => backlog::cmd_draft(
441 repo,
442 &slug,
443 yes,
444 from_mission.as_deref(),
445 cli.dangerously_allow_all,
446 )
447 .await
448 .map_err(augment_limit_hint),
449 Command::Decompose { goal, yes } => {
450 backlog::cmd_decompose(repo, &goal, yes, cli.dangerously_allow_all)
451 .await
452 .map_err(augment_limit_hint)
453 }
454 Command::Exec {
455 file,
456 yes: _,
457 max_cycles,
458 enqueue,
459 enqueue_source,
460 enqueue_external_ref,
461 push,
462 allow_unvalidated,
463 } => crate::exec::cmd_exec(
464 repo,
465 file,
466 crate::exec::ExecOptions {
467 max_cycles,
468 enqueue,
469 enqueue_source: enqueue_source.zip(enqueue_external_ref).map(
470 |(producer, external_ref)| crate::exec::ExternalEnqueueSource {
471 producer,
472 external_ref,
473 },
474 ),
475 push,
476 dangerously_allow_all: cli.dangerously_allow_all,
477 allow_unvalidated,
478 },
479 )
480 .await
481 .map_err(augment_limit_hint),
482 Command::Queue { remove } => match remove {
483 Some(mission_id) => {
484 print!("{}", backlog::cmd_queue_remove(&repo, &mission_id)?);
485 Ok(0)
486 }
487 None => {
488 print!("{}", backlog::cmd_queue(&repo));
489 Ok(0)
490 }
491 },
492 Command::KnowledgeRefresh { json } => cmd_knowledge_refresh(&repo, json),
493 Command::Scan { staged, range } => cmd_scan(&repo, staged, range.as_deref()),
494 Command::DomainLint { seed_config, json } => {
495 cmd_domain_lint(&repo, seed_config.as_deref(), json)
496 }
497 Command::HookGuard { config } => {
498 let mut stdin = std::io::stdin();
499 Ok(crate::hook_guard::run_hook_guard(&config, &mut stdin))
500 }
501 Command::HookStatus { config } => {
502 let mut stdin = std::io::stdin();
503 Ok(crate::hook_status::run_hook_status(
504 &config,
505 &mut stdin,
506 &crate::hook_status::post_signal,
507 )
508 .await)
509 }
510 Command::Ready { json, all } => {
511 if all {
512 let config = kranz_engine::paths::global_config()
513 .ok_or_else(|| anyhow::anyhow!("cannot locate the home directory"))?;
514 let report = crate::ready::assess_all(&config);
515 if json {
516 println!("{}", serde_json::to_string_pretty(&report)?);
517 } else {
518 print!("{}", crate::ready::render_org(&report));
519 }
520 } else {
521 let report = crate::ready::assess(&repo);
522 if json {
523 println!("{}", serde_json::to_string_pretty(&report)?);
524 } else {
525 print!("{}", crate::ready::render(&report));
526 }
527 }
528 Ok(0)
529 }
530 Command::Work { once, expect } => backlog::cmd_work(repo, once, expect)
531 .await
532 .map_err(augment_limit_hint),
533 Command::Serve {
534 port,
535 host,
536 insecure_lan,
537 read_auth,
538 open,
539 dashboard,
540 token,
541 read_token,
542 slack,
543 } => {
544 cmd_serve(
545 repo,
546 host,
547 port,
548 insecure_lan,
549 read_auth,
550 open,
551 dashboard,
552 token,
553 read_token,
554 slack,
555 )
556 .await
557 }
558 Command::Release { url, token } => {
559 let mission = select_mission(&repo, cli.mission.as_deref())?;
560 cmd_release(&repo, &mission, &url, token).await
561 }
562 Command::Config { command } => {
563 crate::config_cmd::cmd_config(&repo, command, cli.mission.as_deref())
564 }
565 Command::Pack { command } => match command {
566 crate::cli::PackCommand::Lint { dir } => cmd_pack_lint(&repo, &dir),
567 },
568 Command::Standards { command } => match command {
569 crate::cli::StandardsCommand::Metrics { json } => {
570 let report = kranz_engine::standards_metrics::compute(&repo)?;
571 if json {
572 println!("{}", serde_json::to_string_pretty(&report)?);
573 } else {
574 print!("{}", crate::output::render_standards_metrics(&report));
575 }
576 Ok(0)
577 }
578 crate::cli::StandardsCommand::Lint { dir, against } => {
579 cmd_standards_lint(&repo, &dir, against.as_deref())
580 }
581 crate::cli::StandardsCommand::Waive {
582 rule,
583 revision,
584 finding,
585 reason,
586 expires,
587 } => {
588 let mission = select_mission(&repo, cli.mission.as_deref())?;
589 cmd_standards_waive(
590 &repo,
591 &mission,
592 &rule,
593 revision,
594 finding.as_deref(),
595 &reason,
596 &expires,
597 lock_force,
598 )
599 }
600 crate::cli::StandardsCommand::Attest { rule, reason } => {
601 let mission = select_mission(&repo, cli.mission.as_deref())?;
602 cmd_standards_attest(&repo, &mission, &rule, &reason, lock_force)
603 }
604 },
605 Command::Otel {
606 endpoint,
607 from_start,
608 } => crate::otel::run_otel(repo, cli.mission.clone(), endpoint, from_start).await,
609 }
610}
611
612fn dispatch_ticket(repo: &Path, command: TicketCommand, mission: Option<&str>) -> Result<i32> {
616 match command {
617 TicketCommand::List => {
618 print!("{}", backlog::cmd_ticket_list(repo));
619 Ok(0)
620 }
621 TicketCommand::Ready { include_deferred } => {
622 print!("{}", backlog::cmd_ticket_ready(repo, include_deferred));
623 Ok(0)
624 }
625 TicketCommand::Show { slug } => {
626 print!("{}", backlog::cmd_ticket_show(repo, &slug)?);
627 Ok(0)
628 }
629 TicketCommand::New { slug, title, goal } => {
630 let path = backlog::cmd_ticket_new(repo, &slug, &title, goal.as_deref())?;
631 println!("created ticket '{slug}' at {}", path.display());
632 Ok(0)
633 }
634 TicketCommand::ImportOpenspec { path, slug } => {
635 let written = crate::openspec::import_change(repo, &path, slug.as_deref())?;
636 println!("imported {} as {}", path.display(), written.display());
637 println!(
638 "acceptance criteria are still prose: replace the placeholder in \
639 '## Acceptance hints' with commands that can fail"
640 );
641 Ok(0)
642 }
643 TicketCommand::Note { slug, text } => {
644 print!(
645 "{}",
646 crate::ticket_notes::cmd_ticket_note(repo, &slug, &text.join(" "))?
647 );
648 Ok(0)
649 }
650 TicketCommand::Notes { slug } => {
651 print!("{}", crate::ticket_notes::cmd_ticket_notes(repo, &slug)?);
652 Ok(0)
653 }
654 TicketCommand::Queue {
655 slug,
656 mission: explicit,
657 force,
658 } => {
659 backlog::cmd_ticket_queue(repo, &slug, explicit.as_deref().or(mission), force)
661 }
662 TicketCommand::Approve {
663 slug,
664 mission: explicit,
665 force,
666 } => {
667 backlog::cmd_ticket_approve(repo, &slug, explicit.as_deref().or(mission), force)
670 }
671 TicketCommand::MigrateState { yes } => backlog::cmd_ticket_migrate_state(repo, yes),
672 }
673}
674
675pub fn load_config(repo: &Path, dangerously_allow_all: bool) -> Result<MissionConfig> {
682 let mut cfg = config::load(repo)?;
683 if dangerously_allow_all {
684 cfg.dangerously_allow_all = true;
685 }
686 config::validate(&cfg)?;
687 Ok(cfg)
688}
689
690pub(crate) fn build_backend(cfg: &MissionConfig) -> Result<Arc<dyn AgentBackend>> {
693 let backend = ClaudeBackend::discover(cfg.claude_binary.as_deref())?;
694 Ok(Arc::new(backend))
695}
696
697pub fn select_mission(repo: &Path, explicit: Option<&str>) -> Result<String> {
700 if let Some(id) = explicit {
701 require_mission(repo, id)?;
702 return Ok(id.to_string());
703 }
704 let ids = MissionPaths::list_missions(repo);
705 match ids.len() {
706 0 => bail!(
707 "no missions found under {}; create one with `kranz plan \"<goal>\"`",
708 repo.join(".kranz").join("missions").display()
709 ),
710 1 => Ok(ids.into_iter().next().expect("len checked")),
711 _ => {
712 let mut best: Option<(SystemTime, String)> = None;
713 for id in ids {
714 let events = MissionPaths::new(repo, &id).events_file();
715 let mtime = std::fs::metadata(&events)
716 .and_then(|m| m.modified())
717 .unwrap_or(SystemTime::UNIX_EPOCH);
718 let newer = match &best {
719 Some((t, _)) => mtime >= *t,
720 None => true,
721 };
722 if newer {
723 best = Some((mtime, id));
724 }
725 }
726 Ok(best.expect("non-empty list").1)
727 }
728 }
729}
730
731pub fn select_control_mission(repo: &Path, explicit: Option<&str>) -> Result<String> {
743 if explicit.is_some() {
744 return Ok(control::resolve_active_mission(repo, explicit)?);
745 }
746 let mission = select_mission(repo, None)?;
747 let status = load_state(repo, &mission)?.mission.status;
748 if mission_catalog::is_terminal_status(status) {
749 bail!(
750 "mission {mission} is {status:?}; control commands apply only to active \
751 missions (a terminal mission's inbox is never drained — see \
752 `kranz missions`)"
753 );
754 }
755 Ok(mission)
756}
757
758pub fn control_queue_hint(repo: &Path, mission_id: &str) -> Option<String> {
763 let paths = MissionPaths::new(repo, mission_id);
764 (!mission_catalog::mission_lock_is_live(&paths)).then(|| {
765 format!(
766 "note: mission {mission_id} is not currently running — the command is \
767 queued and applies when the mission next runs"
768 )
769 })
770}
771
772pub fn select_planning_mission(repo: &Path, explicit: Option<&str>) -> Result<String> {
776 if let Some(id) = explicit {
777 require_mission(repo, id)?;
778 return Ok(id.to_string());
779 }
780 let mut best: Option<(SystemTime, String)> = None;
781 for id in MissionPaths::list_missions(repo) {
782 let Ok(state) = load_state(repo, &id) else {
783 continue; };
785 if state.mission.status != MissionStatus::Planning {
786 continue;
787 }
788 let events = MissionPaths::new(repo, &id).events_file();
789 let mtime = std::fs::metadata(&events)
790 .and_then(|m| m.modified())
791 .unwrap_or(SystemTime::UNIX_EPOCH);
792 if best.as_ref().is_none_or(|(t, _)| mtime >= *t) {
793 best = Some((mtime, id));
794 }
795 }
796 best.map(|(_, id)| id).ok_or_else(|| {
797 anyhow!(
798 "no mission is currently in planning under {} — start one with \
799 `kranz plan \"<goal>\"`",
800 repo.join(".kranz").join("missions").display()
801 )
802 })
803}
804
805pub fn augment_limit_hint(e: anyhow::Error) -> anyhow::Error {
809 let msg = format!("{e:#}").to_ascii_lowercase();
810 if [
811 "session limit",
812 "usage limit",
813 "rate limit",
814 "hit your limit",
815 ]
816 .iter()
817 .any(|s| msg.contains(s))
818 {
819 e.context(
820 "this is your Claude subscription's usage window, not a Kranz failure. \
821 The mission and its conversation are saved: when the limit resets, \
822 `kranz plan` (no goal) resumes planning and `kranz run` resumes execution",
823 )
824 } else {
825 e
826 }
827}
828
829fn require_mission(repo: &Path, mission_id: &str) -> Result<MissionPaths> {
831 if !MissionPaths::is_safe_id(mission_id) {
832 bail!(
833 "invalid mission id '{mission_id}': ids cannot contain path separators, '..', or drive designators"
834 );
835 }
836 let paths = MissionPaths::new(repo, mission_id);
837 if !paths.events_file().is_file() {
838 bail!(
839 "mission '{mission_id}' not found under {} (see `kranz missions`)",
840 paths.missions_dir().display()
841 );
842 }
843 Ok(paths)
844}
845
846pub fn load_state(repo: &Path, mission_id: &str) -> Result<MissionState> {
849 let paths = require_mission(repo, mission_id)?;
850 let events = EventLog::read_events(&paths.events_file())
851 .with_context(|| format!("reading the event log of mission '{mission_id}'"))?;
852 let state = reducer::fold(&events)
853 .with_context(|| format!("folding the event log of mission '{mission_id}'"))?;
854 Ok(state)
855}
856
857pub fn cmd_export_traces(repo: &Path, mission_id: &str) -> Result<String> {
862 let paths = require_mission(repo, mission_id)?;
863 let events = EventLog::read_events(&paths.events_file())
864 .with_context(|| format!("reading the event log of mission '{mission_id}'"))?;
865 let state = reducer::fold(&events)
866 .with_context(|| format!("folding the event log of mission '{mission_id}'"))?;
867 let pairs = trace_export::export_validated_traces(&state, &events);
868 Ok(trace_export::to_jsonl(&pairs))
869}
870
871pub fn cmd_export_traces_all(repo: &Path) -> String {
876 let mut pairs = Vec::new();
877 for mission_id in MissionPaths::list_missions(repo) {
878 let paths = MissionPaths::new(repo, &mission_id);
879 let Ok(events) = EventLog::read_events(&paths.events_file()) else {
880 continue;
881 };
882 let Ok(state) = reducer::fold(&events) else {
883 continue;
884 };
885 pairs.extend(trace_export::export_validated_traces(&state, &events));
886 }
887 trace_export::to_jsonl(&pairs)
888}
889
890pub fn cmd_export_corpus(repo: &Path, mission_id: &str) -> Result<String> {
897 let paths = require_mission(repo, mission_id)?;
898 paths.require_no_follow()?;
899 let events = EventLog::read_events(&paths.events_file())
900 .with_context(|| format!("reading the event log of mission '{mission_id}'"))?;
901 let records = corpus_export::export_corpus(&paths.mission_dir(), mission_id, &events)
902 .with_context(|| format!("deriving the training corpus of mission '{mission_id}'"))?;
903 Ok(corpus_export::to_jsonl(&records))
904}
905
906pub fn cmd_export_corpus_all(repo: &Path) -> String {
912 let mut records = Vec::new();
913 for mission_id in MissionPaths::list_missions(repo) {
914 let paths = MissionPaths::new(repo, &mission_id);
915 if paths.require_no_follow().is_err() {
916 continue;
917 }
918 let Ok(events) = EventLog::read_events(&paths.events_file()) else {
919 continue;
920 };
921 let Ok(mission_records) =
922 corpus_export::export_corpus(&paths.mission_dir(), &mission_id, &events)
923 else {
924 continue;
925 };
926 records.extend(mission_records);
927 }
928 corpus_export::to_jsonl(&records)
929}
930
931pub fn print_danger_banner() {
933 eprintln!(
934 "\n\
935 ============================================================\n\
936 !! --dangerously-allow-all IS SET !!\n\
937 !! !!\n\
938 !! Permission gating is BYPASSED for every agent !!\n\
939 !! session (bypassPermissions). Workers can run ANY !!\n\
940 !! command: file writes, network access, git push, !!\n\
941 !! package publishes, sudo. !!\n\
942 !! !!\n\
943 !! Only use this on a sandboxed, disposable checkout. !!\n\
944 ============================================================\n"
945 );
946}
947
948struct StdinLines {
960 rx: Arc<tokio::sync::Mutex<tokio::sync::mpsc::UnboundedReceiver<String>>>,
961}
962
963impl StdinLines {
964 fn spawn() -> Self {
965 static CHANNEL: std::sync::OnceLock<
966 Arc<tokio::sync::Mutex<tokio::sync::mpsc::UnboundedReceiver<String>>>,
967 > = std::sync::OnceLock::new();
968 let rx = CHANNEL
969 .get_or_init(|| {
970 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
971 std::thread::spawn(move || {
972 let mut buf = String::new();
973 loop {
974 buf.clear();
975 match std::io::stdin().read_line(&mut buf) {
976 Ok(0) | Err(_) => break, Ok(_) => {
978 let line = buf.trim_end_matches(['\r', '\n']).to_string();
979 if tx.send(line).is_err() {
980 break;
981 }
982 }
983 }
984 }
985 });
986 Arc::new(tokio::sync::Mutex::new(rx))
987 })
988 .clone();
989 StdinLines { rx }
990 }
991
992 async fn next(&mut self) -> Option<String> {
994 self.rx.lock().await.recv().await
995 }
996
997 fn drain(&mut self) -> usize {
999 let Ok(mut rx) = self.rx.try_lock() else {
1002 return 0;
1003 };
1004 let mut n = 0;
1005 while rx.try_recv().is_ok() {
1006 n += 1;
1007 }
1008 n
1009 }
1010
1011 fn drain_noisily(&mut self, tty: bool) {
1015 if !std::io::stdin().is_terminal() {
1016 return;
1017 }
1018 let n = self.drain();
1019 if n > 0 {
1020 let (dim, reset) = if tty {
1021 (ansi::DIM, ansi::RESET)
1022 } else {
1023 ("", "")
1024 };
1025 eprintln!(
1026 "{dim}(ignored {n} line(s) typed while the orchestrator was working — \
1027 the prompt below wants fresh input){reset}"
1028 );
1029 }
1030 }
1031}
1032
1033async fn cmd_plan(
1041 repo: PathBuf,
1042 goal: Option<String>,
1043 cfg: MissionConfig,
1044 explicit_mission: Option<&str>,
1045 force_lock: LockForce,
1046) -> Result<i32> {
1047 let backend = build_backend(&cfg)?;
1048 let (mut engine, intro) = match goal {
1049 Some(goal) => {
1050 let engine = MissionEngine::create(backend, repo.clone(), &goal, cfg)?;
1051 let intro = format!("mission {} created (planning)", engine.mission_id());
1052 (engine, intro)
1053 }
1054 None => {
1055 let mission = select_planning_mission(&repo, explicit_mission)?;
1056 let engine = MissionEngine::resume(backend, repo.clone(), &mission, force_lock)?;
1057 if engine.state().mission.status != MissionStatus::Planning {
1058 return Err(anyhow!(
1059 "mission {mission} is {:?}, not in planning — use 'kranz run' \
1060 to execute it, or 'kranz plan \"<goal>\"' to start a new mission",
1061 engine.state().mission.status
1062 ));
1063 }
1064 let intro = format!(
1065 "resuming planning for mission {mission} — the conversation continues \
1066 where it left off"
1067 );
1068 (engine, intro)
1069 }
1070 };
1071
1072 if std::io::stdin().is_terminal() && std::io::stdout().is_terminal() {
1078 let mission_id = engine.mission_id().to_string();
1079 return match crate::planning_tui::run(engine, intro).await? {
1083 PlanningOutcome::ApprovedRun => start_run_after_plan(repo, mission_id).await,
1084 PlanningOutcome::ApprovedExit | PlanningOutcome::NotApproved => Ok(0),
1085 };
1086 }
1087 println!("{intro}");
1088 let tty = std::io::stdout().is_terminal();
1089
1090 let color = std::io::stderr().is_terminal();
1093 let stop = Arc::new(AtomicBool::new(false));
1094 let printer = tokio::spawn(tail::tail_events(
1095 engine.paths().events_file(),
1096 engine.state().last_seq,
1097 EventRenderer::planning(engine.state(), color),
1098 Arc::clone(&stop),
1099 ));
1100 let mut approved = false;
1101 let mut run_now = false;
1102 let mut stdin_lines = StdinLines::spawn();
1103
1104 println!("talk to the orchestrator to shape the plan:");
1105 println!(" /plan request the plan + cost estimate and review it for approval");
1106 println!(" /quit exit planning (Ctrl-D works too)");
1107
1108 loop {
1109 stdin_lines.drain_noisily(tty);
1110 if tty {
1111 print!("you> ");
1112 let _ = std::io::stdout().flush();
1113 }
1114 let Some(line) = stdin_lines.next().await else {
1115 break; };
1117 let line = line.trim().to_string();
1118 if line.is_empty() {
1119 continue;
1120 }
1121 match line.as_str() {
1122 "/quit" => break,
1123 "/plan" => {
1124 let request = engine.request_plan().await;
1125 if let Some(seed) = engine.take_seed_reply() {
1128 print_orchestrator_reply(&seed, tty);
1129 }
1130 let plan = match request {
1131 Ok(PlanRequest::Ready(plan)) => plan,
1132 Ok(PlanRequest::NotReady(text)) => {
1133 print_orchestrator_reply(&text, tty);
1136 println!(
1137 "the orchestrator isn't ready to emit the plan yet — answer it \
1138 above, then /plan again."
1139 );
1140 continue;
1141 }
1142 Ok(PlanRequest::WrongPlan { reason }) => {
1143 print_orchestrator_reply(&reason, tty);
1146 println!(
1147 "the orchestrator believes a plan here is likely WRONG — reframe \
1148 the goal or fix the premise above, then /plan again."
1149 );
1150 continue;
1151 }
1152 Err(e) => {
1153 eprintln!(
1154 "kranz: plan request failed: {:#}",
1155 augment_limit_hint(e.into())
1156 );
1157 continue;
1158 }
1159 };
1160 println!("{}", output::render_plan(&plan));
1161 let calibration = cost::calibrate(&repo);
1164 let estimate = cost::estimate(&plan, &engine.state().config, &calibration.params);
1165 let estimate = cost::apply_shape(estimate, &plan, &calibration);
1166 println!(
1167 "{}",
1168 output::render_cost_estimate(&estimate, calibration.missions_used)
1169 );
1170
1171 stdin_lines.drain_noisily(tty);
1172 print!("approve? [y/N] ");
1173 let _ = std::io::stdout().flush();
1174 let answer = stdin_lines.next().await.unwrap_or_default();
1175 if matches!(answer.trim().to_ascii_lowercase().as_str(), "y" | "yes") {
1176 match engine.approve_plan(plan) {
1177 Ok(()) => {
1178 let branch = engine.state().mission.mission_branch.clone();
1179 approved = true;
1180 if std::io::stdin().is_terminal() {
1185 println!("plan approved and committed on {branch}.");
1186 stdin_lines.drain_noisily(tty);
1187 print!("start execution now? [Y/n] ");
1188 let _ = std::io::stdout().flush();
1189 let reply = stdin_lines.next().await;
1190 if run_now_answer(reply.as_deref()) {
1191 run_now = true;
1192 } else {
1193 println!("run 'kranz run' to execute.");
1194 }
1195 } else {
1196 println!(
1197 "plan approved and committed on {branch}. \
1198 run 'kranz run' to execute."
1199 );
1200 }
1201 break;
1202 }
1203 Err(e) => eprintln!("kranz: plan approval failed: {e}"),
1204 }
1205 } else {
1206 println!("not approved — back to the conversation.");
1207 }
1208 }
1209 _ if line.starts_with('/') => {
1210 println!("unknown command {line}; use /plan or /quit");
1211 }
1212 _ => {
1213 let result = engine.planning_turn(&line).await;
1214 if let Some(seed) = engine.take_seed_reply() {
1217 print_orchestrator_reply(&seed, tty);
1218 }
1219 match result {
1220 Ok(reply) => print_orchestrator_reply(&reply, tty),
1221 Err(e) => eprintln!(
1222 "kranz: orchestrator turn failed: {:#}",
1223 augment_limit_hint(e.into())
1224 ),
1225 }
1226 }
1227 }
1228 }
1229 if !approved {
1230 println!(
1231 "leaving planning; mission {} was not approved. Resume anytime with `kranz plan`.",
1232 engine.mission_id()
1233 );
1234 }
1235 let mission_id = engine.mission_id().to_string();
1236 drop(engine);
1239 stop.store(true, Ordering::Relaxed);
1240 let _ = printer.await;
1241 if run_now {
1242 return start_run_after_plan(repo, mission_id).await;
1245 }
1246 Ok(0)
1247}
1248
1249pub fn run_now_answer(answer: Option<&str>) -> bool {
1254 match answer {
1255 None => false,
1256 Some(text) => !matches!(text.trim().to_ascii_lowercase().as_str(), "n" | "no"),
1257 }
1258}
1259
1260async fn start_run_after_plan(repo: PathBuf, mission_id: String) -> Result<i32> {
1264 println!(
1265 "starting mission {mission_id} — live event feed follows \
1266 (Ctrl-C safe; resume with 'kranz run')"
1267 );
1268 run_mission_loop(repo, mission_id, LockForce::No, true).await
1269}
1270
1271fn print_orchestrator_reply(text: &str, tty: bool) {
1273 let prefix = if tty {
1274 format!("{}orchestrator>{} ", ansi::DIM, ansi::RESET)
1275 } else {
1276 "orchestrator> ".to_string()
1277 };
1278 for line in text.lines() {
1279 println!("{prefix}{line}");
1280 }
1281}
1282
1283async fn cmd_run(
1290 repo: PathBuf,
1291 mission: String,
1292 force_lock: LockForce,
1293 dangerously_allow_all: bool,
1294) -> Result<i32> {
1295 if dangerously_allow_all {
1298 let paths = require_mission(&repo, &mission)?;
1299 control::enqueue(
1300 &paths,
1301 &ControlCommand::ConfigChange {
1302 patch: serde_json::json!({ "dangerouslyAllowAll": true }),
1303 },
1304 )?;
1305 }
1306 run_mission_loop(repo, mission, force_lock, true).await
1307}
1308
1309pub(crate) async fn run_mission_loop(
1313 repo: PathBuf,
1314 mission: String,
1315 force_lock: LockForce,
1316 interactive: bool,
1320) -> Result<i32> {
1321 let cfg = load_config(&repo, false)?;
1322 let backend = build_backend(&cfg)?;
1323 run_mission_loop_with_backend(repo, mission, force_lock, interactive, backend).await
1324}
1325
1326async fn run_mission_loop_with_backend(
1330 repo: PathBuf,
1331 mission: String,
1332 force_lock: LockForce,
1333 interactive: bool,
1334 backend: Arc<dyn AgentBackend>,
1335) -> Result<i32> {
1336 let paths = require_mission(&repo, &mission)?;
1337
1338 loop {
1339 let mut engine =
1340 MissionEngine::resume(Arc::clone(&backend), repo.clone(), &mission, force_lock)?;
1341
1342 let color = std::io::stderr().is_terminal();
1344 let renderer = EventRenderer::seeded(engine.state(), color);
1345 let stop = Arc::new(AtomicBool::new(false));
1346 let printer = tokio::spawn(tail::tail_events(
1347 engine.paths().events_file(),
1348 engine.state().last_seq,
1349 renderer,
1350 Arc::clone(&stop),
1351 ));
1352
1353 let run_result = engine.run().await;
1354 drop(engine);
1357 stop.store(true, Ordering::Relaxed);
1358 let _ = printer.await;
1359
1360 let status = run_result?;
1361 if let Err(e) = kranz_engine::work::reconcile_ticket_for_mission(&repo, &mission) {
1365 eprintln!("kranz run: warning: failed to reconcile linked ticket: {e}");
1366 }
1367
1368 match status {
1369 MissionStatus::Complete => {
1370 println!("mission {mission} COMPLETE");
1371 return Ok(0);
1372 }
1373 MissionStatus::Blocked => {
1374 eprintln!(
1375 "\n\
1376 ==================== MILESTONE BLOCKED ====================\n\
1377 A milestone is blocked (fix-cycle cap reached or blocked by\n\
1378 the orchestrator). Inspect it with `kranz status`.\n\
1379 ==========================================================="
1380 );
1381 if interactive && std::io::stdin().is_terminal() && std::io::stdout().is_terminal()
1385 {
1386 print!(
1387 "guidance for the orchestrator (what to do about the block; \
1388 empty line or Ctrl-D exits)\nguidance> "
1389 );
1390 let _ = std::io::stdout().flush();
1391 let mut lines = StdinLines::spawn();
1392 if let Some(text) = lines.next().await {
1393 let text = text.trim().to_string();
1394 if !text.is_empty() {
1395 control::enqueue(
1396 &paths,
1397 &ControlCommand::Msg {
1398 text,
1399 interrupt: false,
1400 },
1401 )?;
1402 println!("guidance queued — resuming the mission…");
1403 continue;
1404 }
1405 }
1406 }
1407 eprintln!(
1408 "unblock later by sending guidance via `kranz msg \"<text>\"` \
1409 and re-running `kranz run`."
1410 );
1411 println!("mission {mission} BLOCKED");
1412 return Ok(2);
1413 }
1414 MissionStatus::Failed => {
1415 println!("mission {mission} FAILED");
1416 return Ok(1);
1417 }
1418 other => {
1419 println!(
1420 "mission {mission} ended as {}",
1421 output::mission_status_label(other)
1422 );
1423 return Ok(1);
1424 }
1425 }
1426 }
1427}
1428
1429pub fn cmd_pause(repo: &Path, mission_id: &str) -> Result<PathBuf> {
1435 let paths = require_mission(repo, mission_id)?;
1436 Ok(control::enqueue(&paths, &ControlCommand::Pause)?)
1437}
1438
1439pub fn cmd_resume(repo: &Path, mission_id: &str) -> Result<PathBuf> {
1441 let paths = require_mission(repo, mission_id)?;
1442 Ok(control::enqueue(&paths, &ControlCommand::Resume)?)
1443}
1444
1445pub fn cmd_msg(repo: &Path, mission_id: &str, text: &str, interrupt: bool) -> Result<PathBuf> {
1447 let paths = require_mission(repo, mission_id)?;
1448 let cmd = ControlCommand::Msg {
1449 text: text.to_string(),
1450 interrupt,
1451 };
1452 Ok(control::enqueue(&paths, &cmd)?)
1453}
1454
1455pub fn cmd_request_revision(repo: &Path, mission_id: &str, instructions: &str) -> Result<PathBuf> {
1457 let instructions = instructions.trim();
1458 if instructions.is_empty() {
1459 bail!("revision instructions must not be empty");
1460 }
1461 let paths = require_revisable_mission(repo, mission_id)?;
1462 Ok(control::enqueue(
1463 &paths,
1464 &ControlCommand::RequestRevision {
1465 instructions: instructions.to_string(),
1466 },
1467 )?)
1468}
1469
1470pub fn cmd_approve_revision(repo: &Path, mission_id: &str, revision: u32) -> Result<PathBuf> {
1472 let paths = require_pending_revision(repo, mission_id, revision)?;
1473 Ok(control::enqueue(
1474 &paths,
1475 &ControlCommand::ApproveRevision { revision },
1476 )?)
1477}
1478
1479pub fn cmd_reject_revision(repo: &Path, mission_id: &str, revision: u32) -> Result<PathBuf> {
1481 let paths = require_pending_revision(repo, mission_id, revision)?;
1482 Ok(control::enqueue(
1483 &paths,
1484 &ControlCommand::RejectRevision { revision },
1485 )?)
1486}
1487
1488pub fn cmd_permission_answer(
1489 repo: &Path,
1490 mission_id: &str,
1491 request_id: &str,
1492 binding: &str,
1493 allow: bool,
1494) -> Result<PathBuf> {
1495 let paths = require_revisable_mission(repo, mission_id)?;
1496 let state = load_state(repo, mission_id)?;
1497 let record = state
1498 .permissions
1499 .get(request_id)
1500 .ok_or_else(|| anyhow!("unknown live permission"))?;
1501 record.validate_answer(binding, allow, chrono::Utc::now())?;
1502 Ok(control::enqueue(
1503 &paths,
1504 &ControlCommand::ResolvePermission {
1505 resolution: kranz_engine::live_permission::Resolution {
1506 request_id: request_id.into(),
1507 binding_digest: binding.into(),
1508 allow,
1509 actor: kranz_engine::live_permission::Actor::LocalRepositoryAuthority,
1510 reason: "operator answered through the local CLI".into(),
1511 },
1512 },
1513 )?)
1514}
1515
1516pub fn cmd_approve_grant(repo: &Path, mission_id: &str, command: &str) -> Result<PathBuf> {
1518 let paths = require_pending_grant(repo, mission_id, command)?;
1519 Ok(control::enqueue(
1520 &paths,
1521 &ControlCommand::ApproveGrant {
1522 command: command.to_string(),
1523 },
1524 )?)
1525}
1526
1527pub fn cmd_deny_grant(
1529 repo: &Path,
1530 mission_id: &str,
1531 command: &str,
1532 reason: &str,
1533) -> Result<PathBuf> {
1534 let paths = require_pending_grant(repo, mission_id, command)?;
1535 Ok(control::enqueue(
1536 &paths,
1537 &ControlCommand::DenyGrant {
1538 command: command.to_string(),
1539 reason: reason.to_string(),
1540 },
1541 )?)
1542}
1543
1544pub fn cmd_list_questions(repo: &Path, mission_id: &str) -> Result<String> {
1549 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1550 let state = load_state(repo, &mission_id)?;
1551 if state.pending_questions.is_empty() {
1552 return Ok(format!("mission {mission_id} has no open questions\n"));
1553 }
1554 let mut out = String::new();
1555 for q in &state.pending_questions {
1556 out.push_str(&format!(
1557 "{} ({}): {}\n",
1558 q.question_id,
1559 q.feature_id.as_deref().unwrap_or("mission"),
1560 q.text
1561 ));
1562 if q.options.is_empty() {
1563 out.push_str(" free-text answer expected\n");
1564 } else {
1565 for (index, option) in q.options.iter().enumerate() {
1566 out.push_str(&format!(" [{index}] {option}\n"));
1567 }
1568 }
1569 }
1570 Ok(out)
1571}
1572
1573pub fn cmd_answer_question(
1576 repo: &Path,
1577 mission_id: &str,
1578 question_id: &str,
1579 answer: &str,
1580 option: Option<u32>,
1581) -> Result<PathBuf> {
1582 let paths = require_pending_question(repo, mission_id, question_id, option, answer)?;
1583 Ok(control::enqueue(
1584 &paths,
1585 &ControlCommand::AnswerQuestion {
1586 question_id: question_id.to_string(),
1587 answer: answer.to_string(),
1588 option,
1589 },
1590 )?)
1591}
1592
1593fn require_pending_question(
1599 repo: &Path,
1600 mission_id: &str,
1601 question_id: &str,
1602 option: Option<u32>,
1603 answer: &str,
1604) -> Result<MissionPaths> {
1605 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1606 let state = load_state(repo, &mission_id)?;
1607 let Some(pending) = state
1608 .pending_questions
1609 .iter()
1610 .find(|q| q.question_id == question_id)
1611 else {
1612 bail!("mission {mission_id} has no open question '{question_id}'");
1613 };
1614 if let Some(index) = option {
1615 match pending.options.get(index as usize) {
1616 Some(expected) if expected == answer => {}
1617 Some(expected) => bail!(
1618 "answer `{answer}` does not match option {index} (`{expected}`) of question '{question_id}'"
1619 ),
1620 None => bail!(
1621 "question '{question_id}' has no option {index} (it offered {})",
1622 pending.options.len()
1623 ),
1624 }
1625 }
1626 Ok(MissionPaths::new(repo, &mission_id))
1627}
1628
1629fn require_revisable_mission(repo: &Path, mission_id: &str) -> Result<MissionPaths> {
1630 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1631 let state = load_state(repo, &mission_id)?;
1632 if state.mission.status == MissionStatus::Planning {
1633 bail!("mission {mission_id} has no approved plan to revise yet");
1634 }
1635 Ok(MissionPaths::new(repo, &mission_id))
1636}
1637
1638fn require_pending_revision(repo: &Path, mission_id: &str, revision: u32) -> Result<MissionPaths> {
1639 let paths = require_revisable_mission(repo, mission_id)?;
1640 let state = load_state(repo, mission_id)?;
1641 match state.pending_revision {
1642 Some(pending) if pending.revision == revision => Ok(paths),
1643 Some(pending) => bail!(
1644 "mission {mission_id} is awaiting revision {}, not {revision}",
1645 pending.revision
1646 ),
1647 None => bail!("mission {mission_id} has no pending revision"),
1648 }
1649}
1650
1651fn require_pending_grant(repo: &Path, mission_id: &str, command: &str) -> Result<MissionPaths> {
1655 let mission_id = control::resolve_active_mission(repo, Some(mission_id))?;
1656 let state = load_state(repo, &mission_id)?;
1657 match state.pending_grant_request {
1658 Some(pending) if pending.command == command => Ok(MissionPaths::new(repo, &mission_id)),
1659 Some(pending) => bail!(
1660 "mission {mission_id} is awaiting a grant for `{}`, not `{command}`",
1661 pending.command
1662 ),
1663 None => bail!("mission {mission_id} has no pending grant request"),
1664 }
1665}
1666
1667pub fn cmd_knowledge_refresh(repo: &Path, json: bool) -> Result<i32> {
1671 let report = kranz_engine::knowledge::refresh_knowledge(repo);
1672 if json {
1673 println!("{}", serde_json::to_string_pretty(&report)?);
1674 } else if report.findings.is_empty() {
1675 println!("knowledge refresh: no notes");
1676 } else {
1677 for finding in &report.findings {
1678 let labels: Vec<String> = finding
1679 .verdicts
1680 .iter()
1681 .map(|v| match v {
1682 kranz_engine::knowledge::RefreshVerdict::Ok => "ok".into(),
1683 kranz_engine::knowledge::RefreshVerdict::AlreadyStale => "already-stale".into(),
1684 kranz_engine::knowledge::RefreshVerdict::Unverified => "unverified".into(),
1685 kranz_engine::knowledge::RefreshVerdict::InvalidMetadata { field, value } => {
1686 format!(
1687 "invalid-metadata:{field}:{}",
1688 value.as_deref().unwrap_or("missing")
1689 )
1690 }
1691 kranz_engine::knowledge::RefreshVerdict::InvalidCitation { citation } => {
1692 format!("invalid-citation:{citation}")
1693 }
1694 kranz_engine::knowledge::RefreshVerdict::PathMissing { path } => {
1695 format!("path-missing:{path}")
1696 }
1697 kranz_engine::knowledge::RefreshVerdict::PathDrifted { path } => {
1698 format!("path-drifted:{path}")
1699 }
1700 kranz_engine::knowledge::RefreshVerdict::CommandSkipped { command } => {
1701 format!("command-skipped:{command}")
1702 }
1703 kranz_engine::knowledge::RefreshVerdict::ProbeFailed { target, error } => {
1704 format!("probe-failed:{target}:{error}")
1705 }
1706 })
1707 .collect();
1708 println!(
1709 "{} ({}) [{}] {}",
1710 finding.rel_path,
1711 finding.title,
1712 finding.freshness,
1713 labels.join(", ")
1714 );
1715 }
1716 if report.check_needed() {
1717 println!("knowledge refresh: check-needed");
1718 } else {
1719 println!("knowledge refresh: ok");
1720 }
1721 }
1722 Ok(if report.check_needed() { 1 } else { 0 })
1723}
1724
1725pub fn cmd_scan(repo: &Path, staged: bool, range: Option<&str>) -> Result<i32> {
1726 if staged && range.is_some() {
1727 bail!("choose either --staged or --range, not both");
1728 }
1729 let git = kranz_engine::git_ops::GitRepo::open(repo)?;
1730 let diff = if staged {
1731 git.diff_staged()?
1732 } else if let Some(range) = range {
1733 git.diff_range(range)?
1734 } else {
1735 git.diff_range("HEAD")?
1736 };
1737 let allowed = std::fs::read_to_string(repo.join(kranz_engine::scrub::SECRET_ALLOWLIST_PATH))
1738 .ok()
1739 .map(|text| kranz_engine::scrub::read_allowlist_text(&text))
1740 .unwrap_or_default();
1741 let findings = kranz_engine::scrub::filter_allowed(
1742 kranz_engine::scrub::scan_unified_diff(&diff),
1743 &allowed,
1744 );
1745 if findings.is_empty() {
1746 println!("secret scan passed");
1747 Ok(0)
1748 } else {
1749 println!(
1750 "secret scan failed; add a fingerprint to {} only for a reviewed false positive:\n{}",
1751 kranz_engine::scrub::SECRET_ALLOWLIST_PATH,
1752 kranz_engine::scrub::format_findings(&findings)
1753 );
1754 Ok(2)
1755 }
1756}
1757
1758pub fn cmd_domain_lint(repo: &Path, seed_config: Option<&Path>, json: bool) -> Result<i32> {
1769 use kranz_engine::domain_lint as dl;
1770 let config_path = repo.join(dl::DENYLIST_PATH);
1771
1772 if let Some(terms_file) = seed_config {
1773 let terms = std::fs::read_to_string(terms_file)
1774 .with_context(|| format!("read terms file {}", terms_file.display()))?;
1775 let existing = std::fs::read_to_string(&config_path).ok();
1776 let config = dl::seed_config(existing.as_deref(), &terms)?;
1777 if let Some(parent) = config_path.parent() {
1779 std::fs::create_dir_all(parent)
1780 .with_context(|| format!("create {}", parent.display()))?;
1781 }
1782 std::fs::write(&config_path, &config)
1783 .with_context(|| format!("write {}", config_path.display()))?;
1784 let denylist = dl::load_denylist(&config)?;
1785 println!(
1787 "seeded {} ({} hashed terms, {})",
1788 dl::DENYLIST_PATH,
1789 denylist.term_count(),
1790 if existing.is_some() {
1791 "salt preserved"
1792 } else {
1793 "fresh salt"
1794 }
1795 );
1796 warn_if_terms_file_unprotected(repo, terms_file);
1797 return Ok(0);
1798 }
1799
1800 let config_text = std::fs::read_to_string(&config_path).with_context(|| {
1801 format!(
1802 "read {} — seed it with `kranz domain-lint --seed-config <terms-file>` (docs/domain-lint.md)",
1803 dl::DENYLIST_PATH
1804 )
1805 })?;
1806 let denylist = dl::load_denylist(&config_text)?;
1807 let allowed = std::fs::read_to_string(repo.join(dl::ALLOWLIST_PATH))
1808 .ok()
1809 .map(|text| kranz_engine::scrub::read_allowlist_text(&text))
1810 .unwrap_or_default();
1811 let report = dl::lint_tree(repo, &denylist, &allowed)?;
1812
1813 if json {
1814 println!(
1815 "{}",
1816 serde_json::to_string_pretty(&serde_json::json!({
1817 "passed": report.is_clean(),
1818 "filesScanned": report.files_scanned,
1819 "filesSkipped": report.files_skipped,
1820 "findings": report.findings,
1821 }))?
1822 );
1823 }
1824 if report.is_clean() {
1825 if !json {
1826 println!(
1827 "domain lint passed ({} files scanned)",
1828 report.files_scanned
1829 );
1830 }
1831 Ok(0)
1832 } else {
1833 if !json {
1834 println!(
1835 "domain lint failed: {} unwaived hit(s); add a fingerprint to {} only for a reviewed false positive:",
1836 report.findings.len(),
1837 dl::ALLOWLIST_PATH
1838 );
1839 for finding in &report.findings {
1840 println!("{} {}:{}", finding.fingerprint, finding.path, finding.line);
1841 }
1842 }
1843 Ok(1)
1844 }
1845}
1846
1847fn warn_if_terms_file_unprotected(repo: &Path, terms_file: &Path) {
1852 let (Ok(repo), Ok(terms_file)) = (repo.canonicalize(), terms_file.canonicalize()) else {
1853 return;
1854 };
1855 let Ok(relative) = terms_file.strip_prefix(&repo) else {
1856 return; };
1858 let ignored = std::process::Command::new("git")
1859 .args(["check-ignore", "-q", "--"])
1860 .arg(relative)
1861 .current_dir(&repo)
1862 .status()
1863 .map(|status| status.success())
1864 .unwrap_or(true);
1866 if !ignored {
1867 eprintln!(
1868 "warning: {} is inside the repo and NOT gitignored — move it outside the repo or use {} (gitignored)",
1869 relative.display(),
1870 kranz_engine::domain_lint::TERMS_LOCAL_PATH
1871 );
1872 }
1873}
1874
1875fn cmd_pack_lint(repo: &Path, dir: &Path) -> Result<i32> {
1886 let trust = kranz_engine::pack::standards::trust_for_dir(repo, dir);
1887 match kranz_engine::pack::Pack::load_with_trust(dir, trust) {
1888 Ok(Some(pack)) => {
1889 print!("{}", kranz_engine::pack::render_lint(&pack));
1890 Ok(0)
1891 }
1892 Ok(None) => {
1893 println!(
1894 "no pack at {} (no {}) — nothing to lint",
1895 dir.display(),
1896 kranz_engine::pack::PACK_MANIFEST
1897 );
1898 Ok(0)
1899 }
1900 Err(err) => {
1901 eprintln!("invalid pack at {}: {err}", dir.display());
1902 Ok(1)
1903 }
1904 }
1905}
1906
1907fn cmd_standards_lint(repo: &Path, dir: &Path, against: Option<&str>) -> Result<i32> {
1914 let trust = kranz_engine::pack::standards::trust_for_dir(repo, dir);
1915 let pack = match kranz_engine::pack::Pack::load_with_trust(dir, trust) {
1916 Ok(Some(pack)) => pack,
1917 Ok(None) => {
1918 println!(
1919 "no pack at {} (no {}) — nothing to lint",
1920 dir.display(),
1921 kranz_engine::pack::PACK_MANIFEST
1922 );
1923 return Ok(0);
1924 }
1925 Err(err) => {
1926 eprintln!("invalid pack at {}: {err}", dir.display());
1927 return Ok(1);
1928 }
1929 };
1930 let Some(manifest) = &pack.standards else {
1931 println!(
1932 "pack `{}` (schema {}) at {} declares no [standards] root — nothing to lint",
1933 pack.name,
1934 pack.schema,
1935 dir.display()
1936 );
1937 return Ok(0);
1938 };
1939 print!(
1940 "{}",
1941 kranz_engine::pack::standards::render_manifest(manifest, trust)
1942 );
1943 if let Some(refname) = against {
1944 let Some(pack_rel) = kranz_engine::pack::standards::repo_relative_dir(repo, dir) else {
1948 eprintln!(
1949 "--against reads the base pack from tracked git blobs in {}; {} is outside \
1950 the repo — external packs have no base history to compare against (D-A)",
1951 repo.display(),
1952 dir.display()
1953 );
1954 return Ok(1);
1955 };
1956 let git = kranz_engine::git_ops::GitRepo::open(repo)?;
1957 let base = match kranz_engine::pack::standards::load_at_ref(&git, refname, &pack_rel) {
1958 Ok(base) => base,
1959 Err(err) => {
1960 eprintln!("cannot load the base standards at `{refname}`: {err}");
1961 return Ok(1);
1962 }
1963 };
1964 let errors = kranz_engine::pack::standards::check_transitions(base.as_ref(), manifest);
1965 print!(
1966 "{}",
1967 kranz_engine::pack::standards::render_transition_report(
1968 refname,
1969 base.as_ref(),
1970 &errors
1971 )
1972 );
1973 if !errors.is_empty() {
1974 return Ok(1);
1975 }
1976 }
1977 Ok(0)
1978}
1979
1980#[allow(clippy::too_many_arguments)]
1990fn cmd_standards_waive(
1991 repo: &Path,
1992 mission_id: &str,
1993 rule: &str,
1994 revision: Option<u64>,
1995 finding: Option<&str>,
1996 reason: &str,
1997 expires: &str,
1998 force_lock: LockForce,
1999) -> Result<i32> {
2000 let expires_at = match chrono::DateTime::parse_from_rfc3339(expires) {
2001 Ok(parsed) => parsed.with_timezone(&chrono::Utc),
2002 Err(err) => {
2003 eprintln!("waiver refused: --expires must be an RFC 3339 instant: {err}");
2004 return Ok(1);
2005 }
2006 };
2007 let request = kranz_engine::standards_waiver::WaiverRequest {
2008 rule_id: rule.to_string(),
2009 revision,
2010 finding_subject: finding.map(str::to_string),
2011 reason: reason.to_string(),
2012 expires_at,
2013 };
2014 let outcome = match kranz_engine::standards_waiver::approve_standards_waiver(
2015 repo, mission_id, &request, "cli", force_lock,
2016 ) {
2017 Ok(outcome) => outcome,
2018 Err(kranz_engine::error::EngineError::LockHeld(e)) => {
2019 eprintln!(
2020 "waiver refused: an engine still holds mission '{mission_id}'s lock — stop \
2021 the running mission first (a waiver against a live mission would race the \
2022 runner's own appends).\n (underlying: {e})"
2023 );
2024 return Ok(1);
2025 }
2026 Err(e) => {
2027 eprintln!("waiver refused: {e}");
2028 return Ok(1);
2029 }
2030 };
2031 let pinned = &outcome.rule;
2034 println!(
2035 "recorded standards.waiver.approved (seq {})",
2036 outcome.event.seq
2037 );
2038 println!("mission: {mission_id}");
2039 println!(
2040 "rule: {} r{} — {}, {}; checker {}; waivable: {}",
2041 pinned.id,
2042 pinned.revision,
2043 pinned.level,
2044 pinned.effective_status,
2045 pinned.checker.as_deref().unwrap_or("-"),
2046 pinned.waivable
2047 );
2048 println!(" statement: {}", pinned.statement);
2049 println!(
2050 "finding: {} (run {})\n evidence: {}",
2051 outcome.finding_subject, outcome.run_id, outcome.finding_evidence
2052 );
2053 println!(" fingerprint: sha256:{}", outcome.finding_fingerprint);
2054 if pinned.when_paths.is_empty() {
2055 println!(
2056 "affected paths: the whole mission diff ({} path(s)) — the rule is unscoped",
2057 outcome.affected_paths.len()
2058 );
2059 } else if outcome.affected_paths.is_empty() {
2060 println!(
2061 "affected paths: (none — the rule's when-paths match no changed path; the \
2062 waiver binds the empty scoped diff)"
2063 );
2064 } else {
2065 println!("affected paths: {}", outcome.affected_paths.join(", "));
2066 }
2067 println!(
2068 "diff digest: sha256:{} (covers the affected-path diff at the mission branch tip)",
2069 outcome.diff_digest
2070 );
2071 println!(
2072 "approver: {} via cli\nreason: {reason}\nexpires: {}",
2073 kranz_engine::standards_waiver::LOCAL_OPERATOR,
2074 expires_at.to_rfc3339()
2075 );
2076 Ok(0)
2077}
2078
2079fn cmd_standards_attest(
2084 repo: &Path,
2085 mission_id: &str,
2086 rule: &str,
2087 reason: &str,
2088 force_lock: LockForce,
2089) -> Result<i32> {
2090 let record = match kranz_engine::standards_attestation::approve_attestation(
2091 repo, mission_id, rule, reason, "cli", force_lock,
2092 ) {
2093 Ok(record) => record,
2094 Err(kranz_engine::error::EngineError::LockHeld(e)) => {
2095 eprintln!(
2096 "attestation refused: an engine still holds mission '{mission_id}'s lock — stop \
2097 the running mission first.\n (underlying: {e})"
2098 );
2099 return Ok(1);
2100 }
2101 Err(e) => {
2102 eprintln!("attestation refused: {e}");
2103 return Ok(1);
2104 }
2105 };
2106 println!(
2107 "recorded standards.attestation.approved (seq {})",
2108 record.seq
2109 );
2110 println!("mission: {mission_id}");
2111 println!("rule: {} r{}", record.rule_id, record.rule_revision);
2112 if record.paths.is_empty() {
2113 println!("affected paths: (none)");
2114 } else {
2115 println!("affected paths: {}", record.paths.join(", "));
2116 }
2117 println!("diff digest: sha256:{}", record.diff_digest);
2118 println!(
2119 "approver: {} via {}\nreason: {}",
2120 record.approver, record.surface, record.reason
2121 );
2122 Ok(0)
2123}
2124
2125pub fn cmd_missions(repo: &Path) -> Result<String> {
2132 let index_contents =
2133 std::fs::read_to_string(MissionPaths::new(repo, "_").missions_dir().join("index.md"))
2134 .unwrap_or_default();
2135 let mut ids = MissionPaths::list_missions(repo);
2136 for id in mission_catalog::mission_index_ids(&index_contents) {
2137 if !ids.contains(&id) {
2138 ids.push(id);
2139 }
2140 }
2141 ids.sort();
2142 if ids.is_empty() {
2143 return Ok("no missions\n".to_string());
2144 }
2145 let ticket_of: std::collections::HashMap<String, String> =
2148 kranz_engine::ticket::Ticket::list(repo)
2149 .into_iter()
2150 .filter_map(|t| {
2151 kranz_engine::ticket::Ticket::mission_for(repo, &t.slug).map(|m| (m, t.slug))
2152 })
2153 .collect();
2154 let mut out = String::new();
2155 for id in ids {
2156 let ticket = ticket_of
2157 .get(&id)
2158 .map(|s| format!(" [ticket: {s}]"))
2159 .unwrap_or_default();
2160 let paths = MissionPaths::new(repo, &id);
2161 if !paths.events_file().is_file() {
2162 out.push_str(&format!(
2163 "{id} {:<10} deleted mission (no data recorded)\n",
2164 "DELETED"
2165 ));
2166 continue;
2167 }
2168 if let Err(error) = paths.require_no_follow() {
2171 out.push_str(&format!("{id} {:<10} (unreadable: {error})\n", "FAILED"));
2172 continue;
2173 }
2174 match load_state(repo, &id) {
2175 Ok(state) => out.push_str(&format!(
2176 "{id} {:<10} {}{ticket}\n",
2177 output::mission_status_label(state.mission.status),
2178 state.mission.goal
2179 )),
2180 Err(e) => out.push_str(&format!("{id} {:<10} (unreadable: {e:#})\n", "FAILED")),
2181 }
2182 }
2183 Ok(out)
2184}
2185
2186pub fn cmd_abandon(
2194 repo: &Path,
2195 mission_id: &str,
2196 reason: &str,
2197 force_lock: LockForce,
2198) -> Result<()> {
2199 require_mission(repo, mission_id)?;
2200 mission_catalog::abandon_mission(repo, mission_id, reason, force_lock).map_err(|e| {
2201 if matches!(e, kranz_engine::error::EngineError::LockHeld(_)) {
2202 anyhow!(
2203 "cannot abandon mission '{mission_id}' — an engine still holds its lock. \
2204 Stop the running `kranz run` first. If the holder is a crashed leftover, \
2205 pass --force-lock (steals unless the holder is provably alive); a provably \
2206 LIVE holder that you have verified to be a zombie or foreign process \
2207 additionally requires --dangerously-steal-live-lock.\n (underlying: {e})"
2208 )
2209 } else {
2210 anyhow::Error::new(e).context(format!("abandoning mission '{mission_id}'"))
2211 }
2212 })
2213}
2214
2215#[derive(Debug, Clone, PartialEq, Eq)]
2217pub struct CleanEntry {
2218 pub id: String,
2219 pub status_label: String,
2220 pub goal: String,
2221}
2222
2223pub fn select_cleanable(repo: &Path, all: bool) -> Vec<CleanEntry> {
2231 let mut out = Vec::new();
2232 for id in MissionPaths::list_missions(repo) {
2233 let paths = MissionPaths::new(repo, &id);
2234 if mission_catalog::mission_lock_is_live(&paths) {
2236 continue;
2237 }
2238 let Ok(state) = load_state(repo, &id) else {
2239 continue; };
2241 let has_plan = paths.plan_file().is_file();
2242 if mission_catalog::cleanable_class(state.mission.status, has_plan).is_cleaned(all) {
2243 out.push(CleanEntry {
2244 id,
2245 status_label: output::mission_status_label(state.mission.status).to_string(),
2246 goal: state.mission.goal,
2247 });
2248 }
2249 }
2250 out
2251}
2252
2253pub fn render_clean_listing(entries: &[CleanEntry]) -> String {
2256 if entries.is_empty() {
2257 return "nothing to clean\n".to_string();
2258 }
2259 let mut out = String::new();
2260 for e in entries {
2261 out.push_str(&format!("{:<10} {} {}\n", e.status_label, e.id, e.goal));
2262 }
2263 out
2264}
2265
2266fn cmd_clean(repo: &Path, yes: bool, all: bool) -> Result<i32> {
2272 let entries = select_cleanable(repo, all);
2273 if entries.is_empty() {
2274 print!("{}", render_clean_listing(&entries));
2275 return Ok(0);
2276 }
2277
2278 print!("{}", render_clean_listing(&entries));
2279 println!(
2280 "\n{} mission director{} above would be removed.{}",
2281 entries.len(),
2282 if entries.len() == 1 { "y" } else { "ies" },
2283 if all {
2284 ""
2285 } else {
2286 " (Complete missions are kept; pass --all to include them.)"
2287 }
2288 );
2289
2290 if !yes && !confirm_clean()? {
2291 println!("clean aborted; nothing removed.");
2292 return Ok(0);
2293 }
2294
2295 let removed = remove_missions(repo, &entries, true);
2296 println!(
2297 "cleaned {} mission director{}",
2298 removed.len(),
2299 if removed.len() == 1 { "y" } else { "ies" }
2300 );
2301 Ok(0)
2302}
2303
2304pub fn remove_missions(repo: &Path, entries: &[CleanEntry], verbose: bool) -> Vec<String> {
2310 let mut removed = Vec::new();
2311 for e in entries {
2312 let paths = MissionPaths::new(repo, &e.id);
2313 if mission_catalog::mission_lock_is_live(&paths) {
2319 eprintln!("kranz: skipping {} — became live since listing", e.id);
2320 continue;
2321 }
2322 let dir = paths.mission_dir();
2323 match std::fs::remove_dir_all(&dir) {
2324 Ok(()) => {
2325 if verbose {
2326 println!("removed {}", dir.display());
2327 }
2328 mission_catalog::prune_mission_index_file(repo, &e.id);
2329 removed.push(e.id.clone());
2330 }
2331 Err(err) => eprintln!("kranz: could not remove {}: {err}", dir.display()),
2332 }
2333 }
2334 removed
2335}
2336
2337fn confirm_clean() -> Result<bool> {
2340 use std::io::BufRead;
2341 print!("proceed? [y/N] ");
2342 std::io::stdout().flush().ok();
2343 let mut line = String::new();
2344 let n = std::io::stdin().lock().read_line(&mut line)?;
2345 if n == 0 {
2346 return Ok(false); }
2348 Ok(matches!(
2349 line.trim().to_ascii_lowercase().as_str(),
2350 "y" | "yes"
2351 ))
2352}
2353
2354pub(crate) fn refuse_non_loopback_without_insecure_lan(
2364 bind: std::net::IpAddr,
2365 insecure_lan: bool,
2366) -> Result<()> {
2367 if !bind.is_loopback() && !insecure_lan {
2368 anyhow::bail!(
2369 "refusing to bind {bind}: non-loopback binds expose the API on the \
2370 network (reads require the read-only or mutation token; POSTs \
2371 require the mutation token). \
2372 Re-run with `--insecure-lan` if you intentionally trust this \
2373 network (LAN/tailnet), or keep the default `--host 127.0.0.1`."
2374 );
2375 }
2376 Ok(())
2377}
2378
2379pub(crate) fn effective_require_read_token(bind_is_loopback: bool, read_auth: bool) -> bool {
2385 !bind_is_loopback || read_auth
2386}
2387
2388#[allow(clippy::too_many_arguments)]
2398async fn cmd_serve(
2399 repo: PathBuf,
2400 host: String,
2401 port: u16,
2402 insecure_lan: bool,
2403 read_auth: bool,
2404 open: bool,
2405 dashboard: Option<PathBuf>,
2406 token: Option<String>,
2407 read_token: Option<String>,
2408 slack: bool,
2409) -> Result<i32> {
2410 let bind: std::net::IpAddr = host
2411 .parse()
2412 .map_err(|e| anyhow!("--host '{host}' is not an IP address: {e}"))?;
2413 refuse_non_loopback_without_insecure_lan(bind, insecure_lan)?;
2414 let token = token
2415 .or_else(|| std::env::var("KRANZ_TOKEN").ok())
2416 .unwrap_or_else(kranz_server::generate_token);
2417 let read_token = read_token
2418 .or_else(|| std::env::var("KRANZ_READ_TOKEN").ok())
2419 .unwrap_or_else(kranz_server::generate_token);
2420 kranz_server::MutationAuthority::new(token.clone())
2424 .context("--token / KRANZ_TOKEN is invalid")?;
2425 anyhow::ensure!(
2426 !read_token.is_empty() && read_token.bytes().all(|byte| byte.is_ascii_graphic()),
2427 "--read-token / KRANZ_READ_TOKEN must be non-empty visible ASCII without whitespace"
2428 );
2429 anyhow::ensure!(
2430 read_token != token,
2431 "--read-token / KRANZ_READ_TOKEN must differ from the mutation token"
2432 );
2433 if !bind.is_loopback() {
2434 eprintln!(
2435 "WARNING: binding {bind} with --insecure-lan — the API is reachable \
2436 beyond this machine. Reads require the read-only or mutation token; \
2437 POSTs require the mutation token. Use only on a network you trust."
2438 );
2439 }
2440 if read_auth && effective_require_read_token(bind.is_loopback(), read_auth) {
2441 eprintln!(
2442 "--read-auth: GETs and the WS upgrade now require the read-only or \
2443 mutation token, including on loopback; POSTs still require mutation \
2444 authority."
2445 );
2446 }
2447 let listener = kranz_server::bind_listener(bind, port)
2450 .await
2451 .map_err(|e| anyhow!("failed to bind {bind}:{port}: {e}"))?;
2452 let local_addr = listener
2453 .local_addr()
2454 .map_err(|e| anyhow!("failed to read bound address: {e}"))?;
2455 let multi_host = Arc::new(kranz_server::MultiRepoHost::from_global_config(
2459 kranz_engine::paths::global_config().as_deref(),
2460 repo.clone(),
2461 )?);
2462 multi_host.ensure_auto_work_started();
2465
2466 if slack {
2470 if multi_host.uses_operator_catalog() {
2471 let catalog = multi_repo_slack_catalog(&multi_host)?;
2472 tokio::spawn(async move {
2473 let never = std::future::pending::<()>();
2474 if let Err(error) = kranz_slack::serve_slack_catalog(catalog, never).await {
2475 tracing::error!(%error, "slack catalog bridge exited with an error");
2476 }
2477 });
2478 } else {
2479 let context = single_repo_slack_context(&multi_host)?;
2480 let repo_slack = context.root().to_path_buf();
2481 let hosted = context.host().ok_or_else(|| {
2482 anyhow!(
2483 "cannot start Slack bridge for unavailable repository '{}': {}",
2484 context.id(),
2485 context.unavailable_reason().unwrap_or("unavailable")
2486 )
2487 })?;
2488 let bridge_host: kranz_slack::SharedHost =
2489 Arc::new(crate::host_bridge::HostedPlanning(hosted.clone()));
2490 tokio::spawn(async move {
2491 let never = std::future::pending::<()>();
2495 if let Err(e) =
2496 kranz_slack::serve_slack(&repo_slack, Some(bridge_host), never).await
2497 {
2498 tracing::error!(error = %e, "slack bridge exited with an error");
2499 }
2500 });
2501 }
2502 }
2503
2504 let dashboard_assets = resolve_dashboard_assets(&repo, dashboard);
2505 let display_host = match local_addr.ip() {
2508 std::net::IpAddr::V6(v6) => format!("[{v6}]"),
2509 std::net::IpAddr::V4(v4) => v4.to_string(),
2510 };
2511 let url = format!("http://{display_host}:{}/", local_addr.port());
2512 println!("kranz server on {url}");
2513 println!("mutation token: {token}");
2514 println!("read token: {read_token} (GETs/WS only — safe for dashboards and agents)");
2515 if open {
2516 println!(
2517 "opening the bare URL; no token rides in the opener's argv. Paste the \
2518 mutation token above when the dashboard asks for one"
2519 );
2520 }
2521 match &dashboard_assets {
2522 Some(DashboardAssets::Embedded) => println!(
2523 "serving embedded dashboard ({})",
2524 crate::embedded_dashboard::EMBEDDED_DASHBOARD_SOURCE
2525 ),
2526 Some(DashboardAssets::Dir(dir)) => println!("serving dashboard from {}", dir.display()),
2527 None => println!(
2528 "no dashboard build found (--dashboard, $KRANZ_DASHBOARD_DIST, \
2529 <repo>/apps/dashboard/dist, installed asset dirs, or the kranz checkout) \
2530 and no embedded dashboard is available; serving API only"
2531 ),
2532 }
2533
2534 let static_assets = dashboard_assets.map(|assets| match assets {
2535 DashboardAssets::Dir(dir) => kranz_server::DashboardStatic::Dir(dir),
2536 DashboardAssets::Embedded => {
2537 kranz_server::DashboardStatic::Embedded(crate::embedded_dashboard::EMBEDDED_DASHBOARD)
2538 }
2539 });
2540
2541 if open {
2542 let url = url.clone();
2567 tokio::spawn(async move {
2568 tokio::time::sleep(Duration::from_millis(600)).await;
2569 open_browser(&url);
2570 });
2571 }
2572
2573 let shutdown = async {
2574 if let Err(e) = tokio::signal::ctrl_c().await {
2575 tracing::error!(error = %e, "failed to install ctrl-c handler");
2576 }
2577 };
2578 let result = serve_multi_with_token_cleanup(
2579 &repo,
2580 multi_host,
2581 listener,
2582 static_assets,
2583 token,
2584 read_token,
2585 read_auth,
2586 shutdown,
2587 )
2588 .await;
2589 match result {
2590 Ok(()) => Ok(0),
2591 Err(e) => Err(anyhow!("server failed: {e}")),
2592 }
2593}
2594
2595fn single_repo_slack_context(
2596 multi_host: &kranz_server::MultiRepoHost,
2597) -> Result<Arc<kranz_server::RepoContext>> {
2598 multi_host
2599 .compatibility_context()
2600 .ok_or_else(|| anyhow!("--slack requires one healthy configured repository"))
2601}
2602
2603fn multi_repo_slack_catalog(
2604 multi_host: &kranz_server::MultiRepoHost,
2605) -> Result<kranz_slack::SlackCatalog> {
2606 let global_config = kranz_engine::paths::global_config()
2607 .ok_or_else(|| anyhow!("cannot locate the operator config directory for Slack affinity"))?;
2608 let affinity_path = global_config
2609 .parent()
2610 .expect("global config has a parent")
2611 .join("slack")
2612 .join("thread-affinity.json");
2613 let repos = multi_host
2614 .contexts()
2615 .map(|context| {
2616 let host = context.host().map(|host| {
2617 Arc::new(crate::host_bridge::HostedPlanning(host.clone()))
2618 as kranz_slack::SharedHost
2619 });
2620 kranz_slack::SlackRepo {
2621 id: context.id().to_string(),
2622 root: context.root().to_path_buf(),
2623 display_name: context
2624 .config()
2625 .display_name
2626 .clone()
2627 .unwrap_or_else(|| context.id().to_string()),
2628 routes: context
2629 .config()
2630 .slack
2631 .channels
2632 .iter()
2633 .map(|route| kranz_slack::SlackRoute {
2634 team_id: route.team.clone(),
2635 channel_id: route.channel.clone(),
2636 })
2637 .collect(),
2638 allow_users: context.config().slack.allow_users.clone(),
2639 available: context.is_healthy(),
2640 host,
2641 unavailable_reason: context.unavailable_reason().map(str::to_string),
2642 default: multi_host.default_repo() == Some(context.id()),
2643 }
2644 })
2645 .collect();
2646 kranz_slack::SlackCatalog::new(repos, affinity_path)
2647}
2648
2649#[allow(clippy::too_many_arguments)]
2650async fn serve_multi_with_token_cleanup(
2651 repo: &Path,
2652 multi_host: Arc<kranz_server::MultiRepoHost>,
2653 listener: tokio::net::TcpListener,
2654 static_assets: Option<kranz_server::DashboardStatic>,
2655 token: String,
2656 read_token: String,
2657 read_auth: bool,
2658 shutdown: impl std::future::Future<Output = ()> + Send + 'static,
2659) -> anyhow::Result<()> {
2660 let (token_file, read_token_file) = if multi_host.uses_operator_catalog() {
2661 let address = listener
2662 .local_addr()
2663 .context("cannot resolve the bound address for operator token storage")?;
2664 let token_file = write_operator_serve_token(address, &token)
2665 .context("cannot securely store the operator serve token")?;
2666 let read_token_file = write_operator_serve_read_token(address, &read_token)
2667 .context("cannot securely store the operator serve read token")?;
2668 (token_file, read_token_file)
2669 } else {
2670 let token_file =
2671 write_serve_token(repo, &token).context("cannot securely store the serve token")?;
2672 let read_token_file = write_serve_read_token(repo, &read_token)
2673 .context("cannot securely store the serve read token")?;
2674 (token_file, read_token_file)
2675 };
2676 let result = kranz_server::serve_multi_on_listener(
2677 multi_host,
2678 listener,
2679 static_assets,
2680 kranz_server::MutationAuthority::new(token)?,
2681 Some(read_token),
2682 read_auth,
2683 shutdown,
2684 )
2685 .await;
2686 remove_token_file(&token_file);
2687 remove_token_file(&read_token_file);
2688 result
2689}
2690
2691#[cfg(test)]
2697async fn serve_with_token_cleanup(
2698 repo: &Path,
2699 host: Arc<kranz_server::MissionHost>,
2700 listener: tokio::net::TcpListener,
2701 static_assets: Option<kranz_server::DashboardStatic>,
2702 token: String,
2703 read_token: String,
2704 shutdown: impl std::future::Future<Output = ()> + Send + 'static,
2705) -> anyhow::Result<()> {
2706 let token_file = write_serve_token(repo, &token)?;
2710 let read_token_file = write_serve_read_token(repo, &read_token)?;
2711
2712 let result = kranz_server::serve_on_listener(
2713 host,
2714 listener,
2715 static_assets,
2716 kranz_server::MutationAuthority::new(token)?,
2717 shutdown,
2718 )
2719 .await;
2720
2721 remove_token_file(&token_file);
2722 remove_token_file(&read_token_file);
2723
2724 result
2725}
2726
2727fn write_serve_token(repo: &Path, token: &str) -> std::io::Result<PathBuf> {
2733 write_token_file(&repo.join(".kranz").join("serve.token"), token)
2734}
2735
2736fn write_serve_read_token(repo: &Path, read_token: &str) -> std::io::Result<PathBuf> {
2742 write_token_file(&repo.join(".kranz").join("serve.read.token"), read_token)
2743}
2744
2745fn write_operator_serve_token(
2750 address: std::net::SocketAddr,
2751 token: &str,
2752) -> std::io::Result<PathBuf> {
2753 let global = kranz_engine::paths::global_config().ok_or_else(|| {
2754 std::io::Error::new(
2755 std::io::ErrorKind::NotFound,
2756 "cannot resolve operator config directory",
2757 )
2758 })?;
2759 let path = operator_serve_token_path(&global, address);
2760 write_token_file(&path, token)
2761}
2762
2763fn write_operator_serve_read_token(
2765 address: std::net::SocketAddr,
2766 read_token: &str,
2767) -> std::io::Result<PathBuf> {
2768 let global = kranz_engine::paths::global_config().ok_or_else(|| {
2769 std::io::Error::new(
2770 std::io::ErrorKind::NotFound,
2771 "cannot resolve operator config directory",
2772 )
2773 })?;
2774 let path = operator_serve_read_token_path(&global, address);
2775 write_token_file(&path, read_token)
2776}
2777
2778fn operator_serve_read_token_path(global_config: &Path, address: std::net::SocketAddr) -> PathBuf {
2780 operator_serve_token_path(global_config, address).with_extension("read.token")
2781}
2782
2783fn operator_serve_token_path(global_config: &Path, address: std::net::SocketAddr) -> PathBuf {
2784 let endpoint = match address.ip() {
2785 std::net::IpAddr::V4(ip) => format!("v4-{:08x}-{}", u32::from(ip), address.port()),
2786 std::net::IpAddr::V6(ip) => format!("v6-{:032x}-{}", u128::from(ip), address.port()),
2787 };
2788 global_config
2789 .parent()
2790 .unwrap_or_else(|| Path::new("."))
2791 .join("serve")
2792 .join(format!("{endpoint}.token"))
2793}
2794
2795fn legacy_operator_serve_token_path(global_config: &Path, port: u16) -> PathBuf {
2799 global_config
2800 .parent()
2801 .unwrap_or_else(|| Path::new("."))
2802 .join("serve")
2803 .join(format!("{port}.token"))
2804}
2805
2806fn write_token_file(path: &Path, token: &str) -> std::io::Result<PathBuf> {
2807 let dir = path.parent().ok_or_else(|| {
2808 std::io::Error::new(std::io::ErrorKind::InvalidInput, "token path has no parent")
2809 })?;
2810 std::fs::create_dir_all(dir)?;
2811 let tmp = dir.join(format!(
2816 ".{}.tmp-{}",
2817 path.file_name()
2818 .and_then(|name| name.to_str())
2819 .unwrap_or("serve.token"),
2820 uuid::Uuid::new_v4().as_simple()
2821 ));
2822 let write_tmp = || -> std::io::Result<()> {
2823 let mut options = std::fs::OpenOptions::new();
2824 options.create_new(true).write(true);
2825 #[cfg(unix)]
2826 {
2827 use std::os::unix::fs::OpenOptionsExt;
2828 options.mode(0o600);
2829 }
2830 let mut file = options.open(&tmp)?;
2831 file.write_all(token.as_bytes())?;
2832 file.flush()?;
2833 Ok(())
2834 };
2835 if let Err(error) = write_tmp() {
2836 let _ = std::fs::remove_file(&tmp);
2837 return Err(error);
2838 }
2839 #[cfg(unix)]
2840 {
2841 use std::os::unix::fs::PermissionsExt;
2842 std::fs::set_permissions(&tmp, std::fs::Permissions::from_mode(0o600))?;
2843 }
2844 if let Err(error) = std::fs::rename(&tmp, path) {
2845 let _ = std::fs::remove_file(&tmp);
2846 return Err(error);
2847 }
2848 Ok(path.to_path_buf())
2849}
2850
2851fn remove_token_file(path: &Path) {
2852 if let Err(err) = std::fs::remove_file(path) {
2853 if err.kind() != std::io::ErrorKind::NotFound {
2854 eprintln!("kranz: could not remove {}: {err}", path.display());
2855 }
2856 }
2857}
2858
2859#[derive(Debug, Clone, PartialEq, Eq)]
2860enum DashboardAssets {
2861 Dir(PathBuf),
2862 Embedded,
2863}
2864
2865#[derive(Debug, Clone, Default)]
2866struct DashboardResolutionInputs {
2867 env_dist: Option<PathBuf>,
2868 home: Option<PathBuf>,
2869 exe: Option<PathBuf>,
2870 manifest_dir: Option<PathBuf>,
2871 embedded_available: bool,
2872}
2873
2874impl DashboardResolutionInputs {
2875 fn runtime() -> Self {
2876 Self {
2877 env_dist: std::env::var_os("KRANZ_DASHBOARD_DIST").map(PathBuf::from),
2878 home: std::env::var_os(if cfg!(windows) { "USERPROFILE" } else { "HOME" })
2879 .map(PathBuf::from),
2880 exe: std::env::current_exe().ok(),
2881 manifest_dir: Some(PathBuf::from(env!("CARGO_MANIFEST_DIR"))),
2882 embedded_available: !crate::embedded_dashboard::EMBEDDED_DASHBOARD.is_empty(),
2883 }
2884 }
2885}
2886
2887fn resolve_dashboard_assets(repo: &Path, explicit: Option<PathBuf>) -> Option<DashboardAssets> {
2888 resolve_dashboard_assets_from(repo, explicit, &DashboardResolutionInputs::runtime())
2889}
2890
2891fn resolve_dashboard_assets_from(
2899 repo: &Path,
2900 explicit: Option<PathBuf>,
2901 inputs: &DashboardResolutionInputs,
2902) -> Option<DashboardAssets> {
2903 if let Some(d) = explicit {
2904 return Some(DashboardAssets::Dir(d));
2907 }
2908
2909 if let Some(d) = first_dashboard_dir(dashboard_dir_candidates(repo, inputs)) {
2910 return Some(DashboardAssets::Dir(d));
2911 }
2912
2913 inputs
2914 .embedded_available
2915 .then_some(DashboardAssets::Embedded)
2916}
2917
2918pub fn resolve_dashboard_dist(repo: &Path, explicit: Option<PathBuf>) -> Option<PathBuf> {
2921 match resolve_dashboard_assets(repo, explicit) {
2922 Some(DashboardAssets::Dir(dir)) => Some(dir),
2923 Some(DashboardAssets::Embedded) | None => None,
2924 }
2925}
2926
2927fn dashboard_dir_candidates(repo: &Path, inputs: &DashboardResolutionInputs) -> Vec<PathBuf> {
2928 let mut candidates = Vec::new();
2929
2930 candidates.extend(inputs.env_dist.clone());
2931 candidates.push(repo.join("apps").join("dashboard").join("dist"));
2932
2933 if let Some(home) = &inputs.home {
2934 candidates.push(home.join(".kranz").join("dashboard").join("dist"));
2935 candidates.push(home.join(".kranz").join("dashboard"));
2936 }
2937
2938 if let Some(exe) = &inputs.exe {
2939 candidates.extend(installed_dashboard_dirs(exe));
2940 candidates.extend(source_checkout_dist_from_exe(exe));
2941 }
2942
2943 if let Some(manifest_dir) = &inputs.manifest_dir {
2944 candidates.extend(source_checkout_dist_from_manifest(manifest_dir));
2945 candidates.push(manifest_dir.join("assets").join("dashboard").join("dist"));
2946 }
2947
2948 candidates
2949}
2950
2951fn first_dashboard_dir(candidates: Vec<PathBuf>) -> Option<PathBuf> {
2952 candidates
2953 .into_iter()
2954 .find(|d| d.join("index.html").is_file())
2955}
2956
2957fn installed_dashboard_dirs(exe: &Path) -> Vec<PathBuf> {
2958 let Some(bin_dir) = exe.parent() else {
2959 return Vec::new();
2960 };
2961 let mut dirs = vec![
2962 bin_dir.join("dashboard").join("dist"),
2963 bin_dir.join("dashboard"),
2964 ];
2965 if let Some(prefix) = bin_dir.parent() {
2966 dirs.push(
2967 prefix
2968 .join("share")
2969 .join("kranz")
2970 .join("dashboard")
2971 .join("dist"),
2972 );
2973 dirs.push(prefix.join("share").join("kranz").join("dashboard"));
2974 }
2975 dirs
2976}
2977
2978fn source_checkout_dist_from_exe(exe: &Path) -> Option<PathBuf> {
2979 let profile_dir = exe.parent()?;
2980 let target_dir = profile_dir.parent()?;
2981 if target_dir.file_name()? != "target" {
2982 return None;
2983 }
2984 Some(
2985 target_dir
2986 .parent()?
2987 .join("apps")
2988 .join("dashboard")
2989 .join("dist"),
2990 )
2991}
2992
2993fn source_checkout_dist_from_manifest(manifest_dir: &Path) -> Option<PathBuf> {
2994 Some(
2995 manifest_dir
2996 .parent()?
2997 .parent()?
2998 .join("apps")
2999 .join("dashboard")
3000 .join("dist"),
3001 )
3002}
3003
3004fn open_browser(url: &str) {
3006 #[cfg(target_os = "macos")]
3007 let mut command = {
3008 let mut c = std::process::Command::new("open");
3009 c.arg(url);
3010 c
3011 };
3012 #[cfg(target_os = "windows")]
3013 let mut command = {
3014 let mut c = std::process::Command::new("cmd");
3015 c.args(["/C", "start", "", url]);
3017 c
3018 };
3019 #[cfg(not(any(target_os = "macos", target_os = "windows")))]
3020 let mut command = {
3021 let mut c = std::process::Command::new("xdg-open");
3022 c.arg(url);
3023 c
3024 };
3025
3026 match command.spawn() {
3027 Ok(mut child) => {
3028 std::thread::spawn(move || {
3030 let _ = child.wait();
3031 });
3032 }
3033 Err(e) => eprintln!("kranz: could not open the browser: {e}"),
3034 }
3035}
3036
3037fn resolve_release_token(
3049 repo: &Path,
3050 url: &str,
3051 flag: Option<String>,
3052) -> Result<String, ReleaseTokenError> {
3053 let parsed = reqwest::Url::parse(url).ok();
3054 let operator_lookup = parsed
3055 .as_ref()
3056 .and_then(|url| {
3057 kranz_engine::paths::global_config()
3058 .as_deref()
3059 .map(|global| operator_token_for_url(global, url))
3060 })
3061 .unwrap_or(OperatorTokenLookup::Absent);
3062 let allow_repo_fallback = parsed.as_ref().is_some_and(automatic_repo_token_allowed);
3063 resolve_release_token_from_sources(repo, operator_lookup, allow_repo_fallback, flag)
3064}
3065
3066#[derive(Debug, Clone, PartialEq, Eq)]
3071enum OperatorTokenLookup {
3072 Absent,
3073 Found(String),
3074 Legacy(String),
3075 Ambiguous,
3076}
3077
3078#[derive(Debug, Clone, PartialEq, Eq)]
3079enum ReleaseTokenError {
3080 Absent,
3081 Ambiguous,
3082}
3083
3084fn resolve_release_token_from_sources(
3085 repo: &Path,
3086 operator_lookup: OperatorTokenLookup,
3087 allow_repo_fallback: bool,
3088 flag: Option<String>,
3089) -> Result<String, ReleaseTokenError> {
3090 if let Some(flag) = flag {
3091 return Ok(flag);
3092 }
3093 if let Ok(env) = std::env::var("KRANZ_TOKEN") {
3094 return Ok(env);
3095 }
3096 match operator_lookup {
3097 OperatorTokenLookup::Found(authority) => Ok(authority),
3098 OperatorTokenLookup::Legacy(authority) => {
3103 let repo_authority = allow_repo_fallback
3104 .then(|| read_token_file(&repo.join(".kranz").join("serve.token")))
3105 .flatten();
3106 Ok(repo_authority.unwrap_or(authority))
3107 }
3108 OperatorTokenLookup::Ambiguous => Err(ReleaseTokenError::Ambiguous),
3109 OperatorTokenLookup::Absent => allow_repo_fallback
3110 .then(|| read_token_file(&repo.join(".kranz").join("serve.token")))
3111 .flatten()
3112 .ok_or(ReleaseTokenError::Absent),
3113 }
3114}
3115
3116fn scan_operator_token_addresses(
3117 global_config: &Path,
3118 addresses: &[std::net::SocketAddr],
3119) -> OperatorTokenLookup {
3120 let mut authority = None;
3121 for address in addresses {
3122 if let Some(found) = read_token_file(&operator_serve_token_path(global_config, *address)) {
3123 if authority.is_some() {
3124 return OperatorTokenLookup::Ambiguous;
3127 }
3128 authority = Some(found);
3129 }
3130 }
3131 match authority {
3132 Some(authority) => OperatorTokenLookup::Found(authority),
3133 None => OperatorTokenLookup::Absent,
3134 }
3135}
3136
3137fn operator_token_for_url(global_config: &Path, url: &reqwest::Url) -> OperatorTokenLookup {
3138 let Some(port) = url.port_or_known_default() else {
3139 return OperatorTokenLookup::Absent;
3140 };
3141 let Some(host) = normalized_url_host(url) else {
3142 return OperatorTokenLookup::Absent;
3143 };
3144 let endpoint_lookup = if host.eq_ignore_ascii_case("localhost") {
3148 let primary = [
3149 std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, port)),
3150 std::net::SocketAddr::from((std::net::Ipv6Addr::LOCALHOST, port)),
3151 ];
3152 match scan_operator_token_addresses(global_config, &primary) {
3153 OperatorTokenLookup::Absent => scan_operator_token_addresses(
3154 global_config,
3155 &[
3156 std::net::SocketAddr::from((std::net::Ipv4Addr::UNSPECIFIED, port)),
3157 std::net::SocketAddr::from((std::net::Ipv6Addr::UNSPECIFIED, port)),
3158 ],
3159 ),
3160 other => other,
3161 }
3162 } else {
3163 let Ok(ip) = host.parse::<std::net::IpAddr>() else {
3164 return OperatorTokenLookup::Absent;
3165 };
3166 if !ip.is_loopback() {
3169 return OperatorTokenLookup::Absent;
3170 }
3171 let primary = [std::net::SocketAddr::new(ip, port)];
3172 match scan_operator_token_addresses(global_config, &primary) {
3173 OperatorTokenLookup::Absent => {
3174 let unspecified = match ip {
3175 std::net::IpAddr::V4(_) => {
3176 std::net::IpAddr::V4(std::net::Ipv4Addr::UNSPECIFIED)
3177 }
3178 std::net::IpAddr::V6(_) => {
3179 std::net::IpAddr::V6(std::net::Ipv6Addr::UNSPECIFIED)
3180 }
3181 };
3182 scan_operator_token_addresses(
3183 global_config,
3184 &[std::net::SocketAddr::new(unspecified, port)],
3185 )
3186 }
3187 other => other,
3188 }
3189 };
3190 match endpoint_lookup {
3191 OperatorTokenLookup::Absent => {
3192 match read_token_file(&legacy_operator_serve_token_path(global_config, port)) {
3194 Some(token) => OperatorTokenLookup::Legacy(token),
3195 None => OperatorTokenLookup::Absent,
3196 }
3197 }
3198 other => other,
3199 }
3200}
3201
3202fn automatic_repo_token_allowed(url: &reqwest::Url) -> bool {
3203 normalized_url_host(url).is_some_and(|host| {
3204 host.eq_ignore_ascii_case("localhost")
3205 || host
3206 .parse::<std::net::IpAddr>()
3207 .is_ok_and(|ip| ip.is_loopback())
3208 })
3209}
3210
3211fn normalized_url_host(url: &reqwest::Url) -> Option<&str> {
3212 let host = url.host_str()?;
3213 Some(
3214 host.strip_prefix('[')
3215 .and_then(|host| host.strip_suffix(']'))
3216 .unwrap_or(host),
3217 )
3218}
3219
3220fn read_token_file(path: &Path) -> Option<String> {
3221 let contents = std::fs::read_to_string(path).ok()?;
3222 let authority = contents.trim_end().to_string();
3223 if authority.is_empty() {
3224 None
3225 } else {
3226 Some(authority)
3227 }
3228}
3229
3230fn release_http_client() -> Result<reqwest::Client> {
3231 reqwest::Client::builder()
3232 .redirect(reqwest::redirect::Policy::none())
3233 .connect_timeout(Duration::from_secs(5))
3234 .timeout(Duration::from_secs(30))
3235 .build()
3236 .context("building release HTTP client")
3237}
3238
3239fn release_repo_id_is_valid(id: &str) -> bool {
3240 let mut chars = id.chars();
3241 let valid_first = chars.next().is_some_and(|ch| ch.is_ascii_alphanumeric());
3242 let valid_rest = chars.all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_'));
3243 valid_first && valid_rest
3244}
3245
3246async fn cmd_release(
3247 repo: &Path,
3248 mission_id: &str,
3249 url: &str,
3250 token: Option<String>,
3251) -> Result<i32> {
3252 let token = match resolve_release_token(repo, url, token) {
3253 Ok(token) => token,
3254 Err(ReleaseTokenError::Ambiguous) => {
3255 bail!(
3256 "multiple ~/.kranz/serve/<endpoint>.token files match {url}; \
3257 pass --token for the intended serve instead of guessing"
3258 );
3259 }
3260 Err(ReleaseTokenError::Absent) => {
3261 bail!(
3262 "no mutation token available — pass --token, set $KRANZ_TOKEN, or use a local URL matching \
3263 a live ~/.kranz/serve/<endpoint>.token / single-repo .kranz/serve.token \
3264 (the token `kranz serve` prints on startup)"
3265 );
3266 }
3267 };
3268
3269 let client = release_http_client()?;
3270 let repo_id = resolve_release_repo_id(repo, url, &client, &token).await?;
3271 if let Some(id) = repo_id.as_deref() {
3272 if !release_repo_id_is_valid(id) {
3273 bail!("live serve returned an invalid repository id '{id}'; refusing release");
3274 }
3275 }
3276 let endpoint = release_endpoint(url, mission_id, repo_id.as_deref());
3277
3278 let response = client
3279 .post(&endpoint)
3280 .header("x-kranz-token", token)
3281 .json(&serde_json::json!({}))
3282 .send()
3283 .await
3284 .map_err(|e| {
3285 if e.is_connect() {
3286 anyhow!("no kranz serve reachable at {url} — is it running?")
3287 } else {
3288 anyhow::Error::new(e).context(format!("releasing mission '{mission_id}'"))
3289 }
3290 })?;
3291
3292 match response.status() {
3293 reqwest::StatusCode::OK => {
3294 let body: serde_json::Value = response.json().await.unwrap_or_default();
3295 let released = body
3296 .get("released")
3297 .and_then(serde_json::Value::as_bool)
3298 .unwrap_or(false);
3299 if released {
3300 println!("mission {mission_id} released — the lock is now free");
3301 } else {
3302 println!("mission {mission_id} was already free (no lock held)");
3303 }
3304 Ok(0)
3305 }
3306 reqwest::StatusCode::CONFLICT => {
3307 eprintln!("kranz: mission {mission_id} has a turn in flight — try again shortly");
3308 Ok(1)
3309 }
3310 reqwest::StatusCode::NOT_FOUND => {
3311 eprintln!("kranz: unknown mission '{mission_id}' at {url}");
3312 Ok(1)
3313 }
3314 reqwest::StatusCode::UNAUTHORIZED => {
3315 eprintln!(
3316 "kranz: the token was missing or invalid — check --token / $KRANZ_TOKEN \
3317 against the token `kranz serve` printed on startup"
3318 );
3319 Ok(1)
3320 }
3321 other => {
3322 let body = response.text().await.unwrap_or_default();
3323 eprintln!("kranz: release failed ({other}): {body}");
3324 Ok(1)
3325 }
3326 }
3327}
3328
3329async fn resolve_release_repo_id(
3334 repo: &Path,
3335 url: &str,
3336 client: &reqwest::Client,
3337 token: &str,
3338) -> Result<Option<String>> {
3339 live_release_repo_id(repo, url, client, token).await
3340}
3341
3342fn release_repo_id_from_summaries(
3343 repo: &Path,
3344 repos: &[kranz_server::RepoSummary],
3345) -> Result<Option<String>> {
3346 if repos.is_empty() {
3347 bail!("the live serve returned an empty repository catalog; refusing an unscoped release");
3348 }
3349 let root = std::fs::canonicalize(repo).unwrap_or_else(|_| repo.to_path_buf());
3350 let root_str = root.to_string_lossy();
3351 let matches: Vec<&kranz_server::RepoSummary> = repos
3352 .iter()
3353 .filter(|entry| {
3354 let entry_root = PathBuf::from(&entry.root);
3355 let entry_canon =
3356 std::fs::canonicalize(&entry_root).unwrap_or_else(|_| entry_root.clone());
3357 entry_canon == root || entry.root == root_str
3358 })
3359 .collect();
3360 match matches.as_slice() {
3361 [one] => Ok(Some(one.id.clone())),
3362 [] => Err(anyhow!(
3363 "selected repository '{}' is not present in the live serve catalog; refusing an unscoped release",
3364 repo.display()
3365 )),
3366 _ => Err(anyhow!(
3367 "selected repository '{}' matches multiple live serve catalog entries; refusing release",
3368 repo.display()
3369 )),
3370 }
3371}
3372
3373async fn live_release_repo_id(
3374 repo: &Path,
3375 url: &str,
3376 client: &reqwest::Client,
3377 token: &str,
3378) -> Result<Option<String>> {
3379 let base = url.trim_end_matches('/');
3380 let mut request = client.get(format!("{base}/api/repos"));
3384 let loopback_catalog = reqwest::Url::parse(url)
3385 .ok()
3386 .is_some_and(|parsed| automatic_repo_token_allowed(&parsed));
3387 if !loopback_catalog {
3388 request = request.header("x-kranz-token", token);
3389 }
3390 let send_error = |e: reqwest::Error| {
3391 if e.is_connect() {
3392 anyhow!("no kranz serve reachable at {url} — is it running?")
3393 } else {
3394 anyhow::Error::new(e).context("listing live serve repositories")
3395 }
3396 };
3397 let mut response = request.send().await.map_err(send_error)?;
3398 if loopback_catalog && response.status() == reqwest::StatusCode::UNAUTHORIZED {
3403 response = client
3404 .get(format!("{base}/api/repos"))
3405 .header("x-kranz-token", token)
3406 .send()
3407 .await
3408 .map_err(send_error)?;
3409 }
3410 if !response.status().is_success() {
3411 bail!(
3412 "cannot list repositories at {url}: HTTP {}",
3413 response.status()
3414 );
3415 }
3416 let repos: Vec<kranz_server::RepoSummary> = response
3417 .json()
3418 .await
3419 .context("parsing GET /api/repos response")?;
3420 release_repo_id_from_summaries(repo, &repos)
3421}
3422
3423fn release_endpoint(url: &str, mission_id: &str, repo_id: Option<&str>) -> String {
3424 let base = url.trim_end_matches('/');
3425 match repo_id {
3426 Some(repo_id) => {
3427 format!("{base}/api/repos/{repo_id}/missions/{mission_id}/release")
3428 }
3429 None => format!("{base}/api/missions/{mission_id}/release"),
3430 }
3431}
3432
3433#[cfg(test)]
3434mod tests {
3435 use super::*;
3436 use std::fs;
3437
3438 static KRANZ_TOKEN_ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
3443
3444 #[tokio::test]
3447 #[allow(clippy::await_holding_lock)]
3451 async fn release_without_a_token_errors_clearly() {
3452 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3453 unsafe {
3455 std::env::remove_var("KRANZ_TOKEN");
3456 }
3457 let tmp = tempfile::tempdir().unwrap();
3458 let repo = tmp.path().to_path_buf();
3459 let err = cmd_release(&repo, "m-1", "http://192.0.2.1:4560", None)
3463 .await
3464 .unwrap_err();
3465 let msg = err.to_string();
3466 assert!(msg.contains("--token"), "{msg}");
3467 assert!(msg.contains("KRANZ_TOKEN"), "{msg}");
3468 }
3469
3470 #[test]
3473 fn release_reads_token_file_when_flag_and_env_absent() {
3474 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3475 unsafe {
3477 std::env::remove_var("KRANZ_TOKEN");
3478 }
3479 let tmp = tempfile::tempdir().unwrap();
3480 let repo = tmp.path().to_path_buf();
3481 write_serve_token(&repo, "file-token").unwrap();
3482
3483 assert_eq!(
3484 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Absent, true, None),
3485 Ok("file-token".to_string())
3486 );
3487 }
3488
3489 #[test]
3490 fn release_reads_token_file_but_flag_overrides() {
3491 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3492 unsafe {
3494 std::env::remove_var("KRANZ_TOKEN");
3495 }
3496 let tmp = tempfile::tempdir().unwrap();
3497 let repo = tmp.path().to_path_buf();
3498 write_serve_token(&repo, "file-token").unwrap();
3499
3500 assert_eq!(
3501 resolve_release_token_from_sources(
3502 &repo,
3503 OperatorTokenLookup::Absent,
3504 true,
3505 Some("flag-token".to_string()),
3506 ),
3507 Ok("flag-token".to_string())
3508 );
3509 }
3510
3511 #[test]
3512 fn release_reads_token_file_but_env_overrides() {
3513 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3514 unsafe {
3516 std::env::set_var("KRANZ_TOKEN", "env-token");
3517 }
3518 let tmp = tempfile::tempdir().unwrap();
3519 let repo = tmp.path().to_path_buf();
3520 write_serve_token(&repo, "file-token").unwrap();
3521
3522 let result =
3523 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Absent, true, None);
3524 unsafe {
3525 std::env::remove_var("KRANZ_TOKEN");
3526 }
3527 assert_eq!(result, Ok("env-token".to_string()));
3528 }
3529
3530 #[test]
3533 fn release_reads_token_file_none_present_yields_none() {
3534 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3535 unsafe {
3537 std::env::remove_var("KRANZ_TOKEN");
3538 }
3539 let tmp = tempfile::tempdir().unwrap();
3540 let repo = tmp.path().to_path_buf();
3541
3542 assert_eq!(
3543 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Absent, true, None),
3544 Err(ReleaseTokenError::Absent)
3545 );
3546 }
3547
3548 #[test]
3549 fn release_prefers_operator_process_token_over_repo_compatibility_token() {
3550 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3551 unsafe {
3552 std::env::remove_var("KRANZ_TOKEN");
3553 }
3554 let tmp = tempfile::tempdir().unwrap();
3555 let repo = tmp.path().join("repo");
3556 write_serve_token(&repo, "repo-token").unwrap();
3557 let global = tmp.path().join("operator").join("config.json");
3558 let address = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3559 let operator = operator_serve_token_path(&global, address);
3560 write_token_file(&operator, "operator-token").unwrap();
3561 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3562
3563 assert_eq!(
3564 resolve_release_token_from_sources(
3565 &repo,
3566 operator_token_for_url(&global, &url),
3567 true,
3568 None,
3569 ),
3570 Ok("operator-token".to_string())
3571 );
3572 }
3573
3574 #[test]
3575 fn ambiguous_operator_token_does_not_fall_through_to_repo_file() {
3576 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3577 unsafe {
3578 std::env::remove_var("KRANZ_TOKEN");
3579 }
3580 let tmp = tempfile::tempdir().unwrap();
3581 let repo = tmp.path().join("repo");
3582 write_serve_token(&repo, "repo-token").unwrap();
3583
3584 assert_eq!(
3585 resolve_release_token_from_sources(&repo, OperatorTokenLookup::Ambiguous, true, None,),
3586 Err(ReleaseTokenError::Ambiguous)
3587 );
3588 }
3589
3590 #[test]
3591 fn operator_token_discovery_is_ambiguous_for_dual_localhost_endpoint_files() {
3592 let tmp = tempfile::tempdir().unwrap();
3593 let global = tmp.path().join(".kranz").join("config.json");
3594 let v4 = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3595 let v6 = std::net::SocketAddr::from((std::net::Ipv6Addr::LOCALHOST, 4560));
3596 write_token_file(&operator_serve_token_path(&global, v4), "v4-token").unwrap();
3597 write_token_file(&operator_serve_token_path(&global, v6), "v6-token").unwrap();
3598 let url = reqwest::Url::parse("http://localhost:4560").unwrap();
3599
3600 assert_eq!(
3601 operator_token_for_url(&global, &url),
3602 OperatorTokenLookup::Ambiguous
3603 );
3604 }
3605
3606 #[test]
3607 fn operator_token_path_is_scoped_by_full_bound_endpoint() {
3608 let global = Path::new("/operator/.kranz/config.json");
3609 let a =
3610 operator_serve_token_path(global, std::net::SocketAddr::from(([127, 0, 0, 1], 4560)));
3611 let b =
3612 operator_serve_token_path(global, std::net::SocketAddr::from(([127, 0, 0, 2], 4560)));
3613 assert_ne!(a, b);
3614 assert_eq!(
3615 a,
3616 Path::new("/operator/.kranz/serve/v4-7f000001-4560.token")
3617 );
3618 }
3619
3620 #[test]
3621 fn operator_token_discovery_handles_ipv6_url_brackets() {
3622 let tmp = tempfile::tempdir().unwrap();
3623 let global = tmp.path().join(".kranz").join("config.json");
3624 let address = std::net::SocketAddr::from((std::net::Ipv6Addr::LOCALHOST, 4560));
3625 write_token_file(&operator_serve_token_path(&global, address), "ipv6-token").unwrap();
3626 let url = reqwest::Url::parse("http://[::1]:4560").unwrap();
3627
3628 assert_eq!(
3629 operator_token_for_url(&global, &url),
3630 OperatorTokenLookup::Found("ipv6-token".to_string())
3631 );
3632 }
3633
3634 #[test]
3635 fn automatic_token_discovery_refuses_remote_domain_urls() {
3636 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3637 unsafe {
3638 std::env::remove_var("KRANZ_TOKEN");
3639 }
3640 let tmp = tempfile::tempdir().unwrap();
3641 let repo = tmp.path().join("repo");
3642 write_serve_token(&repo, "repo-token").unwrap();
3643
3644 assert_eq!(
3645 resolve_release_token(&repo, "https://example.com:4560", None),
3646 Err(ReleaseTokenError::Absent)
3647 );
3648 }
3649
3650 #[test]
3651 fn operator_token_discovery_refuses_non_loopback_ip_literals() {
3652 let tmp = tempfile::tempdir().unwrap();
3653 let global = tmp.path().join(".kranz").join("config.json");
3654 let address = std::net::SocketAddr::from(([203, 0, 113, 10], 4560));
3655 write_token_file(&operator_serve_token_path(&global, address), "remote-token").unwrap();
3656 let url = reqwest::Url::parse("http://203.0.113.10:4560").unwrap();
3657
3658 assert_eq!(
3659 operator_token_for_url(&global, &url),
3660 OperatorTokenLookup::Absent
3661 );
3662 }
3663
3664 #[test]
3665 fn operator_token_discovery_prefers_exact_loopback_over_stale_unspecified() {
3666 let tmp = tempfile::tempdir().unwrap();
3667 let global = tmp.path().join(".kranz").join("config.json");
3668 let loopback = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3669 let unspecified = std::net::SocketAddr::from((std::net::Ipv4Addr::UNSPECIFIED, 4560));
3670 write_token_file(&operator_serve_token_path(&global, loopback), "live-token").unwrap();
3671 write_token_file(
3672 &operator_serve_token_path(&global, unspecified),
3673 "stale-token",
3674 )
3675 .unwrap();
3676 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3677
3678 assert_eq!(
3679 operator_token_for_url(&global, &url),
3680 OperatorTokenLookup::Found("live-token".to_string())
3681 );
3682 }
3683
3684 #[test]
3685 fn operator_token_discovery_falls_back_to_legacy_port_token() {
3686 let tmp = tempfile::tempdir().unwrap();
3687 let global = tmp.path().join(".kranz").join("config.json");
3688 write_token_file(
3689 &legacy_operator_serve_token_path(&global, 4560),
3690 "legacy-token",
3691 )
3692 .unwrap();
3693 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3694
3695 assert_eq!(
3696 operator_token_for_url(&global, &url),
3697 OperatorTokenLookup::Legacy("legacy-token".to_string())
3698 );
3699 }
3700
3701 #[test]
3702 fn stale_legacy_operator_token_does_not_mask_live_repo_token() {
3703 let _guard = KRANZ_TOKEN_ENV_LOCK.lock().unwrap();
3704 unsafe {
3705 std::env::remove_var("KRANZ_TOKEN");
3706 }
3707 let tmp = tempfile::tempdir().unwrap();
3708 let repo = tmp.path().join("repo");
3709 write_serve_token(&repo, "live-repo-token").unwrap();
3710 let global = tmp.path().join(".kranz").join("config.json");
3711 write_token_file(
3712 &legacy_operator_serve_token_path(&global, 4560),
3713 "stale-legacy-token",
3714 )
3715 .unwrap();
3716 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3717 let operator = operator_token_for_url(&global, &url);
3718
3719 assert_eq!(
3720 resolve_release_token_from_sources(&repo, operator, true, None),
3721 Ok("live-repo-token".to_string())
3722 );
3723 }
3724
3725 #[test]
3726 fn operator_token_discovery_prefers_endpoint_file_over_legacy_port_token() {
3727 let tmp = tempfile::tempdir().unwrap();
3728 let global = tmp.path().join(".kranz").join("config.json");
3729 let loopback = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 4560));
3730 write_token_file(
3731 &operator_serve_token_path(&global, loopback),
3732 "endpoint-token",
3733 )
3734 .unwrap();
3735 write_token_file(
3736 &legacy_operator_serve_token_path(&global, 4560),
3737 "legacy-token",
3738 )
3739 .unwrap();
3740 let url = reqwest::Url::parse("http://127.0.0.1:4560").unwrap();
3741
3742 assert_eq!(
3743 operator_token_for_url(&global, &url),
3744 OperatorTokenLookup::Found("endpoint-token".to_string())
3745 );
3746 }
3747
3748 #[test]
3749 fn release_repo_id_from_summaries_refuses_duplicate_root() {
3750 let tmp = tempfile::tempdir().unwrap();
3751 let repo = tmp.path().join("repo");
3752 std::fs::create_dir_all(&repo).unwrap();
3753 let root = repo.to_string_lossy().into_owned();
3754 let summaries = vec![
3755 kranz_server::RepoSummary {
3756 id: "alpha".to_string(),
3757 root: root.clone(),
3758 display_name: "alpha".to_string(),
3759 group: None,
3760 pinned: false,
3761 is_default: true,
3762 status: "healthy".to_string(),
3763 error: None,
3764 activity: kranz_server::RepoActivity::default(),
3765 },
3766 kranz_server::RepoSummary {
3767 id: "beta".to_string(),
3768 root,
3769 display_name: "beta".to_string(),
3770 group: None,
3771 pinned: false,
3772 is_default: false,
3773 status: "healthy".to_string(),
3774 error: None,
3775 activity: kranz_server::RepoActivity::default(),
3776 },
3777 ];
3778 let error = release_repo_id_from_summaries(&repo, &summaries).unwrap_err();
3779 assert!(error.to_string().contains("matches multiple"));
3780 }
3781
3782 #[test]
3783 fn slack_operator_catalog_is_not_rejected_by_legacy_context_helper() {
3784 fn init_git(root: &Path) {
3785 std::fs::create_dir_all(root).unwrap();
3786 let status = std::process::Command::new("git")
3787 .args(["init", "-q"])
3788 .arg(root)
3789 .status()
3790 .unwrap();
3791 assert!(status.success());
3792 }
3793
3794 let tmp = tempfile::tempdir().unwrap();
3795 let a = tmp.path().join("a");
3796 init_git(&a);
3797 let multi = kranz_server::MultiRepoHost::from_config(kranz_server::HostConfig {
3798 default_repo: Some("a".to_string()),
3799 max_concurrent_repos: 1,
3800 repos: vec![kranz_server::RepoConfig {
3801 id: "a".to_string(),
3802 root: a,
3803 display_name: None,
3804 group: None,
3805 pinned: false,
3806 slack: kranz_server::RepoSlackConfig::default(),
3807 }],
3808 })
3809 .unwrap();
3810 let context = single_repo_slack_context(&multi).unwrap();
3811 assert_eq!(context.id(), "a");
3812 }
3813
3814 #[test]
3815 fn release_endpoint_is_repo_scoped_when_catalog_id_is_known() {
3816 assert_eq!(
3817 release_endpoint("http://127.0.0.1:4560/", "same-id", Some("repo-b")),
3818 "http://127.0.0.1:4560/api/repos/repo-b/missions/same-id/release"
3819 );
3820 }
3821
3822 #[test]
3823 fn release_repo_id_from_summaries_matches_live_catalog_root() {
3824 let tmp = tempfile::tempdir().unwrap();
3825 let repo = tmp.path().join("repo");
3826 std::fs::create_dir_all(&repo).unwrap();
3827 let summaries = vec![kranz_server::RepoSummary {
3828 id: "alpha".to_string(),
3829 root: repo.to_string_lossy().into_owned(),
3830 display_name: "alpha".to_string(),
3831 group: None,
3832 pinned: false,
3833 is_default: true,
3834 status: "healthy".to_string(),
3835 error: None,
3836 activity: kranz_server::RepoActivity::default(),
3837 }];
3838 assert_eq!(
3839 release_repo_id_from_summaries(&repo, &summaries)
3840 .unwrap()
3841 .as_deref(),
3842 Some("alpha")
3843 );
3844 }
3845
3846 #[test]
3847 fn release_repo_id_from_summaries_refuses_unmatched_root() {
3848 let tmp = tempfile::tempdir().unwrap();
3849 let repo = tmp.path().join("local");
3850 let other = tmp.path().join("other");
3851 std::fs::create_dir_all(&repo).unwrap();
3852 std::fs::create_dir_all(&other).unwrap();
3853 let summaries = vec![kranz_server::RepoSummary {
3854 id: "other".to_string(),
3855 root: other.to_string_lossy().into_owned(),
3856 display_name: "other".to_string(),
3857 group: None,
3858 pinned: false,
3859 is_default: false,
3860 status: "healthy".to_string(),
3861 error: None,
3862 activity: kranz_server::RepoActivity::default(),
3863 }];
3864 let error = release_repo_id_from_summaries(&repo, &summaries).unwrap_err();
3865 assert!(error
3866 .to_string()
3867 .contains("not present in the live serve catalog"));
3868 }
3869
3870 #[test]
3871 fn release_repo_id_from_summaries_refuses_empty_catalog() {
3872 let tmp = tempfile::tempdir().unwrap();
3873 let error = release_repo_id_from_summaries(tmp.path(), &[]).unwrap_err();
3874 assert!(error.to_string().contains("empty repository catalog"));
3875 }
3876
3877 #[tokio::test]
3878 async fn live_release_repo_lookup_authenticates_protected_catalog() {
3879 let tmp = tempfile::tempdir().unwrap();
3880 let repo = tmp.path().join("repo");
3881 std::fs::create_dir_all(&repo).unwrap();
3882 let status = std::process::Command::new("git")
3883 .args(["init", "-q"])
3884 .arg(&repo)
3885 .status()
3886 .unwrap();
3887 assert!(status.success());
3888
3889 let catalog = Arc::new(kranz_server::MultiRepoHost::single(repo.clone()).unwrap());
3890 let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
3891 .await
3892 .unwrap();
3893 let address = listener.local_addr().unwrap();
3894 let app = kranz_server::router_with_multi_repo_host_and_addr(
3896 catalog,
3897 None,
3898 kranz_server::MutationAuthority::new("catalog-token").unwrap(),
3899 Some(address),
3900 true,
3901 false,
3902 );
3903 let server = tokio::spawn(async move {
3904 axum::serve(listener, app).await.unwrap();
3905 });
3906 let client = release_http_client().unwrap();
3907 let url = format!("http://{address}");
3908
3909 let repo_id = live_release_repo_id(&repo, &url, &client, "catalog-token")
3910 .await
3911 .unwrap();
3912
3913 server.abort();
3914 let _ = server.await;
3915 assert_eq!(repo_id.as_deref(), Some("repo"));
3916 }
3917
3918 #[tokio::test]
3919 async fn live_release_repo_lookup_retries_with_token_when_read_gated() {
3920 let tmp = tempfile::tempdir().unwrap();
3921 let repo = tmp.path().join("repo");
3922 std::fs::create_dir_all(&repo).unwrap();
3923 let status = std::process::Command::new("git")
3924 .args(["init", "-q"])
3925 .arg(&repo)
3926 .status()
3927 .unwrap();
3928 assert!(status.success());
3929
3930 let catalog = Arc::new(kranz_server::MultiRepoHost::single(repo.clone()).unwrap());
3931 let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
3932 .await
3933 .unwrap();
3934 let address = listener.local_addr().unwrap();
3935 let app = kranz_server::router_with_multi_repo_host_and_addr(
3940 catalog,
3941 None,
3942 kranz_server::MutationAuthority::new("catalog-token").unwrap(),
3943 Some(address),
3944 false,
3945 true,
3946 );
3947 let server = tokio::spawn(async move {
3948 axum::serve(listener, app).await.unwrap();
3949 });
3950 let client = release_http_client().unwrap();
3951 let url = format!("http://{address}");
3952
3953 let repo_id = live_release_repo_id(&repo, &url, &client, "catalog-token")
3954 .await
3955 .unwrap();
3956
3957 server.abort();
3958 let _ = server.await;
3959 assert_eq!(repo_id.as_deref(), Some("repo"));
3960 }
3961
3962 #[test]
3963 fn empty_token_files_are_ignored() {
3964 let tmp = tempfile::tempdir().unwrap();
3965 let path = tmp.path().join("serve.token");
3966 std::fs::write(&path, "").unwrap();
3967 assert_eq!(read_token_file(&path), None);
3968 }
3969
3970 #[cfg(unix)]
3971 #[test]
3972 fn serve_token_file_is_written_with_owner_only_permissions() {
3973 use std::os::unix::fs::PermissionsExt;
3974 let tmp = tempfile::tempdir().unwrap();
3975 let repo = tmp.path().to_path_buf();
3976 let path = write_serve_token(&repo, "secret").unwrap();
3977 let mode = std::fs::metadata(&path).unwrap().permissions().mode();
3978 assert_eq!(mode & 0o777, 0o600);
3979 }
3980
3981 #[cfg(unix)]
3982 #[test]
3983 fn serve_read_token_file_is_written_with_owner_only_permissions() {
3984 use std::os::unix::fs::PermissionsExt;
3985 let tmp = tempfile::tempdir().unwrap();
3986 let repo = tmp.path().to_path_buf();
3987 let path = write_serve_read_token(&repo, "read-secret").unwrap();
3988 assert!(path.ends_with(".kranz/serve.read.token"));
3989 let mode = std::fs::metadata(&path).unwrap().permissions().mode();
3990 assert_eq!(mode & 0o777, 0o600);
3991 }
3992
3993 #[test]
3994 fn operator_read_token_path_sits_beside_the_operator_token() {
3995 let global = Path::new("/home/op/.kranz/config.json");
3996 let address = std::net::SocketAddr::from(([127, 0, 0, 1], 4560));
3997 let mutation = operator_serve_token_path(global, address);
3998 let read = operator_serve_read_token_path(global, address);
3999 assert_eq!(read, mutation.with_extension("read.token"));
4000 assert!(read.to_string_lossy().ends_with(".read.token"));
4001 }
4002
4003 #[cfg(unix)]
4004 #[test]
4005 fn existing_token_permissions_are_hardened_before_replacement() {
4006 use std::os::unix::fs::PermissionsExt;
4007 let tmp = tempfile::tempdir().unwrap();
4008 let path = tmp.path().join("serve.token");
4009 std::fs::write(&path, "old-token").unwrap();
4010 std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o644)).unwrap();
4011
4012 write_token_file(&path, "new-token").unwrap();
4013
4014 assert_eq!(std::fs::read_to_string(&path).unwrap(), "new-token");
4015 let mode = std::fs::metadata(&path).unwrap().permissions().mode();
4016 assert_eq!(mode & 0o777, 0o600);
4017 }
4018
4019 #[tokio::test]
4020 async fn serve_token_file_is_removed_after_graceful_shutdown() {
4021 let tmp = tempfile::tempdir().unwrap();
4022 let repo = tmp.path().to_path_buf();
4023 let path = repo.join(".kranz").join("serve.token");
4024 let read_path = repo.join(".kranz").join("serve.read.token");
4025
4026 let host = std::sync::Arc::new(kranz_server::MissionHost::new(repo.clone()));
4027 let listener =
4028 kranz_server::bind_listener(std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST), 0)
4029 .await
4030 .unwrap();
4031 serve_with_token_cleanup(
4032 &repo,
4033 host,
4034 listener,
4035 None,
4036 "tok".to_string(),
4037 "read-tok".to_string(),
4038 std::future::ready(()),
4039 )
4040 .await
4041 .unwrap();
4042
4043 assert!(!path.exists());
4044 assert!(!read_path.exists());
4045 }
4046
4047 fn dashboard_at(path: PathBuf) -> PathBuf {
4048 fs::create_dir_all(&path).unwrap();
4049 fs::write(path.join("index.html"), "<!doctype html>").unwrap();
4050 path
4051 }
4052
4053 fn inputs() -> DashboardResolutionInputs {
4054 DashboardResolutionInputs {
4055 embedded_available: false,
4056 ..Default::default()
4057 }
4058 }
4059
4060 #[test]
4061 fn dashboard_resolution_honors_explicit_path_verbatim() {
4062 let tmp = tempfile::tempdir().unwrap();
4063 let repo = tmp.path().join("repo");
4064 let explicit = tmp.path().join("missing-dashboard");
4065
4066 assert_eq!(
4067 resolve_dashboard_assets_from(&repo, Some(explicit.clone()), &inputs()),
4068 Some(DashboardAssets::Dir(explicit))
4069 );
4070 }
4071
4072 #[test]
4073 fn dashboard_resolution_env_precedes_repo_and_invalid_env_is_skipped() {
4074 let tmp = tempfile::tempdir().unwrap();
4075 let repo = tmp.path().join("repo");
4076 let repo_dist = dashboard_at(repo.join("apps").join("dashboard").join("dist"));
4077 let env_dist = dashboard_at(tmp.path().join("env-dist"));
4078
4079 let mut with_env = inputs();
4080 with_env.env_dist = Some(env_dist.clone());
4081 assert_eq!(
4082 resolve_dashboard_assets_from(&repo, None, &with_env),
4083 Some(DashboardAssets::Dir(env_dist))
4084 );
4085
4086 let mut with_invalid_env = inputs();
4087 with_invalid_env.env_dist = Some(tmp.path().join("missing-env-dist"));
4088 assert_eq!(
4089 resolve_dashboard_assets_from(&repo, None, &with_invalid_env),
4090 Some(DashboardAssets::Dir(repo_dist))
4091 );
4092 }
4093
4094 #[test]
4095 fn dashboard_resolution_finds_installed_asset_dirs() {
4096 let tmp = tempfile::tempdir().unwrap();
4097 let repo = tmp.path().join("repo");
4098 let exe = tmp.path().join("prefix").join("bin").join("kranz");
4099 let installed = dashboard_at(
4100 tmp.path()
4101 .join("prefix")
4102 .join("share")
4103 .join("kranz")
4104 .join("dashboard")
4105 .join("dist"),
4106 );
4107
4108 let mut inputs = inputs();
4109 inputs.exe = Some(exe);
4110 assert_eq!(
4111 resolve_dashboard_assets_from(&repo, None, &inputs),
4112 Some(DashboardAssets::Dir(installed))
4113 );
4114 }
4115
4116 #[test]
4117 fn dashboard_resolution_finds_checkout_used_to_build_installed_binary() {
4118 let tmp = tempfile::tempdir().unwrap();
4119 let repo = tmp.path().join("mission-repo");
4120 let checkout = tmp.path().join("kranz");
4121 let manifest_dir = checkout.join("crates").join("cli");
4122 let checkout_dist = dashboard_at(checkout.join("apps").join("dashboard").join("dist"));
4123
4124 let mut inputs = inputs();
4125 inputs.manifest_dir = Some(manifest_dir);
4126 inputs.exe = Some(tmp.path().join("cargo-home").join("bin").join("kranz"));
4127 assert_eq!(
4128 resolve_dashboard_assets_from(&repo, None, &inputs),
4129 Some(DashboardAssets::Dir(checkout_dist))
4130 );
4131 }
4132
4133 #[test]
4134 fn dashboard_resolution_falls_back_to_embedded_assets() {
4135 let tmp = tempfile::tempdir().unwrap();
4136 let repo = tmp.path().join("repo");
4137 let mut inputs = inputs();
4138 inputs.embedded_available = true;
4139
4140 assert_eq!(
4141 resolve_dashboard_assets_from(&repo, None, &inputs),
4142 Some(DashboardAssets::Embedded)
4143 );
4144 }
4145
4146 #[test]
4147 fn embedded_dashboard_bundle_contains_index() {
4148 assert!(
4149 crate::embedded_dashboard::EMBEDDED_DASHBOARD
4150 .iter()
4151 .any(|file| file.path == "index.html"),
4152 "embedded dashboard source: {}",
4153 crate::embedded_dashboard::EMBEDDED_DASHBOARD_SOURCE
4154 );
4155 }
4156
4157 #[test]
4158 fn serve_refuses_non_loopback_without_insecure_lan() {
4159 let bind: std::net::IpAddr = "0.0.0.0".parse().unwrap();
4160 let err = refuse_non_loopback_without_insecure_lan(bind, false).unwrap_err();
4161 let msg = err.to_string();
4162 assert!(
4163 msg.contains("refusing to bind") && msg.contains("--insecure-lan"),
4164 "unexpected error: {msg}"
4165 );
4166 }
4167
4168 #[test]
4169 fn serve_allows_non_loopback_with_insecure_lan() {
4170 let bind: std::net::IpAddr = "0.0.0.0".parse().unwrap();
4171 refuse_non_loopback_without_insecure_lan(bind, true).unwrap();
4172 }
4173
4174 #[test]
4175 fn serve_allows_loopback_without_insecure_lan() {
4176 let bind: std::net::IpAddr = "127.0.0.1".parse().unwrap();
4177 refuse_non_loopback_without_insecure_lan(bind, false).unwrap();
4178 }
4179
4180 #[test]
4181 fn read_auth_on_loopback_requires_read_token() {
4182 assert!(effective_require_read_token(true, true));
4183 }
4184
4185 #[test]
4186 fn read_auth_off_loopback_bind_does_not_require_read_token() {
4187 assert!(!effective_require_read_token(true, false));
4188 }
4189
4190 #[test]
4191 fn read_auth_non_loopback_always_requires_read_token() {
4192 assert!(effective_require_read_token(false, true));
4193 assert!(effective_require_read_token(false, false));
4194 }
4195
4196 fn reconcile_turn(reply: &str) -> Vec<kranz_engine::backend::AgentEvent> {
4201 vec![
4202 kranz_engine::backend_mock::mock_text(reply),
4203 kranz_engine::backend_mock::mock_result_text(reply),
4204 ]
4205 }
4206
4207 fn reconcile_worker_pass() -> kranz_engine::backend_mock::MockScript {
4208 kranz_engine::backend_mock::MockScript::single_shot_json(&serde_json::json!({
4209 "result": "pass",
4210 "summary": "implemented and tested",
4211 "filesTouched": ["delivered.txt"],
4212 "testsAdded": [],
4213 "testEvidence": "all green",
4214 "commits": []
4215 }))
4216 .writes_file("delivered.txt", "delivered by the mock worker\n")
4217 }
4218
4219 fn reconcile_plan_json() -> serde_json::Value {
4220 serde_json::json!({
4221 "goal": "ship the demo",
4222 "validationContract": [],
4223 "milestones": [{
4224 "title": "M1",
4225 "features": [{
4226 "title": "F1",
4227 "spec": "build the thing",
4228 "validationCriteria": ["it works"]
4229 }]
4230 }]
4231 })
4232 }
4233
4234 #[tokio::test]
4243 async fn reconcile_on_terminal_after_cli_run_marks_ticket_done() {
4244 let tmp = tempfile::tempdir().unwrap();
4245 let repo = tmp.path().to_path_buf();
4246 let status = std::process::Command::new("git")
4247 .args(["init", "-b", "main"])
4248 .current_dir(&repo)
4249 .status()
4250 .unwrap();
4251 assert!(status.success());
4252 std::process::Command::new("git")
4253 .args(["config", "user.name", "test"])
4254 .current_dir(&repo)
4255 .status()
4256 .unwrap();
4257 std::process::Command::new("git")
4258 .args(["config", "user.email", "test@example.com"])
4259 .current_dir(&repo)
4260 .status()
4261 .unwrap();
4262 std::fs::write(repo.join("README.md"), "seed\n").unwrap();
4263 std::process::Command::new("git")
4264 .args(["add", "-A"])
4265 .current_dir(&repo)
4266 .status()
4267 .unwrap();
4268 std::process::Command::new("git")
4269 .args(["commit", "-m", "seed"])
4270 .current_dir(&repo)
4271 .status()
4272 .unwrap();
4273 let repo = std::fs::canonicalize(&repo).unwrap();
4274
4275 let judgement = serde_json::json!({
4276 "decision": "complete",
4277 "guidance": "",
4278 "summary": "worker did the job"
4279 });
4280 let orch_setup = kranz_engine::backend_mock::MockScript::streaming(vec![
4284 kranz_engine::backend_mock::mock_init("orch-session"),
4285 kranz_engine::backend_mock::mock_result_text("seed-hi"),
4286 ])
4287 .responding(vec![
4288 reconcile_turn("let's scope the demo"),
4289 reconcile_turn(&reconcile_plan_json().to_string()),
4290 ]);
4291 let orch_run = kranz_engine::backend_mock::MockScript::streaming(vec![
4292 kranz_engine::backend_mock::mock_init("orch-session-2"),
4293 kranz_engine::backend_mock::mock_result_text("ack"),
4294 ])
4295 .responding(vec![
4296 reconcile_turn("ack"),
4297 reconcile_turn(
4298 &serde_json::json!({"action": "commit-as-is", "note": "worker delivered files"})
4299 .to_string(),
4300 ),
4301 reconcile_turn(&judgement.to_string()),
4302 reconcile_turn("NONE"),
4303 reconcile_turn("NONE"),
4307 reconcile_turn("NONE"),
4308 reconcile_turn("NONE"),
4309 ]);
4310 let backend: Arc<dyn AgentBackend> =
4311 Arc::new(kranz_engine::backend_mock::MockBackend::with_scripts(vec![
4312 orch_setup,
4313 kranz_engine::backend_mock::MockScript::single_shot("ok"),
4318 reconcile_worker_pass(),
4321 orch_run,
4322 ]));
4323
4324 let cfg = MissionConfig {
4325 skip_scrutiny: true,
4326 skip_functional: true,
4327 ..Default::default()
4328 };
4329 let mut engine =
4330 MissionEngine::create(Arc::clone(&backend), repo.clone(), "ship the demo", cfg)
4331 .unwrap();
4332 let mission_id = engine.mission_id().to_string();
4333 engine.planning_turn("ship the demo").await.unwrap();
4334 let request = engine.request_plan().await.unwrap();
4335 let plan = match request {
4336 PlanRequest::Ready(plan) => plan,
4337 PlanRequest::NotReady(text) => panic!("expected a ready plan, got: {text}"),
4338 PlanRequest::WrongPlan { reason } => {
4339 panic!("expected a ready plan, got a wrong-plan escalation: {reason}")
4340 }
4341 };
4342 engine.approve_plan(plan).unwrap();
4343 drop(engine);
4344
4345 kranz_engine::ticket::Ticket::record_mission(&repo, "my-ticket", &mission_id).unwrap();
4348 kranz_engine::ticket::Ticket::write_state(
4349 &repo,
4350 "my-ticket",
4351 kranz_engine::ticket::TicketState::Failed,
4352 None,
4353 )
4354 .unwrap();
4355 assert_eq!(
4356 kranz_engine::ticket::Ticket::read_state(&repo, "my-ticket"),
4357 kranz_engine::ticket::TicketState::Failed
4358 );
4359
4360 let exit_code =
4361 run_mission_loop_with_backend(repo.clone(), mission_id, LockForce::No, false, backend)
4362 .await
4363 .unwrap();
4364 assert_eq!(exit_code, 0);
4365
4366 assert_eq!(
4367 kranz_engine::ticket::Ticket::read_state(&repo, "my-ticket"),
4368 kranz_engine::ticket::TicketState::Done,
4369 "run_mission_loop must reconcile the linked ticket to Done on Complete"
4370 );
4371 }
4372}