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/// 常驻工作的状态取值(对应模型 `StandingWork.status`)。
28///
29/// `superseded` 与 `revoked` 都不出现在 `WorkGrant.standing` 里:前者是「被新版本取代」,
30/// 后者是「授权被撤」。两者都要留痕,所以还是要有状态而不是删行。
31pub const STANDING_WORK_STATUSES: [&str; 4] = ["active", "paused", "superseded", "revoked"];
32
33/// 一次性工作的状态取值(对应模型 `OneShotWork.status`)。
34pub 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
46/// Agent **允许上报**的一次性工作状态([`ONE_SHOT_WORK_STATUSES`] 的真子集)。
47///
48/// 为什么不是全集:`dispatched` 是网关自己写的(派下去那一刻),`paused` / `canceled` / `expired`
49/// 归**控制面与期限**管(运维暂停撤回、或过了截止)—— agent 无权把它们写回去。
50/// 剩下的问题只有一种:「这件活做完了没有」,答案就这三种。
51pub const AGENT_REPORTABLE_WORK_STATUSES: [&str; 3] = ["running", "succeeded", "failed"];
52
53/// 一次性工作的**终态**:到了这几个状态就了结了,不再出现在快照里。
54pub const ONE_SHOT_TERMINAL_STATUSES: [&str; 5] =
55    ["succeeded", "failed", "timed_out", "canceled", "expired"];
56
57/// 工作类型:常驻(持续到被替换或撤回)与一次性(有期限与终态)。
58#[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    /// 常驻工作:按**采集面**授权,一个面一份。
63    Standing,
64    /// 一次性工作:按**动作**授权。
65    OneShot,
66}
67
68/// 常驻工作:网关声明该 Agent 应当持续执行的工作。
69///
70/// 粒度是**采集面**(一个面 = 一份工作),不是 capability:一份 `collect_logs` 会裹住十几个面,
71/// 那样「按面暂停 / 限流 / 审计」全都无从下手。面由模板的 `family_scope` 展开而来,
72/// **不携带 capability** —— 该面由哪个采集器承接,由 `spec` 里的单元各自决定。
73#[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    /// 采集面(`CollectionFamily`)。
80    pub family: String,
81    /// 工作参数:由采集目录的条目组合而成(**不是自由文本**),随 `plan_version` 整体替换。
82    pub spec: String,
83    /// 本工作按哪一版目录展开:目录换版**不追改**已授权工作(要跟新版得走新提案 + 审定)。
84    pub catalog_version: i64,
85    /// 生效依据:指向已批准的提案(人工直填 spec 时为空)。
86    #[serde(default)]
87    pub proposal_id: Option<String>,
88    /// 期望版本:网关每次改动 +1;agentd 回报实际版本,与它比对即得漂移。
89    pub plan_version: i64,
90    pub effective_from: String,
91    /// 见 [`STANDING_WORK_STATUSES`]。
92    pub status: String,
93    pub updated_by: String,
94    pub updated_at: String,
95}
96
97/// 一次性工作:有计划开始时间与期限,有明确终态。
98#[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    /// 动作面:upgrade / snapshot / exec / ...(执行载体是 Reporting 域的 ActionPlan)。
105    pub action: String,
106    pub spec: String,
107    /// 计划开始时间:到点前不应执行(与「立即派发」区分开)。
108    pub scheduled_at: String,
109    /// 绝对截止:业务要求的时间点,**暂停也照走**(不由暂停顺延)。
110    pub deadline_at: String,
111    /// 执行预算(秒):**只在实际执行时消耗**,暂停期间不计。
112    pub timeout_seconds: i64,
113    /// 可中断性:只有可中断的动作才允许运行中暂停。
114    pub interruptible: bool,
115    /// 见 [`ONE_SHOT_WORK_STATUSES`]。
116    pub status: String,
117    /// 当前暂停的起点(运行中暂停才有;恢复后清空)。
118    #[serde(default)]
119    pub paused_at: Option<String>,
120    /// 累计暂停时长(秒):用于审计与「预算未被暂停消耗」的核对。
121    pub paused_total_seconds: i64,
122    /// 步级断点:已完成步骤保留、`current_step` 在恢复时重做。
123    #[serde(default)]
124    pub current_step: Option<String>,
125    #[serde(default)]
126    pub completed_steps: Vec<String>,
127    /// 已尝试次数(含恢复后的重做)。
128    pub attempt: i64,
129    pub issued_by: String,
130    pub issued_at: String,
131}
132
133impl OneShotWork {
134    /// 是否**未了结**(快照里只带未了结的活)。
135    pub fn is_outstanding(&self) -> bool {
136        !ONE_SHOT_TERMINAL_STATUSES.contains(&self.status.as_str())
137    }
138
139    /// 是否允许在运行期间暂停。
140    ///
141    /// 两条都要满足:动作声明了 `interruptible`,且此刻确实在做(`accepted`/`running`)。
142    /// 不可中断的动作**拒绝**而不是「尽力暂停」—— 挂起半个升级进程比不暂停更危险。
143    pub fn can_pause(&self) -> bool {
144        self.interruptible && matches!(self.status.as_str(), "accepted" | "running")
145    }
146}
147
148/// 管理面授权或撤回工作的回执:一份工作一次。
149#[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    /// accepted | paused | resumed | revoked | rejected。
157    pub status: String,
158    pub plan_version: i64,
159    pub created_at: String,
160}
161
162// ─────────────────────────────────────────────────────────────────────────────
163// 工作参数(`StandingWork.spec` 的内容)
164// ─────────────────────────────────────────────────────────────────────────────
165//
166// 模型里 `StandingWork.spec` 是一个 `String`,语义是「工作参数:由采集目录的条目
167// 组合而成,**不是自由文本**」。这里的三个类型就是它的**编码**:一串已物化的采集单元。
168//
169// 为什么必须带上来源与规则标识,而不是只给一串 `unit_id`:agentd 拿到工作要能
170// **直接照做**。单元的采集来源(`sources`)与数据面规则标识(`rule_ref`)本来都只
171// 存在于网关的采集目录里,只发 id 等于发了一张自己去不了的地址 —— 于是要么再去网关
172// 拉一次目录(多一条必须鉴权的路径),要么两边各维护一份目录(必然漂移)。
173//
174// 为什么不把它们塞成 `StandingWork` 的字段:那会把「工作」与「内容目录」的边界糊掉;
175// 模型里这个字段就是**不透明的工作参数**,保持它不透明是两侧能独立演进的前提。
176
177/// 一份常驻工作的工作参数。
178#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
179#[serde(deny_unknown_fields)]
180pub struct WorkSpec {
181    /// 本工作包含的采集单元(已按该机器的事实裁剪过)。
182    #[serde(default)]
183    pub units: Vec<WorkSpecUnit>,
184}
185
186/// 一个已物化的采集单元:只说「采什么、怎么落地」,不复述策展元信息。
187///
188/// 刻意不带 `status` / `match` / `catalog_version`:那些是网关策展与裁剪的输入,
189/// 工作一旦发出去就已经裁剪完了。带上它们会让 agentd 有「再判断一次」的空间,
190/// 而 agentd 不是第二个策展器。
191#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192#[serde(deny_unknown_fields)]
193pub struct WorkSpecUnit {
194    pub unit_id: String,
195    /// `collect_logs` | `collect_metrics`:该单元由哪类采集器承接。
196    pub capability: String,
197    /// 数据面 rule/oml 标识(空 = 该单元还没接上规则,理论不会入工作)。
198    #[serde(default)]
199    pub rule_ref: String,
200    /// `none` | `root` | `fda`:要采到这东西得有什么权限 —— 缺权限时应**说清缺什么**,
201    /// 而不是安静地采不到。
202    #[serde(default)]
203    pub requires_privilege: String,
204    #[serde(default)]
205    pub sources: Vec<WorkSpecSource>,
206}
207
208/// 采集来源:`kind` ∈ `FileGlob` | `Exporter` | `UnifiedLogPredicate` | `MetricInterval`。
209///
210/// 与模型 `CollectionSource` 同形:`target` 的含义由 `kind` 决定
211/// (路径通配 / 导出器标识 / 谓词 / 周期)。
212#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
213#[serde(deny_unknown_fields)]
214pub struct WorkSpecSource {
215    pub kind: String,
216    pub target: String,
217    /// 怎么读这条来源:`none`(一行一条)| `indented`(行首缩进是上一条的续行)。
218    ///
219    /// 只对可 tail 的 `FileGlob` 有意义,其余 kind 恒为 `none`。**不是**"要不要多行"这个偏好,
220    /// 而是这条日志的**固有格式**:说错了要么把多行记录拆散、要么把独立记录粘成一条。
221    #[serde(default = "default_multiline")]
222    pub multiline: String,
223}
224
225fn default_multiline() -> String {
226    "none".to_string()
227}
228
229impl WorkSpec {
230    /// 解析 `StandingWork.spec`。
231    ///
232    /// 解析失败**不当作空工作**:空工作会让「参数坏了」退化成「没事可做」,
233    /// 两种情形的运维动作完全不同。
234    pub fn parse(spec: &str) -> Result<Self, serde_json::Error> {
235        serde_json::from_str(spec)
236    }
237
238    /// 序列化(网关写 `spec` 用)。字段顺序固定 → 同一份工作每次编码都一样,
239    /// `plan_version` 之外的字节也稳定,方便对账。
240    pub fn encode(&self) -> Result<String, serde_json::Error> {
241        serde_json::to_string(self)
242    }
243
244    /// 周期类来源声明的间隔(秒):取所有 `MetricInterval` 里最密的那一个。
245    ///
246    /// 取最密而不是取第一个:一份工作里若有两个指标单元,采慢的那个会
247    /// 静默压掉采密的要求,而“要得最急”的需求才是上送频率的下界。
248    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    /// 是否要采日志(有 `collect_logs` 单元)。
258    pub fn collects_logs(&self) -> bool {
259        self.units
260            .iter()
261            .any(|unit| unit.capability == "collect_logs")
262    }
263}
264
265/// agentd **有可能真的执行**的来源 kind(其余 kind 一律 `unsupported`,不当没看见)。
266///
267/// 注意:kind 对了**还不够** —— 目标形态也有限制,见 [`is_executable_source`]:
268/// `FileGlob` 要是显式绝对路径、`Exporter` 要是**已知导出器 ID**。只按 kind 判会把
269/// 「通配目标」「写错的导出器」当成可采,而那在采集端落不了地。
270pub const EXECUTABLE_SOURCE_KINDS: &[&str] = &["FileGlob", "MetricInterval", "Exporter"];
271
272/// 目标里有没有通配元字符 —— 有就说明它不是**显式单路径**。
273pub fn has_glob_meta(target: &str) -> bool {
274    target.contains(['*', '?', '['])
275}
276
277/// 这个目标是不是采集端今天**真能打开**的路径:**绝对路径、无通配**。
278///
279/// 采集端(agentd 第一版)只实现显式单路径输入:**不做 glob 展开、也不展开 `~`**
280/// (见 `wist/wist-agentd/docs/design/log-file-input-spec.md` §2 的能力边界)。
281/// 所以 `/var/log/wifi.log*`、`~/Library/Logs/Homebrew/*` 这类目标今天都采不到。
282pub fn is_explicit_path(target: &str) -> bool {
283    target.starts_with('/') && !has_glob_meta(target)
284}
285
286/// 定时导出器(`Exporter`)的**已知 ID**:协议层词表,网关 / agentd / web 三侧共用。
287///
288/// 为什么词表放这里:`is_executable_source` 要判「这条 `Exporter` 目标今天真能不能采」,
289/// 而网关的采集就绪度与 agentd 的「接不接得了」必须是**同一个判据**(见 [`is_executable_source`])。
290/// 词表只收**已实现**的导出器 —— 收录一个 ID 就等于承诺 agentd 能跑它;未知 ID(旧网关发新 ID、
291/// 或写错)一律当不可采,不要「写了就算」。新增能力 = 四侧(本 crate + 网关 + agentd + web)
292/// 一起改 + 知识库把对应单元开 `active`。
293pub 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
304/// 把一个 `Exporter` 目标拆成 `(id, arg)`:`"dmesg:panic"` → `("dmesg", Some("panic"))`。
305///
306/// 无 `:` 时 `arg` 为 `None`;空 id 或含空白的 id 视为无效(返回 `None`)。`arg` 是**声明式的窄选项**
307/// (如 `dmesg` 的过滤键),不是命令行片段 —— 实现侧各自解释,绝不拼进 shell。
308pub 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
320/// 这条 `Exporter` 目标的 ID 是不是**已知且已实现**的导出器。
321pub 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
328/// 一条来源今天**能不能被采**:`kind` 要 agentd 能执行,`target` 还要是显式路径。
329///
330/// 为什么把两件事放在一个函数里:网关的采集就绪度(「这个面能不能派下去」)与 agentd 的
331/// 「这条来源接不接得了」必须是**同一个判据**,否则会出现「网关说可采、agent 拿到后报
332/// unsupported」这种没人能发现的矛盾。加一种 kind 支持时**只改这里**。
333pub fn is_executable_source(kind: &str, target: &str) -> bool {
334    match kind {
335        // 本地文件采集:只支持显式绝对路径(无 glob、无 `~`)。
336        "FileGlob" => is_explicit_path(target),
337        // 指标:`target` 是采集周期(如 `15s`、`60s`),不是路径。
338        "MetricInterval" => true,
339        // 定时导出器:`target` 要是**已知导出器 ID**(可带 `:arg`)。未知 ID 不可执行。
340        "Exporter" => is_known_exporter(target),
341        _ => false,
342    }
343}
344
345impl WorkSpecUnit {
346    /// 该单元里可以交给本地文件采集器的来源(`FileGlob`)。
347    ///
348    /// 其它来源(导出器 / 统一日志谓词)不是本地 tail 一个文件能承接的:
349    /// 调用方应当把它们当成**还没支持的单元**如实报出来,而不是当没看见。
350    pub fn file_sources(&self) -> Vec<&WorkSpecSource> {
351        self.sources
352            .iter()
353            .filter(|source| source.kind == "FileGlob")
354            .collect()
355    }
356
357    /// 本单元有没有 agentd 目前接不了的来源(kind 不认识,或目标不是显式路径)。
358    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    /// 本单元至少有一条来源是 agentd **今天真能采**的。
366    ///
367    /// 「能不能采」只看这个:与解析规则(`rule_ref`)无关 —— 采原文不需要规则,
368    /// 规则只决定采下来的东西能不能被归类、抽字段。
369    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
376/// 解析 `15s` / `60s` / `5m` 这类周期写法(与目录里 `MetricInterval` 的写法一致)。
377///
378/// 只认秒与分两种后缀:目录里的值是人写的,多一个单位就多一种笔误的可能,
379/// 而解析不了时返回 `None` 会让调用方回退到自己的默认值(不静默取 0)。
380pub 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        // `checked_mul`:`123456789012345678m` 这类输入会让 `minutes * 60` 溢出 i64。
388        // 契约是「解析不了返回 None」(调用方回退到自己的默认值),不能让它 panic(debug)
389        // 或悄悄回绕成一个错的间隔(release)。
390        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        // 未了结的活才进快照:终态留在库里供审计,但不必再发给 agent。
423        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        // 不可中断:拒绝,而不是「尽力暂停」。
434        assert!(!one_shot("running", false).can_pause());
435        // 可中断但没在做(还没到点 / 已经了结):也没什么可暂停的。
436        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        // 采集端(agentd 第一版)只实现**显式单路径**:不做 glob 展开、不展开 `~`。
457        // 这个判据被两侧共用(网关的采集就绪度 + agentd 的接受逻辑),所以在这里钉死。
458        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        // 指标:`target` 是周期不是路径,不管显式路径那一条。
468        assert!(is_executable_source("MetricInterval", "15s"));
469        // 文件:kind 对了,目标还得是显式路径 —— 否则网关会把采不到的面报成“可采”。
470        assert!(is_executable_source("FileGlob", "/var/log/app.log"));
471        assert!(!is_executable_source("FileGlob", "/var/log/app*"));
472        // 导出器:kind 对了,ID 还得是**已知的**(`last,lastb` 是 macOS 侧、尚未实现)。
473        assert!(is_executable_source("Exporter", "smartctl"));
474        assert!(is_executable_source("Exporter", "dmesg:panic"));
475        assert!(!is_executable_source("Exporter", "last,lastb"));
476        // 还没实现的 kind:一律不可执行。
477        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        // 空 id / 含空白的 id 不是合法 ID。
488        assert_eq!(parse_exporter_target(":panic"), None);
489        assert_eq!(parse_exporter_target(" dm esg"), None);
490        // 判定只看冒号前的 ID,`arg` 不参与。
491        assert!(is_known_exporter("dmesg:nvidia-xid"));
492        assert!(!is_known_exporter("nope"));
493    }
494
495    // ── 工作参数(`spec`)的编码 ──
496
497    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                            // 显式单路径:采集端今天只能执行这种(通配未实现)。
525                            target: "/var/log/tccd.log".to_string(),
526                            // 多行格式是这条日志的固有属性,跟着来源走。
527                            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        // 取了最密的周期:要得最急的才是上送频率的下界。
543        assert_eq!(decoded.metric_interval_seconds(), Some(15));
544        // 可本地 tail 的来源挑得出来,接不了的来源也报得出来(而不是当没看见)。
545        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        // 类型是 FileGlob("可 tail 的文件")不等于可采:**目标形态**也得是采集端能执行的。
556        // 否则网关会把一个根本采不到的面报成“可采”,而 agentd 那边只会报“路径不存在”。
557        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        // 老网关发来的 spec 没有 `multiline`:默认必须是一行一条。
576        // 取错默认值会把多行日志**粘**成一条 —— 那是无声的内容损坏,不是格式问题。
577        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        // 再编码出去会把默认值写实:新网关发的 spec 字段总是齐的。
581        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        // 参数坏了与「没事可做」是两回事:前者要人去看,后者什么都不用做。
590        assert!(WorkSpec::parse("mac-host-metrics").is_err());
591        assert!(WorkSpec::parse("{").is_err());
592        // 空工作本身是合法的(能表示「这个面暂时没东西可采」)。
593        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        // 认不出来的回 None(调用方回退到自己的默认值),而不是静默当 0。
608        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        // 契约是「解析不了返回 None」,不是 panic:溢出在 debug 下 panic、release 下回绕成一个
616        // 错误的间隔 —— 两个后果都不能接受。这是两侧都要解析的字节,输入可能来自远端。
617        assert_eq!(parse_interval_seconds("922337203685477580m"), None);
618        assert_eq!(parse_interval_seconds("153722867280912931m"), None);
619        // 秒分支本来就靠 parse 失败兜底。
620        assert_eq!(parse_interval_seconds("99999999999999999999s"), None);
621        // 合法但很大的分钟值仍要算得出来(不误伤)。
622        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        // 认不出的周期不能静默当 0(0 会被当成「无限密」),而要如实跳过,只取能解析的最密值。
628        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}