Skip to main content

wist_contracts/
agent_uplink.rs

1//! **数据面上送启用**:控制面决定「这个 Agent 现在该不该向数据面上送」,agentd 拉取并在内存生效。
2//!
3//! ## 为什么是「现算的下发」而不是一条配置
4//!
5//! 上送目标(`host:port`)是**控制面的事实**:只有网关知道数据面(warp-parse)在哪。
6//! 把它写进 Agent 的静态配置意味着「改一次地址 = 重装所有 Agent」,而重装又会把配置
7//! 重置回模板默认值 —— 那正是「授权前待命」被实现成 `kind = "file"` 之后反复踩到的坑。
8//!
9//! 所以这里只下发一个**派生结果**:网关每次被问到时,用「该 Agent 有没有生效的工作」
10//! 加上管理面已设的「数据面上送地址」现算出来。没有新状态、没有推送通道,也不需要
11//! 管理面多一个动作 —— 派活(写一条 standing work)下一个 poll 就变成 `enabled = true`,
12//! 撤回即变回 `false`。
13//!
14//! ## 为什么另开端点,不塞进 [`crate::work::WorkGrant`]
15//!
16//! `WorkGrant` 带 `#[serde(deny_unknown_fields)]`。往里加字段会让**新网关 + 旧 agentd**
17//! 直接解析失败:旧 agentd 连工作授权都收不到,会停在最后一次应用的工作上 —— 这是
18//! 舰队级的静默停摆。独立端点对两个方向都安全:旧 agentd 从不调它;新 agentd 遇到
19//! 旧网关得到 404,按「无下发」回落本机配置。
20//!
21//! ## 三方语义(agentd 侧生效规则,见 `agentd` 的 `effective_output`)
22//!
23//! * 未下发(未入网 / 旧网关 404 / 网络失败) → 用**本机** `[telemetry.logs.output]`;
24//! * `enabled = false` → **待命**:不采、不写文件、不上送日志/指标(因此不会报错);
25//!   **但事实摘要仍上报** —— 它不是主机内容,而是让平台能推断「这台机器是什么」的最小
26//!   元数据(进程列表 / 监听端口 / os / arch)。没有它,新装机器在网关侧一片空白,
27//!   连「该派什么活」都定不下来(下发的目标指到哪,事实就发到哪,仅 tcp 能承载);
28//! * `enabled = true` → 若带 `host`/`port` 则**强制** `tcp(host, port)`(覆盖本机 `kind`),
29//!   否则沿用本机 `kind`。
30//!
31//! 授权优先于本机:托管模式下目标与开关都由控制面说了算,本机配置只作 standalone 的兜底。
32
33use serde::{Deserialize, Serialize};
34
35/// agentd → 网关:拉取数据面上送启用的 envelope kind。
36pub const POLL_AGENT_UPLINK_KIND: &str = "poll_agent_uplink";
37
38/// agentd → 网关:拉取数据面上送启用。
39///
40/// 与 [`crate::work::PollWork`] 同形(同一套 agent 凭据、同一份实例标识),因为它是
41/// 同一类「拉期望状态」的动作;只是期望状态的内容不同。
42#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
43#[jumo(
44    kind = "message",
45    role = "command",
46    domain = "Control",
47    module = "Control.AgentApp.FacingInterface"
48)]
49#[serde(deny_unknown_fields)]
50pub struct PollAgentUplink {
51    pub api_version: String,
52    pub kind: String,
53    pub agent_id: String,
54    pub instance_id: String,
55    pub requested_at: String,
56}
57
58/// 网关 → agentd:数据面上送的当前期望状态。
59///
60/// 与 [`crate::work::WorkGrant`] 同类:描述的是「这个 Agent 的上送现在应当是什么样」这个
61/// 领域事实(幂等、可重复拉取),而不是一次协议动作。它**派生自工作授权**,所以与
62/// `WorkGrant` 同属 `Control.Agent.Work`。
63#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
64#[jumo(kind = "struct", domain = "Control", module = "Control.Agent.Work")]
65#[serde(deny_unknown_fields)]
66pub struct AgentUplinkGrant {
67    /// 是否启用数据面上送的**主机内容**(日志 / 指标)。`false` = 待命:不采集、不上送日志与
68    /// 指标。**但事实摘要不受它管** —— 待命期仍上报进程列表等推断元数据(见模块文档)。
69    pub enabled: bool,
70    /// 数据面地址。`enabled = true` 时给出即**覆盖**本机 `tcp.addr`;
71    /// 缺省表示「沿用本机配置的目标」(本机 `kind = "file"` 时就仍是本地文件)。
72    #[serde(default, skip_serializing_if = "Option::is_none")]
73    pub host: Option<String>,
74    #[serde(default, skip_serializing_if = "Option::is_none")]
75    pub port: Option<u16>,
76    pub granted_at: String,
77}
78
79impl AgentUplinkGrant {
80    /// 待命:明确关掉。网关在「没有生效工作」或「管理面还没设上送地址」时下发它。
81    pub fn standby(granted_at: String) -> Self {
82        Self {
83            enabled: false,
84            host: None,
85            port: None,
86            granted_at,
87        }
88    }
89
90    /// 启用:带上目标。`host`/`port` 由管理面的「数据面上送地址」给。
91    pub fn enabled_at(host: String, port: u16, granted_at: String) -> Self {
92        Self {
93            enabled: true,
94            host: Some(host),
95            port: Some(port),
96            granted_at,
97        }
98    }
99
100    /// 要覆盖的本机目标(`enabled` 且绑定齐了 `host`/`port` 才有)。
101    ///
102    /// `host` 会 `trim`:空白主机(`" "`)不是目标 —— 把它当目标只会拼出连不上的
103    /// `" :9000"`,而「拼出一个错地址」比「不覆盖、回落本机」难查得多。
104    pub fn target(&self) -> Option<(&str, u16)> {
105        match (self.enabled, self.host.as_deref(), self.port) {
106            (true, Some(host), Some(port)) => {
107                let host = host.trim();
108                (!host.is_empty()).then_some((host, port))
109            }
110            _ => None,
111        }
112    }
113}
114
115/// Agent **实际生效**的采集输出状态(agent → 网关,自报)。
116///
117/// 与 [`AgentUplinkGrant`] 是一对:grant 说「控制面要它怎样」,这里说「它实际成了怎样」。
118/// 为什么必须由 agent 上报:网关知道自己**下发**了什么,但不知道 agent **生效**成什么 ——
119/// grant 可能还没拉到、可能被本机总闸拦住、可能目标连不上。这与
120/// `AgentStatusReport::discovery_policy_version` 是同一类事实(「我发布了哪一版」≠「它生效了哪一版」)。
121#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ::jumo_derive::Jumo)]
122#[jumo(kind = "struct", domain = "Reporting", module = "Reporting.Protocol")]
123#[serde(deny_unknown_fields)]
124pub struct AgentUplinkState {
125    /// 生效输出的总闸:`false` = 待命(不读源、不写本地采集输出、不上送)。
126    pub enabled: bool,
127    /// 生效输出类型:`file` | `tcp`(grant 覆盖本机 `kind` 之后的**结果**)。
128    pub kind: String,
129    /// 生效目标(`tcp` 且有目标时)。`None` = 本地文件输出,或 tcp 但没有目标。
130    #[serde(default, skip_serializing_if = "Option::is_none")]
131    pub target: Option<String>,
132    /// 这份生效状态的来源:`grant`(控制面下发)| `local`(未下发,用本机配置)。
133    ///
134    /// 有了它,平台才能区分「控制面没启用」与「启用了但 agent 没听从」(后者才是故障)。
135    pub source: String,
136    /// 最近的出口写失败是否**尚未恢复**(细节在 agent 日志里:目标与原因)。
137    pub output_write_failing: bool,
138}
139
140#[cfg(test)]
141mod state_tests {
142    use super::*;
143
144    #[test]
145    fn a_local_file_output_state_omits_the_target() {
146        let state = AgentUplinkState {
147            enabled: true,
148            kind: "file".to_string(),
149            target: None,
150            source: "local".to_string(),
151            output_write_failing: false,
152        };
153        let json = serde_json::to_string(&state).expect("encode");
154        assert!(!json.contains("target"), "{json}");
155        let back: AgentUplinkState = serde_json::from_str(&json).expect("decode");
156        assert_eq!(back, state);
157    }
158
159    #[test]
160    fn a_grant_driven_tcp_state_carries_the_target() {
161        let state = AgentUplinkState {
162            enabled: true,
163            kind: "tcp".to_string(),
164            target: Some("10.0.1.9:9000".to_string()),
165            source: "grant".to_string(),
166            output_write_failing: true,
167        };
168        let json = serde_json::to_string(&state).expect("encode");
169        let back: AgentUplinkState = serde_json::from_str(&json).expect("decode");
170        assert_eq!(back, state);
171        // 与其它契约同口径:多出来的键要硬失败,别把分叉静默掉。
172        assert!(
173            serde_json::from_str::<AgentUplinkState>(
174                r#"{"enabled":true,"kind":"tcp",
175               "source":"local","output_write_failing":false,"extra":1}"#
176            )
177            .is_err()
178        );
179    }
180}
181
182#[cfg(test)]
183mod tests {
184    use super::*;
185
186    #[test]
187    fn standby_round_trips_and_omits_absent_target() {
188        // 待命帧不该把空目标序列化出去:`host`/`port` 缺失就是「沿用本机」的判据。
189        let grant = AgentUplinkGrant::standby("2026-09-26T00:00:00Z".to_string());
190        let json = serde_json::to_string(&grant).expect("encode");
191        assert!(!json.contains("host"), "{json}");
192        assert!(!json.contains("port"), "{json}");
193        let back: AgentUplinkGrant = serde_json::from_str(&json).expect("decode");
194        assert_eq!(back, grant);
195        assert_eq!(back.target(), None);
196    }
197
198    #[test]
199    fn an_enabled_grant_carries_the_target_to_override() {
200        let grant = AgentUplinkGrant::enabled_at(
201            "c-001.gateway.example".to_string(),
202            9000,
203            "2026-09-26T00:00:00Z".to_string(),
204        );
205        assert_eq!(grant.target(), Some(("c-001.gateway.example", 9000)));
206        let json = serde_json::to_string(&grant).expect("encode");
207        let back: AgentUplinkGrant = serde_json::from_str(&json).expect("decode");
208        assert_eq!(back, grant);
209    }
210
211    #[test]
212    fn an_enabled_grant_without_a_target_falls_back_to_the_local_kind() {
213        // 允许「只开不指目标」:本机 kind 说了算。`target()` 必须为 None 而不是拼出空地址。
214        let grant = AgentUplinkGrant {
215            enabled: true,
216            host: None,
217            port: None,
218            granted_at: "t".to_string(),
219        };
220        assert_eq!(grant.target(), None);
221
222        // 空 host 也不能被当成目标:那只会拼出连不上的 ":9000"。
223        let blank = AgentUplinkGrant {
224            enabled: true,
225            host: Some(String::new()),
226            port: Some(9000),
227            granted_at: "t".to_string(),
228        };
229        assert_eq!(blank.target(), None);
230
231        // 只有空白的 host 同理(网关写入时会 trim,这里是第二道)。
232        let whitespace = AgentUplinkGrant {
233            enabled: true,
234            host: Some("   ".to_string()),
235            port: Some(9000),
236            granted_at: "t".to_string(),
237        };
238        assert_eq!(whitespace.target(), None);
239
240        // 前后带空白的真实主机要被裁成可用目标,而不是被判为无效。
241        let padded = AgentUplinkGrant {
242            enabled: true,
243            host: Some(" gw.example ".to_string()),
244            port: Some(9000),
245            granted_at: "t".to_string(),
246        };
247        assert_eq!(padded.target(), Some(("gw.example", 9000)));
248
249        // 有 host 但缺 port 也不算目标:缺一半就只能拼出 `host:` 这种残地址。
250        let host_only = AgentUplinkGrant {
251            enabled: true,
252            host: Some("gw.example".to_string()),
253            port: None,
254            granted_at: "t".to_string(),
255        };
256        assert_eq!(host_only.target(), None);
257
258        // 未启用的 grant 即使目标齐全也不构成覆盖。
259        let standby_with_target = AgentUplinkGrant {
260            enabled: false,
261            host: Some("gw.example".to_string()),
262            port: Some(9000),
263            granted_at: "t".to_string(),
264        };
265        assert_eq!(standby_with_target.target(), None);
266    }
267
268    #[test]
269    fn a_newer_gateway_field_is_rejected_rather_than_ignored() {
270        // 契约是两侧都解析的字节:多出来的键必须硬失败,否则「网关发了 agentd 没读的字段」
271        // 会变成静默分叉(与 work.rs 同一约定)。
272        let json = r#"{"enabled":true,"granted_at":"t","extra":1}"#;
273        assert!(serde_json::from_str::<AgentUplinkGrant>(json).is_err());
274    }
275
276    #[test]
277    fn poll_round_trips_and_rejects_unknown_fields() {
278        let poll = PollAgentUplink {
279            api_version: crate::API_VERSION_V1.to_string(),
280            kind: POLL_AGENT_UPLINK_KIND.to_string(),
281            agent_id: "agent-1".to_string(),
282            instance_id: "inst-1".to_string(),
283            requested_at: "2026-09-26T00:00:00Z".to_string(),
284        };
285        let json = serde_json::to_string(&poll).expect("encode");
286        let back: PollAgentUplink = serde_json::from_str(&json).expect("decode");
287        assert_eq!(back, poll);
288
289        let bad = r#"{"api_version":"v1","kind":"poll_agent_uplink","agent_id":"a",
290                      "instance_id":"i","requested_at":"t","extra":1}"#;
291        assert!(serde_json::from_str::<PollAgentUplink>(bad).is_err());
292    }
293}