nx_mqtt_sink_for_telemetry.nx source
↩ module page · 207 lines · 8484 B
1// nx_mqtt_sink_for_telemetry.nx -- MQTT sink for telemetry events.
2//
3// Converts an NxTelemetryEvent into an MQTT PUBLISH frame, so the
4// substrate's anomaly verdicts (pred-maint commit d3037668 +
5// telemetry bus 469f822f) become operator-visible on whatever MQTT
6// broker the operator points to.
7//
8// Architectural note (V1 vs V2):
9// V1 ships a STANDALONE sink module. The operator calls
10// nx_mqtt_sink_emit alongside the bus's stdout-JSON sink in a
11// fanout pattern (manual wiring):
12// nx_tm_emit(bus, event) -- writes stdout NDJSON
13// nx_mqtt_sink_emit(cfg, event, ...) -- builds PUBLISH frame
14// that operator transmits
15// over socket
16// V2 will extend NxTelemetrySink with a sealed sink-kind enum
17// (STDOUT_JSON / MQTT / SYSLOG / FILE / FANOUT) so the bus
18// dispatches automatically; this V1 module survives as the
19// MQTT-specific frame builder either way.
20//
21// Status: SEED v0.1.0. 2026-05-27.
22// WINNER-TIER: WINNER-A CANDIDATE (preserved from telemetry bus;
23// this commit adds the MQTT-output capability without
24// changing the underlying value-add claim)
25// INCUMBENTS: Eclipse Paho's MQTTClient_publishMessage,
26// mqtt.js client.publish(), AWS IoT SDK publish(),
27// Mosquitto's mosquitto_pub CLI
28// NUMBERS: V1 ships the frame-builder; latency vs Paho
29// publish() pending paired bench on real broker
30// GAP: Paho buffers + retries + auto-reconnects internally;
31// V1 is purely the wire-format builder. Operator
32// owns socket lifecycle for V1; V2 adds nx_tcp_shim
33// integration when that ships.
34// PLAN: M-next: paired publish-latency bench vs Paho; add
35// socket-write integration when nx_tcp_shim ships
36// EXEMPTION REASON: n/a; provisional pending measurement
37//
38// Topic shape (operator-customizable via topic_prefix):
39// <prefix>/<source>
40// where <source> is the NxTelemetryEvent.source_str (e.g.,
41// "nx_pred_maint_consumer"). Default prefix: "nishi/telemetry".
42//
43// Payload shape (compact JSON; same schema as stdout-JSON sink for
44// downstream-consumer uniformity):
45// {"kind":N,"sev":N,"ts":N,"v0":N,"v1":N,"v2":N}
46// Source field is OMITTED from payload (it's in the topic).
47
48import "nx_syscalls.nx"
49import "nx_telemetry.nx"
50import "shims/nx_mqtt_shim.nx"
51
52// ===== Sealed verdict surface =================================================
53const NX_MQTT_SINK_OK: i64 = 0
54const NX_MQTT_SINK_BAD_INPUT: i64 = 1
55const NX_MQTT_SINK_NOT_INITIALIZED: i64 = 2
56const NX_MQTT_SINK_BAD_TOPIC: i64 = 3
57const NX_MQTT_SINK_PAYLOAD_OVERFLOW: i64 = 4
58const NX_MQTT_SINK_FRAME_OVERFLOW: i64 = 5
59
60// ===== Config =================================================
61
62struct NxMqttSinkConfig {
63 shim: *NxMqttShim // already-initialised MQTT shim
64 topic_prefix: *u8 // e.g., "nishi/telemetry"
65 topic_prefix_len: i64
66 qos: i64 // 0 or 1 (V1)
67 valid: i64
68}
69
70func nx_mqtt_sink_init(cfg: *NxMqttSinkConfig,
71 shim: *NxMqttShim,
72 topic_prefix: *u8, topic_prefix_len: i64,
73 qos: i64) -> i64 {
74 if (cfg as i64) == 0 { return 0 - NX_MQTT_SINK_BAD_INPUT }
75 if (shim as i64) == 0 { return 0 - NX_MQTT_SINK_BAD_INPUT }
76 if topic_prefix_len < 1 { return 0 - NX_MQTT_SINK_BAD_TOPIC }
77 if qos < 0 { return 0 - NX_MQTT_SINK_BAD_INPUT }
78 if qos > 1 { return 0 - NX_MQTT_SINK_BAD_INPUT } // V1 limit
79 cfg.shim = shim
80 cfg.topic_prefix = topic_prefix
81 cfg.topic_prefix_len = topic_prefix_len
82 cfg.qos = qos
83 cfg.valid = 1
84 return NX_MQTT_SINK_OK
85}
86
87// ===== Helper: compact JSON payload =================================================
88//
89// Reuses the substrate's int-emit + str-emit primitives from
90// nx_telemetry's stdout-JSON sink. V1 inlines a minimal version
91// to avoid taking an additional dependency; V2 can refactor to
92// share helpers once a "compact-JSON" common module lands.
93
94func nx_mqtt_sink_emit_int(out: *u8, off: i64, value: i64) -> i64 {
95 var p: i64 = off
96 var n: i64 = value
97 if n < 0 {
98 out[p] = 0x2d as u8
99 p = p + 1
100 n = 0 - n
101 }
102 if n == 0 {
103 out[p] = 0x30 as u8
104 return p + 1
105 }
106 let digits: *u8 = sys_mmap(24)
107 var k: i64 = 0
108 while n > 0 {
109 digits[k] = (0x30 + (n % 10)) as u8
110 n = n / 10
111 k = k + 1
112 }
113 var i: i64 = k - 1
114 while i >= 0 {
115 out[p] = digits[i]
116 p = p + 1
117 i = i - 1
118 }
119 return p
120}
121
122func nx_mqtt_sink_emit_str(out: *u8, off: i64, src: *u8, src_len: i64) -> i64 {
123 var p: i64 = off
124 var i: i64 = 0
125 while i < src_len {
126 out[p] = src[i]
127 p = p + 1
128 i = i + 1
129 }
130 return p
131}
132
133// ===== Build topic name: <prefix>/<source> =================================================
134
135func nx_mqtt_sink_build_topic(cfg: *NxMqttSinkConfig,
136 source: *u8, source_len: i64,
137 out: *u8, out_cap: i64) -> i64 {
138 let total: i64 = cfg.topic_prefix_len + 1 + source_len
139 if total > out_cap { return 0 - NX_MQTT_SINK_BAD_TOPIC }
140 var p: i64 = nx_mqtt_sink_emit_str(out, 0, cfg.topic_prefix, cfg.topic_prefix_len)
141 out[p] = 0x2f as u8 // '/'
142 p = p + 1
143 p = nx_mqtt_sink_emit_str(out, p, source, source_len)
144 return p
145}
146
147// ===== Build compact-JSON payload =================================================
148
149func nx_mqtt_sink_build_payload(event: *NxTelemetryEvent,
150 out: *u8, out_cap: i64) -> i64 {
151 if out_cap < 128 { return 0 - NX_MQTT_SINK_PAYLOAD_OVERFLOW }
152 var p: i64 = 0
153 p = nx_mqtt_sink_emit_str(out, p, "{\"kind\":" as *u8, 8)
154 p = nx_mqtt_sink_emit_int(out, p, event.kind)
155 p = nx_mqtt_sink_emit_str(out, p, ",\"sev\":" as *u8, 7)
156 p = nx_mqtt_sink_emit_int(out, p, event.severity)
157 p = nx_mqtt_sink_emit_str(out, p, ",\"ts\":" as *u8, 6)
158 p = nx_mqtt_sink_emit_int(out, p, event.ts_ns)
159 p = nx_mqtt_sink_emit_str(out, p, ",\"v0\":" as *u8, 6)
160 p = nx_mqtt_sink_emit_int(out, p, event.value0)
161 p = nx_mqtt_sink_emit_str(out, p, ",\"v1\":" as *u8, 6)
162 p = nx_mqtt_sink_emit_int(out, p, event.value1)
163 p = nx_mqtt_sink_emit_str(out, p, ",\"v2\":" as *u8, 6)
164 p = nx_mqtt_sink_emit_int(out, p, event.value2)
165 p = nx_mqtt_sink_emit_str(out, p, "}" as *u8, 1)
166 return p
167}
168
169// ===== Top-level emit: build MQTT PUBLISH frame for event =================================================
170//
171// Caller supplies output buffer (recommended: >= 512 bytes). Sink
172// writes the complete PUBLISH frame (fixed header + variable header
173// + payload). Caller transmits the bytes over their socket.
174//
175// V2 will integrate nx_tcp_shim's TX call so the operator only
176// needs one call instead of two; V1 keeps wire I/O out-of-scope
177// per architectural separation between framing + transport.
178//
179// Returns frame size in bytes on success, or 0 - verdict.
180
181func nx_mqtt_sink_emit(cfg: *NxMqttSinkConfig,
182 event: *NxTelemetryEvent,
183 out: *u8, out_cap: i64) -> i64 {
184 if cfg.valid != 1 { return 0 - NX_MQTT_SINK_NOT_INITIALIZED }
185 if (event as i64) == 0 { return 0 - NX_MQTT_SINK_BAD_INPUT }
186
187 // Topic = prefix + "/" + source.
188 let topic_buf: *u8 = sys_mmap(256)
189 let topic_len: i64 = nx_mqtt_sink_build_topic(cfg, event.source_str,
190 event.source_len,
191 topic_buf, 256)
192 if topic_len < 0 { return topic_len }
193
194 // Payload = compact JSON.
195 let payload_buf: *u8 = sys_mmap(256)
196 let payload_len: i64 = nx_mqtt_sink_build_payload(event, payload_buf, 256)
197 if payload_len < 0 { return payload_len }
198
199 // Build PUBLISH frame via the MQTT shim's builder.
200 let frame_len: i64 = nx_mqtt_build_publish(cfg.shim,
201 topic_buf, topic_len,
202 payload_buf, payload_len,
203 cfg.qos,
204 out, out_cap)
205 if frame_len < 0 { return 0 - NX_MQTT_SINK_FRAME_OVERFLOW }
206 return frame_len
207}