Skip to main content

helix_driver_host/metrics/
business.rs

1use helix_core::effect::DomainEventBytes;
2
3use super::{AsyncMetricSink, LabelKey, MetricEvent, MetricId, MetricLabels};
4
5const MESSAGE_V3_EVENT_PREFIX: &[u8] = b"{\"event\":\"";
6
7/// ClientTerminalObservation 是边缘 UI 只读事实;所有字段都必须来自固定低基数枚举。
8#[derive(Clone, Copy, Debug, PartialEq, Eq)]
9pub struct ClientTerminalObservation {
10    pub terminal: &'static str,
11    pub path: &'static str,
12    pub role: &'static str,
13    pub result: &'static str,
14    pub platform: &'static str,
15    pub recipient_scope: &'static str,
16}
17
18/// record_message_client_terminal 将客户端终点旁路写入有界 sink,不参与业务 Tick 队列。
19pub fn record_message_client_terminal(
20    metrics: &dyn AsyncMetricSink,
21    observation: ClientTerminalObservation,
22) {
23    let labels = MetricLabels::one(LabelKey::Terminal, observation.terminal)
24        .with(LabelKey::Path, observation.path)
25        .with(LabelKey::Role, observation.role)
26        .with(LabelKey::Result, observation.result)
27        .with(LabelKey::Platform, observation.platform)
28        .with(LabelKey::RecipientScope, observation.recipient_scope);
29    let _ = metrics.try_record(MetricEvent::counter(
30        MetricId::ImMessageClientTerminalTotal,
31        1.0,
32        labels,
33    ));
34}
35
36/// record_message_e2e_view_updated 记录同一客户端单调时钟产生的真实端到端样本。
37pub fn record_message_e2e_view_updated(
38    metrics: &dyn AsyncMetricSink,
39    seconds: f64,
40    path: &'static str,
41    result: &'static str,
42    platform: &'static str,
43) {
44    let labels = MetricLabels::one(LabelKey::Path, path)
45        .with(LabelKey::Result, result)
46        .with(LabelKey::Platform, platform);
47    let _ = metrics.try_record(MetricEvent::histogram(
48        MetricId::ImMessageE2eViewUpdatedDurationSeconds,
49        seconds,
50        labels,
51    ));
52}
53
54/// record_seq_observation 记录已由本地同步状态机判定的固定 Seq 结果。
55pub fn record_seq_observation(
56    metrics: &dyn AsyncMetricSink,
57    outcome: &'static str,
58    path: &'static str,
59    operation: &'static str,
60) {
61    let labels = MetricLabels::one(LabelKey::State, outcome)
62        .with(LabelKey::Path, path)
63        .with(LabelKey::Operation, operation);
64    let _ = metrics.try_record(MetricEvent::counter(
65        MetricId::ImSeqObservationTotal,
66        1.0,
67        labels,
68    ));
69}
70
71/// record_gap_duration 仅在 gap lifecycle 形成稳定终态时记录持续时间。
72pub fn record_gap_duration(
73    metrics: &dyn AsyncMetricSink,
74    seconds: f64,
75    terminal: &'static str,
76    path: &'static str,
77    result: &'static str,
78) {
79    let labels = MetricLabels::one(LabelKey::Terminal, terminal)
80        .with(LabelKey::Path, path)
81        .with(LabelKey::Result, result);
82    let _ = metrics.try_record(MetricEvent::histogram(
83        MetricId::ImGapDurationSeconds,
84        seconds,
85        labels,
86    ));
87}
88
89/// record_tick_stage_duration 记录已有 Tick 生命周期边界计算出的阶段样本。
90pub fn record_tick_stage_duration(
91    metrics: &dyn AsyncMetricSink,
92    seconds: f64,
93    stage: &'static str,
94    result: &'static str,
95) {
96    let labels = MetricLabels::one(LabelKey::Stage, stage).with(LabelKey::Result, result);
97    let _ = metrics.try_record(MetricEvent::histogram(
98        MetricId::ImTickStageDurationSeconds,
99        seconds,
100        labels,
101    ));
102}
103
104/// record_im_business_event 在 host Emit 边界记录投影、同步与恢复终态,不把远端推送冒充本地 Command。
105pub fn record_im_business_event(metrics: &dyn AsyncMetricSink, event: &DomainEventBytes) {
106    if !metrics.is_enabled() {
107        return;
108    }
109    let Some(event_name) = message_v3_event_name(event.0.as_ref()) else {
110        return;
111    };
112
113    match event_name {
114        b"im:post:received" | b"im:post:sent" => {
115            record_terminal(
116                metrics,
117                MetricId::ImProjectionTerminalTotal,
118                "send_message",
119                "success",
120                "none",
121            );
122            record_terminal(
123                metrics,
124                MetricId::ImMessageCorrectnessTotal,
125                "message_projection",
126                "success",
127                "none",
128            );
129        }
130        b"im:post:send-failed" => {}
131        b"im:post:client-ack-succeeded" => record_terminal(
132            metrics,
133            MetricId::ImClientAckTerminalTotal,
134            "client_ack",
135            "success",
136            "none",
137        ),
138        b"im:post:client-ack-failed" => record_terminal(
139            metrics,
140            MetricId::ImClientAckTerminalTotal,
141            "client_ack",
142            "error",
143            "ack_failed",
144        ),
145        b"im:post:increment-failed" => {
146            record_terminal(
147                metrics,
148                MetricId::ImProjectionTerminalTotal,
149                "channel_increment",
150                "error",
151                "increment_failed",
152            );
153            record_terminal(
154                metrics,
155                MetricId::ImSyncAnomalyTotal,
156                "channel_increment",
157                "error",
158                "increment_failed",
159            );
160        }
161        b"im:channel-sync-complete" => {
162            record_terminal(
163                metrics,
164                MetricId::ImProjectionTerminalTotal,
165                "channel_sync",
166                "success",
167                "none",
168            );
169            record_terminal(
170                metrics,
171                MetricId::ImSyncSessionTotal,
172                "channel_sync",
173                "success",
174                "none",
175            );
176        }
177        b"im:sync:loaded" => record_terminal(
178            metrics,
179            MetricId::ImSyncSessionTotal,
180            "startup_sync",
181            "success",
182            "none",
183        ),
184        b"im:sync:recovered" => {
185            let payload = serde_json::from_slice::<serde_json::Value>(event.0.as_ref()).ok();
186            let state = payload
187                .as_ref()
188                .and_then(|v| v.get("data"))
189                .and_then(|v| v.get("state"))
190                .and_then(serde_json::Value::as_str);
191            let (status, error) = match state {
192                Some("recovered") => ("success", "none"),
193                Some("failed") => ("error", "recovery_failed"),
194                _ => return,
195            };
196            record_terminal(
197                metrics,
198                MetricId::ImSyncSessionTotal,
199                "offline_recovery",
200                status,
201                error,
202            );
203            record_terminal(
204                metrics,
205                MetricId::ImRecoverySessionTotal,
206                "offline_recovery",
207                status,
208                error,
209            );
210        }
211        b"im:sync:gap-repaired" => {
212            record_terminal(
213                metrics,
214                MetricId::ImRecoverySessionTotal,
215                "gap_repair",
216                "success",
217                "none",
218            );
219            record_terminal(
220                metrics,
221                MetricId::ImMessageCorrectnessTotal,
222                "gap_repair",
223                "success",
224                "gap_repaired",
225            );
226        }
227        b"im:sync:channel-hydrated" => record_terminal(
228            metrics,
229            MetricId::ImRecoverySessionTotal,
230            "channel_hydration",
231            "success",
232            "none",
233        ),
234        b"im:sync:too_long" => {
235            record_terminal(
236                metrics,
237                MetricId::ImSyncSessionTotal,
238                "channel_sync",
239                "error",
240                "too_long",
241            );
242            record_terminal(
243                metrics,
244                MetricId::ImSyncAnomalyTotal,
245                "channel_sync",
246                "error",
247                "too_long",
248            );
249            record_terminal(
250                metrics,
251                MetricId::ImMessageCorrectnessTotal,
252                "channel_sync",
253                "error",
254                "gap",
255            );
256        }
257        b"im:post:updated" | b"im:post:updates" | b"im:post:batch-updated" => {
258            record_projection(metrics, "update_message")
259        }
260        b"im:post:revoke" => record_projection(metrics, "revoke_message"),
261        b"im:post:deleted" => record_projection(metrics, "delete_message"),
262        b"im:post:read" | b"im:post:readers" | b"im:channel:read_echo" => {
263            record_projection(metrics, "mark_read")
264        }
265        b"im:channel:created" => record_projection(metrics, "create_channel"),
266        b"im:channel:closed" => record_projection(metrics, "close_channel"),
267        b"im:channel:schedule-created" => record_projection(metrics, "create_schedule"),
268        b"im:channel:schedule-canceled" => record_projection(metrics, "cancel_schedule"),
269        b"im:channel:member-updated" | b"im:channel:member-nickname" => {
270            record_projection(metrics, "update_channel_member")
271        }
272        b"im:channel:settings-updated" => record_projection(metrics, "update_channel_settings"),
273        b"im:todo:updated" => record_projection(metrics, "update_todo"),
274        b"im:post_chain:publish"
275        | b"im:post_chain:upsert"
276        | b"im:post_chain:close"
277        | b"im:post_chain:retract"
278        | b"im:post_chain:read_cursor" => record_projection(metrics, "post_chain"),
279        b"im:post_chain:append_rejected" => {}
280        b"im:read:result" => record_unread_reconcile(metrics, event),
281        _ => {}
282    }
283}
284
285/// record_unread_reconcile 只接受 hydration 终态中同一 authority snapshot 的双侧绝对值结果。
286fn record_unread_reconcile(metrics: &dyn AsyncMetricSink, event: &DomainEventBytes) {
287    let Ok(payload) = serde_json::from_slice::<serde_json::Value>(event.0.as_ref()) else {
288        return;
289    };
290    let Some(status) = payload
291        .pointer("/data/body/unreadReconcile/status")
292        .and_then(serde_json::Value::as_str)
293    else {
294        tracing::debug!("unread reconcile terminal absent from read result");
295        return;
296    };
297    tracing::info!(status, "unread reconcile terminal recorded");
298    match status {
299        "match" => record_terminal(
300            metrics,
301            MetricId::ImUnreadReconcileTotal,
302            "unread_reconcile",
303            "match",
304            "none",
305        ),
306        "mismatch" => record_terminal(
307            metrics,
308            MetricId::ImUnreadReconcileTotal,
309            "unread_reconcile",
310            "mismatch",
311            "unread_mismatch",
312        ),
313        _ => {}
314    }
315}
316
317/// record_im_command_terminal_event 仅在 Engine 已证明本地 Command lineage 时记录首个稳定终态。
318pub(crate) fn record_im_command_terminal_event(
319    metrics: &dyn AsyncMetricSink,
320    event: &DomainEventBytes,
321) -> bool {
322    let Some(event_name) = message_v3_event_name(event.0.as_ref()) else {
323        return false;
324    };
325    let Some((operation, status, error_kind)) = command_terminal(event_name) else {
326        return false;
327    };
328    record_terminal(
329        metrics,
330        MetricId::ImCommandTerminalTotal,
331        operation,
332        status,
333        error_kind,
334    );
335    true
336}
337
338/// command_terminal 把固定 MessageV3 终态事件映射成低基数 Command 结果。
339fn command_terminal(event_name: &[u8]) -> Option<(&'static str, &'static str, &'static str)> {
340    match event_name {
341        b"im:post:received" | b"im:post:sent" => Some(("send_message", "success", "none")),
342        b"im:post:send-failed" => Some(("send_message", "error", "send_failed")),
343        b"im:post:updated" | b"im:post:updates" | b"im:post:batch-updated" => {
344            Some(("update_message", "success", "none"))
345        }
346        b"im:post:revoke" => Some(("revoke_message", "success", "none")),
347        b"im:post:deleted" => Some(("delete_message", "success", "none")),
348        b"im:post:read" | b"im:post:readers" | b"im:channel:read_echo" => {
349            Some(("mark_read", "success", "none"))
350        }
351        b"im:channel:created" => Some(("create_channel", "success", "none")),
352        b"im:channel:closed" => Some(("close_channel", "success", "none")),
353        b"im:channel:schedule-created" => Some(("create_schedule", "success", "none")),
354        b"im:channel:schedule-canceled" => Some(("cancel_schedule", "success", "none")),
355        b"im:channel:member-updated" | b"im:channel:member-nickname" => {
356            Some(("update_channel_member", "success", "none"))
357        }
358        b"im:channel:settings-updated" => Some(("update_channel_settings", "success", "none")),
359        b"im:todo:updated" => Some(("update_todo", "success", "none")),
360        b"im:post_chain:publish"
361        | b"im:post_chain:upsert"
362        | b"im:post_chain:close"
363        | b"im:post_chain:retract"
364        | b"im:post_chain:read_cursor" => Some(("post_chain", "success", "none")),
365        b"im:post_chain:append_rejected" => Some(("post_chain", "error", "append_rejected")),
366        _ => None,
367    }
368}
369
370/// message_v3_event_name 只读取 canonical envelope 的有界 event 前缀,拒绝非标准或截断输入。
371fn message_v3_event_name(bytes: &[u8]) -> Option<&[u8]> {
372    let rest = bytes.strip_prefix(MESSAGE_V3_EVENT_PREFIX)?;
373    let end = rest.iter().position(|byte| *byte == b'"')?;
374    (end <= 96).then_some(&rest[..end])
375}
376
377/// record_projection 为 durable projection 记录唯一投影终态,Command 终态由 Engine lineage 单独闭合。
378fn record_projection(metrics: &dyn AsyncMetricSink, operation: &'static str) {
379    record_terminal(
380        metrics,
381        MetricId::ImProjectionTerminalTotal,
382        operation,
383        "success",
384        "none",
385    );
386}
387
388/// record_terminal 统一写入 operation/status/error_kind 三个静态低基数标签。
389fn record_terminal(
390    metrics: &dyn AsyncMetricSink,
391    id: MetricId,
392    operation: &'static str,
393    status: &'static str,
394    error_kind: &'static str,
395) {
396    let labels = MetricLabels::one(LabelKey::Operation, operation)
397        .with(LabelKey::Status, status)
398        .with(LabelKey::ErrorKind, error_kind);
399    let _ = metrics.try_record(MetricEvent::counter(id, 1.0, labels));
400}