1use helix_core::effect::DomainEventBytes;
2
3use super::{AsyncMetricSink, LabelKey, MetricEvent, MetricId, MetricLabels};
4
5const MESSAGE_V3_EVENT_PREFIX: &[u8] = b"{\"event\":\"";
6
7#[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
18pub 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
36pub 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
54pub 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
71pub 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
89pub 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
104pub 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
285fn 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
317pub(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
338fn 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
370fn 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
377fn 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
388fn 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}