Skip to main content

wist_contracts/
local_work.rs

1//! agentd 上报的**本机工作内容视图**(`state/work.json` 的同形子集)。
2//!
3//! 为什么要有它:网关知道自己**授权**了什么(`StandingWork.spec`),但那不等于「agent 真的在采
4//! 哪些文件」—— 本机配置里手工加的输入、工作被暂停、某条来源今天接不接得了,都只有 agent 自己
5//! 知道。`state/work.json` 就是这份本机事实,但它只落在目标机上;把它的**必要子集**随状态上报
6//! 带上来,页面上才能回答「这台机器现在到底在干什么」。
7//!
8//! 刻意**不带**工作内容(`units`):那是网关发下去的,网关自己有 —— 本机再抄一份只会诱使人拿它
9//! 当第二份真相(与 `state_store::work` 的取舍一致)。
10
11use serde::{Deserialize, Serialize};
12
13/// 本机一条**采集任务**:一个任务 id 盯一个路径。
14#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
15#[jumo(kind = "struct", domain = "Reporting", module = "Reporting.Protocol")]
16#[serde(deny_unknown_fields)]
17pub struct AgentLocalTask {
18    /// 任务 id:同时是本机 `state/logs/file_inputs/<input_id>/` 与 spool 的目录名。
19    pub input_id: String,
20    /// 盯的文件路径(今天只实现显式绝对路径)。
21    pub path: String,
22    /// `head` | `tail`:起读位置。授权采集一律 `tail`(不重放历史)。
23    pub startup_position: String,
24}
25
26/// 一份**常驻工作**在本机的样子。
27#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
28#[jumo(kind = "struct", domain = "Reporting", module = "Reporting.Protocol")]
29#[serde(deny_unknown_fields)]
30pub struct AgentLocalStandingWork {
31    pub work_id: String,
32    /// 采集面(`CollectionFamily`)。
33    pub family: String,
34    /// `active` | `paused`(被撤回 / 被取代的不算「我手里的工作」)。
35    pub status: String,
36    /// 网关的期望版本。
37    pub plan_version: i64,
38    /// 我回报过的版本;`None` = 还没回报(网关那边看到的就是漂移)。
39    #[serde(default)]
40    pub acknowledged_version: Option<i64>,
41    pub effective_from: String,
42    /// 折算成的本机采集任务(指标类工作没有任务)。
43    #[serde(default)]
44    pub tasks: Vec<AgentLocalTask>,
45}
46
47/// 一件**一次性工作**在本机的样子。
48///
49/// 两个状态轴分开:`status` 是**网关侧**的派发状态(dispatched / accepted / …),`execution` 是
50/// **本机执行状态**(现在一律 `unexecuted`)。
51#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
52#[jumo(kind = "struct", domain = "Reporting", module = "Reporting.Protocol")]
53#[serde(deny_unknown_fields)]
54pub struct AgentLocalOneShotWork {
55    pub work_id: String,
56    /// 动作面:upgrade / snapshot / exec / …。
57    pub action: String,
58    pub status: String,
59    /// `unexecuted`(本机尚未执行)。
60    pub execution: String,
61    pub scheduled_at: String,
62    pub deadline_at: String,
63    pub timeout_seconds: i64,
64}
65
66/// agentd 上报的**本机工作视图**。
67#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
68#[jumo(kind = "struct", domain = "Reporting", module = "Reporting.Protocol")]
69#[serde(deny_unknown_fields)]
70pub struct AgentLocalWork {
71    /// 这份视图是什么时候算出来的(**本机时钟**)。
72    pub recorded_at: String,
73    /// 溯源:本视图对应网关的授权序号。**不是授权副本** —— 只用来把本机视图与网关那一版对上号。
74    pub gateway_sequence: i64,
75    /// 手里生效或暂停的常驻工作。
76    #[serde(default)]
77    pub standing: Vec<AgentLocalStandingWork>,
78    /// 手里未了结的一次性工作。
79    #[serde(default)]
80    pub one_shot: Vec<AgentLocalOneShotWork>,
81    /// 本机配置里手工加的日志输入(运维逃生舱):**不来自任何采集面**,网关无从得知。
82    #[serde(default)]
83    pub local_inputs: Vec<AgentLocalTask>,
84    /// 汇总:指标上送周期(秒);`None` = 不上送。
85    #[serde(default)]
86    pub metrics_interval_seconds: Option<i64>,
87}
88
89#[cfg(test)]
90mod tests {
91    use super::*;
92
93    #[test]
94    fn a_minimal_report_decodes_with_empty_collections() {
95        // 只报「我算过一份视图」时,四类列表与周期都不写也应能解析。
96        let json = r#"{"recorded_at":"2026-09-27T00:00:00Z","gateway_sequence":4}"#;
97        let work: AgentLocalWork = serde_json::from_str(json).expect("decode");
98        assert!(work.standing.is_empty());
99        assert!(work.one_shot.is_empty());
100        assert!(work.local_inputs.is_empty());
101        assert_eq!(work.metrics_interval_seconds, None);
102    }
103
104    #[test]
105    fn a_full_report_round_trips_and_rejects_unknown_fields() {
106        let work = AgentLocalWork {
107            recorded_at: "2026-09-27T00:00:00Z".to_string(),
108            gateway_sequence: 9,
109            standing: vec![AgentLocalStandingWork {
110                work_id: "work-1".to_string(),
111                family: "ServiceLifecycle".to_string(),
112                status: "active".to_string(),
113                plan_version: 2,
114                acknowledged_version: Some(2),
115                effective_from: "2026-09-27T00:00:00Z".to_string(),
116                tasks: vec![AgentLocalTask {
117                    input_id: "app".to_string(),
118                    path: "/var/log/app.log".to_string(),
119                    startup_position: "tail".to_string(),
120                }],
121            }],
122            one_shot: vec![AgentLocalOneShotWork {
123                work_id: "work-2".to_string(),
124                action: "upgrade".to_string(),
125                status: "dispatched".to_string(),
126                execution: "unexecuted".to_string(),
127                scheduled_at: "2026-09-27T00:00:00Z".to_string(),
128                deadline_at: "2026-09-28T00:00:00Z".to_string(),
129                timeout_seconds: 600,
130            }],
131            local_inputs: Vec::new(),
132            metrics_interval_seconds: Some(15),
133        };
134        let json = serde_json::to_string(&work).expect("encode");
135        let back: AgentLocalWork = serde_json::from_str(&json).expect("decode");
136        assert_eq!(back, work);
137
138        let mutated = json.replacen('{', "{\"extra\":1,", 1);
139        assert!(serde_json::from_str::<AgentLocalWork>(&mutated).is_err());
140    }
141}