Skip to main content

wist_contracts/
work.rs

1//! **工作授权快照**:网关授权(grant)、agentd 拉取(`WorkGrant`)并确认。
2//!
3//! ## 为什么是「快照」而不是「指令流」
4//!
5//! 常驻工作(`StandingWork`)表达的是**期望状态**:某个采集面应当持续执行什么。
6//! 期望状态天然是幂等的(同一份重复拉取不产生副作用),断了网、重启了进程,
7//! 再拉一次就回到期望 —— 不必让网关记住「上次推到哪条」。
8//! 这跟 `PollControlCommands` 的长轮询指令流是两种东西:那条要求**不重放**,这条要求**可重拉**。
9//!
10//! 一次性工作(`OneShotWork`)是命令式的,但它也搭这份快照回来:
11//! 决定「现在还该做吗」的是它的状态(未了结才出现在快照里),不是「推没推过」。
12//!
13//! ## 为什么在契约 crate
14//!
15//! 这是**两侧都要解析**的字节:网关写授权、agentd 读并执行。各写一份结构体,
16//! 迟早出现「网关发了 `plan_version`、agentd 读的是 `version`」这种只能在真机联调时才发现的分叉。
17//! 与 `DiscoveryAspectPolicySet` 同一个理由。
18//!
19//! ## 与本地保护暂停的区别
20//!
21//! agentd 自己也会暂停(如 `spool over limit`),那是**本地保护**:自动、临时、本机可见。
22//! 这里下发的 `paused` 是**授权层状态**:人工、持久、跨重启。两者在观测里必须能分辨,
23//! 不要合并成一个 paused —— 合并之后「谁把采集停了」就再也说不清。
24
25use serde::{Deserialize, Serialize};
26
27/// agentd → 网关:拉取工作授权快照的 envelope kind。
28pub const POLL_WORK_KIND: &str = "poll_work";
29/// agentd → 网关:确认收到工作的 envelope kind。
30pub const ACK_WORK_KIND: &str = "ack_work";
31/// agentd → 网关:上报一次性工作执行结果的 envelope kind。
32pub const REPORT_WORK_RESULT_KIND: &str = "report_work_result";
33
34/// 常驻工作的状态取值(对应模型 `StandingWork.status`)。
35///
36/// `superseded` 与 `revoked` 都不出现在 `WorkGrant.standing` 里:前者是「被新版本取代」,
37/// 后者是「授权被撤」。两者都要留痕,所以还是要有状态而不是删行。
38pub const STANDING_WORK_STATUSES: [&str; 4] = ["active", "paused", "superseded", "revoked"];
39
40/// 一次性工作的状态取值(对应模型 `OneShotWork.status`)。
41pub 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
53/// Agent **允许上报**的一次性工作状态([`ONE_SHOT_WORK_STATUSES`] 的真子集)。
54///
55/// 为什么不是全集:`dispatched` 是网关自己写的(派下去那一刻),`paused` / `canceled` / `expired`
56/// 归**控制面与期限**管(运维暂停撤回、或过了截止)—— agent 无权把它们写回去。
57/// 剩下的问题只有一种:「这件活做完了没有」,答案就这三种。
58pub const AGENT_REPORTABLE_WORK_STATUSES: [&str; 3] = ["running", "succeeded", "failed"];
59
60/// 一次性工作的**终态**:到了这几个状态就了结了,不再出现在快照里。
61pub const ONE_SHOT_TERMINAL_STATUSES: [&str; 5] =
62    ["succeeded", "failed", "timed_out", "canceled", "expired"];
63
64/// 工作类型:常驻(持续到被替换或撤回)与一次性(有期限与终态)。
65#[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    /// 常驻工作:按**采集面**授权,一个面一份。
70    Standing,
71    /// 一次性工作:按**动作**授权。
72    OneShot,
73}
74
75/// 常驻工作:网关声明该 Agent 应当持续执行的工作。
76///
77/// 粒度是**采集面**(一个面 = 一份工作),不是 capability:一份 `collect_logs` 会裹住十几个面,
78/// 那样「按面暂停 / 限流 / 审计」全都无从下手。面由模板的 `family_scope` 展开而来,
79/// **不携带 capability** —— 该面由哪个采集器承接,由 `spec` 里的单元各自决定。
80#[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    /// 采集面(`CollectionFamily`)。
87    pub family: String,
88    /// 工作参数:由采集目录的条目组合而成(**不是自由文本**),随 `plan_version` 整体替换。
89    pub spec: String,
90    /// 本工作按哪一版目录展开:目录换版**不追改**已授权工作(要跟新版得走新提案 + 审定)。
91    pub catalog_version: i64,
92    /// 生效依据:指向已批准的提案(人工直填 spec 时为空)。
93    #[serde(default)]
94    pub proposal_id: Option<String>,
95    /// 期望版本:网关每次改动 +1;agentd 回报实际版本,与它比对即得漂移。
96    pub plan_version: i64,
97    pub effective_from: String,
98    /// 见 [`STANDING_WORK_STATUSES`]。
99    pub status: String,
100    pub updated_by: String,
101    pub updated_at: String,
102}
103
104/// 一次性工作:有计划开始时间与期限,有明确终态。
105#[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    /// 动作面:upgrade / snapshot / exec / ...(执行载体是 Reporting 域的 ActionPlan)。
112    pub action: String,
113    pub spec: String,
114    /// 计划开始时间:到点前不应执行(与「立即派发」区分开)。
115    pub scheduled_at: String,
116    /// 绝对截止:业务要求的时间点,**暂停也照走**(不由暂停顺延)。
117    pub deadline_at: String,
118    /// 执行预算(秒):**只在实际执行时消耗**,暂停期间不计。
119    pub timeout_seconds: i64,
120    /// 可中断性:只有可中断的动作才允许运行中暂停。
121    pub interruptible: bool,
122    /// 见 [`ONE_SHOT_WORK_STATUSES`]。
123    pub status: String,
124    /// 当前暂停的起点(运行中暂停才有;恢复后清空)。
125    #[serde(default)]
126    pub paused_at: Option<String>,
127    /// 累计暂停时长(秒):用于审计与「预算未被暂停消耗」的核对。
128    pub paused_total_seconds: i64,
129    /// 步级断点:已完成步骤保留、`current_step` 在恢复时重做。
130    #[serde(default)]
131    pub current_step: Option<String>,
132    #[serde(default)]
133    pub completed_steps: Vec<String>,
134    /// 已尝试次数(含恢复后的重做)。
135    pub attempt: i64,
136    pub issued_by: String,
137    pub issued_at: String,
138}
139
140impl OneShotWork {
141    /// 是否**未了结**(快照里只带未了结的活)。
142    pub fn is_outstanding(&self) -> bool {
143        !ONE_SHOT_TERMINAL_STATUSES.contains(&self.status.as_str())
144    }
145
146    /// 是否允许在运行期间暂停。
147    ///
148    /// 两条都要满足:动作声明了 `interruptible`,且此刻确实在做(`accepted`/`running`)。
149    /// 不可中断的动作**拒绝**而不是「尽力暂停」—— 挂起半个升级进程比不暂停更危险。
150    pub fn can_pause(&self) -> bool {
151        self.interruptible && matches!(self.status.as_str(), "accepted" | "running")
152    }
153}
154
155/// 工作授权快照:常驻工作的当前生效版本 + 未了结的一次性工作。
156///
157/// 幂等、可重复拉取;`sequence` 只用来让 agentd 判断「这份跟我手上的有没有变」,
158/// **不承担「指令重放」的语义**(那是控制指令流的事)。
159#[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    /// 每个面一条。
165    #[serde(default)]
166    pub standing: Vec<StandingWork>,
167    #[serde(default)]
168    pub one_shot: Vec<OneShotWork>,
169    /// 授权序号(单调递增,每次授权/撤回/暂停/继续都 +1)。
170    pub sequence: i64,
171    pub granted_at: String,
172}
173
174/// 管理面授权或撤回工作的回执:一份工作一次。
175#[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    /// accepted | paused | resumed | revoked | rejected。
183    pub status: String,
184    pub plan_version: i64,
185    pub created_at: String,
186}
187
188/// agentd → 网关:拉取工作授权快照。
189///
190/// 带 `last_seen_sequence`(与本机手上那份的序号),网关可以据此在没变化时短路;
191/// 带 `wait_ms` 是为了允许将来的长轮询(现在是立即返回,字段先留着,免得改协议)。
192#[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/// agentd → 网关:确认收到某份工作。
211///
212/// 常驻工作在 `plan_version` 变化后**也必须**确认:网关据此判断「期望的版本真到了吗」,
213/// 一直没确认的就是漂移。
214#[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/// 网关对 [`AckWork`] 的回应。
233///
234/// 是**工作域的结构**而不是协议消息(与 [`WorkGrant`] 同类):它描述的是
235/// 「工作已被确认」这个领域事实,也要能被用例当成 outcome 引用。
236#[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    /// accepted | stale | unknown。
242    pub status: String,
243    pub accepted_at: String,
244}
245
246/// agentd → 网关:上报一次性工作的**执行结果**(进度与终态)。
247///
248/// 与 [`AckWork`] 的分工:确认回答「我收到了」,本消息回答「我做得怎么样了」。
249/// 两者分开是因为它们的**失败代价不同**:确认丢了只是页面晚一拍,结果丢了则意味着
250/// 「一件改变机器状态的活做完了,而控制面永远不知道它成没成」。
251///
252/// 为什么必须由 agent 上报而不是网关自己推:网关只知道派了什么;一件升级是否真的换上了、
253/// 失败了有没有回滚,只有在机器上的那个执行体知道。
254#[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    /// 见 [`AGENT_REPORTABLE_WORK_STATUSES`]。
269    pub status: String,
270    /// 人看的说明:失败原因**原样带上**(如「摘要不符」「新版 60s 没起来,已回滚到 0.1.3」)。
271    /// 不写清原因,页面上就只剩一个无法解释的「失败」。
272    #[serde(default)]
273    pub detail: String,
274    pub reported_at: String,
275}
276
277/// 网关对 [`ReportWorkResult`] 的回应。
278#[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    /// accepted | stale | unknown。
284    ///
285    /// `stale` = 这件活已经到终态了(被撤回、超期,或已经报过终态),后到的结果**不覆盖**它:
286    /// 一件活只能有一个终态,写第二次只会把历史抹掉。
287    pub status: String,
288    pub accepted_at: String,
289}
290
291// ─────────────────────────────────────────────────────────────────────────────
292// 工作参数(`StandingWork.spec` 的内容)
293// ─────────────────────────────────────────────────────────────────────────────
294//
295// 模型里 `StandingWork.spec` 是一个 `String`,语义是「工作参数:由采集目录的条目
296// 组合而成,**不是自由文本**」。这里的三个类型就是它的**编码**:一串已物化的采集单元。
297//
298// 为什么必须带上来源与规则标识,而不是只给一串 `unit_id`:agentd 拿到工作要能
299// **直接照做**。单元的采集来源(`sources`)与数据面规则标识(`rule_ref`)本来都只
300// 存在于网关的采集目录里,只发 id 等于发了一张自己去不了的地址 —— 于是要么再去网关
301// 拉一次目录(多一条必须鉴权的路径),要么两边各维护一份目录(必然漂移)。
302//
303// 为什么不把它们塞成 `StandingWork` 的字段:那会把「工作」与「内容目录」的边界糊掉;
304// 模型里这个字段就是**不透明的工作参数**,保持它不透明是两侧能独立演进的前提。
305
306/// 一份常驻工作的工作参数。
307#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
308#[serde(deny_unknown_fields)]
309pub struct WorkSpec {
310    /// 本工作包含的采集单元(已按该机器的事实裁剪过)。
311    #[serde(default)]
312    pub units: Vec<WorkSpecUnit>,
313}
314
315/// 一个已物化的采集单元:只说「采什么、怎么落地」,不复述策展元信息。
316///
317/// 刻意不带 `status` / `match` / `catalog_version`:那些是网关策展与裁剪的输入,
318/// 工作一旦发出去就已经裁剪完了。带上它们会让 agentd 有「再判断一次」的空间,
319/// 而 agentd 不是第二个策展器。
320#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
321#[serde(deny_unknown_fields)]
322pub struct WorkSpecUnit {
323    pub unit_id: String,
324    /// `collect_logs` | `collect_metrics`:该单元由哪类采集器承接。
325    pub capability: String,
326    /// 数据面 rule/oml 标识(空 = 该单元还没接上规则,理论不会入工作)。
327    #[serde(default)]
328    pub rule_ref: String,
329    /// `none` | `root` | `fda`:要采到这东西得有什么权限 —— 缺权限时应**说清缺什么**,
330    /// 而不是安静地采不到。
331    #[serde(default)]
332    pub requires_privilege: String,
333    #[serde(default)]
334    pub sources: Vec<WorkSpecSource>,
335}
336
337/// 采集来源:`kind` ∈ `FileGlob` | `Exporter` | `UnifiedLogPredicate` | `MetricInterval`。
338///
339/// 与模型 `CollectionSource` 同形:`target` 的含义由 `kind` 决定
340/// (路径通配 / 导出器标识 / 谓词 / 周期)。
341#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
342#[serde(deny_unknown_fields)]
343pub struct WorkSpecSource {
344    pub kind: String,
345    pub target: String,
346    /// 怎么读这条来源:`none`(一行一条)| `indented`(行首缩进是上一条的续行)。
347    ///
348    /// 只对可 tail 的 `FileGlob` 有意义,其余 kind 恒为 `none`。**不是**"要不要多行"这个偏好,
349    /// 而是这条日志的**固有格式**:说错了要么把多行记录拆散、要么把独立记录粘成一条。
350    #[serde(default = "default_multiline")]
351    pub multiline: String,
352}
353
354fn default_multiline() -> String {
355    "none".to_string()
356}
357
358impl WorkSpec {
359    /// 解析 `StandingWork.spec`。
360    ///
361    /// 解析失败**不当作空工作**:空工作会让「参数坏了」退化成「没事可做」,
362    /// 两种情形的运维动作完全不同。
363    pub fn parse(spec: &str) -> Result<Self, serde_json::Error> {
364        serde_json::from_str(spec)
365    }
366
367    /// 序列化(网关写 `spec` 用)。字段顺序固定 → 同一份工作每次编码都一样,
368    /// `plan_version` 之外的字节也稳定,方便对账。
369    pub fn encode(&self) -> Result<String, serde_json::Error> {
370        serde_json::to_string(self)
371    }
372
373    /// 周期类来源声明的间隔(秒):取所有 `MetricInterval` 里最密的那一个。
374    ///
375    /// 取最密而不是取第一个:一份工作里若有两个指标单元,采慢的那个会
376    /// 静默压掉采密的要求,而“要得最急”的需求才是上送频率的下界。
377    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    /// 是否要采日志(有 `collect_logs` 单元)。
387    pub fn collects_logs(&self) -> bool {
388        self.units
389            .iter()
390            .any(|unit| unit.capability == "collect_logs")
391    }
392}
393
394/// agentd **有可能真的执行**的来源 kind(其余 kind 一律 `unsupported`,不当没看见)。
395///
396/// 注意:kind 对了**还不够** —— 目标形态也有限制,见 [`is_executable_source`]:
397/// `FileGlob` 要是显式绝对路径、`Exporter` 要是**已知导出器 ID**。只按 kind 判会把
398/// 「通配目标」「写错的导出器」当成可采,而那在采集端落不了地。
399pub const EXECUTABLE_SOURCE_KINDS: &[&str] = &["FileGlob", "MetricInterval", "Exporter"];
400
401/// 目标里有没有通配元字符 —— 有就说明它不是**显式单路径**。
402pub fn has_glob_meta(target: &str) -> bool {
403    target.contains(['*', '?', '['])
404}
405
406/// 这个目标是不是采集端今天**真能打开**的路径:**绝对路径、无通配**。
407///
408/// 采集端(agentd 第一版)只实现显式单路径输入:**不做 glob 展开、也不展开 `~`**
409/// (见 `wist/wist-agentd/docs/design/log-file-input-spec.md` §2 的能力边界)。
410/// 所以 `/var/log/wifi.log*`、`~/Library/Logs/Homebrew/*` 这类目标今天都采不到。
411pub fn is_explicit_path(target: &str) -> bool {
412    target.starts_with('/') && !has_glob_meta(target)
413}
414
415/// 定时导出器(`Exporter`)的**已知 ID**:协议层词表,网关 / agentd / web 三侧共用。
416///
417/// 为什么词表放这里:`is_executable_source` 要判「这条 `Exporter` 目标今天真能不能采」,
418/// 而网关的采集就绪度与 agentd 的「接不接得了」必须是**同一个判据**(见 [`is_executable_source`])。
419/// 词表只收**已实现**的导出器 —— 收录一个 ID 就等于承诺 agentd 能跑它;未知 ID(旧网关发新 ID、
420/// 或写错)一律当不可采,不要「写了就算」。新增能力 = 四侧(本 crate + 网关 + agentd + web)
421/// 一起改 + 知识库把对应单元开 `active`。
422pub 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
433/// 把一个 `Exporter` 目标拆成 `(id, arg)`:`"dmesg:panic"` → `("dmesg", Some("panic"))`。
434///
435/// 无 `:` 时 `arg` 为 `None`;空 id 或含空白的 id 视为无效(返回 `None`)。`arg` 是**声明式的窄选项**
436/// (如 `dmesg` 的过滤键),不是命令行片段 —— 实现侧各自解释,绝不拼进 shell。
437pub 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
449/// 这条 `Exporter` 目标的 ID 是不是**已知且已实现**的导出器。
450pub 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
457/// 一条来源今天**能不能被采**:`kind` 要 agentd 能执行,`target` 还要是显式路径。
458///
459/// 为什么把两件事放在一个函数里:网关的采集就绪度(「这个面能不能派下去」)与 agentd 的
460/// 「这条来源接不接得了」必须是**同一个判据**,否则会出现「网关说可采、agent 拿到后报
461/// unsupported」这种没人能发现的矛盾。加一种 kind 支持时**只改这里**。
462pub fn is_executable_source(kind: &str, target: &str) -> bool {
463    match kind {
464        // 本地文件采集:只支持显式绝对路径(无 glob、无 `~`)。
465        "FileGlob" => is_explicit_path(target),
466        // 指标:`target` 是采集周期(如 `15s`、`60s`),不是路径。
467        "MetricInterval" => true,
468        // 定时导出器:`target` 要是**已知导出器 ID**(可带 `:arg`)。未知 ID 不可执行。
469        "Exporter" => is_known_exporter(target),
470        _ => false,
471    }
472}
473
474impl WorkSpecUnit {
475    /// 该单元里可以交给本地文件采集器的来源(`FileGlob`)。
476    ///
477    /// 其它来源(导出器 / 统一日志谓词)不是本地 tail 一个文件能承接的:
478    /// 调用方应当把它们当成**还没支持的单元**如实报出来,而不是当没看见。
479    pub fn file_sources(&self) -> Vec<&WorkSpecSource> {
480        self.sources
481            .iter()
482            .filter(|source| source.kind == "FileGlob")
483            .collect()
484    }
485
486    /// 本单元有没有 agentd 目前接不了的来源(kind 不认识,或目标不是显式路径)。
487    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    /// 本单元至少有一条来源是 agentd **今天真能采**的。
495    ///
496    /// 「能不能采」只看这个:与解析规则(`rule_ref`)无关 —— 采原文不需要规则,
497    /// 规则只决定采下来的东西能不能被归类、抽字段。
498    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
505/// 解析 `15s` / `60s` / `5m` 这类周期写法(与目录里 `MetricInterval` 的写法一致)。
506///
507/// 只认秒与分两种后缀:目录里的值是人写的,多一个单位就多一种笔误的可能,
508/// 而解析不了时返回 `None` 会让调用方回退到自己的默认值(不静默取 0)。
509pub 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        // `checked_mul`:`123456789012345678m` 这类输入会让 `minutes * 60` 溢出 i64。
517        // 契约是「解析不了返回 None」(调用方回退到自己的默认值),不能让它 panic(debug)
518        // 或悄悄回绕成一个错的间隔(release)。
519        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        // 未了结的活才进快照:终态留在库里供审计,但不必再发给 agent。
568        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        // 不可中断:拒绝,而不是「尽力暂停」。
579        assert!(!one_shot("running", false).can_pause());
580        // 可中断但没在做(还没到点 / 已经了结):也没什么可暂停的。
581        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        // 采集端(agentd 第一版)只实现**显式单路径**:不做 glob 展开、不展开 `~`。
602        // 这个判据被两侧共用(网关的采集就绪度 + agentd 的接受逻辑),所以在这里钉死。
603        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        // 指标:`target` 是周期不是路径,不管显式路径那一条。
613        assert!(is_executable_source("MetricInterval", "15s"));
614        // 文件:kind 对了,目标还得是显式路径 —— 否则网关会把采不到的面报成“可采”。
615        assert!(is_executable_source("FileGlob", "/var/log/app.log"));
616        assert!(!is_executable_source("FileGlob", "/var/log/app*"));
617        // 导出器:kind 对了,ID 还得是**已知的**(`last,lastb` 是 macOS 侧、尚未实现)。
618        assert!(is_executable_source("Exporter", "smartctl"));
619        assert!(is_executable_source("Exporter", "dmesg:panic"));
620        assert!(!is_executable_source("Exporter", "last,lastb"));
621        // 还没实现的 kind:一律不可执行。
622        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        // 空 id / 含空白的 id 不是合法 ID。
633        assert_eq!(parse_exporter_target(":panic"), None);
634        assert_eq!(parse_exporter_target(" dm esg"), None);
635        // 判定只看冒号前的 ID,`arg` 不参与。
636        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        // `proposal_id` / `paused_at` / `current_step` 缺省时必须能解析:
657        // 手写 JSON(页面、联调)不该被迫填一堆 null。
658        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        // 两侧各自演进时,多出来的字段必须是响亮的错误:静默忽略会让
668        // 「网关发了新字段、agentd 装作没看见」变成长期无声的语义分叉。
669        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    // ── 工作参数(`spec`)的编码 ──
675
676    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                            // 显式单路径:采集端今天只能执行这种(通配未实现)。
704                            target: "/var/log/tccd.log".to_string(),
705                            // 多行格式是这条日志的固有属性,跟着来源走。
706                            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        // 取了最密的周期:要得最急的才是上送频率的下界。
722        assert_eq!(decoded.metric_interval_seconds(), Some(15));
723        // 可本地 tail 的来源挑得出来,接不了的来源也报得出来(而不是当没看见)。
724        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        // 类型是 FileGlob("可 tail 的文件")不等于可采:**目标形态**也得是采集端能执行的。
735        // 否则网关会把一个根本采不到的面报成“可采”,而 agentd 那边只会报“路径不存在”。
736        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        // 老网关发来的 spec 没有 `multiline`:默认必须是一行一条。
755        // 取错默认值会把多行日志**粘**成一条 —— 那是无声的内容损坏,不是格式问题。
756        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        // 再编码出去会把默认值写实:新网关发的 spec 字段总是齐的。
760        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        // 参数坏了与「没事可做」是两回事:前者要人去看,后者什么都不用做。
769        assert!(WorkSpec::parse("mac-host-metrics").is_err());
770        assert!(WorkSpec::parse("{").is_err());
771        // 空工作本身是合法的(能表示「这个面暂时没东西可采」)。
772        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        // 认不出来的回 None(调用方回退到自己的默认值),而不是静默当 0。
787        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        // 契约是「解析不了返回 None」,不是 panic:溢出在 debug 下 panic、release 下回绕成一个
795        // 错误的间隔 —— 两个后果都不能接受。这是两侧都要解析的字节,输入可能来自远端。
796        assert_eq!(parse_interval_seconds("922337203685477580m"), None);
797        assert_eq!(parse_interval_seconds("153722867280912931m"), None);
798        // 秒分支本来就靠 parse 失败兜底。
799        assert_eq!(parse_interval_seconds("99999999999999999999s"), None);
800        // 合法但很大的分钟值仍要算得出来(不误伤)。
801        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        // 认不出的周期不能静默当 0(0 会被当成「无限密」),而要如实跳过,只取能解析的最密值。
807        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}