1use serde::{Deserialize, Serialize};
4
5use crate::deps::ProgramDep;
6use crate::parallel::ParallelConfig;
7
8pub const CONCIERGE_AGENT: &str = "mur";
9
10pub const RUN_ID_ENV: &str = "MUR_RUN_ID";
15
16#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
17pub struct Fleet {
18 pub name: String,
19 #[serde(default)]
20 pub display_name: String,
21 #[serde(default)]
22 pub goal: String,
23 #[serde(default, skip_serializing_if = "Option::is_none")]
24 pub router: Option<String>,
25 #[serde(default, skip_serializing_if = "Option::is_none")]
29 pub team_id: Option<String>,
30 #[serde(default)]
31 pub members: Vec<String>,
32 pub channel_id: String,
33 #[serde(default, skip_serializing_if = "Vec::is_empty")]
34 pub rules: Vec<String>,
35 #[serde(default, skip_serializing_if = "Vec::is_empty")]
36 pub skills: Vec<String>,
37 #[serde(default, rename = "loop", skip_serializing_if = "Option::is_none")]
38 pub loop_cfg: Option<FleetLoop>,
39 #[serde(default, skip_serializing_if = "Option::is_none")]
40 pub parallel: Option<ParallelConfig>,
41 #[serde(default, skip_serializing_if = "Option::is_none")]
44 pub hitl: Option<FleetHitl>,
45 #[serde(default, skip_serializing_if = "Vec::is_empty")]
48 pub requires_programs: Vec<ProgramDep>,
49 #[serde(default, skip_serializing_if = "Option::is_none")]
53 pub limits: Option<crate::limits::Limits>,
54 #[serde(default, skip_serializing_if = "Vec::is_empty")]
60 pub needs: Vec<String>,
61 #[serde(default, skip_serializing_if = "Vec::is_empty")]
67 pub procedure: Vec<FleetStep>,
68}
69
70pub const FLEET_PROCEDURE_MAX_STEPS: usize = 32;
72
73#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
75pub struct FleetStep {
76 pub id: String,
79 pub member: String,
81 #[serde(default)]
83 pub task: String,
84 #[serde(default, skip_serializing_if = "Vec::is_empty")]
86 pub depends_on: Vec<String>,
87}
88
89impl Fleet {
90 pub fn validate_procedure(&self) -> Result<(), String> {
94 if self.procedure.is_empty() {
95 return Ok(());
96 }
97 if self.parallel.is_some() {
98 return Err(
99 "fleet sets both `procedure:` and `parallel:` — pick one. `parallel:` would win \
100 and the procedure would be silently ignored."
101 .into(),
102 );
103 }
104 if self.procedure.len() > FLEET_PROCEDURE_MAX_STEPS {
105 return Err(format!(
106 "procedure has {} steps; the limit is {FLEET_PROCEDURE_MAX_STEPS}",
107 self.procedure.len()
108 ));
109 }
110 let mut ids = std::collections::HashSet::new();
111 for s in &self.procedure {
112 if s.id.trim().is_empty() {
113 return Err(format!("procedure step for `{}` has an empty id", s.member));
114 }
115 if !ids.insert(s.id.as_str()) {
116 return Err(format!("procedure step id `{}` appears twice", s.id));
117 }
118 if !self.members.iter().any(|m| m == &s.member) {
119 return Err(format!(
120 "procedure step `{}` names `{}`, which is not a fleet member ({})",
121 s.id,
122 s.member,
123 self.members.join(", ")
124 ));
125 }
126 }
127 for s in &self.procedure {
128 if let Some(d) = s.depends_on.iter().find(|d| !ids.contains(d.as_str())) {
129 return Err(format!(
130 "procedure step `{}` depends on `{d}`, which is not a step id",
131 s.id
132 ));
133 }
134 }
135 Ok(())
136 }
137}
138
139#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
144pub struct FleetHitl {
145 #[serde(default, skip_serializing_if = "Option::is_none")]
150 pub mode: Option<crate::hitl::Unanswered>,
151 #[serde(default, skip_serializing_if = "Vec::is_empty")]
160 pub auto_approve_tiers: Vec<crate::hitl::RiskTier>,
161}
162
163impl FleetHitl {
164 pub fn validate(&self) -> Result<(), String> {
170 let bad: Vec<String> = self
171 .auto_approve_tiers
172 .iter()
173 .filter(|t| !crate::hitl::tier_may_be_granted(**t))
174 .map(|t| format!("{t:?}").to_lowercase())
175 .collect();
176 if bad.is_empty() {
177 return Ok(());
178 }
179 Err(format!(
180 "hitl.auto_approve_tiers may not include {} — standing approval is capped at `write`. \
181 Those actions have to be approved per-run (`mur channel approve`), because their cost \
182 cannot be undone by noticing afterwards.",
183 bad.join(", ")
184 ))
185 }
186}
187
188#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
189pub struct FleetLoop {
190 #[serde(default = "default_trigger")]
191 pub trigger: String,
192 #[serde(default)]
195 pub max_iterations: u32,
196 #[serde(default)]
197 pub budget_usd: f64,
198 #[serde(default)]
199 pub deadline: String,
200 #[serde(default)]
201 pub done_when: String,
202}
203
204fn default_trigger() -> String {
205 "manual".to_string()
206}
207
208impl Fleet {
209 pub fn limits_or_legacy(&self) -> Option<crate::limits::Limits> {
214 if let Some(l) = &self.limits {
215 return Some(l.clone());
216 }
217 let lc = self.loop_cfg.as_ref()?;
218 let deadline = (!lc.deadline.trim().is_empty()).then(|| lc.deadline.trim().to_string());
219 let cost_usd = (lc.budget_usd > 0.0).then_some(lc.budget_usd);
220 if deadline.is_none() && cost_usd.is_none() {
221 return None;
222 }
223 Some(crate::limits::Limits {
224 deadline,
225 stuck: None,
226 cost_usd,
227 })
228 }
229
230 pub fn router_or_concierge(&self) -> &str {
231 self.router.as_deref().unwrap_or(CONCIERGE_AGENT)
232 }
233}
234
235pub fn valid_fleet_name(name: &str) -> bool {
238 !name.is_empty()
239 && name.len() <= 64
240 && name
241 .chars()
242 .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-' || c == '_')
243}
244
245pub const CHANNEL_PREFIX: &str = "fleet-";
247
248pub fn fleet_name_from_channel_id(channel_id: &str) -> Option<&str> {
255 let name = channel_id.strip_prefix(CHANNEL_PREFIX)?;
256 valid_fleet_name(name).then_some(name)
257}
258
259#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
265#[serde(rename_all = "lowercase")]
266pub enum JobStatus {
267 Queued,
268 Running,
269 Blocked,
270 Done,
271 Failed,
272 Canceled,
273}
274
275impl JobStatus {
276 pub fn is_terminal(&self) -> bool {
278 matches!(
279 self,
280 JobStatus::Done | JobStatus::Failed | JobStatus::Canceled
281 )
282 }
283
284 pub fn as_str(&self) -> &'static str {
288 match self {
289 JobStatus::Queued => "queued",
290 JobStatus::Running => "running",
291 JobStatus::Blocked => "blocked",
292 JobStatus::Done => "done",
293 JobStatus::Failed => "failed",
294 JobStatus::Canceled => "canceled",
295 }
296 }
297}
298
299impl std::fmt::Display for JobStatus {
300 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
301 f.write_str(self.as_str())
302 }
303}
304
305#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
308pub struct Job {
309 pub id: String,
310 pub text: String,
311 pub source: String,
313 pub status: JobStatus,
314 pub created_at: String,
316 #[serde(default, skip_serializing_if = "Option::is_none")]
317 pub started_at: Option<String>,
318 #[serde(default, skip_serializing_if = "Option::is_none")]
319 pub finished_at: Option<String>,
320 #[serde(default, skip_serializing_if = "Option::is_none")]
322 pub run_id: Option<String>,
323 #[serde(default, skip_serializing_if = "Option::is_none")]
324 pub result: Option<String>,
325 #[serde(default, skip_serializing_if = "Option::is_none")]
326 pub error: Option<String>,
327}
328
329#[cfg(test)]
330mod tests {
331 use super::*;
332
333 #[test]
334 fn env_name_is_stable() {
335 assert_eq!(RUN_ID_ENV, "MUR_RUN_ID");
339 }
340
341 #[test]
342 fn valid_fleet_name_accepts_and_rejects() {
343 assert!(valid_fleet_name("dev"));
345 assert!(valid_fleet_name("dev-team"));
346 assert!(valid_fleet_name("dev_1"));
347 assert!(valid_fleet_name("ab12"));
348 assert!(!valid_fleet_name("")); assert!(!valid_fleet_name("../x")); assert!(!valid_fleet_name("a/b")); assert!(!valid_fleet_name("a\\b")); assert!(!valid_fleet_name("Dev")); assert!(!valid_fleet_name("a b")); assert!(!valid_fleet_name(".hidden")); }
357
358 #[test]
359 fn fleet_name_from_channel_id_extracts_and_validates() {
360 assert_eq!(fleet_name_from_channel_id("fleet-dev"), Some("dev"));
362 assert_eq!(
363 fleet_name_from_channel_id("fleet-my-squad"),
364 Some("my-squad")
365 );
366 assert_eq!(fleet_name_from_channel_id("fleet-ab12"), Some("ab12"));
367 assert_eq!(fleet_name_from_channel_id("dev"), None);
369 assert_eq!(fleet_name_from_channel_id("agent:foo:uuid"), None);
370 assert_eq!(fleet_name_from_channel_id("fleet-"), None); assert_eq!(fleet_name_from_channel_id("fleet-../etc"), None); assert_eq!(fleet_name_from_channel_id("fleet-a/b"), None); assert_eq!(fleet_name_from_channel_id("fleet-Dev"), None); }
376
377 #[test]
378 fn fleet_minimal_yaml_deserializes_with_defaults() {
379 let f: Fleet = serde_yaml::from_str("name: dev\nchannel_id: fleet-dev\n").unwrap();
380 assert_eq!(f.name, "dev");
381 assert_eq!(f.channel_id, "fleet-dev");
382 assert!(f.members.is_empty());
383 assert_eq!(f.router_or_concierge(), CONCIERGE_AGENT);
384 assert!(f.loop_cfg.is_none());
385 }
386
387 #[test]
391 fn fleet_hitl_mode_parses_and_defaults_to_absent() {
392 let bare: Fleet = serde_yaml::from_str("name: dev\nchannel_id: fleet-dev\n").unwrap();
393 assert!(bare.hitl.is_none(), "no policy stated ⇒ auto");
394
395 let declared: Fleet =
396 serde_yaml::from_str("name: dev\nchannel_id: fleet-dev\nhitl:\n mode: deny\n")
397 .unwrap();
398 assert_eq!(
399 declared.hitl.as_ref().and_then(|h| h.mode),
400 Some(crate::hitl::Unanswered::Deny)
401 );
402
403 for (yaml, want) in [
404 ("defer", crate::hitl::Unanswered::Defer),
405 ("wait", crate::hitl::Unanswered::Wait),
406 ("deny", crate::hitl::Unanswered::Deny),
407 ] {
408 let f: Fleet = serde_yaml::from_str(&format!(
409 "name: dev\nchannel_id: fleet-dev\nhitl:\n mode: {yaml}\n"
410 ))
411 .unwrap();
412 assert_eq!(f.hitl.and_then(|h| h.mode), Some(want));
413 }
414
415 assert!(
417 serde_yaml::from_str::<Fleet>(
418 "name: dev\nchannel_id: fleet-dev\nhitl:\n mode: yolo\n"
419 )
420 .is_err()
421 );
422 }
423
424 #[test]
429 fn auto_approve_tiers_parse_and_validate_against_the_ceiling() {
430 let f: Fleet = serde_yaml::from_str(
431 "name: dev\nchannel_id: fleet-dev\nhitl:\n auto_approve_tiers: [read, write]\n",
432 )
433 .unwrap();
434 let hitl = f.hitl.as_ref().unwrap();
435 assert_eq!(
436 hitl.auto_approve_tiers,
437 vec![crate::hitl::RiskTier::Read, crate::hitl::RiskTier::Write]
438 );
439 assert!(hitl.validate().is_ok());
440
441 let f: Fleet = serde_yaml::from_str(
442 "name: dev\nchannel_id: fleet-dev\nhitl:\n auto_approve_tiers: [write, destructive]\n",
443 )
444 .unwrap();
445 let err = f.hitl.as_ref().unwrap().validate().unwrap_err();
446 assert!(err.contains("destructive"), "{err}");
447 assert!(err.contains("write"), "the ceiling must be named: {err}");
448 }
449
450 #[test]
451 fn fleet_yaml_roundtrip_and_router_default() {
452 let f = Fleet {
453 name: "dev".into(),
454 display_name: "Dev Team".into(),
455 goal: "ship it".into(),
456 router: None,
457 team_id: None,
458 members: vec!["pm".into(), "qa".into()],
459 channel_id: "fleet-dev".into(),
460 procedure: vec![],
461 rules: vec![],
462 skills: vec![],
463 loop_cfg: None,
464 parallel: None,
465 hitl: None,
466 requires_programs: vec![],
467 limits: None,
468 needs: vec![],
469 };
470 assert_eq!(f.router_or_concierge(), CONCIERGE_AGENT);
471 let yaml = serde_yaml::to_string(&f).unwrap();
472 let back: Fleet = serde_yaml::from_str(&yaml).unwrap();
473 assert_eq!(back, f);
474 let with_loop: Fleet = serde_yaml::from_str(
476 "name: dev\ndisplay_name: Dev\ngoal: test\nchannel_id: fleet-dev\nrules: []\nskills: []\nmembers: []\nloop:\n trigger: manual\n max_iterations: 3\n budget_usd: 1.0\n deadline: '2026-12-31'\n done_when: 'all_tasks_done'\n",
477 ).unwrap();
478 assert_eq!(with_loop.loop_cfg.unwrap().max_iterations, 3);
479 }
480
481 #[test]
482 fn minimal_loop_block_deserializes_with_defaults() {
483 let f: Fleet = serde_yaml::from_str(
486 "name: dev\nchannel_id: fleet-dev\nloop:\n trigger: \"interval:1h\"\n",
487 )
488 .unwrap();
489 let l = f.loop_cfg.unwrap();
490 assert_eq!(l.trigger, "interval:1h");
491 assert_eq!(l.max_iterations, 0);
492 assert_eq!(l.budget_usd, 0.0);
493 }
494
495 #[test]
496 fn job_status_serde_is_lowercase_and_terminal_predicate() {
497 assert_eq!(
498 serde_yaml::to_string(&JobStatus::Queued).unwrap().trim(),
499 "queued"
500 );
501 assert_eq!(
502 serde_yaml::to_string(&JobStatus::Done).unwrap().trim(),
503 "done"
504 );
505 assert!(!JobStatus::Queued.is_terminal());
506 assert!(!JobStatus::Running.is_terminal());
507 assert!(JobStatus::Done.is_terminal());
508 assert!(JobStatus::Failed.is_terminal());
509 assert!(JobStatus::Canceled.is_terminal());
510 }
511
512 #[test]
513 fn job_status_as_str_and_display_match_serde_for_all_variants() {
514 for s in [
515 JobStatus::Queued,
516 JobStatus::Running,
517 JobStatus::Done,
518 JobStatus::Failed,
519 JobStatus::Canceled,
520 ] {
521 let serde = serde_yaml::to_string(&s).unwrap();
522 assert_eq!(
523 serde.trim(),
524 s.as_str(),
525 "as_str must match serde for {s:?}"
526 );
527 assert_eq!(s.to_string(), s.as_str(), "Display must delegate to as_str");
528 }
529 }
530
531 #[test]
532 fn job_yaml_roundtrip_with_optional_fields_skipped() {
533 let j = Job {
534 id: "0190f3a2-0000-7000-8000-000000000000".into(),
535 text: "ship it".into(),
536 source: "cli".into(),
537 status: JobStatus::Queued,
538 created_at: "2026-06-24T00:00:00Z".into(),
539 started_at: None,
540 finished_at: None,
541 run_id: None,
542 result: None,
543 error: None,
544 };
545 let yaml = serde_yaml::to_string(&j).unwrap();
546 assert!(
547 !yaml.contains("started_at"),
548 "None optionals must be skipped: {yaml}"
549 );
550 let back: Job = serde_yaml::from_str(&yaml).unwrap();
551 assert_eq!(back, j);
552 }
553}
554
555#[cfg(test)]
556mod limits_tests {
557 use super::*;
558
559 fn fleet(loop_cfg: Option<FleetLoop>, limits: Option<crate::limits::Limits>) -> Fleet {
560 Fleet {
561 name: "dev".into(),
562 display_name: String::new(),
563 goal: "g".into(),
564 router: None,
565 team_id: None,
566 members: vec![],
567 channel_id: "fleet-dev".into(),
568 procedure: vec![],
569 rules: vec![],
570 skills: vec![],
571 loop_cfg,
572 parallel: None,
573 hitl: None,
574 requires_programs: vec![],
575 limits,
576 needs: vec![],
577 }
578 }
579
580 #[test]
584 fn legacy_loop_fields_are_read_as_limits_only_when_the_block_is_absent() {
585 let lc = FleetLoop {
586 trigger: "manual".into(),
587 max_iterations: 8,
588 budget_usd: 5.0,
589 deadline: "2h".into(),
590 done_when: String::new(),
591 };
592 let legacy = fleet(Some(lc.clone()), None)
593 .limits_or_legacy()
594 .expect("derived");
595 assert_eq!(legacy.deadline.as_deref(), Some("2h"));
596 assert_eq!(legacy.cost_usd, Some(5.0));
597 assert_eq!(legacy.stuck, None, "the old loop had no stuck setting");
598
599 let zero = FleetLoop {
601 budget_usd: 0.0,
602 deadline: String::new(),
603 ..lc.clone()
604 };
605 assert_eq!(fleet(Some(zero), None).limits_or_legacy(), None);
606
607 let explicit = fleet(Some(lc), Some(crate::limits::Limits::default())).limits_or_legacy();
609 assert_eq!(explicit, Some(crate::limits::Limits::default()));
610
611 let yaml = serde_yaml_ng::to_string(&fleet(None, None)).unwrap();
613 assert!(!yaml.contains("limits"), "{yaml}");
614 }
615
616 #[test]
619 fn needs_round_trip_and_stay_absent_when_empty() {
620 let mut f = fleet(None, None);
621 assert!(!serde_yaml_ng::to_string(&f).unwrap().contains("needs"));
622 f.needs = vec!["write_file".into(), "bash".into()];
623 let yaml = serde_yaml_ng::to_string(&f).unwrap();
624 assert!(yaml.contains("needs:"), "{yaml}");
625 let back: Fleet = serde_yaml_ng::from_str(&yaml).unwrap();
626 assert_eq!(
627 back.needs,
628 vec!["write_file".to_string(), "bash".to_string()]
629 );
630 }
631
632 fn proc_fleet(extra: &str) -> Fleet {
633 serde_yaml::from_str(&format!(
634 "name: council\nchannel_id: fleet-council\nmembers: [a, b]\n{extra}"
635 ))
636 .unwrap()
637 }
638
639 #[test]
640 fn procedure_absent_is_empty_and_not_serialized() {
641 let f = proc_fleet("");
642 assert!(f.procedure.is_empty());
643 assert!(f.validate_procedure().is_ok());
644 assert!(!serde_yaml::to_string(&f).unwrap().contains("procedure"));
645 }
646
647 #[test]
648 fn procedure_parses_and_keeps_dep_order() {
649 let f = proc_fleet(
650 "procedure:\n - {id: p1, member: a, task: one}\n - {id: p2, member: b}\n - {id: r, member: a, task: rev, depends_on: [p2, p1]}\n",
651 );
652 assert_eq!(f.procedure.len(), 3);
653 assert_eq!(f.procedure[1].task, "");
654 assert_eq!(f.procedure[2].depends_on, vec!["p2", "p1"]);
655 assert!(f.validate_procedure().is_ok());
656 }
657
658 #[test]
659 fn procedure_rejects_unknown_member() {
660 let f = proc_fleet("procedure:\n - {id: p1, member: ghost}\n");
661 let err = f.validate_procedure().unwrap_err();
662 assert!(err.contains("ghost"), "{err}");
663 }
664
665 #[test]
666 fn procedure_rejects_duplicate_and_empty_ids() {
667 let f = proc_fleet("procedure:\n - {id: p1, member: a}\n - {id: p1, member: b}\n");
668 assert!(f.validate_procedure().unwrap_err().contains("p1"));
669 let f = proc_fleet("procedure:\n - {id: ' ', member: a}\n");
670 assert!(f.validate_procedure().is_err());
671 }
672
673 #[test]
674 fn procedure_rejects_missing_dep() {
675 let f = proc_fleet("procedure:\n - {id: p1, member: a, depends_on: [nope]}\n");
676 let err = f.validate_procedure().unwrap_err();
677 assert!(err.contains("nope"), "{err}");
678 }
679
680 #[test]
681 fn procedure_rejects_oversize() {
682 let steps: String = (0..=FLEET_PROCEDURE_MAX_STEPS)
683 .map(|i| format!(" - {{id: s{i}, member: a}}\n"))
684 .collect();
685 let f = proc_fleet(&format!("procedure:\n{steps}"));
686 assert!(f.validate_procedure().is_err());
687 }
688
689 #[test]
690 fn procedure_and_parallel_are_exclusive() {
691 let mut f = proc_fleet("procedure:\n - {id: p1, member: a}\n");
692 f.parallel = Some(serde_yaml::from_str("judge:\n model: m\ntracks: []\n").unwrap());
693 let err = f.validate_procedure().unwrap_err();
694 assert!(err.contains("parallel"), "{err}");
695 }
696}