ProtoCore v0.0.2
Deterministic, zero-heap network stack for embedded targets
Loading...
Searching...
No Matches
nats.cpp
Go to the documentation of this file.
1// Copyright (C) 2026 Douglas Quigg (dstroy0) <dquigg123@gmail.com>
2// SPDX-License-Identifier: AGPL-3.0-or-later
3
4/**
5 * @file nats.cpp
6 * @brief NATS client protocol builder + parser (pure, host-tested).
7 */
8
10
11#if PC_ENABLE_NATS
12
13#include <string.h>
14
15// A tiny bounded append cursor; sets ok=false on overflow and stops.
16struct Buf
17{
18 char *p;
19 size_t cap;
20 size_t pos;
21 bool ok;
22};
23
24static void put_str(Buf *b, const char *s)
25{
26 if (!b->ok || !s) // GCOVR_EXCL_BR_LINE the !s half is unreachable: every put_str arg is a literal or a
27 // checked/guarded non-null string (see below)
28 {
29 if (!s) // GCOVR_EXCL_BR_LINE unreachable for the same reason as the condition above
30 {
31 b->ok = false; // GCOVR_EXCL_LINE unreachable: every put_str arg is a literal or a checked/guarded non-null
32 // string
33 }
34 return;
35 }
36 size_t n = strnlen(s, b->cap + 1);
37 if (b->pos + n > b->cap)
38 {
39 b->ok = false;
40 return;
41 }
42 memcpy(b->p + b->pos, s, n);
43 b->pos += n;
44}
45
46static void put_bytes(Buf *b, const uint8_t *d, size_t n)
47{
48 if (!b->ok)
49 {
50 return;
51 }
52 if (b->pos + n > b->cap)
53 {
54 b->ok = false;
55 return;
56 }
57 if (n)
58 {
59 memcpy(b->p + b->pos, d, n);
60 }
61 b->pos += n;
62}
63
64static void put_ch(Buf *b, char c)
65{
66 if (!b->ok)
67 {
68 return;
69 }
70 if (b->pos + 1 > b->cap)
71 {
72 b->ok = false;
73 return;
74 }
75 b->p[b->pos++] = c;
76}
77
78static void put_uint(Buf *b, uint64_t v)
79{
80 char tmp[20];
81 size_t n = 0;
82 char rev[20];
83 size_t r = 0;
84 if (v == 0)
85 {
86 rev[r++] = '0';
87 }
88 while (v)
89 {
90 rev[r++] = (char)('0' + (int)(v % 10));
91 v /= 10;
92 }
93 while (r)
94 {
95 tmp[n++] = rev[--r];
96 }
97 put_bytes(b, (const uint8_t *)tmp, n);
98}
99
100static size_t finish(Buf *b)
101{
102 if (!b->ok)
103 {
104 return 0;
105 }
106 if (b->pos < b->cap) // NUL-terminate when there's room (the returned length excludes it)
107 {
108 b->p[b->pos] = '\0';
109 }
110 return b->pos;
111}
112
113size_t pc_nats_build_connect(char *buf, size_t cap, const char *options_json)
114{
115 if (!buf || !options_json)
116 {
117 return 0;
118 }
119 Buf b = {buf, cap, 0, true};
120 put_str(&b, "CONNECT ");
121 put_str(&b, options_json);
122 put_str(&b, "\r\n");
123 return finish(&b);
124}
125
126size_t pc_nats_build_pub(char *buf, size_t cap, const char *subject, const char *reply_to, const uint8_t *payload,
127 size_t payload_len)
128{
129 if (!buf || !subject || (payload_len && !payload))
130 {
131 return 0;
132 }
133 Buf b = {buf, cap, 0, true};
134 put_str(&b, "PUB ");
135 put_str(&b, subject);
136 if (reply_to)
137 {
138 put_ch(&b, ' ');
139 put_str(&b, reply_to);
140 }
141 put_ch(&b, ' ');
142 put_uint(&b, payload_len);
143 put_str(&b, "\r\n");
144 put_bytes(&b, payload, payload_len);
145 put_str(&b, "\r\n");
146 return finish(&b);
147}
148
149size_t pc_nats_build_hpub(char *buf, size_t cap, const char *subject, const char *reply_to, const char *headers,
150 size_t headers_len, const uint8_t *payload, size_t payload_len)
151{
152 if (!buf || !subject || !headers || headers_len == 0 || (payload_len && !payload))
153 {
154 return 0;
155 }
156 Buf b = {buf, cap, 0, true};
157 put_str(&b, "HPUB ");
158 put_str(&b, subject);
159 if (reply_to)
160 {
161 put_ch(&b, ' ');
162 put_str(&b, reply_to);
163 }
164 put_ch(&b, ' ');
165 put_uint(&b, headers_len); // hdr_len
166 put_ch(&b, ' ');
167 put_uint(&b, headers_len + payload_len); // total_len = headers + payload
168 put_str(&b, "\r\n");
169 put_bytes(&b, (const uint8_t *)headers, headers_len);
170 put_bytes(&b, payload, payload_len);
171 put_str(&b, "\r\n");
172 return finish(&b);
173}
174
175size_t pc_nats_build_sub(char *buf, size_t cap, const char *subject, const char *queue, const char *sid)
176{
177 if (!buf || !subject || !sid)
178 {
179 return 0;
180 }
181 Buf b = {buf, cap, 0, true};
182 put_str(&b, "SUB ");
183 put_str(&b, subject);
184 if (queue)
185 {
186 put_ch(&b, ' ');
187 put_str(&b, queue);
188 }
189 put_ch(&b, ' ');
190 put_str(&b, sid);
191 put_str(&b, "\r\n");
192 return finish(&b);
193}
194
195size_t pc_nats_build_unsub(char *buf, size_t cap, const char *sid, uint32_t max_msgs, bool with_max)
196{
197 if (!buf || !sid)
198 {
199 return 0;
200 }
201 Buf b = {buf, cap, 0, true};
202 put_str(&b, "UNSUB ");
203 put_str(&b, sid);
204 if (with_max)
205 {
206 put_ch(&b, ' ');
207 put_uint(&b, max_msgs);
208 }
209 put_str(&b, "\r\n");
210 return finish(&b);
211}
212
213size_t pc_nats_build_ping(char *buf, size_t cap)
214{
215 Buf b = {buf, cap, 0, true};
216 put_str(&b, "PING\r\n");
217 return finish(&b);
218}
219
220size_t pc_nats_build_pong(char *buf, size_t cap)
221{
222 Buf b = {buf, cap, 0, true};
223 put_str(&b, "PONG\r\n");
224 return finish(&b);
225}
226
227// Find the CRLF that ends the control line; returns the index of '\r', or len if absent.
228static size_t find_crlf(const char *buf, size_t len)
229{
230 for (size_t i = 0; i + 1 < len; i++)
231 {
232 if (buf[i] == '\r' && buf[i + 1] == '\n')
233 {
234 return i;
235 }
236 }
237 return len;
238}
239
240// Decimal parse of [s, s+n); false on a non-digit.
241static bool parse_uint(const char *s, size_t n, size_t *out)
242{
243 if (n == 0) // GCOVR_EXCL_BR_LINE unreachable: the sole caller passes a space-delimited MSG token, always
244 // length>=1 (see below)
245 {
246 return false; // GCOVR_EXCL_LINE unreachable: the sole caller passes a space-delimited MSG token, always
247 // length>=1
248 }
249 size_t v = 0;
250 for (size_t i = 0; i < n; i++)
251 {
252 if (s[i] < '0' || s[i] > '9')
253 {
254 return false;
255 }
256 v = v * 10 + (size_t)(s[i] - '0');
257 }
258 *out = v;
259 return true;
260}
261
262// True when the verb token at buf matches op (and is followed by a space or end-of-token).
263static bool verb_is(const char *buf, size_t line_len, const char *op)
264{
265 size_t n = strnlen(op, line_len + 1);
266 if (line_len < n)
267 {
268 return false;
269 }
270 if (memcmp(buf, op, n) != 0)
271 {
272 return false;
273 }
274 return line_len == n || buf[n] == ' ' || buf[n] == '\t';
275}
276
277bool pc_nats_parse(const char *buf, size_t len, NatsMsg *out, size_t *consumed)
278{
279 if (!buf || !out || !consumed)
280 {
281 return false;
282 }
283 size_t crlf = find_crlf(buf, len);
284 if (crlf == len)
285 {
286 return false; // control line not fully buffered
287 }
288 size_t line_len = crlf;
289 size_t after_line = crlf + 2;
290
291 out->subject = out->sid = out->reply = out->arg = nullptr;
292 out->subject_len = out->sid_len = out->reply_len = out->arg_len = 0;
293 out->payload = nullptr;
294 out->payload_len = 0;
295 out->headers = nullptr;
296 out->headers_len = 0;
297
298 if (verb_is(buf, line_len, "PING"))
299 {
300 out->type = NatsMsgType::NATS_PING;
301 *consumed = after_line;
302 return true;
303 }
304 if (verb_is(buf, line_len, "PONG"))
305 {
306 out->type = NatsMsgType::NATS_PONG;
307 *consumed = after_line;
308 return true;
309 }
310 if (verb_is(buf, line_len, "+OK"))
311 {
312 out->type = NatsMsgType::NATS_OK;
313 *consumed = after_line;
314 return true;
315 }
316 if (verb_is(buf, line_len, "-ERR") || verb_is(buf, line_len, "INFO"))
317 {
318 out->type = (buf[0] == '-') ? NatsMsgType::NATS_ERR : NatsMsgType::NATS_INFO;
319 size_t a = 4; // skip the verb: both "-ERR" and "INFO" are 4 chars
320 while (a < line_len && (buf[a] == ' ' || buf[a] == '\t'))
321 {
322 a++;
323 }
324 out->arg = buf + a;
325 out->arg_len = line_len - a;
326 *consumed = after_line;
327 return true;
328 }
329 if (verb_is(buf, line_len, "MSG"))
330 {
331 // MSG <subject> <sid> [reply-to] <#bytes>
332 const char *tok[4];
333 size_t tlen[4];
334 size_t ntok = 0;
335 size_t i = 3; // past "MSG"
336 while (i < line_len && ntok < 4)
337 {
338 while (i < line_len && (buf[i] == ' ' || buf[i] == '\t'))
339 {
340 i++;
341 }
342 if (i >= line_len)
343 {
344 break;
345 }
346 size_t start = i;
347 while (i < line_len && buf[i] != ' ' && buf[i] != '\t')
348 {
349 i++;
350 }
351 tok[ntok] = buf + start;
352 tlen[ntok] = i - start;
353 ntok++;
354 }
355 if (ntok != 3 && ntok != 4) // subject sid [reply] size
356 {
357 return false;
358 }
359 size_t size;
360 if (!parse_uint(tok[ntok - 1], tlen[ntok - 1], &size)) // the last token is the byte count
361 {
362 return false;
363 }
364 // Bound the byte count against the remaining capacity without adding it (a 32-bit
365 // size_t would wrap if we computed after_line + size + 2 first).
366 if (after_line + 2 > len || size > len - after_line - 2)
367 {
368 return false; // payload + trailing CRLF not fully buffered
369 }
370 size_t total = after_line + size + 2;
371 out->type = NatsMsgType::NATS_MSG;
372 out->subject = tok[0];
373 out->subject_len = tlen[0];
374 out->sid = tok[1];
375 out->sid_len = tlen[1];
376 if (ntok == 4)
377 {
378 out->reply = tok[2];
379 out->reply_len = tlen[2];
380 }
381 out->payload = (const uint8_t *)(buf + after_line);
382 out->payload_len = size;
383 *consumed = total;
384 return true;
385 }
386 if (verb_is(buf, line_len, "HMSG"))
387 {
388 // HMSG <subject> <sid> [reply-to] <hdr_len> <total_len>
389 const char *tok[5];
390 size_t tlen[5];
391 size_t ntok = 0;
392 size_t i = 4; // past "HMSG"
393 while (i < line_len && ntok < 5)
394 {
395 while (i < line_len && (buf[i] == ' ' || buf[i] == '\t'))
396 {
397 i++;
398 }
399 if (i >= line_len)
400 {
401 break;
402 }
403 size_t start = i;
404 while (i < line_len && buf[i] != ' ' && buf[i] != '\t')
405 {
406 i++;
407 }
408 tok[ntok] = buf + start;
409 tlen[ntok] = i - start;
410 ntok++;
411 }
412 if (ntok != 4 && ntok != 5) // subject sid [reply] hdr_len total_len
413 {
414 return false;
415 }
416 size_t hdr_len, total_size;
417 if (!parse_uint(tok[ntok - 2], tlen[ntok - 2], &hdr_len) ||
418 !parse_uint(tok[ntok - 1], tlen[ntok - 1], &total_size))
419 {
420 return false;
421 }
422 if (hdr_len > total_size) // the header block cannot exceed the header+payload total
423 {
424 return false;
425 }
426 // Bound the total against the remaining capacity without overflowing size_t.
427 if (after_line + 2 > len || total_size > len - after_line - 2)
428 {
429 return false; // headers + payload + trailing CRLF not fully buffered
430 }
431 size_t total = after_line + total_size + 2;
432 out->type = NatsMsgType::NATS_MSG;
433 out->subject = tok[0];
434 out->subject_len = tlen[0];
435 out->sid = tok[1];
436 out->sid_len = tlen[1];
437 if (ntok == 5)
438 {
439 out->reply = tok[2];
440 out->reply_len = tlen[2];
441 }
442 out->headers = buf + after_line;
443 out->headers_len = hdr_len;
444 out->payload = (const uint8_t *)(buf + after_line + hdr_len);
445 out->payload_len = total_size - hdr_len;
446 *consumed = total;
447 return true;
448 }
449
450 out->type = NatsMsgType::NATS_UNKNOWN;
451 *consumed = after_line;
452 return true;
453}
454
455#endif // PC_ENABLE_NATS
NATS client protocol codec (PC_ENABLE_NATS) - zero-heap builder + parser for the text-based NATS pub/...