code wiki / (root) / nx_mqtt_sink_for_telemetry.nx

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}