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