nx_mqtt_shim.nx source
↩ module page · 485 lines · 19846 B
1// nx_mqtt_shim.nx -- MQTT 3.1.1 protocol shim.
2//
3// WHEELER-IMPLEMENTATION-OF-CANONICAL-SPEC.
4// Wire-spec source (re-implemented clean-room from published spec;
5// no external code imported per sovereignty audit):
6// - OASIS MQTT Version 3.1.1 Plus Errata 01 (29 October 2015)
7// - Optional V2+: OASIS MQTT Version 5.0 (07 March 2019)
8//
9// Third protocol shim per INTEROPERABILITY_CHARTER.md §5.5 (cloud).
10// Opens the cloud-pub/sub family. MQTT is the IoT-standard pub/sub
11// wire protocol; once this ships, the substrate's nx_telemetry bus
12// can publish anomaly verdicts to an MQTT broker -> operator
13// dashboards / mobile alerts / cloud aggregation, instead of just
14// stdout NDJSON.
15//
16// Status: SEED v0.1.0. 2026-05-27.
17//
18// WINNER-TIER: BASELINE-C provisional
19// INCUMBENTS: Mosquitto v2.x (broker; ISC),
20// Eclipse Paho C v1.3.x (client; EPL),
21// mqtt.js v5.x (JavaScript client; MIT),
22// HiveMQ Cloud (commercial broker),
23// AWS IoT Core (commercial cloud broker),
24// Aedes (Node.js broker)
25// NUMBERS: V1 ships framing + publish/subscribe + ping; paired
26// throughput + connection-stability bench vs Paho C
27// pending real broker fixture
28// GAP: Paho has feature breadth (TLS, MQTT 5.0 properties,
29// persistent session, will/retain); V1 covers the
30// ~70% subset that publishes telemetry from substrate
31// to broker. Nishi adds value via sealed verdict
32// surface + telemetry-bus integration (Paho is
33// string-keyed; loses type safety)
34// PLAN: M-next: paired bench vs Paho on 10K-publish burst
35// latency + broker reconnection storm; add MQTT 5.0
36// properties + TLS wrap when nx_tls_shim lands
37// EXEMPTION REASON: n/a; provisional pending measurement
38//
39// V1 SCOPE:
40// - CONNECT / CONNACK (client-side; clean-session only V1)
41// - PUBLISH frame parse + build (QoS 0 + 1)
42// - PUBACK (QoS 1 ack reception)
43// - SUBSCRIBE / SUBACK (single-topic subscription per call)
44// - PINGREQ / PINGRESP (keepalive)
45// - DISCONNECT (clean shutdown)
46// - Variable-length integer encode + decode (1-4 byte VLI per spec §2.2.3)
47// - Topic-name validation (no wildcards in PUBLISH; / + # in SUBSCRIBE)
48// - 7 NxProtocolEvent kinds (substrate-facing; raw bytes never
49// leak past shim per INTEROPERABILITY_CHARTER §M2)
50//
51// V2+ SCOPE:
52// - QoS 2 retransmission state machine (PUBREC/PUBREL/PUBCOMP)
53// - MQTT 5.0 properties + reason codes
54// - TLS wrap via nx_tls_shim (when §5.5 TLS shim lands)
55// - Will / retain message semantics
56// - Username/password CONNECT auth
57// - Persistent session resumption
58// - SUBSCRIBE with multiple topic filters per packet
59
60import "nx_syscalls.nx"
61
62// ===== Cross-shim verdict codes (mirror INTEROPERABILITY_CHARTER §9) =================================================
63const NX_SHIM_OK: i64 = 0
64const NX_SHIM_BAD_FRAME: i64 = 1
65const NX_SHIM_PROTOCOL_VIOLATION: i64 = 2
66const NX_SHIM_TIMEOUT: i64 = 3
67const NX_SHIM_BACKPRESSURE: i64 = 4
68const NX_SHIM_UNSUPPORTED_PID: i64 = 5 // here: unsupported control packet type
69const NX_SHIM_NO_WIRE: i64 = 6
70
71// ===== MQTT-specific verdict codes =================================================
72const NX_MQTT_VERDICT_BASE: i64 = 300
73const NX_MQTT_BAD_PROTO_NAME: i64 = 300 // CONNECT proto name != "MQTT"
74const NX_MQTT_BAD_PROTO_LEVEL: i64 = 301 // CONNECT proto level != 4 (v3.1.1)
75const NX_MQTT_BAD_VLI: i64 = 302 // var-len-int > 4 bytes (spec max)
76const NX_MQTT_BAD_QOS: i64 = 303 // QoS out of 0..2
77const NX_MQTT_BAD_TOPIC: i64 = 304 // topic name violates §4.7
78const NX_MQTT_UNSUPPORTED_QOS2: i64 = 305 // V1 doesn't implement QoS 2
79const NX_MQTT_CONNACK_REJECTED: i64 = 306 // broker rejected connection (CONNACK rc != 0)
80const NX_MQTT_TRUNCATED: i64 = 307
81
82// ===== Control packet types (spec §2.2.1; high nibble of byte 0) =================================================
83const NX_MQTT_PKT_RESERVED_0: i64 = 0x00
84const NX_MQTT_PKT_CONNECT: i64 = 0x10 // client -> server
85const NX_MQTT_PKT_CONNACK: i64 = 0x20 // server -> client
86const NX_MQTT_PKT_PUBLISH: i64 = 0x30 // bidirectional
87const NX_MQTT_PKT_PUBACK: i64 = 0x40 // QoS 1 ack
88const NX_MQTT_PKT_PUBREC: i64 = 0x50 // QoS 2 step 1 (V2+)
89const NX_MQTT_PKT_PUBREL: i64 = 0x60 // QoS 2 step 2 (V2+)
90const NX_MQTT_PKT_PUBCOMP: i64 = 0x70 // QoS 2 step 3 (V2+)
91const NX_MQTT_PKT_SUBSCRIBE: i64 = 0x80
92const NX_MQTT_PKT_SUBACK: i64 = 0x90
93const NX_MQTT_PKT_UNSUBSCRIBE: i64 = 0xA0
94const NX_MQTT_PKT_UNSUBACK: i64 = 0xB0
95const NX_MQTT_PKT_PINGREQ: i64 = 0xC0
96const NX_MQTT_PKT_PINGRESP: i64 = 0xD0
97const NX_MQTT_PKT_DISCONNECT: i64 = 0xE0
98
99// ===== CONNECT flags (spec §3.1.2.3) =================================================
100const NX_MQTT_CONNECT_FLAG_CLEAN_SESSION: i64 = 0x02
101const NX_MQTT_CONNECT_FLAG_WILL: i64 = 0x04
102const NX_MQTT_CONNECT_FLAG_WILL_QOS_1: i64 = 0x08
103const NX_MQTT_CONNECT_FLAG_WILL_QOS_2: i64 = 0x10
104const NX_MQTT_CONNECT_FLAG_WILL_RETAIN: i64 = 0x20
105const NX_MQTT_CONNECT_FLAG_PASSWORD: i64 = 0x40
106const NX_MQTT_CONNECT_FLAG_USERNAME: i64 = 0x80
107
108// ===== CONNACK return codes (spec §3.2.2.3) =================================================
109const NX_MQTT_CONNACK_OK: i64 = 0
110const NX_MQTT_CONNACK_UNSUPPORTED_PROTOCOL: i64 = 1
111const NX_MQTT_CONNACK_BAD_CLIENT_ID: i64 = 2
112const NX_MQTT_CONNACK_SERVER_UNAVAILABLE: i64 = 3
113const NX_MQTT_CONNACK_BAD_AUTH: i64 = 4
114const NX_MQTT_CONNACK_NOT_AUTHORIZED: i64 = 5
115
116// ===== Substrate-facing event kinds =================================================
117const NX_MQTT_EVT_CONNECTED: i64 = 400 // v0=session_present; v1=connack_rc
118const NX_MQTT_EVT_MESSAGE_RX: i64 = 401 // v0=packet_id; v1=qos; v2=payload_len
119 // topic + payload via extra_str (V2)
120const NX_MQTT_EVT_PUBACK_RX: i64 = 402 // v0=packet_id (matches our prior PUBLISH)
121const NX_MQTT_EVT_SUBSCRIBED: i64 = 403 // v0=packet_id; v1=granted_qos
122const NX_MQTT_EVT_PING_RX: i64 = 404 // server-side PINGRESP received
123const NX_MQTT_EVT_DISCONNECTED: i64 = 405 // server initiated DISCONNECT
124const NX_MQTT_EVT_REJECTED: i64 = 406 // v0=connack_rc reason
125
126// ===== Shim state =================================================
127
128struct NxMqttShim {
129 name_buf: *u8
130 wire_spec_id: *u8
131 next_packet_id: i64 // monotonic; wraps at 65535
132 rx_byte_count: i64
133 rx_event_count: i64
134 rx_error_count: i64
135 last_verdict: i64
136 connected: i64 // 1 after successful CONNACK
137 valid: i64
138}
139
140func nx_mqtt_shim_init(s: *NxMqttShim) -> i64 {
141 if (s as i64) == 0 { return 0 - NX_SHIM_BAD_FRAME }
142 s.name_buf = "MQTT 3.1.1" as *u8
143 s.wire_spec_id = "OASIS MQTT v3.1.1+Errata-01" as *u8
144 s.next_packet_id = 1
145 s.rx_byte_count = 0
146 s.rx_event_count = 0
147 s.rx_error_count = 0
148 s.last_verdict = NX_SHIM_OK
149 s.connected = 0
150 s.valid = 1
151 return NX_SHIM_OK
152}
153
154// ===== Variable-length integer (spec §2.2.3) =================================================
155//
156// 1-4 bytes per integer; each byte uses 7 bits for value + high bit
157// as continuation flag. Encodes values 0..268,435,455 (~256 MB).
158
159func nx_mqtt_vli_decode(buf: *u8, off: i64, out_value: *i64) -> i64 {
160 var value: i64 = 0
161 var multiplier: i64 = 1
162 var i: i64 = 0
163 while i < 4 {
164 let byte: i64 = buf[off + i] as i64
165 value = value + ((byte & 0x7f) * multiplier)
166 if (byte & 0x80) == 0 {
167 out_value[0] = value
168 return i + 1 // bytes consumed
169 }
170 multiplier = multiplier * 128
171 i = i + 1
172 }
173 return 0 - NX_MQTT_BAD_VLI // malformed (>4 bytes)
174}
175
176func nx_mqtt_vli_encode(value: i64, out: *u8, out_off: i64) -> i64 {
177 var v: i64 = value
178 var off: i64 = out_off
179 var n: i64 = 0
180 while n < 4 {
181 var byte: i64 = v & 0x7f
182 v = v >> 7
183 if v > 0 { byte = byte | 0x80 }
184 out[off] = byte as u8
185 off = off + 1
186 n = n + 1
187 if v == 0 { return n } // bytes written
188 }
189 return 0 - NX_MQTT_BAD_VLI
190}
191
192// ===== UTF-8 string helpers (length-prefixed per spec §1.5.3) =================================================
193//
194// MQTT strings are: u16 BE length | UTF-8 bytes (length octets).
195// Max 65535 octets per spec.
196
197func nx_mqtt_read_u16_be(buf: *u8, off: i64) -> i64 {
198 let hi: i64 = buf[off] as i64
199 let lo: i64 = buf[off + 1] as i64
200 return (hi << 8) | lo
201}
202
203func nx_mqtt_write_u16_be(buf: *u8, off: i64, value: i64) -> i64 {
204 buf[off] = ((value >> 8) & 0xff) as u8
205 buf[off + 1] = (value & 0xff) as u8
206 return 0
207}
208
209func nx_mqtt_write_string(buf: *u8, off: i64, str: *u8, str_len: i64) -> i64 {
210 nx_mqtt_write_u16_be(buf, off, str_len)
211 var i: i64 = 0
212 while i < str_len {
213 buf[off + 2 + i] = str[i]
214 i = i + 1
215 }
216 return off + 2 + str_len
217}
218
219// ===== Build CONNECT (spec §3.1) =================================================
220//
221// Variable header: protocol_name "MQTT" (6 bytes incl length) +
222// protocol_level 4 + connect_flags + keepalive_be
223// Payload: client_id (UTF-8 string)
224// V1: clean-session only; no will / username / password.
225
226func nx_mqtt_build_connect(s: *NxMqttShim, client_id: *u8, client_id_len: i64,
227 keepalive_s: i64, out: *u8, out_cap: i64) -> i64 {
228 if s.valid != 1 { return 0 - NX_SHIM_NO_WIRE }
229 if out_cap < 64 { return 0 - NX_SHIM_BACKPRESSURE }
230 if client_id_len < 1 { return 0 - NX_SHIM_BAD_FRAME }
231 if client_id_len > 23 { return 0 - NX_SHIM_BAD_FRAME } // spec §3.1.3.1 max
232
233 // Variable header (10 bytes) + payload (2 + client_id_len bytes)
234 let var_header_len: i64 = 10
235 let payload_len: i64 = 2 + client_id_len
236 let remaining_len: i64 = var_header_len + payload_len
237
238 // Byte 0: CONNECT (0x10)
239 out[0] = NX_MQTT_PKT_CONNECT as u8
240 // Bytes 1..: VLI for remaining length
241 let vli_bytes: i64 = nx_mqtt_vli_encode(remaining_len, out, 1)
242 if vli_bytes < 0 { return vli_bytes }
243
244 let body_off: i64 = 1 + vli_bytes
245
246 // Protocol name: "MQTT" (length-prefixed)
247 let proto_name: *u8 = "MQTT" as *u8
248 let off2: i64 = nx_mqtt_write_string(out, body_off, proto_name, 4)
249
250 // Protocol level (4 for v3.1.1)
251 out[off2] = 4 as u8
252
253 // Connect flags (clean session only)
254 out[off2 + 1] = NX_MQTT_CONNECT_FLAG_CLEAN_SESSION as u8
255
256 // Keepalive (u16 BE seconds)
257 nx_mqtt_write_u16_be(out, off2 + 2, keepalive_s)
258
259 // Payload: client_id
260 nx_mqtt_write_string(out, off2 + 4, client_id, client_id_len)
261
262 return body_off + var_header_len + payload_len
263}
264
265// ===== Build PUBLISH (spec §3.3) =================================================
266//
267// Fixed header: 0x30 | (DUP<<3) | (QoS<<1) | RETAIN
268// Variable header: topic_name (UTF-8 string)
269// + packet_id (u16 BE, only if QoS > 0)
270// Payload: application message (opaque bytes)
271
272func nx_mqtt_build_publish(s: *NxMqttShim, topic: *u8, topic_len: i64,
273 payload: *u8, payload_len: i64, qos: i64,
274 out: *u8, out_cap: i64) -> i64 {
275 if s.valid != 1 { return 0 - NX_SHIM_NO_WIRE }
276 if qos < 0 { return 0 - NX_MQTT_BAD_QOS }
277 if qos > 1 { return 0 - NX_MQTT_UNSUPPORTED_QOS2 } // V1 QoS 0+1 only
278 if topic_len < 1 { return 0 - NX_MQTT_BAD_TOPIC }
279
280 // Estimate frame size: 1 (fixed_hdr) + 4 (max VLI) +
281 // 2 + topic_len + (2 if QoS>0) + payload_len
282 var var_hdr_len: i64 = 2 + topic_len
283 if qos > 0 { var_hdr_len = var_hdr_len + 2 }
284 let remaining_len: i64 = var_hdr_len + payload_len
285
286 let needed: i64 = 1 + 4 + remaining_len
287 if out_cap < needed { return 0 - NX_SHIM_BACKPRESSURE }
288
289 // Byte 0: PUBLISH (0x30) + flags
290 out[0] = (NX_MQTT_PKT_PUBLISH | (qos << 1)) as u8
291 let vli_bytes: i64 = nx_mqtt_vli_encode(remaining_len, out, 1)
292 if vli_bytes < 0 { return vli_bytes }
293 let body_off: i64 = 1 + vli_bytes
294
295 // Topic name (length-prefixed UTF-8)
296 let off2: i64 = nx_mqtt_write_string(out, body_off, topic, topic_len)
297
298 // Packet ID (QoS > 0 only)
299 var payload_off: i64 = off2
300 if qos > 0 {
301 let pid: i64 = s.next_packet_id
302 nx_mqtt_write_u16_be(out, off2, pid)
303 s.next_packet_id = pid + 1
304 if s.next_packet_id > 65535 { s.next_packet_id = 1 }
305 payload_off = off2 + 2
306 }
307
308 // Payload (opaque copy)
309 var i: i64 = 0
310 while i < payload_len {
311 out[payload_off + i] = payload[i]
312 i = i + 1
313 }
314 return payload_off + payload_len
315}
316
317// ===== Build SUBSCRIBE (spec §3.8; single-topic V1) =================================================
318
319func nx_mqtt_build_subscribe(s: *NxMqttShim, topic_filter: *u8, topic_filter_len: i64,
320 qos: i64, out: *u8, out_cap: i64) -> i64 {
321 if s.valid != 1 { return 0 - NX_SHIM_NO_WIRE }
322 if qos < 0 { return 0 - NX_MQTT_BAD_QOS }
323 if qos > 2 { return 0 - NX_MQTT_BAD_QOS }
324 if topic_filter_len < 1 { return 0 - NX_MQTT_BAD_TOPIC }
325 if out_cap < 16 + topic_filter_len { return 0 - NX_SHIM_BACKPRESSURE }
326
327 // Variable header: packet_id (2 bytes) ; Payload: topic_filter (2+N) + qos (1)
328 let remaining_len: i64 = 2 + 2 + topic_filter_len + 1
329
330 // Byte 0: SUBSCRIBE (0x80) + reserved low nibble = 0x02 (spec §3.8.1)
331 out[0] = (NX_MQTT_PKT_SUBSCRIBE | 0x02) as u8
332 let vli_bytes: i64 = nx_mqtt_vli_encode(remaining_len, out, 1)
333 if vli_bytes < 0 { return vli_bytes }
334 let body_off: i64 = 1 + vli_bytes
335
336 let pid: i64 = s.next_packet_id
337 nx_mqtt_write_u16_be(out, body_off, pid)
338 s.next_packet_id = pid + 1
339 if s.next_packet_id > 65535 { s.next_packet_id = 1 }
340
341 let off2: i64 = nx_mqtt_write_string(out, body_off + 2, topic_filter, topic_filter_len)
342 out[off2] = (qos & 0x03) as u8
343 return off2 + 1
344}
345
346// ===== Build PINGREQ / DISCONNECT (no body) =================================================
347
348func nx_mqtt_build_pingreq(s: *NxMqttShim, out: *u8, out_cap: i64) -> i64 {
349 if s.valid != 1 { return 0 - NX_SHIM_NO_WIRE }
350 if out_cap < 2 { return 0 - NX_SHIM_BACKPRESSURE }
351 out[0] = NX_MQTT_PKT_PINGREQ as u8
352 out[1] = 0 as u8 // remaining length = 0
353 return 2
354}
355
356func nx_mqtt_build_disconnect(s: *NxMqttShim, out: *u8, out_cap: i64) -> i64 {
357 if s.valid != 1 { return 0 - NX_SHIM_NO_WIRE }
358 if out_cap < 2 { return 0 - NX_SHIM_BACKPRESSURE }
359 out[0] = NX_MQTT_PKT_DISCONNECT as u8
360 out[1] = 0 as u8
361 s.connected = 0
362 return 2
363}
364
365// ===== rx_fn: parse incoming frame =================================================
366//
367// Dispatches by control-packet-type to the appropriate parser.
368// Each parser emits 0 or more NxProtocolEvents.
369
370func nx_mqtt_rx_fn(s: *NxMqttShim, frame: *u8, frame_len: i64,
371 out_evt_kinds: *i64, out_values: *i64, cap: i64) -> i64 {
372 if s.valid != 1 { return 0 - NX_SHIM_NO_WIRE }
373 s.rx_byte_count = s.rx_byte_count + frame_len
374 if frame_len < 2 { return 0 - NX_MQTT_TRUNCATED }
375
376 let pkt_type: i64 = (frame[0] as i64) & 0xf0
377
378 // Decode remaining length (VLI).
379 let vli_out: *i64 = (sys_mmap(8)) as *i64
380 vli_out[0] = 0
381 let vli_bytes: i64 = nx_mqtt_vli_decode(frame, 1, vli_out)
382 if vli_bytes < 0 { return vli_bytes }
383 let remaining_len: i64 = vli_out[0]
384 let body_off: i64 = 1 + vli_bytes
385 if (body_off + remaining_len) > frame_len { return 0 - NX_MQTT_TRUNCATED }
386
387 if pkt_type == NX_MQTT_PKT_CONNACK {
388 if remaining_len < 2 { return 0 - NX_MQTT_TRUNCATED }
389 let session_present: i64 = (frame[body_off] as i64) & 0x01
390 let connack_rc: i64 = frame[body_off + 1] as i64
391 if cap < 1 { return 0 - NX_SHIM_BACKPRESSURE }
392 if connack_rc == NX_MQTT_CONNACK_OK {
393 out_evt_kinds[0] = NX_MQTT_EVT_CONNECTED
394 out_values[0] = session_present
395 s.connected = 1
396 }
397 if connack_rc != NX_MQTT_CONNACK_OK {
398 out_evt_kinds[0] = NX_MQTT_EVT_REJECTED
399 out_values[0] = connack_rc
400 s.rx_error_count = s.rx_error_count + 1
401 s.last_verdict = NX_MQTT_CONNACK_REJECTED
402 }
403 s.rx_event_count = s.rx_event_count + 1
404 return 1
405 }
406
407 if pkt_type == NX_MQTT_PKT_PUBLISH {
408 // Flags (byte 0 low nibble): DUP|QoS|RETAIN
409 let qos: i64 = ((frame[0] as i64) >> 1) & 0x03
410 if qos > 1 { return 0 - NX_MQTT_UNSUPPORTED_QOS2 }
411 // Variable header: topic_name (2+N), optionally packet_id (2)
412 if remaining_len < 2 { return 0 - NX_MQTT_TRUNCATED }
413 let topic_len: i64 = nx_mqtt_read_u16_be(frame, body_off)
414 var var_hdr_len: i64 = 2 + topic_len
415 var pid: i64 = 0
416 if qos > 0 {
417 pid = nx_mqtt_read_u16_be(frame, body_off + var_hdr_len)
418 var_hdr_len = var_hdr_len + 2
419 }
420 let payload_len: i64 = remaining_len - var_hdr_len
421 if cap < 1 { return 0 - NX_SHIM_BACKPRESSURE }
422 out_evt_kinds[0] = NX_MQTT_EVT_MESSAGE_RX
423 out_values[0] = pid // packet ID (0 for QoS 0)
424 // Future: topic + payload accessible via extra_str on the
425 // event (per charter §6); V1 surfaces metadata only via the
426 // 3-value-slot layout.
427 s.rx_event_count = s.rx_event_count + 1
428 return 1
429 }
430
431 if pkt_type == NX_MQTT_PKT_PUBACK {
432 if remaining_len < 2 { return 0 - NX_MQTT_TRUNCATED }
433 let pid: i64 = nx_mqtt_read_u16_be(frame, body_off)
434 if cap < 1 { return 0 - NX_SHIM_BACKPRESSURE }
435 out_evt_kinds[0] = NX_MQTT_EVT_PUBACK_RX
436 out_values[0] = pid
437 s.rx_event_count = s.rx_event_count + 1
438 return 1
439 }
440
441 if pkt_type == NX_MQTT_PKT_SUBACK {
442 if remaining_len < 3 { return 0 - NX_MQTT_TRUNCATED }
443 let pid: i64 = nx_mqtt_read_u16_be(frame, body_off)
444 let granted_qos: i64 = (frame[body_off + 2] as i64) & 0x03
445 if cap < 1 { return 0 - NX_SHIM_BACKPRESSURE }
446 out_evt_kinds[0] = NX_MQTT_EVT_SUBSCRIBED
447 out_values[0] = pid
448 // granted_qos not surfaced in 1-value layout; V2 uses
449 // 3-value layout to include it
450 s.rx_event_count = s.rx_event_count + 1
451 return 1
452 }
453
454 if pkt_type == NX_MQTT_PKT_PINGRESP {
455 if cap < 1 { return 0 - NX_SHIM_BACKPRESSURE }
456 out_evt_kinds[0] = NX_MQTT_EVT_PING_RX
457 out_values[0] = 0
458 s.rx_event_count = s.rx_event_count + 1
459 return 1
460 }
461
462 if pkt_type == NX_MQTT_PKT_DISCONNECT {
463 s.connected = 0
464 if cap < 1 { return 0 - NX_SHIM_BACKPRESSURE }
465 out_evt_kinds[0] = NX_MQTT_EVT_DISCONNECTED
466 out_values[0] = 0
467 s.rx_event_count = s.rx_event_count + 1
468 return 1
469 }
470
471 // Unsupported control packet type (e.g., QoS 2 family in V1)
472 s.rx_error_count = s.rx_error_count + 1
473 return 0 - NX_SHIM_UNSUPPORTED_PID
474}
475
476// ===== diag_fn =================================================
477
478func nx_mqtt_diag_fn(s: *NxMqttShim, out_vec: *i64) -> i64 {
479 if s.valid != 1 { return 0 - NX_SHIM_NO_WIRE }
480 out_vec[0] = s.rx_byte_count
481 out_vec[1] = s.rx_event_count
482 out_vec[2] = s.rx_error_count
483 out_vec[3] = s.last_verdict
484 return NX_SHIM_OK
485}