code wiki / shims / nx_mqtt_shim.nx

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}