nx_telemetry_bus.nx source
↩ module page · 564 lines · 22709 B
1// nx_telemetry.nx -- substrate telemetry bus.
2//
3// Sealed-enum-typed event bus for substrate-internal observability.
4// Closes the OBD predictive-maintenance pipeline loop (commit
5// 87dca0a0): consumer emits anomaly verdicts -> telemetry bus ->
6// pluggable sink (stdout JSON V1; future MQTT shim / syslog / file
7// when those sinks land).
8//
9// Per INTEROPERABILITY_CHARTER.md the same telemetry bus consumes
10// events from EVERY protocol shim's diag_fn + EVERY substrate-side
11// consumer + the silicon-feedback hot-report. One uniform format
12// across the whole substrate.
13//
14// Status: SEED v0.1.0. 2026-05-27.
15//
16// WINNER-TIER: WINNER-A CANDIDATE
17// INCUMBENTS: syslog (RFC 5424), Prometheus client libs,
18// OpenTelemetry (OTLP), structlog (Python), statsd
19// protocol, Vector.dev event router
20// NUMBERS: V1 ships the bus + stdout-JSON sink; throughput
21// measurement vs structlog/Prometheus pending paired
22// bench
23// GAP: vs OpenTelemetry the gap is feature coverage (no
24// distributed-trace context yet); vs syslog the gap is
25// transport (V1 is in-process; future MQTT/syslog
26// sinks close it); vs structlog Nishi WINS on
27// sealed-enum typing (structlog is string-keyed and
28// loses type safety)
29// PLAN: M-next: paired bench vs structlog on 1M-event burst;
30// re-rate; add distributed-trace context for
31// OpenTelemetry parity
32// EXEMPTION REASON: n/a; provisional pending measurement
33//
34// V1 SCOPE:
35// - 8 sealed event kinds (covers anomaly / perf / health /
36// shim_diag / silicon_feedback / substrate_panic / user / debug)
37// - 4 sealed severity levels (INFO / WARN / ERROR / CRITICAL)
38// - NxTelemetryEvent struct (kind / severity / source-string /
39// ts_ns / 3 value fields for structured numeric payload)
40// - Pluggable NxTelemetrySink (caller-supplied emit function)
41// - Built-in stdout-JSON sink (V1 default; one line per event)
42// - Token-bucket rate limiter per source (avoids flood on
43// storm conditions; configurable cap + refill rate)
44// - Sealed verdict surface (5 codes)
45//
46// V2+ SCOPE:
47// - MQTT sink (publishes to broker; uses nx_mqtt_shim once that
48// lands in the ยง5.5 interop family)
49// - syslog sink (RFC 5424 over UDP per nx_syslog_shim)
50// - file sink (append-only NDJSON; rotates by size)
51// - Distributed-trace context propagation (parent_span_id etc;
52// OpenTelemetry parity)
53// - Per-event-kind sampling rates (full INFO sampling vs 100%
54// CRITICAL retention)
55// - Histogram metrics (latency distributions for shim_diag events)
56
57import "nx_syscalls.nx"
58
59// ===== Sealed verdict surface =================================================
60const NX_TM_OK: i64 = 0
61const NX_TM_BAD_KIND: i64 = 1 // event kind out of sealed enum
62const NX_TM_BAD_SEVERITY: i64 = 2
63const NX_TM_RATE_LIMITED: i64 = 3 // token bucket exhausted
64const NX_TM_SINK_FAULT: i64 = 4 // configured sink fn returned non-OK
65const NX_TM_BAD_SINK_KIND: i64 = 5 // sink kind out of sealed enum
66const NX_TM_FANOUT_OVERFLOW: i64 = 6 // fanout array exceeds NX_TM_FANOUT_MAX
67
68// ===== Sealed sink-kind enum (V2 fanout dispatch) =================================================
69//
70// Bus dispatches per sink.kind in nx_tm_emit. Adding a kind
71// requires MINOR bump per migration policy. Cardinal: sealed-enum
72// dispatch over function-pointer-callback to preserve compile-time
73// type safety + avoid circular imports between bus + per-sink modules.
74
75const NX_TM_SINK_KIND_STDOUT_JSON: i64 = 0 // V1; built-in NDJSON sink
76const NX_TM_SINK_KIND_NULL: i64 = 1 // no-op (test/silence)
77const NX_TM_SINK_KIND_FANOUT: i64 = 2 // dispatches to N sub-sinks
78const NX_TM_SINK_KIND_BUFFER: i64 = 3 // accumulates events for operator drain
79 // (use case: MQTT-via-operator's
80 // socket write; bus doesn't import
81 // the MQTT module, avoiding the
82 // circular-import problem)
83const NX_TM_SINK_KIND_N: i64 = 4
84
85func nx_tm_sink_kind_is_valid(k: i64) -> i64 {
86 if k < 0 { return 0 }
87 if k >= NX_TM_SINK_KIND_N { return 0 }
88 return 1
89}
90
91// Fanout cap: a sealed-cap V1 keeps the dispatch loop bounded
92// per [[feedback-loop-discipline-cardinal-2026-05-21]].
93const NX_TM_FANOUT_MAX: i64 = 8
94
95// ===== Sealed event-kind enum =================================================
96//
97// Adding a new kind requires a MINOR version bump per the standard
98// migration policy. Eight kinds cover every substrate-internal
99// observability source today.
100
101const NX_TM_KIND_ANOMALY: i64 = 0 // pred-maint anomaly verdict
102const NX_TM_KIND_PERF_SAMPLE: i64 = 1 // bench / silicon-feedback per-PC sample
103const NX_TM_KIND_HEALTH: i64 = 2 // substrate process / kernel health metric
104const NX_TM_KIND_SHIM_DIAG: i64 = 3 // protocol shim diag_fn output
105const NX_TM_KIND_SILICON_FB: i64 = 4 // hot-report / multi-tier ROI verdict
106const NX_TM_KIND_PANIC: i64 = 5 // substrate panic / unrecoverable verdict
107const NX_TM_KIND_USER: i64 = 6 // operator-emitted (debug / annotation)
108const NX_TM_KIND_DEBUG: i64 = 7 // dev-only; suppressed in production builds
109const NX_TM_KIND_N: i64 = 8
110
111func nx_tm_kind_is_valid(k: i64) -> i64 {
112 if k < 0 { return 0 }
113 if k >= NX_TM_KIND_N { return 0 }
114 return 1
115}
116
117// ===== Sealed severity enum =================================================
118const NX_TM_SEV_INFO: i64 = 0
119const NX_TM_SEV_WARN: i64 = 1
120const NX_TM_SEV_ERROR: i64 = 2
121const NX_TM_SEV_CRITICAL: i64 = 3
122const NX_TM_SEV_N: i64 = 4
123
124func nx_tm_sev_is_valid(s: i64) -> i64 {
125 if s < 0 { return 0 }
126 if s >= NX_TM_SEV_N { return 0 }
127 return 1
128}
129
130// ===== Event payload =================================================
131//
132// Fixed-width struct (caller-friendly; no per-event alloc). Three
133// value fields cover the common case (single scalar, two-tuple,
134// or three-tuple); structured payloads beyond that pass a
135// caller-allocated string in `extra_str`.
136
137struct NxTelemetryEvent {
138 kind: i64 // NX_TM_KIND_*
139 severity: i64 // NX_TM_SEV_*
140 source_str: *u8 // caller-owned UTF-8 (e.g., "nx_obd2_shim")
141 source_len: i64
142 ts_ns: i64 // ns since epoch (or boot for embedded targets)
143 value0: i64 // scalar payload slot 0
144 value1: i64
145 value2: i64
146 extra_str: *u8 // optional structured payload (NULL = absent)
147 extra_len: i64
148}
149
150// ===== Sink contract =================================================
151//
152// A sink is a (caller-supplied via init) emit function pointer.
153// V1 ships one built-in sink: stdout-JSON. Future commits add
154// MQTT / syslog / file sinks once the corresponding shims land.
155
156struct NxTelemetrySink {
157 kind: i64 // NX_TM_SINK_KIND_* (sealed enum)
158 rate_cap: i64 // tokens before refill (events/sec budget)
159 // Fanout-kind fields (used iff kind == NX_TM_SINK_KIND_FANOUT)
160 fanout_sinks: *i64 // array of *NxTelemetrySink (stored as i64
161 // handles to avoid pointer-of-pointer
162 // syntax friction)
163 fanout_count: i64
164 // Buffer-kind fields (used iff kind == NX_TM_SINK_KIND_BUFFER)
165 buf_storage: *NxTelemetryEvent // ring of recent events
166 buf_cap: i64
167 buf_write_idx: i64
168 buf_n_filled: i64
169 valid: i64
170}
171
172// ===== Sink-kind init helpers =================================================
173//
174// Each kind has a focused init function so the operator never has
175// to populate the wrong field for the wrong kind.
176
177func nx_tm_sink_init_stdout_json(sink: *NxTelemetrySink, rate_cap: i64) -> i64 {
178 if (sink as i64) == 0 { return 0 - NX_TM_BAD_SINK_KIND }
179 sink.kind = NX_TM_SINK_KIND_STDOUT_JSON
180 sink.rate_cap = rate_cap
181 sink.fanout_sinks = (0 as i64) as *i64
182 sink.fanout_count = 0
183 sink.buf_storage = (0 as i64) as *NxTelemetryEvent
184 sink.buf_cap = 0
185 sink.buf_write_idx = 0
186 sink.buf_n_filled = 0
187 sink.valid = 1
188 return NX_TM_OK
189}
190
191func nx_tm_sink_init_null(sink: *NxTelemetrySink) -> i64 {
192 if (sink as i64) == 0 { return 0 - NX_TM_BAD_SINK_KIND }
193 sink.kind = NX_TM_SINK_KIND_NULL
194 sink.rate_cap = 0
195 sink.fanout_sinks = (0 as i64) as *i64
196 sink.fanout_count = 0
197 sink.buf_storage = (0 as i64) as *NxTelemetryEvent
198 sink.buf_cap = 0
199 sink.buf_write_idx = 0
200 sink.buf_n_filled = 0
201 sink.valid = 1
202 return NX_TM_OK
203}
204
205func nx_tm_sink_init_fanout(sink: *NxTelemetrySink,
206 sub_sinks: *i64, n_subs: i64) -> i64 {
207 if (sink as i64) == 0 { return 0 - NX_TM_BAD_SINK_KIND }
208 if n_subs < 0 { return 0 - NX_TM_BAD_SINK_KIND }
209 if n_subs > NX_TM_FANOUT_MAX { return 0 - NX_TM_FANOUT_OVERFLOW }
210 sink.kind = NX_TM_SINK_KIND_FANOUT
211 sink.rate_cap = 0
212 sink.fanout_sinks = sub_sinks
213 sink.fanout_count = n_subs
214 sink.buf_storage = (0 as i64) as *NxTelemetryEvent
215 sink.buf_cap = 0
216 sink.buf_write_idx = 0
217 sink.buf_n_filled = 0
218 sink.valid = 1
219 return NX_TM_OK
220}
221
222func nx_tm_sink_init_buffer(sink: *NxTelemetrySink,
223 buf_storage: *NxTelemetryEvent, buf_cap: i64) -> i64 {
224 if (sink as i64) == 0 { return 0 - NX_TM_BAD_SINK_KIND }
225 if (buf_storage as i64) == 0 { return 0 - NX_TM_BAD_SINK_KIND }
226 if buf_cap < 1 { return 0 - NX_TM_BAD_SINK_KIND }
227 sink.kind = NX_TM_SINK_KIND_BUFFER
228 sink.rate_cap = 0
229 sink.fanout_sinks = (0 as i64) as *i64
230 sink.fanout_count = 0
231 sink.buf_storage = buf_storage
232 sink.buf_cap = buf_cap
233 sink.buf_write_idx = 0
234 sink.buf_n_filled = 0
235 sink.valid = 1
236 return NX_TM_OK
237}
238
239// ===== Buffer-kind drain (for operator's external transport) =================================================
240//
241// The BUFFER sink accumulates events in a ring; operator polls
242// via nx_tm_sink_buffer_drain to copy out + reset write position.
243// Used for the MQTT-via-operator-socket pattern: bus pushes events
244// into the buffer; operator's tick loop drains + transmits.
245
246func nx_tm_sink_buffer_drain(sink: *NxTelemetrySink,
247 out_events: *NxTelemetryEvent, out_cap: i64) -> i64 {
248 if sink.valid != 1 { return 0 - NX_TM_BAD_SINK_KIND }
249 if sink.kind != NX_TM_SINK_KIND_BUFFER { return 0 - NX_TM_BAD_SINK_KIND }
250 let n: i64 = sink.buf_n_filled
251 let copy_n: i64 = if n < out_cap then n else out_cap
252 // Copy from ring starting at oldest entry (V1: write_idx-N mod cap;
253 // for simplicity, drain all of buf_n_filled in write order).
254 var i: i64 = 0
255 while i < copy_n {
256 let src_idx: i64 = (sink.buf_write_idx + NX_TM_FANOUT_MAX - n + i) % sink.buf_cap
257 out_events[i].kind = sink.buf_storage[src_idx].kind
258 out_events[i].severity = sink.buf_storage[src_idx].severity
259 out_events[i].source_str = sink.buf_storage[src_idx].source_str
260 out_events[i].source_len = sink.buf_storage[src_idx].source_len
261 out_events[i].ts_ns = sink.buf_storage[src_idx].ts_ns
262 out_events[i].value0 = sink.buf_storage[src_idx].value0
263 out_events[i].value1 = sink.buf_storage[src_idx].value1
264 out_events[i].value2 = sink.buf_storage[src_idx].value2
265 out_events[i].extra_str = sink.buf_storage[src_idx].extra_str
266 out_events[i].extra_len = sink.buf_storage[src_idx].extra_len
267 i = i + 1
268 }
269 sink.buf_n_filled = 0
270 sink.buf_write_idx = 0
271 return copy_n
272}
273
274// ===== Token-bucket rate limiter =================================================
275//
276// Per-source token bucket. Each source string hashes to a slot;
277// each slot has (tokens, last_refill_ns). When emit attempted:
278// - refill tokens at rate_cap per second since last_refill_ns
279// - if tokens > 0: emit + tokens--
280// - if tokens == 0: return NX_TM_RATE_LIMITED
281//
282// V1 fixed 64-slot table; future commit grows to operator-tunable.
283
284const NX_TM_RL_SLOTS: i64 = 64
285
286struct NxTelemetryRateLimiter {
287 slot_keys: *i64 // 64 i64 hashes
288 slot_tokens: *i64 // 64 i64 current token count
289 slot_last_refill: *i64 // 64 i64 ns timestamps
290 default_cap: i64 // default tokens per source
291 valid: i64
292}
293
294// Tiny FNV-1a hash for source-string -> slot index.
295func nx_tm_hash_source(src: *u8, len: i64) -> i64 {
296 var h: i64 = 0xcbf29ce484222325
297 var i: i64 = 0
298 while i < len {
299 h = h ^ (src[i] as i64)
300 h = h * 1099511628211
301 i = i + 1
302 }
303 return h
304}
305
306func nx_tm_rl_init(rl: *NxTelemetryRateLimiter,
307 keys_buf: *i64, tokens_buf: *i64, refill_buf: *i64,
308 default_cap: i64) -> i64 {
309 if (rl as i64) == 0 { return 0 - NX_TM_BAD_KIND }
310 rl.slot_keys = keys_buf
311 rl.slot_tokens = tokens_buf
312 rl.slot_last_refill = refill_buf
313 rl.default_cap = default_cap
314 rl.valid = 1
315 var i: i64 = 0
316 while i < NX_TM_RL_SLOTS {
317 keys_buf[i] = 0
318 tokens_buf[i] = default_cap
319 refill_buf[i] = 0
320 i = i + 1
321 }
322 return NX_TM_OK
323}
324
325// Try to consume one token for source. Returns NX_TM_OK or
326// NX_TM_RATE_LIMITED.
327func nx_tm_rl_try_consume(rl: *NxTelemetryRateLimiter,
328 source_str: *u8, source_len: i64,
329 now_ns: i64) -> i64 {
330 if rl.valid != 1 { return 0 - NX_TM_BAD_KIND }
331 let h: i64 = nx_tm_hash_source(source_str, source_len)
332 let slot: i64 = (h & 0x7fffffffffffffff) % NX_TM_RL_SLOTS
333 // Refill: tokens += elapsed_seconds * default_cap, cap at default_cap
334 let elapsed_ns: i64 = now_ns - rl.slot_last_refill[slot]
335 if elapsed_ns > 0 {
336 // refill_amount = (elapsed_ns / 1_000_000_000) * default_cap
337 let refill: i64 = (elapsed_ns / 1000000000) * rl.default_cap
338 if refill > 0 {
339 let new_tokens: i64 = rl.slot_tokens[slot] + refill
340 if new_tokens > rl.default_cap {
341 rl.slot_tokens[slot] = rl.default_cap
342 }
343 if new_tokens <= rl.default_cap {
344 rl.slot_tokens[slot] = new_tokens
345 }
346 rl.slot_last_refill[slot] = now_ns
347 }
348 }
349 if rl.slot_tokens[slot] > 0 {
350 rl.slot_tokens[slot] = rl.slot_tokens[slot] - 1
351 return NX_TM_OK
352 }
353 return 0 - NX_TM_RATE_LIMITED
354}
355
356// ===== Built-in stdout-JSON sink =================================================
357//
358// Writes one NDJSON line per event to fd 1 (stdout). Format:
359// {"kind":N,"sev":N,"source":"...","ts":N,"v0":N,"v1":N,"v2":N}
360//
361// Operator-tunable: future commit adds pretty-print + per-kind
362// formatter overrides. V1 is fixed-format for parser stability.
363
364func nx_tm_emit_int(out_buf: *u8, out_off: i64, value: i64) -> i64 {
365 // Render i64 to out_buf starting at out_off; returns new off.
366 var off: i64 = out_off
367 var n: i64 = value
368 if n < 0 {
369 out_buf[off] = 0x2d as u8 // '-'
370 off = off + 1
371 n = 0 - n
372 }
373 if n == 0 {
374 out_buf[off] = 0x30 as u8 // '0'
375 return off + 1
376 }
377 // Build digits backward then reverse-copy into out_buf.
378 let digits: *u8 = sys_mmap(24)
379 var k: i64 = 0
380 while n > 0 {
381 digits[k] = (0x30 + (n % 10)) as u8
382 n = n / 10
383 k = k + 1
384 }
385 var i: i64 = k - 1
386 while i >= 0 {
387 out_buf[off] = digits[i]
388 off = off + 1
389 i = i - 1
390 }
391 return off
392}
393
394func nx_tm_emit_str(out_buf: *u8, out_off: i64, src: *u8, src_len: i64) -> i64 {
395 var off: i64 = out_off
396 var i: i64 = 0
397 while i < src_len {
398 out_buf[off] = src[i]
399 off = off + 1
400 i = i + 1
401 }
402 return off
403}
404
405func nx_tm_sink_stdout_json(event: *NxTelemetryEvent) -> i64 {
406 let buf: *u8 = sys_mmap(512)
407 var off: i64 = 0
408 off = nx_tm_emit_str(buf, off, "{\"kind\":" as *u8, 8)
409 off = nx_tm_emit_int(buf, off, event.kind)
410 off = nx_tm_emit_str(buf, off, ",\"sev\":" as *u8, 7)
411 off = nx_tm_emit_int(buf, off, event.severity)
412 off = nx_tm_emit_str(buf, off, ",\"source\":\"" as *u8, 11)
413 off = nx_tm_emit_str(buf, off, event.source_str, event.source_len)
414 off = nx_tm_emit_str(buf, off, "\",\"ts\":" as *u8, 7)
415 off = nx_tm_emit_int(buf, off, event.ts_ns)
416 off = nx_tm_emit_str(buf, off, ",\"v0\":" as *u8, 6)
417 off = nx_tm_emit_int(buf, off, event.value0)
418 off = nx_tm_emit_str(buf, off, ",\"v1\":" as *u8, 6)
419 off = nx_tm_emit_int(buf, off, event.value1)
420 off = nx_tm_emit_str(buf, off, ",\"v2\":" as *u8, 6)
421 off = nx_tm_emit_int(buf, off, event.value2)
422 off = nx_tm_emit_str(buf, off, "}\n" as *u8, 2)
423 sys_write(1, buf, off)
424 return NX_TM_OK
425}
426
427// ===== Top-level emit =================================================
428//
429// Public entry point: callers (pred-maint consumer, shim diag_fn,
430// silicon-feedback report) call this once per observable event.
431// V1 dispatches to the configured sink directly; future variants
432// add async queue + worker thread.
433
434struct NxTelemetryBus {
435 sink: *NxTelemetrySink
436 rl: *NxTelemetryRateLimiter
437 emit_count: i64
438 drop_count: i64
439 valid: i64
440}
441
442func nx_tm_bus_init(bus: *NxTelemetryBus,
443 sink: *NxTelemetrySink, rl: *NxTelemetryRateLimiter) -> i64 {
444 if (bus as i64) == 0 { return 0 - NX_TM_BAD_KIND }
445 bus.sink = sink
446 bus.rl = rl
447 bus.emit_count = 0
448 bus.drop_count = 0
449 bus.valid = 1
450 return NX_TM_OK
451}
452
453// ===== Per-sink dispatch (V2 sealed-kind switch) =================================================
454//
455// Helper called from nx_tm_emit (bus path) AND recursively from the
456// FANOUT kind. Recursion bounded by NX_TM_FANOUT_MAX per cardinal
457// loop-discipline.
458
459func nx_tm_dispatch_to_sink(sink: *NxTelemetrySink,
460 event: *NxTelemetryEvent) -> i64 {
461 if sink.valid != 1 { return 0 - NX_TM_BAD_SINK_KIND }
462 if nx_tm_sink_kind_is_valid(sink.kind) != 1 { return 0 - NX_TM_BAD_SINK_KIND }
463
464 if sink.kind == NX_TM_SINK_KIND_STDOUT_JSON {
465 return nx_tm_sink_stdout_json(event)
466 }
467 if sink.kind == NX_TM_SINK_KIND_NULL {
468 // No-op (event silently swallowed). Useful for test
469 // harnesses + temporarily silencing a noisy source.
470 return NX_TM_OK
471 }
472 if sink.kind == NX_TM_SINK_KIND_FANOUT {
473 var i: i64 = 0
474 var last_rc: i64 = NX_TM_OK
475 while i < sink.fanout_count {
476 let sub_handle: i64 = sink.fanout_sinks[i]
477 let sub_sink: *NxTelemetrySink = sub_handle as *NxTelemetrySink
478 // Best-effort: continue on per-sub failure; record last
479 // verdict for caller diagnostics. V2+ could escalate to
480 // bus.fanout_failure_count.
481 let sub_rc: i64 = nx_tm_dispatch_to_sink(sub_sink, event)
482 if sub_rc != NX_TM_OK { last_rc = sub_rc }
483 i = i + 1
484 }
485 return last_rc
486 }
487 if sink.kind == NX_TM_SINK_KIND_BUFFER {
488 if sink.buf_cap < 1 { return 0 - NX_TM_SINK_FAULT }
489 let slot: i64 = sink.buf_write_idx
490 sink.buf_storage[slot].kind = event.kind
491 sink.buf_storage[slot].severity = event.severity
492 sink.buf_storage[slot].source_str = event.source_str
493 sink.buf_storage[slot].source_len = event.source_len
494 sink.buf_storage[slot].ts_ns = event.ts_ns
495 sink.buf_storage[slot].value0 = event.value0
496 sink.buf_storage[slot].value1 = event.value1
497 sink.buf_storage[slot].value2 = event.value2
498 sink.buf_storage[slot].extra_str = event.extra_str
499 sink.buf_storage[slot].extra_len = event.extra_len
500 sink.buf_write_idx = (slot + 1) % sink.buf_cap
501 if sink.buf_n_filled < sink.buf_cap {
502 sink.buf_n_filled = sink.buf_n_filled + 1
503 }
504 return NX_TM_OK
505 }
506 return 0 - NX_TM_BAD_SINK_KIND
507}
508
509func nx_tm_emit(bus: *NxTelemetryBus, event: *NxTelemetryEvent) -> i64 {
510 if bus.valid != 1 { return 0 - NX_TM_BAD_KIND }
511 if nx_tm_kind_is_valid(event.kind) != 1 { return 0 - NX_TM_BAD_KIND }
512 if nx_tm_sev_is_valid(event.severity) != 1 { return 0 - NX_TM_BAD_SEVERITY }
513
514 // Rate-limit check (bus-level; per-source token bucket).
515 let rl_rc: i64 = nx_tm_rl_try_consume(bus.rl,
516 event.source_str, event.source_len,
517 event.ts_ns)
518 if rl_rc != NX_TM_OK {
519 bus.drop_count = bus.drop_count + 1
520 return rl_rc
521 }
522
523 // V2: dispatch via sealed-kind switch in nx_tm_dispatch_to_sink.
524 // The FANOUT kind recursively walks an array of sub-sinks so
525 // operator wires the bus ONCE; bus delivers to all sinks
526 // automatically (replaces the V1 manual-fanout pattern).
527 let sink_rc: i64 = nx_tm_dispatch_to_sink(bus.sink, event)
528 if sink_rc != NX_TM_OK {
529 return 0 - NX_TM_SINK_FAULT
530 }
531
532 bus.emit_count = bus.emit_count + 1
533 return NX_TM_OK
534}
535
536// ===== Convenience: emit anomaly verdict from pred-maint =================================================
537//
538// Wrapper for the most common caller (nx_pred_maint_consumer) so
539// the consumer doesn't have to build the full NxTelemetryEvent
540// struct per emit. Captures the operator's "X days early"
541// metric as value2 for the operator dashboard.
542
543const NX_TM_SOURCE_PRED_MAINT: *u8 = "nx_pred_maint_consumer" as *u8
544const NX_TM_SOURCE_PRED_MAINT_LEN: i64 = 22
545
546func nx_tm_emit_anomaly(bus: *NxTelemetryBus, pid: i64, anom_kind: i64,
547 days_early: i64, ts_ns: i64) -> i64 {
548 let evt: *NxTelemetryEvent = (sys_mmap(96)) as *NxTelemetryEvent
549 evt.kind = NX_TM_KIND_ANOMALY
550 // Severity: 3-sigma anomalies are WARN; drift-detection that
551 // crosses threshold within 7 days is ERROR.
552 var sev: i64 = NX_TM_SEV_WARN
553 if anom_kind >= 3 { sev = NX_TM_SEV_ERROR } // ANOM_DRIFT_UP/DOWN
554 evt.severity = sev
555 evt.source_str = NX_TM_SOURCE_PRED_MAINT
556 evt.source_len = NX_TM_SOURCE_PRED_MAINT_LEN
557 evt.ts_ns = ts_ns
558 evt.value0 = pid // OBD PID number
559 evt.value1 = anom_kind // NX_PM_ANOM_*
560 evt.value2 = days_early // 0 if 3-sigma now; N if drift predicts in N days
561 evt.extra_str = (0 as i64) as *u8
562 evt.extra_len = 0
563 return nx_tm_emit(bus, evt)
564}