use serde::{Deserialize, Serialize};
pub const POLL_WORK_KIND: &str = "poll_work";
pub const ACK_WORK_KIND: &str = "ack_work";
pub const REPORT_WORK_RESULT_KIND: &str = "report_work_result";
pub const STANDING_WORK_STATUSES: [&str; 4] = ["active", "paused", "superseded", "revoked"];
pub const ONE_SHOT_WORK_STATUSES: [&str; 9] = [
"dispatched",
"accepted",
"running",
"paused",
"succeeded",
"failed",
"timed_out",
"canceled",
"expired",
];
pub const AGENT_REPORTABLE_WORK_STATUSES: [&str; 3] = ["running", "succeeded", "failed"];
pub const ONE_SHOT_TERMINAL_STATUSES: [&str; 5] =
["succeeded", "failed", "timed_out", "canceled", "expired"];
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(kind = "state", domain = "Control", module = "Control.Agent.Work")]
#[serde(rename_all = "PascalCase")]
pub enum WorkKind {
Standing,
OneShot,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
#[serde(deny_unknown_fields)]
pub struct StandingWork {
pub work_id: String,
pub agent_id: String,
pub family: String,
pub spec: String,
pub catalog_version: i64,
#[serde(default)]
pub proposal_id: Option<String>,
pub plan_version: i64,
pub effective_from: String,
pub status: String,
pub updated_by: String,
pub updated_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
#[serde(deny_unknown_fields)]
pub struct OneShotWork {
pub work_id: String,
pub agent_id: String,
pub action: String,
pub spec: String,
pub scheduled_at: String,
pub deadline_at: String,
pub timeout_seconds: i64,
pub interruptible: bool,
pub status: String,
#[serde(default)]
pub paused_at: Option<String>,
pub paused_total_seconds: i64,
#[serde(default)]
pub current_step: Option<String>,
#[serde(default)]
pub completed_steps: Vec<String>,
pub attempt: i64,
pub issued_by: String,
pub issued_at: String,
}
impl OneShotWork {
pub fn is_outstanding(&self) -> bool {
!ONE_SHOT_TERMINAL_STATUSES.contains(&self.status.as_str())
}
pub fn can_pause(&self) -> bool {
self.interruptible && matches!(self.status.as_str(), "accepted" | "running")
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
#[serde(deny_unknown_fields)]
pub struct WorkGrant {
pub agent_id: String,
#[serde(default)]
pub standing: Vec<StandingWork>,
#[serde(default)]
pub one_shot: Vec<OneShotWork>,
pub sequence: i64,
pub granted_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
#[serde(deny_unknown_fields)]
pub struct WorkReceipt {
pub work_id: String,
pub agent_id: String,
pub work_kind: WorkKind,
pub status: String,
pub plan_version: i64,
pub created_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(
kind = "message",
role = "command",
domain = "Control",
module = "Control.AgentApp.FacingInterface"
)]
#[serde(deny_unknown_fields)]
pub struct PollWork {
pub api_version: String,
pub kind: String,
pub agent_id: String,
pub instance_id: String,
pub last_seen_sequence: i64,
pub wait_ms: i64,
pub requested_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(
kind = "message",
role = "command",
domain = "Control",
module = "Control.AgentApp.FacingInterface"
)]
#[serde(deny_unknown_fields)]
pub struct AckWork {
pub api_version: String,
pub kind: String,
pub agent_id: String,
pub instance_id: String,
pub work_id: String,
pub plan_version: i64,
pub acknowledged_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
#[serde(deny_unknown_fields)]
pub struct WorkAccepted {
pub work_id: String,
pub status: String,
pub accepted_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(
kind = "message",
role = "command",
domain = "Control",
module = "Control.AgentApp.FacingInterface"
)]
#[serde(deny_unknown_fields)]
pub struct ReportWorkResult {
pub api_version: String,
pub kind: String,
pub agent_id: String,
pub instance_id: String,
pub work_id: String,
pub status: String,
#[serde(default)]
pub detail: String,
pub reported_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
#[serde(deny_unknown_fields)]
pub struct WorkResultAccepted {
pub work_id: String,
pub status: String,
pub accepted_at: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct WorkSpec {
#[serde(default)]
pub units: Vec<WorkSpecUnit>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct WorkSpecUnit {
pub unit_id: String,
pub capability: String,
#[serde(default)]
pub rule_ref: String,
#[serde(default)]
pub requires_privilege: String,
#[serde(default)]
pub sources: Vec<WorkSpecSource>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct WorkSpecSource {
pub kind: String,
pub target: String,
#[serde(default = "default_multiline")]
pub multiline: String,
}
fn default_multiline() -> String {
"none".to_string()
}
impl WorkSpec {
pub fn parse(spec: &str) -> Result<Self, serde_json::Error> {
serde_json::from_str(spec)
}
pub fn encode(&self) -> Result<String, serde_json::Error> {
serde_json::to_string(self)
}
pub fn metric_interval_seconds(&self) -> Option<i64> {
self.units
.iter()
.flat_map(|unit| unit.sources.iter())
.filter(|source| source.kind == "MetricInterval")
.filter_map(|source| parse_interval_seconds(&source.target))
.min()
}
pub fn collects_logs(&self) -> bool {
self.units
.iter()
.any(|unit| unit.capability == "collect_logs")
}
}
pub const EXECUTABLE_SOURCE_KINDS: &[&str] = &["FileGlob", "MetricInterval", "Exporter"];
pub fn has_glob_meta(target: &str) -> bool {
target.contains(['*', '?', '['])
}
pub fn is_explicit_path(target: &str) -> bool {
target.starts_with('/') && !has_glob_meta(target)
}
pub const EXPORTER_IDS: &[&str] = &[
"journalctl-unit",
"journalctl-shutdown",
"last-reboot",
"nft-ruleset",
"iptables-save",
"smartctl",
"dmesg",
"auditd-execve",
];
pub fn parse_exporter_target(target: &str) -> Option<(&str, Option<&str>)> {
let (id, arg) = match target.split_once(':') {
Some((id, arg)) => (id, Some(arg)),
None => (target, None),
};
let id = id.trim();
if id.is_empty() || id.contains(char::is_whitespace) {
return None;
}
Some((id, arg))
}
pub fn is_known_exporter(target: &str) -> bool {
match parse_exporter_target(target) {
Some((id, _)) => EXPORTER_IDS.contains(&id),
None => false,
}
}
pub fn is_executable_source(kind: &str, target: &str) -> bool {
match kind {
"FileGlob" => is_explicit_path(target),
"MetricInterval" => true,
"Exporter" => is_known_exporter(target),
_ => false,
}
}
impl WorkSpecUnit {
pub fn file_sources(&self) -> Vec<&WorkSpecSource> {
self.sources
.iter()
.filter(|source| source.kind == "FileGlob")
.collect()
}
pub fn unsupported_sources(&self) -> Vec<&WorkSpecSource> {
self.sources
.iter()
.filter(|source| !is_executable_source(&source.kind, &source.target))
.collect()
}
pub fn has_executable_source(&self) -> bool {
self.sources
.iter()
.any(|source| is_executable_source(&source.kind, &source.target))
}
}
pub fn parse_interval_seconds(raw: &str) -> Option<i64> {
let raw = raw.trim();
if let Some(value) = raw.strip_suffix('s') {
return value.trim().parse::<i64>().ok().filter(|value| *value > 0);
}
if let Some(value) = raw.strip_suffix('m') {
let minutes = value.trim().parse::<i64>().ok()?;
return minutes.checked_mul(60).filter(|value| *value > 0);
}
None
}
#[cfg(test)]
mod tests {
use super::*;
fn standing(work_id: &str, family: &str, plan_version: i64) -> StandingWork {
StandingWork {
work_id: work_id.to_string(),
agent_id: "agent-1".to_string(),
family: family.to_string(),
spec: "unit-a,unit-b".to_string(),
catalog_version: 1,
proposal_id: None,
plan_version,
effective_from: "2026-09-23T00:00:00Z".to_string(),
status: "active".to_string(),
updated_by: "admin".to_string(),
updated_at: "2026-09-23T00:00:00Z".to_string(),
}
}
fn one_shot(status: &str, interruptible: bool) -> OneShotWork {
OneShotWork {
work_id: "work-1".to_string(),
agent_id: "agent-1".to_string(),
action: "upgrade".to_string(),
spec: "0.1.4".to_string(),
scheduled_at: "2026-09-23T00:00:00Z".to_string(),
deadline_at: "2026-09-24T00:00:00Z".to_string(),
timeout_seconds: 600,
interruptible,
status: status.to_string(),
paused_at: None,
paused_total_seconds: 0,
current_step: None,
completed_steps: vec![],
attempt: 0,
issued_by: "admin".to_string(),
issued_at: "2026-09-23T00:00:00Z".to_string(),
}
}
#[test]
fn terminal_one_shot_work_is_not_outstanding() {
for status in ONE_SHOT_TERMINAL_STATUSES {
assert!(!one_shot(status, true).is_outstanding(), "{status}");
}
for status in ["dispatched", "accepted", "running", "paused"] {
assert!(one_shot(status, true).is_outstanding(), "{status}");
}
}
#[test]
fn only_interruptible_running_work_can_pause() {
assert!(!one_shot("running", false).can_pause());
assert!(!one_shot("dispatched", true).can_pause());
assert!(!one_shot("succeeded", true).can_pause());
assert!(one_shot("accepted", true).can_pause());
assert!(one_shot("running", true).can_pause());
}
#[test]
fn work_kind_serializes_as_the_model_names() {
assert_eq!(
serde_json::to_string(&WorkKind::Standing).expect("serialize"),
"\"Standing\""
);
assert_eq!(
serde_json::to_string(&WorkKind::OneShot).expect("serialize"),
"\"OneShot\""
);
}
#[test]
fn only_explicit_absolute_paths_are_collectable_file_targets() {
assert!(is_explicit_path("/var/log/install.log"));
assert!(!is_explicit_path("/var/log/wifi.log*"));
assert!(!is_explicit_path("/Library/Logs/DiagnosticReports/*.ips"));
assert!(!is_explicit_path("~/Library/Logs/Homebrew/*"));
assert!(!is_explicit_path("relative/app.log"));
}
#[test]
fn executability_needs_both_a_supported_kind_and_a_supported_target() {
assert!(is_executable_source("MetricInterval", "15s"));
assert!(is_executable_source("FileGlob", "/var/log/app.log"));
assert!(!is_executable_source("FileGlob", "/var/log/app*"));
assert!(is_executable_source("Exporter", "smartctl"));
assert!(is_executable_source("Exporter", "dmesg:panic"));
assert!(!is_executable_source("Exporter", "last,lastb"));
assert!(!is_executable_source("UnifiedLogPredicate", "syspolicyd"));
}
#[test]
fn exporter_targets_split_into_a_known_id_and_an_optional_arg() {
assert_eq!(parse_exporter_target("smartctl"), Some(("smartctl", None)));
assert_eq!(
parse_exporter_target("dmesg:panic"),
Some(("dmesg", Some("panic")))
);
assert_eq!(parse_exporter_target(":panic"), None);
assert_eq!(parse_exporter_target(" dm esg"), None);
assert!(is_known_exporter("dmesg:nvidia-xid"));
assert!(!is_known_exporter("nope"));
}
#[test]
fn grant_round_trips_with_serde() {
let grant = WorkGrant {
agent_id: "agent-1".to_string(),
standing: vec![standing("work-a", "LoginSession", 2)],
one_shot: vec![one_shot("running", true)],
sequence: 7,
granted_at: "2026-09-23T00:00:00Z".to_string(),
};
let json = serde_json::to_string(&grant).expect("serialize");
let decoded: WorkGrant = serde_json::from_str(&json).expect("deserialize");
assert_eq!(decoded, grant);
}
#[test]
fn grant_omits_absent_optional_fields_and_still_decodes() {
let json = r#"{"agent_id":"a","standing":[],"one_shot":[],"sequence":0,
"granted_at":"t"}"#;
let grant: WorkGrant = serde_json::from_str(json).expect("deserialize");
assert!(grant.standing.is_empty());
assert!(grant.one_shot.is_empty());
}
#[test]
fn rejects_unknown_fields() {
let json = r#"{"agent_id":"a","standing":[],"one_shot":[],"sequence":0,
"granted_at":"t","extra":1}"#;
assert!(serde_json::from_str::<WorkGrant>(json).is_err());
}
fn spec() -> WorkSpec {
WorkSpec {
units: vec![
WorkSpecUnit {
unit_id: "mac-host-metrics".to_string(),
capability: "collect_metrics".to_string(),
rule_ref: "agent_uplink".to_string(),
requires_privilege: "none".to_string(),
sources: vec![WorkSpecSource {
kind: "MetricInterval".to_string(),
target: "15s".to_string(),
multiline: "none".to_string(),
}],
},
WorkSpecUnit {
unit_id: "mac-privacy-tcc".to_string(),
capability: "collect_logs".to_string(),
rule_ref: "macos/tcc".to_string(),
requires_privilege: "fda".to_string(),
sources: vec![
WorkSpecSource {
kind: "Exporter".to_string(),
target: "sqlite-snapshot(TCC.db)".to_string(),
multiline: "none".to_string(),
},
WorkSpecSource {
kind: "FileGlob".to_string(),
target: "/var/log/tccd.log".to_string(),
multiline: "indented".to_string(),
},
],
},
],
}
}
#[test]
fn spec_round_trips_and_reports_what_agentd_can_do() {
let encoded = spec().encode().expect("encode");
let decoded = WorkSpec::parse(&encoded).expect("parse");
assert_eq!(decoded, spec());
assert!(decoded.collects_logs());
assert_eq!(decoded.metric_interval_seconds(), Some(15));
let files = decoded.units[1].file_sources();
assert_eq!(files.len(), 1);
assert_eq!(files[0].target, "/var/log/tccd.log");
assert_eq!(files[0].multiline, "indented");
assert_eq!(decoded.units[1].unsupported_sources().len(), 1);
assert_eq!(decoded.units[1].unsupported_sources()[0].kind, "Exporter");
}
#[test]
fn a_glob_file_source_counts_as_unsupported() {
let unit = WorkSpecUnit {
unit_id: "mac-wifi".to_string(),
capability: "collect_logs".to_string(),
rule_ref: String::new(),
requires_privilege: "root".to_string(),
sources: vec![WorkSpecSource {
kind: "FileGlob".to_string(),
target: "/var/log/wifi.log*".to_string(),
multiline: "none".to_string(),
}],
};
assert_eq!(unit.file_sources().len(), 1, "它仍是一条文件来源");
assert_eq!(unit.unsupported_sources().len(), 1, "但今天采不了");
assert!(!unit.has_executable_source());
}
#[test]
fn a_source_without_multiline_decodes_as_single_line() {
let json = r#"{"units":[{"unit_id":"u","capability":"collect_logs","sources":[{"kind":"FileGlob","target":"/a/*"}]}]}"#;
let spec = WorkSpec::parse(json).expect("parse");
assert_eq!(spec.units[0].sources[0].multiline, "none");
assert_eq!(
spec.encode().expect("encode"),
r#"{"units":[{"unit_id":"u","capability":"collect_logs","rule_ref":"","requires_privilege":"","sources":[{"kind":"FileGlob","target":"/a/*","multiline":"none"}]}]}"#
);
}
#[test]
fn a_broken_spec_is_an_error_not_an_empty_work() {
assert!(WorkSpec::parse("mac-host-metrics").is_err());
assert!(WorkSpec::parse("{").is_err());
assert!(WorkSpec::parse(r#"{"units":[]}"#).unwrap().units.is_empty());
assert_eq!(
WorkSpec::parse(r#"{"units":[]}"#)
.unwrap()
.metric_interval_seconds(),
None
);
}
#[test]
fn interval_parsing_accepts_seconds_and_minutes_only() {
assert_eq!(parse_interval_seconds("15s"), Some(15));
assert_eq!(parse_interval_seconds(" 5s "), Some(5));
assert_eq!(parse_interval_seconds("5m"), Some(300));
assert_eq!(parse_interval_seconds("15"), None);
assert_eq!(parse_interval_seconds("0s"), None);
assert_eq!(parse_interval_seconds(""), None);
}
#[test]
fn interval_parsing_never_overflows_on_huge_values() {
assert_eq!(parse_interval_seconds("922337203685477580m"), None);
assert_eq!(parse_interval_seconds("153722867280912931m"), None);
assert_eq!(parse_interval_seconds("99999999999999999999s"), None);
assert_eq!(parse_interval_seconds("600m"), Some(36_000));
}
#[test]
fn metric_interval_ignores_unparseable_targets_instead_of_treating_them_as_zero() {
let unit = |target: &str| WorkSpecUnit {
unit_id: "u".to_string(),
capability: "collect_metrics".to_string(),
rule_ref: String::new(),
requires_privilege: String::new(),
sources: vec![WorkSpecSource {
kind: "MetricInterval".to_string(),
target: target.to_string(),
multiline: "none".to_string(),
}],
};
let mixed = WorkSpec {
units: vec![unit("garbage"), unit("30s")],
};
assert_eq!(mixed.metric_interval_seconds(), Some(30));
let all_bad = WorkSpec {
units: vec![unit("nope")],
};
assert_eq!(all_bad.metric_interval_seconds(), None);
}
}