1use serde::{Deserialize, Serialize};
26
27pub const POLL_WORK_KIND: &str = "poll_work";
29pub const ACK_WORK_KIND: &str = "ack_work";
31pub const REPORT_WORK_RESULT_KIND: &str = "report_work_result";
33
34pub const STANDING_WORK_STATUSES: [&str; 4] = ["active", "paused", "superseded", "revoked"];
39
40pub const ONE_SHOT_WORK_STATUSES: [&str; 9] = [
42 "dispatched",
43 "accepted",
44 "running",
45 "paused",
46 "succeeded",
47 "failed",
48 "timed_out",
49 "canceled",
50 "expired",
51];
52
53pub const AGENT_REPORTABLE_WORK_STATUSES: [&str; 3] = ["running", "succeeded", "failed"];
59
60pub const ONE_SHOT_TERMINAL_STATUSES: [&str; 5] =
62 ["succeeded", "failed", "timed_out", "canceled", "expired"];
63
64#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
66#[jumo(kind = "state", domain = "Control", module = "Control.Agent.Work")]
67#[serde(rename_all = "PascalCase")]
68pub enum WorkKind {
69 Standing,
71 OneShot,
73}
74
75#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
81#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
82#[serde(deny_unknown_fields)]
83pub struct StandingWork {
84 pub work_id: String,
85 pub agent_id: String,
86 pub family: String,
88 pub spec: String,
90 pub catalog_version: i64,
92 #[serde(default)]
94 pub proposal_id: Option<String>,
95 pub plan_version: i64,
97 pub effective_from: String,
98 pub status: String,
100 pub updated_by: String,
101 pub updated_at: String,
102}
103
104#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
106#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
107#[serde(deny_unknown_fields)]
108pub struct OneShotWork {
109 pub work_id: String,
110 pub agent_id: String,
111 pub action: String,
113 pub spec: String,
114 pub scheduled_at: String,
116 pub deadline_at: String,
118 pub timeout_seconds: i64,
120 pub interruptible: bool,
122 pub status: String,
124 #[serde(default)]
126 pub paused_at: Option<String>,
127 pub paused_total_seconds: i64,
129 #[serde(default)]
131 pub current_step: Option<String>,
132 #[serde(default)]
133 pub completed_steps: Vec<String>,
134 pub attempt: i64,
136 pub issued_by: String,
137 pub issued_at: String,
138}
139
140impl OneShotWork {
141 pub fn is_outstanding(&self) -> bool {
143 !ONE_SHOT_TERMINAL_STATUSES.contains(&self.status.as_str())
144 }
145
146 pub fn can_pause(&self) -> bool {
151 self.interruptible && matches!(self.status.as_str(), "accepted" | "running")
152 }
153}
154
155#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
160#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
161#[serde(deny_unknown_fields)]
162pub struct WorkGrant {
163 pub agent_id: String,
164 #[serde(default)]
166 pub standing: Vec<StandingWork>,
167 #[serde(default)]
168 pub one_shot: Vec<OneShotWork>,
169 pub sequence: i64,
171 pub granted_at: String,
172}
173
174#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
176#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
177#[serde(deny_unknown_fields)]
178pub struct WorkReceipt {
179 pub work_id: String,
180 pub agent_id: String,
181 pub work_kind: WorkKind,
182 pub status: String,
184 pub plan_version: i64,
185 pub created_at: String,
186}
187
188#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
193#[jumo(
194 kind = "message",
195 role = "command",
196 domain = "Control",
197 module = "Control.AgentApp.FacingInterface"
198)]
199#[serde(deny_unknown_fields)]
200pub struct PollWork {
201 pub api_version: String,
202 pub kind: String,
203 pub agent_id: String,
204 pub instance_id: String,
205 pub last_seen_sequence: i64,
206 pub wait_ms: i64,
207 pub requested_at: String,
208}
209
210#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
215#[jumo(
216 kind = "message",
217 role = "command",
218 domain = "Control",
219 module = "Control.AgentApp.FacingInterface"
220)]
221#[serde(deny_unknown_fields)]
222pub struct AckWork {
223 pub api_version: String,
224 pub kind: String,
225 pub agent_id: String,
226 pub instance_id: String,
227 pub work_id: String,
228 pub plan_version: i64,
229 pub acknowledged_at: String,
230}
231
232#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
237#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
238#[serde(deny_unknown_fields)]
239pub struct WorkAccepted {
240 pub work_id: String,
241 pub status: String,
243 pub accepted_at: String,
244}
245
246#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
255#[jumo(
256 kind = "message",
257 role = "command",
258 domain = "Control",
259 module = "Control.AgentApp.FacingInterface"
260)]
261#[serde(deny_unknown_fields)]
262pub struct ReportWorkResult {
263 pub api_version: String,
264 pub kind: String,
265 pub agent_id: String,
266 pub instance_id: String,
267 pub work_id: String,
268 pub status: String,
270 #[serde(default)]
273 pub detail: String,
274 pub reported_at: String,
275}
276
277#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
279#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
280#[serde(deny_unknown_fields)]
281pub struct WorkResultAccepted {
282 pub work_id: String,
283 pub status: String,
288 pub accepted_at: String,
289}
290
291#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
308#[serde(deny_unknown_fields)]
309pub struct WorkSpec {
310 #[serde(default)]
312 pub units: Vec<WorkSpecUnit>,
313}
314
315#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
321#[serde(deny_unknown_fields)]
322pub struct WorkSpecUnit {
323 pub unit_id: String,
324 pub capability: String,
326 #[serde(default)]
328 pub rule_ref: String,
329 #[serde(default)]
332 pub requires_privilege: String,
333 #[serde(default)]
334 pub sources: Vec<WorkSpecSource>,
335}
336
337#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
342#[serde(deny_unknown_fields)]
343pub struct WorkSpecSource {
344 pub kind: String,
345 pub target: String,
346 #[serde(default = "default_multiline")]
351 pub multiline: String,
352}
353
354fn default_multiline() -> String {
355 "none".to_string()
356}
357
358impl WorkSpec {
359 pub fn parse(spec: &str) -> Result<Self, serde_json::Error> {
364 serde_json::from_str(spec)
365 }
366
367 pub fn encode(&self) -> Result<String, serde_json::Error> {
370 serde_json::to_string(self)
371 }
372
373 pub fn metric_interval_seconds(&self) -> Option<i64> {
378 self.units
379 .iter()
380 .flat_map(|unit| unit.sources.iter())
381 .filter(|source| source.kind == "MetricInterval")
382 .filter_map(|source| parse_interval_seconds(&source.target))
383 .min()
384 }
385
386 pub fn collects_logs(&self) -> bool {
388 self.units
389 .iter()
390 .any(|unit| unit.capability == "collect_logs")
391 }
392}
393
394pub const EXECUTABLE_SOURCE_KINDS: &[&str] = &["FileGlob", "MetricInterval", "Exporter"];
400
401pub fn has_glob_meta(target: &str) -> bool {
403 target.contains(['*', '?', '['])
404}
405
406pub fn is_explicit_path(target: &str) -> bool {
412 target.starts_with('/') && !has_glob_meta(target)
413}
414
415pub const EXPORTER_IDS: &[&str] = &[
423 "journalctl-unit",
424 "journalctl-shutdown",
425 "last-reboot",
426 "nft-ruleset",
427 "iptables-save",
428 "smartctl",
429 "dmesg",
430 "auditd-execve",
431];
432
433pub fn parse_exporter_target(target: &str) -> Option<(&str, Option<&str>)> {
438 let (id, arg) = match target.split_once(':') {
439 Some((id, arg)) => (id, Some(arg)),
440 None => (target, None),
441 };
442 let id = id.trim();
443 if id.is_empty() || id.contains(char::is_whitespace) {
444 return None;
445 }
446 Some((id, arg))
447}
448
449pub fn is_known_exporter(target: &str) -> bool {
451 match parse_exporter_target(target) {
452 Some((id, _)) => EXPORTER_IDS.contains(&id),
453 None => false,
454 }
455}
456
457pub fn is_executable_source(kind: &str, target: &str) -> bool {
463 match kind {
464 "FileGlob" => is_explicit_path(target),
466 "MetricInterval" => true,
468 "Exporter" => is_known_exporter(target),
470 _ => false,
471 }
472}
473
474impl WorkSpecUnit {
475 pub fn file_sources(&self) -> Vec<&WorkSpecSource> {
480 self.sources
481 .iter()
482 .filter(|source| source.kind == "FileGlob")
483 .collect()
484 }
485
486 pub fn unsupported_sources(&self) -> Vec<&WorkSpecSource> {
488 self.sources
489 .iter()
490 .filter(|source| !is_executable_source(&source.kind, &source.target))
491 .collect()
492 }
493
494 pub fn has_executable_source(&self) -> bool {
499 self.sources
500 .iter()
501 .any(|source| is_executable_source(&source.kind, &source.target))
502 }
503}
504
505pub fn parse_interval_seconds(raw: &str) -> Option<i64> {
510 let raw = raw.trim();
511 if let Some(value) = raw.strip_suffix('s') {
512 return value.trim().parse::<i64>().ok().filter(|value| *value > 0);
513 }
514 if let Some(value) = raw.strip_suffix('m') {
515 let minutes = value.trim().parse::<i64>().ok()?;
516 return minutes.checked_mul(60).filter(|value| *value > 0);
520 }
521 None
522}
523
524#[cfg(test)]
525mod tests {
526 use super::*;
527
528 fn standing(work_id: &str, family: &str, plan_version: i64) -> StandingWork {
529 StandingWork {
530 work_id: work_id.to_string(),
531 agent_id: "agent-1".to_string(),
532 family: family.to_string(),
533 spec: "unit-a,unit-b".to_string(),
534 catalog_version: 1,
535 proposal_id: None,
536 plan_version,
537 effective_from: "2026-09-23T00:00:00Z".to_string(),
538 status: "active".to_string(),
539 updated_by: "admin".to_string(),
540 updated_at: "2026-09-23T00:00:00Z".to_string(),
541 }
542 }
543
544 fn one_shot(status: &str, interruptible: bool) -> OneShotWork {
545 OneShotWork {
546 work_id: "work-1".to_string(),
547 agent_id: "agent-1".to_string(),
548 action: "upgrade".to_string(),
549 spec: "0.1.4".to_string(),
550 scheduled_at: "2026-09-23T00:00:00Z".to_string(),
551 deadline_at: "2026-09-24T00:00:00Z".to_string(),
552 timeout_seconds: 600,
553 interruptible,
554 status: status.to_string(),
555 paused_at: None,
556 paused_total_seconds: 0,
557 current_step: None,
558 completed_steps: vec![],
559 attempt: 0,
560 issued_by: "admin".to_string(),
561 issued_at: "2026-09-23T00:00:00Z".to_string(),
562 }
563 }
564
565 #[test]
566 fn terminal_one_shot_work_is_not_outstanding() {
567 for status in ONE_SHOT_TERMINAL_STATUSES {
569 assert!(!one_shot(status, true).is_outstanding(), "{status}");
570 }
571 for status in ["dispatched", "accepted", "running", "paused"] {
572 assert!(one_shot(status, true).is_outstanding(), "{status}");
573 }
574 }
575
576 #[test]
577 fn only_interruptible_running_work_can_pause() {
578 assert!(!one_shot("running", false).can_pause());
580 assert!(!one_shot("dispatched", true).can_pause());
582 assert!(!one_shot("succeeded", true).can_pause());
583 assert!(one_shot("accepted", true).can_pause());
584 assert!(one_shot("running", true).can_pause());
585 }
586
587 #[test]
588 fn work_kind_serializes_as_the_model_names() {
589 assert_eq!(
590 serde_json::to_string(&WorkKind::Standing).expect("serialize"),
591 "\"Standing\""
592 );
593 assert_eq!(
594 serde_json::to_string(&WorkKind::OneShot).expect("serialize"),
595 "\"OneShot\""
596 );
597 }
598
599 #[test]
600 fn only_explicit_absolute_paths_are_collectable_file_targets() {
601 assert!(is_explicit_path("/var/log/install.log"));
604 assert!(!is_explicit_path("/var/log/wifi.log*"));
605 assert!(!is_explicit_path("/Library/Logs/DiagnosticReports/*.ips"));
606 assert!(!is_explicit_path("~/Library/Logs/Homebrew/*"));
607 assert!(!is_explicit_path("relative/app.log"));
608 }
609
610 #[test]
611 fn executability_needs_both_a_supported_kind_and_a_supported_target() {
612 assert!(is_executable_source("MetricInterval", "15s"));
614 assert!(is_executable_source("FileGlob", "/var/log/app.log"));
616 assert!(!is_executable_source("FileGlob", "/var/log/app*"));
617 assert!(is_executable_source("Exporter", "smartctl"));
619 assert!(is_executable_source("Exporter", "dmesg:panic"));
620 assert!(!is_executable_source("Exporter", "last,lastb"));
621 assert!(!is_executable_source("UnifiedLogPredicate", "syspolicyd"));
623 }
624
625 #[test]
626 fn exporter_targets_split_into_a_known_id_and_an_optional_arg() {
627 assert_eq!(parse_exporter_target("smartctl"), Some(("smartctl", None)));
628 assert_eq!(
629 parse_exporter_target("dmesg:panic"),
630 Some(("dmesg", Some("panic")))
631 );
632 assert_eq!(parse_exporter_target(":panic"), None);
634 assert_eq!(parse_exporter_target(" dm esg"), None);
635 assert!(is_known_exporter("dmesg:nvidia-xid"));
637 assert!(!is_known_exporter("nope"));
638 }
639
640 #[test]
641 fn grant_round_trips_with_serde() {
642 let grant = WorkGrant {
643 agent_id: "agent-1".to_string(),
644 standing: vec![standing("work-a", "LoginSession", 2)],
645 one_shot: vec![one_shot("running", true)],
646 sequence: 7,
647 granted_at: "2026-09-23T00:00:00Z".to_string(),
648 };
649 let json = serde_json::to_string(&grant).expect("serialize");
650 let decoded: WorkGrant = serde_json::from_str(&json).expect("deserialize");
651 assert_eq!(decoded, grant);
652 }
653
654 #[test]
655 fn grant_omits_absent_optional_fields_and_still_decodes() {
656 let json = r#"{"agent_id":"a","standing":[],"one_shot":[],"sequence":0,
659 "granted_at":"t"}"#;
660 let grant: WorkGrant = serde_json::from_str(json).expect("deserialize");
661 assert!(grant.standing.is_empty());
662 assert!(grant.one_shot.is_empty());
663 }
664
665 #[test]
666 fn rejects_unknown_fields() {
667 let json = r#"{"agent_id":"a","standing":[],"one_shot":[],"sequence":0,
670 "granted_at":"t","extra":1}"#;
671 assert!(serde_json::from_str::<WorkGrant>(json).is_err());
672 }
673
674 fn spec() -> WorkSpec {
677 WorkSpec {
678 units: vec![
679 WorkSpecUnit {
680 unit_id: "mac-host-metrics".to_string(),
681 capability: "collect_metrics".to_string(),
682 rule_ref: "agent_uplink".to_string(),
683 requires_privilege: "none".to_string(),
684 sources: vec![WorkSpecSource {
685 kind: "MetricInterval".to_string(),
686 target: "15s".to_string(),
687 multiline: "none".to_string(),
688 }],
689 },
690 WorkSpecUnit {
691 unit_id: "mac-privacy-tcc".to_string(),
692 capability: "collect_logs".to_string(),
693 rule_ref: "macos/tcc".to_string(),
694 requires_privilege: "fda".to_string(),
695 sources: vec![
696 WorkSpecSource {
697 kind: "Exporter".to_string(),
698 target: "sqlite-snapshot(TCC.db)".to_string(),
699 multiline: "none".to_string(),
700 },
701 WorkSpecSource {
702 kind: "FileGlob".to_string(),
703 target: "/var/log/tccd.log".to_string(),
705 multiline: "indented".to_string(),
707 },
708 ],
709 },
710 ],
711 }
712 }
713
714 #[test]
715 fn spec_round_trips_and_reports_what_agentd_can_do() {
716 let encoded = spec().encode().expect("encode");
717 let decoded = WorkSpec::parse(&encoded).expect("parse");
718 assert_eq!(decoded, spec());
719
720 assert!(decoded.collects_logs());
721 assert_eq!(decoded.metric_interval_seconds(), Some(15));
723 let files = decoded.units[1].file_sources();
725 assert_eq!(files.len(), 1);
726 assert_eq!(files[0].target, "/var/log/tccd.log");
727 assert_eq!(files[0].multiline, "indented");
728 assert_eq!(decoded.units[1].unsupported_sources().len(), 1);
729 assert_eq!(decoded.units[1].unsupported_sources()[0].kind, "Exporter");
730 }
731
732 #[test]
733 fn a_glob_file_source_counts_as_unsupported() {
734 let unit = WorkSpecUnit {
737 unit_id: "mac-wifi".to_string(),
738 capability: "collect_logs".to_string(),
739 rule_ref: String::new(),
740 requires_privilege: "root".to_string(),
741 sources: vec![WorkSpecSource {
742 kind: "FileGlob".to_string(),
743 target: "/var/log/wifi.log*".to_string(),
744 multiline: "none".to_string(),
745 }],
746 };
747 assert_eq!(unit.file_sources().len(), 1, "它仍是一条文件来源");
748 assert_eq!(unit.unsupported_sources().len(), 1, "但今天采不了");
749 assert!(!unit.has_executable_source());
750 }
751
752 #[test]
753 fn a_source_without_multiline_decodes_as_single_line() {
754 let json = r#"{"units":[{"unit_id":"u","capability":"collect_logs","sources":[{"kind":"FileGlob","target":"/a/*"}]}]}"#;
757 let spec = WorkSpec::parse(json).expect("parse");
758 assert_eq!(spec.units[0].sources[0].multiline, "none");
759 assert_eq!(
761 spec.encode().expect("encode"),
762 r#"{"units":[{"unit_id":"u","capability":"collect_logs","rule_ref":"","requires_privilege":"","sources":[{"kind":"FileGlob","target":"/a/*","multiline":"none"}]}]}"#
763 );
764 }
765
766 #[test]
767 fn a_broken_spec_is_an_error_not_an_empty_work() {
768 assert!(WorkSpec::parse("mac-host-metrics").is_err());
770 assert!(WorkSpec::parse("{").is_err());
771 assert!(WorkSpec::parse(r#"{"units":[]}"#).unwrap().units.is_empty());
773 assert_eq!(
774 WorkSpec::parse(r#"{"units":[]}"#)
775 .unwrap()
776 .metric_interval_seconds(),
777 None
778 );
779 }
780
781 #[test]
782 fn interval_parsing_accepts_seconds_and_minutes_only() {
783 assert_eq!(parse_interval_seconds("15s"), Some(15));
784 assert_eq!(parse_interval_seconds(" 5s "), Some(5));
785 assert_eq!(parse_interval_seconds("5m"), Some(300));
786 assert_eq!(parse_interval_seconds("15"), None);
788 assert_eq!(parse_interval_seconds("0s"), None);
789 assert_eq!(parse_interval_seconds(""), None);
790 }
791
792 #[test]
793 fn interval_parsing_never_overflows_on_huge_values() {
794 assert_eq!(parse_interval_seconds("922337203685477580m"), None);
797 assert_eq!(parse_interval_seconds("153722867280912931m"), None);
798 assert_eq!(parse_interval_seconds("99999999999999999999s"), None);
800 assert_eq!(parse_interval_seconds("600m"), Some(36_000));
802 }
803
804 #[test]
805 fn metric_interval_ignores_unparseable_targets_instead_of_treating_them_as_zero() {
806 let unit = |target: &str| WorkSpecUnit {
808 unit_id: "u".to_string(),
809 capability: "collect_metrics".to_string(),
810 rule_ref: String::new(),
811 requires_privilege: String::new(),
812 sources: vec![WorkSpecSource {
813 kind: "MetricInterval".to_string(),
814 target: target.to_string(),
815 multiline: "none".to_string(),
816 }],
817 };
818 let mixed = WorkSpec {
819 units: vec![unit("garbage"), unit("30s")],
820 };
821 assert_eq!(mixed.metric_interval_seconds(), Some(30));
822
823 let all_bad = WorkSpec {
824 units: vec![unit("nope")],
825 };
826 assert_eq!(all_bad.metric_interval_seconds(), None);
827 }
828}