1use serde::{Deserialize, Serialize};
26
27pub const STANDING_WORK_STATUSES: [&str; 4] = ["active", "paused", "superseded", "revoked"];
32
33pub const ONE_SHOT_WORK_STATUSES: [&str; 9] = [
35 "dispatched",
36 "accepted",
37 "running",
38 "paused",
39 "succeeded",
40 "failed",
41 "timed_out",
42 "canceled",
43 "expired",
44];
45
46pub const AGENT_REPORTABLE_WORK_STATUSES: [&str; 3] = ["running", "succeeded", "failed"];
52
53pub const ONE_SHOT_TERMINAL_STATUSES: [&str; 5] =
55 ["succeeded", "failed", "timed_out", "canceled", "expired"];
56
57#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
59#[jumo(kind = "state", domain = "Control", module = "Control.Agent.Work")]
60#[serde(rename_all = "PascalCase")]
61pub enum WorkKind {
62 Standing,
64 OneShot,
66}
67
68#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
74#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
75#[serde(deny_unknown_fields)]
76pub struct StandingWork {
77 pub work_id: String,
78 pub agent_id: String,
79 pub family: String,
81 pub spec: String,
83 pub catalog_version: i64,
85 #[serde(default)]
87 pub proposal_id: Option<String>,
88 pub plan_version: i64,
90 pub effective_from: String,
91 pub status: String,
93 pub updated_by: String,
94 pub updated_at: String,
95}
96
97#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
99#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
100#[serde(deny_unknown_fields)]
101pub struct OneShotWork {
102 pub work_id: String,
103 pub agent_id: String,
104 pub action: String,
106 pub spec: String,
107 pub scheduled_at: String,
109 pub deadline_at: String,
111 pub timeout_seconds: i64,
113 pub interruptible: bool,
115 pub status: String,
117 #[serde(default)]
119 pub paused_at: Option<String>,
120 pub paused_total_seconds: i64,
122 #[serde(default)]
124 pub current_step: Option<String>,
125 #[serde(default)]
126 pub completed_steps: Vec<String>,
127 pub attempt: i64,
129 pub issued_by: String,
130 pub issued_at: String,
131}
132
133impl OneShotWork {
134 pub fn is_outstanding(&self) -> bool {
136 !ONE_SHOT_TERMINAL_STATUSES.contains(&self.status.as_str())
137 }
138
139 pub fn can_pause(&self) -> bool {
144 self.interruptible && matches!(self.status.as_str(), "accepted" | "running")
145 }
146}
147
148#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
150#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
151#[serde(deny_unknown_fields)]
152pub struct WorkReceipt {
153 pub work_id: String,
154 pub agent_id: String,
155 pub work_kind: WorkKind,
156 pub status: String,
158 pub plan_version: i64,
159 pub created_at: String,
160}
161
162#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
179#[serde(deny_unknown_fields)]
180pub struct WorkSpec {
181 #[serde(default)]
183 pub units: Vec<WorkSpecUnit>,
184}
185
186#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192#[serde(deny_unknown_fields)]
193pub struct WorkSpecUnit {
194 pub unit_id: String,
195 pub capability: String,
197 #[serde(default)]
199 pub rule_ref: String,
200 #[serde(default)]
203 pub requires_privilege: String,
204 #[serde(default)]
205 pub sources: Vec<WorkSpecSource>,
206}
207
208#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
213#[serde(deny_unknown_fields)]
214pub struct WorkSpecSource {
215 pub kind: String,
216 pub target: String,
217 #[serde(default = "default_multiline")]
222 pub multiline: String,
223}
224
225fn default_multiline() -> String {
226 "none".to_string()
227}
228
229impl WorkSpec {
230 pub fn parse(spec: &str) -> Result<Self, serde_json::Error> {
235 serde_json::from_str(spec)
236 }
237
238 pub fn encode(&self) -> Result<String, serde_json::Error> {
241 serde_json::to_string(self)
242 }
243
244 pub fn metric_interval_seconds(&self) -> Option<i64> {
249 self.units
250 .iter()
251 .flat_map(|unit| unit.sources.iter())
252 .filter(|source| source.kind == "MetricInterval")
253 .filter_map(|source| parse_interval_seconds(&source.target))
254 .min()
255 }
256
257 pub fn collects_logs(&self) -> bool {
259 self.units
260 .iter()
261 .any(|unit| unit.capability == "collect_logs")
262 }
263}
264
265pub const EXECUTABLE_SOURCE_KINDS: &[&str] = &["FileGlob", "MetricInterval", "Exporter"];
271
272pub fn has_glob_meta(target: &str) -> bool {
274 target.contains(['*', '?', '['])
275}
276
277pub fn is_explicit_path(target: &str) -> bool {
283 target.starts_with('/') && !has_glob_meta(target)
284}
285
286pub const EXPORTER_IDS: &[&str] = &[
294 "journalctl-unit",
295 "journalctl-shutdown",
296 "last-reboot",
297 "nft-ruleset",
298 "iptables-save",
299 "smartctl",
300 "dmesg",
301 "auditd-execve",
302];
303
304pub fn parse_exporter_target(target: &str) -> Option<(&str, Option<&str>)> {
309 let (id, arg) = match target.split_once(':') {
310 Some((id, arg)) => (id, Some(arg)),
311 None => (target, None),
312 };
313 let id = id.trim();
314 if id.is_empty() || id.contains(char::is_whitespace) {
315 return None;
316 }
317 Some((id, arg))
318}
319
320pub fn is_known_exporter(target: &str) -> bool {
322 match parse_exporter_target(target) {
323 Some((id, _)) => EXPORTER_IDS.contains(&id),
324 None => false,
325 }
326}
327
328pub fn is_executable_source(kind: &str, target: &str) -> bool {
334 match kind {
335 "FileGlob" => is_explicit_path(target),
337 "MetricInterval" => true,
339 "Exporter" => is_known_exporter(target),
341 _ => false,
342 }
343}
344
345impl WorkSpecUnit {
346 pub fn file_sources(&self) -> Vec<&WorkSpecSource> {
351 self.sources
352 .iter()
353 .filter(|source| source.kind == "FileGlob")
354 .collect()
355 }
356
357 pub fn unsupported_sources(&self) -> Vec<&WorkSpecSource> {
359 self.sources
360 .iter()
361 .filter(|source| !is_executable_source(&source.kind, &source.target))
362 .collect()
363 }
364
365 pub fn has_executable_source(&self) -> bool {
370 self.sources
371 .iter()
372 .any(|source| is_executable_source(&source.kind, &source.target))
373 }
374}
375
376pub fn parse_interval_seconds(raw: &str) -> Option<i64> {
381 let raw = raw.trim();
382 if let Some(value) = raw.strip_suffix('s') {
383 return value.trim().parse::<i64>().ok().filter(|value| *value > 0);
384 }
385 if let Some(value) = raw.strip_suffix('m') {
386 let minutes = value.trim().parse::<i64>().ok()?;
387 return minutes.checked_mul(60).filter(|value| *value > 0);
391 }
392 None
393}
394
395#[cfg(test)]
396mod tests {
397 use super::*;
398
399 fn one_shot(status: &str, interruptible: bool) -> OneShotWork {
400 OneShotWork {
401 work_id: "work-1".to_string(),
402 agent_id: "agent-1".to_string(),
403 action: "upgrade".to_string(),
404 spec: "0.1.4".to_string(),
405 scheduled_at: "2026-09-23T00:00:00Z".to_string(),
406 deadline_at: "2026-09-24T00:00:00Z".to_string(),
407 timeout_seconds: 600,
408 interruptible,
409 status: status.to_string(),
410 paused_at: None,
411 paused_total_seconds: 0,
412 current_step: None,
413 completed_steps: vec![],
414 attempt: 0,
415 issued_by: "admin".to_string(),
416 issued_at: "2026-09-23T00:00:00Z".to_string(),
417 }
418 }
419
420 #[test]
421 fn terminal_one_shot_work_is_not_outstanding() {
422 for status in ONE_SHOT_TERMINAL_STATUSES {
424 assert!(!one_shot(status, true).is_outstanding(), "{status}");
425 }
426 for status in ["dispatched", "accepted", "running", "paused"] {
427 assert!(one_shot(status, true).is_outstanding(), "{status}");
428 }
429 }
430
431 #[test]
432 fn only_interruptible_running_work_can_pause() {
433 assert!(!one_shot("running", false).can_pause());
435 assert!(!one_shot("dispatched", true).can_pause());
437 assert!(!one_shot("succeeded", true).can_pause());
438 assert!(one_shot("accepted", true).can_pause());
439 assert!(one_shot("running", true).can_pause());
440 }
441
442 #[test]
443 fn work_kind_serializes_as_the_model_names() {
444 assert_eq!(
445 serde_json::to_string(&WorkKind::Standing).expect("serialize"),
446 "\"Standing\""
447 );
448 assert_eq!(
449 serde_json::to_string(&WorkKind::OneShot).expect("serialize"),
450 "\"OneShot\""
451 );
452 }
453
454 #[test]
455 fn only_explicit_absolute_paths_are_collectable_file_targets() {
456 assert!(is_explicit_path("/var/log/install.log"));
459 assert!(!is_explicit_path("/var/log/wifi.log*"));
460 assert!(!is_explicit_path("/Library/Logs/DiagnosticReports/*.ips"));
461 assert!(!is_explicit_path("~/Library/Logs/Homebrew/*"));
462 assert!(!is_explicit_path("relative/app.log"));
463 }
464
465 #[test]
466 fn executability_needs_both_a_supported_kind_and_a_supported_target() {
467 assert!(is_executable_source("MetricInterval", "15s"));
469 assert!(is_executable_source("FileGlob", "/var/log/app.log"));
471 assert!(!is_executable_source("FileGlob", "/var/log/app*"));
472 assert!(is_executable_source("Exporter", "smartctl"));
474 assert!(is_executable_source("Exporter", "dmesg:panic"));
475 assert!(!is_executable_source("Exporter", "last,lastb"));
476 assert!(!is_executable_source("UnifiedLogPredicate", "syspolicyd"));
478 }
479
480 #[test]
481 fn exporter_targets_split_into_a_known_id_and_an_optional_arg() {
482 assert_eq!(parse_exporter_target("smartctl"), Some(("smartctl", None)));
483 assert_eq!(
484 parse_exporter_target("dmesg:panic"),
485 Some(("dmesg", Some("panic")))
486 );
487 assert_eq!(parse_exporter_target(":panic"), None);
489 assert_eq!(parse_exporter_target(" dm esg"), None);
490 assert!(is_known_exporter("dmesg:nvidia-xid"));
492 assert!(!is_known_exporter("nope"));
493 }
494
495 fn spec() -> WorkSpec {
498 WorkSpec {
499 units: vec![
500 WorkSpecUnit {
501 unit_id: "mac-host-metrics".to_string(),
502 capability: "collect_metrics".to_string(),
503 rule_ref: "agent_uplink".to_string(),
504 requires_privilege: "none".to_string(),
505 sources: vec![WorkSpecSource {
506 kind: "MetricInterval".to_string(),
507 target: "15s".to_string(),
508 multiline: "none".to_string(),
509 }],
510 },
511 WorkSpecUnit {
512 unit_id: "mac-privacy-tcc".to_string(),
513 capability: "collect_logs".to_string(),
514 rule_ref: "macos/tcc".to_string(),
515 requires_privilege: "fda".to_string(),
516 sources: vec![
517 WorkSpecSource {
518 kind: "Exporter".to_string(),
519 target: "sqlite-snapshot(TCC.db)".to_string(),
520 multiline: "none".to_string(),
521 },
522 WorkSpecSource {
523 kind: "FileGlob".to_string(),
524 target: "/var/log/tccd.log".to_string(),
526 multiline: "indented".to_string(),
528 },
529 ],
530 },
531 ],
532 }
533 }
534
535 #[test]
536 fn spec_round_trips_and_reports_what_agentd_can_do() {
537 let encoded = spec().encode().expect("encode");
538 let decoded = WorkSpec::parse(&encoded).expect("parse");
539 assert_eq!(decoded, spec());
540
541 assert!(decoded.collects_logs());
542 assert_eq!(decoded.metric_interval_seconds(), Some(15));
544 let files = decoded.units[1].file_sources();
546 assert_eq!(files.len(), 1);
547 assert_eq!(files[0].target, "/var/log/tccd.log");
548 assert_eq!(files[0].multiline, "indented");
549 assert_eq!(decoded.units[1].unsupported_sources().len(), 1);
550 assert_eq!(decoded.units[1].unsupported_sources()[0].kind, "Exporter");
551 }
552
553 #[test]
554 fn a_glob_file_source_counts_as_unsupported() {
555 let unit = WorkSpecUnit {
558 unit_id: "mac-wifi".to_string(),
559 capability: "collect_logs".to_string(),
560 rule_ref: String::new(),
561 requires_privilege: "root".to_string(),
562 sources: vec![WorkSpecSource {
563 kind: "FileGlob".to_string(),
564 target: "/var/log/wifi.log*".to_string(),
565 multiline: "none".to_string(),
566 }],
567 };
568 assert_eq!(unit.file_sources().len(), 1, "它仍是一条文件来源");
569 assert_eq!(unit.unsupported_sources().len(), 1, "但今天采不了");
570 assert!(!unit.has_executable_source());
571 }
572
573 #[test]
574 fn a_source_without_multiline_decodes_as_single_line() {
575 let json = r#"{"units":[{"unit_id":"u","capability":"collect_logs","sources":[{"kind":"FileGlob","target":"/a/*"}]}]}"#;
578 let spec = WorkSpec::parse(json).expect("parse");
579 assert_eq!(spec.units[0].sources[0].multiline, "none");
580 assert_eq!(
582 spec.encode().expect("encode"),
583 r#"{"units":[{"unit_id":"u","capability":"collect_logs","rule_ref":"","requires_privilege":"","sources":[{"kind":"FileGlob","target":"/a/*","multiline":"none"}]}]}"#
584 );
585 }
586
587 #[test]
588 fn a_broken_spec_is_an_error_not_an_empty_work() {
589 assert!(WorkSpec::parse("mac-host-metrics").is_err());
591 assert!(WorkSpec::parse("{").is_err());
592 assert!(WorkSpec::parse(r#"{"units":[]}"#).unwrap().units.is_empty());
594 assert_eq!(
595 WorkSpec::parse(r#"{"units":[]}"#)
596 .unwrap()
597 .metric_interval_seconds(),
598 None
599 );
600 }
601
602 #[test]
603 fn interval_parsing_accepts_seconds_and_minutes_only() {
604 assert_eq!(parse_interval_seconds("15s"), Some(15));
605 assert_eq!(parse_interval_seconds(" 5s "), Some(5));
606 assert_eq!(parse_interval_seconds("5m"), Some(300));
607 assert_eq!(parse_interval_seconds("15"), None);
609 assert_eq!(parse_interval_seconds("0s"), None);
610 assert_eq!(parse_interval_seconds(""), None);
611 }
612
613 #[test]
614 fn interval_parsing_never_overflows_on_huge_values() {
615 assert_eq!(parse_interval_seconds("922337203685477580m"), None);
618 assert_eq!(parse_interval_seconds("153722867280912931m"), None);
619 assert_eq!(parse_interval_seconds("99999999999999999999s"), None);
621 assert_eq!(parse_interval_seconds("600m"), Some(36_000));
623 }
624
625 #[test]
626 fn metric_interval_ignores_unparseable_targets_instead_of_treating_them_as_zero() {
627 let unit = |target: &str| WorkSpecUnit {
629 unit_id: "u".to_string(),
630 capability: "collect_metrics".to_string(),
631 rule_ref: String::new(),
632 requires_privilege: String::new(),
633 sources: vec![WorkSpecSource {
634 kind: "MetricInterval".to_string(),
635 target: target.to_string(),
636 multiline: "none".to_string(),
637 }],
638 };
639 let mixed = WorkSpec {
640 units: vec![unit("garbage"), unit("30s")],
641 };
642 assert_eq!(mixed.metric_interval_seconds(), Some(30));
643
644 let all_bad = WorkSpec {
645 units: vec![unit("nope")],
646 };
647 assert_eq!(all_bad.metric_interval_seconds(), None);
648 }
649}