From c7d98f3714b1080eabcb9dfe37e9862f812d5470 Mon Sep 17 00:00:00 2001 From: ATH User Date: Sun, 4 Oct 2026 15:10:56 -0400 Subject: [PATCH 1/2] fix(q): harden codec and connection handling --- embed/rayforce_q.c | 18 ++-- q.c | 158 ++++++++++++++++++++++++----- q.h | 10 +- q_server.c | 161 +++++++++++++++++++++++++++--- q_server.h | 3 + test/driver.c | 243 +++++++++++++++++++++++++++++++++++++++++++-- 6 files changed, 531 insertions(+), 62 deletions(-) diff --git a/embed/rayforce_q.c b/embed/rayforce_q.c index 95d81d3..60abc0e 100644 --- a/embed/rayforce_q.c +++ b/embed/rayforce_q.c @@ -131,8 +131,8 @@ static ray_t *qb_connect(ray_t **args, int64_t n) { return ray_error("type", ".q.connect: password must be a string"); char user[128], password[128]; - ray_t *arg_err = q_copy_str_arg(n >= 3 ? args[2] : NULL, "user", user, - sizeof user); + ray_t *arg_err = + q_copy_str_arg(n >= 3 ? args[2] : NULL, "user", user, sizeof user); if (arg_err) return arg_err; arg_err = q_copy_str_arg(n >= 4 ? args[3] : NULL, "password", password, @@ -146,8 +146,7 @@ static ray_t *qb_connect(ray_t **args, int64_t n) { if (!tok) return ray_error("type", ".q.connect: timeout must be an integer"); if (timeout_value < 0 || timeout_value > INT_MAX) - return ray_error("range", - ".q.connect: timeout must be in 0..INT_MAX ms"); + return ray_error("range", ".q.connect: timeout must be in 0..INT_MAX ms"); timeout_ms = (int)timeout_value; } @@ -185,6 +184,8 @@ static ray_t *qb_send(ray_t *handle, ray_t *msg) { ray_poll_t *poll = q_poll(); if (poll != NULL) return q_conn_send(poll, fd, msg); + if (fd > INT_MAX) + return ray_error("range", ".q.send: handle exceeds the socket range"); char err[128] = {0}; ray_t *res = q_send((int)fd, msg, err, sizeof err); @@ -208,9 +209,11 @@ static ray_t *qb_close(ray_t *handle) { ray_poll_t *poll = q_poll(); if (poll != NULL) { - if (ray_poll_get(poll, fd) == NULL) - return ray_error("handle", ".q.close: not an open connection"); + if (!q_conn_is_handle(poll, fd)) + return ray_error("handle", ".q.close: not an open Q connection"); q_conn_close(poll, fd); + } else if (fd > INT_MAX) { + return ray_error("range", ".q.close: handle exceeds the socket range"); } else if (q_close((int)fd) < 0) { return ray_error("handle", ".q.close: not an open connection"); } @@ -252,7 +255,8 @@ int64_t q_serve_from_args(ray_poll_t *poll, int argc, char **argv) { break; if (strcmp(argv[i], "-q") == 0 || strcmp(argv[i], "--q-serve") == 0) { if (i + 1 >= argc) { - fprintf(stderr, "q: missing port after %s (expected 1..65535)\n", argv[i]); + fprintf(stderr, "q: missing port after %s (expected 1..65535)\n", + argv[i]); return -1; } int port = 0; diff --git a/q.c b/q.c index 596e45f..9fd25a5 100644 --- a/q.c +++ b/q.c @@ -171,8 +171,10 @@ static ssize_t q_recv_all(int fd, void *buf, size_t n) { uint8_t *p = (uint8_t *)buf; while (total < n) { ssize_t r = recv(fd, p + total, n - total, 0); - if (r == 0) + if (r == 0) { + errno = ECONNRESET; return -1; + } if (r < 0) { if (errno == EINTR) continue; @@ -328,9 +330,14 @@ static int8_t q_type_of(int8_t ray_type) { } } -static int64_t q_size_obj(ray_t *obj); +/* Keep all recursive codec walks bounded. A frame-size limit alone does not + * prevent a tiny, deeply nested list from exhausting the C stack. */ +#define Q_MAX_NESTING 128 +#define Q_SIZE_OVERFLOW (-1) + +static int64_t q_size_obj(ray_t *obj, int depth); static int64_t q_ser_obj(uint8_t *buf, ray_t *obj); -static ray_t *q_des_obj(uint8_t **buf, int64_t *len); +static ray_t *q_des_obj(uint8_t **buf, int64_t *len, int depth); /* Resolve element i of a SYM vector to its string atom. Splayed/mmap SYM * vectors carry narrow storage widths (attrs low bits, RAY_SYM_W8..W64) @@ -376,7 +383,25 @@ static ray_t *q_make_table(ray_t *keys, ray_t *vals) { return tbl; } -static int64_t q_size_obj(ray_t *obj) { +static int q_add_size(int64_t *size, int64_t part) { + if (part < 0 || *size > INT64_MAX - part) + return -1; + *size += part; + return 0; +} + +static int q_is_keyed_table(ray_t *obj) { + if (obj == NULL || obj->type != RAY_LIST || !(obj->attrs & RAY_ATTR_DICT) || + obj->len != 2) + return 0; + ray_t **items = (ray_t **)ray_data(obj); + return items[0] != NULL && items[1] != NULL && items[0]->type == RAY_TABLE && + items[1]->type == RAY_TABLE; +} + +static int64_t q_size_obj(ray_t *obj, int depth) { + if (depth > Q_MAX_NESTING) + return Q_SIZE_OVERFLOW; if (obj == NULL || obj == RAY_NULL_OBJ) return 1 + 1; /* identity type + primitive code */ @@ -412,27 +437,44 @@ static int64_t q_size_obj(ray_t *obj) { case RAY_STR: { /* Single char goes as -KC atom; multi-char as KC vector. */ int64_t n = (int64_t)ray_str_len(obj); + if (n > UINT32_MAX) + return Q_SIZE_OVERFLOW; if (n == 1) return 1 + 1; return 1 + 1 + 4 + n; } } - return 0; + return Q_SIZE_OVERFLOW; } /* Vectors / containers */ if (t == RAY_LIST) { + if (obj->len < 0 || obj->len > UINT32_MAX) + return Q_SIZE_OVERFLOW; + if (q_is_keyed_table(obj)) { + ray_t **items = (ray_t **)ray_data(obj); + int64_t keys = q_size_obj(items[0], depth + 1); + int64_t vals = q_size_obj(items[1], depth + 1); + if (keys < 0 || vals < 0 || keys > INT64_MAX - vals - 1) + return Q_SIZE_OVERFLOW; + return 1 + keys + vals; + } int64_t size = 1 + 1 + 4; ray_t **elems = (ray_t **)ray_data(obj); - for (int64_t i = 0; i < obj->len; i++) - size += q_size_obj(elems[i]); + for (int64_t i = 0; i < obj->len; i++) { + if (q_add_size(&size, q_size_obj(elems[i], depth + 1)) < 0) + return Q_SIZE_OVERFLOW; + } return size; } if (t == RAY_SYM) { + if (obj->len < 0 || obj->len > UINT32_MAX) + return Q_SIZE_OVERFLOW; int64_t size = 1 + 1 + 4; for (int64_t i = 0; i < obj->len; i++) { ray_t *s = q_sym_vec_str(obj, i); /* borrowed */ - size += s ? (int64_t)ray_str_len(s) + 1 : 1; + if (q_add_size(&size, s ? (int64_t)ray_str_len(s) + 1 : 1) < 0) + return Q_SIZE_OVERFLOW; } return size; } @@ -441,20 +483,37 @@ static int64_t q_size_obj(ray_t *obj) { * list-0 of column vectors. RAY_TABLE is opaque - extract via the * public accessors instead of indexing ray_data directly. */ int64_t ncols = ray_table_ncols(obj); + if (ncols < 0 || ncols > UINT32_MAX) + return Q_SIZE_OVERFLOW; int64_t names = 1 + 1 + 4; /* KS type + attrs + count */ for (int64_t i = 0; i < ncols; i++) { ray_t *s = ray_sym_str(ray_table_col_name(obj, i)); - names += s ? (int64_t)ray_str_len(s) + 1 : 1; + if (q_add_size(&names, s ? (int64_t)ray_str_len(s) + 1 : 1) < 0) { + if (s) + ray_release(s); + return Q_SIZE_OVERFLOW; + } if (s) ray_release(s); } int64_t cols = 1 + 1 + 4; /* list type + attrs + count */ - for (int64_t i = 0; i < ncols; i++) - cols += q_size_obj(ray_table_get_col_idx(obj, i)); - return 3 + names + cols; /* XT + attrs + XD */ + for (int64_t i = 0; i < ncols; i++) { + if (q_add_size(&cols, + q_size_obj(ray_table_get_col_idx(obj, i), depth + 1)) < 0) + return Q_SIZE_OVERFLOW; + } + int64_t total = 3; + if (q_add_size(&total, names) < 0 || q_add_size(&total, cols) < 0) + return Q_SIZE_OVERFLOW; + return total; /* XT + attrs + XD */ + } + if (t == RAY_DICT) { + int64_t keys = q_size_obj(ray_dict_keys(obj), depth + 1); + int64_t vals = q_size_obj(ray_dict_vals(obj), depth + 1); + if (keys < 0 || vals < 0 || keys > INT64_MAX - vals - 1) + return Q_SIZE_OVERFLOW; + return 1 + keys + vals; } - if (t == RAY_DICT) - return 1 + q_size_obj(ray_dict_keys(obj)) + q_size_obj(ray_dict_vals(obj)); if (t == RAY_ERROR) { const char *msg = ray_err_code(obj); int64_t n = msg ? (int64_t)strlen(msg) : 0; @@ -464,18 +523,22 @@ static int64_t q_size_obj(ray_t *obj) { return 1 + 1 + 4; if (t == RAY_STR) { + if (obj->len < 0 || obj->len > UINT32_MAX) + return Q_SIZE_OVERFLOW; int64_t size = 1 + 1 + 4; /* list type + attrs + count */ for (int64_t i = 0; i < obj->len; i++) { size_t slen = 0; ray_str_vec_get(obj, i, &slen); - size += 1 + 1 + 4 + (int64_t)slen; /* KC vec header + chars */ + if (q_add_size(&size, 1 + 1 + 4 + (int64_t)slen) < 0) + return Q_SIZE_OVERFLOW; } return size; } int esz = (int)ray_scalar_elem_size(t); - if (esz == 0) - return 0; + if (esz == 0 || obj->len < 0 || obj->len > UINT32_MAX || + obj->len > (INT64_MAX - 6) / esz) + return Q_SIZE_OVERFLOW; return 1 + 1 + 4 + obj->len * esz; } @@ -489,7 +552,22 @@ static int64_t q_ser_obj(uint8_t *buf, ray_t *obj) { } int8_t t = obj->type; - *buf++ = (uint8_t)q_type_of(t); + if (t == RAY_LIST && q_is_keyed_table(obj)) { + ray_t **items = (ray_t **)ray_data(obj); + *buf++ = Q_XD; + int64_t r = q_ser_obj(buf, items[0]); + if (r < 0) + return -1; + buf += r; + r = q_ser_obj(buf, items[1]); + if (r < 0) + return -1; + return buf + r - start; + } + int8_t wire_type = q_type_of(t); + if (wire_type == 0 && t != RAY_LIST) + return -1; + *buf++ = (uint8_t)wire_type; if (t < 0) { int abs_t = -t; @@ -917,7 +995,9 @@ static ray_t *q_des_vec_i32_conv(uint8_t **buf, int64_t *len, int8_t ray_type, return vec; } -static ray_t *q_des_obj(uint8_t **buf, int64_t *len) { +static ray_t *q_des_obj(uint8_t **buf, int64_t *len, int depth) { + if (depth > Q_MAX_NESTING) + return ray_error("q: maximum nesting exceeded", NULL); if (*len < 1) return ray_error("q: buffer underflow (type)", NULL); @@ -1110,6 +1190,8 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) { return ray_error("q: buffer underflow", NULL); if (n < 0) return ray_error("q: negative symbol-vec length", NULL); + if ((int64_t)n > *len) + return ray_error("q: symbol vector exceeds remaining frame", NULL); ray_t *vec = ray_sym_vec_new(RAY_SYM_W64, n); if (vec == NULL || RAY_IS_ERR(vec)) { if (vec) @@ -1146,7 +1228,7 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) { return ray_error("q: negative list length", NULL); ray_t *list = ray_list_new(0); for (int32_t i = 0; i < n; i++) { - ray_t *elem = q_des_obj(buf, len); + ray_t *elem = q_des_obj(buf, len, depth + 1); if (elem == NULL || RAY_IS_ERR(elem)) { ray_release(list); return elem ? elem : ray_error("q: list element decode failed", NULL); @@ -1165,10 +1247,10 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) { return ray_error("q: malformed table marker", NULL); (*buf) += 2; *len -= 2; - ray_t *keys = q_des_obj(buf, len); + ray_t *keys = q_des_obj(buf, len, depth + 1); if (keys == NULL || RAY_IS_ERR(keys)) return keys; - ray_t *vals = q_des_obj(buf, len); + ray_t *vals = q_des_obj(buf, len, depth + 1); if (vals == NULL || RAY_IS_ERR(vals)) { ray_release(keys); return vals; @@ -1177,10 +1259,10 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) { } case Q_XD: { /* dict = keys + values; could be a keyed table */ - ray_t *keys = q_des_obj(buf, len); + ray_t *keys = q_des_obj(buf, len, depth + 1); if (keys == NULL || RAY_IS_ERR(keys)) return keys; - ray_t *vals = q_des_obj(buf, len); + ray_t *vals = q_des_obj(buf, len, depth + 1); if (vals == NULL || RAY_IS_ERR(vals)) { ray_release(keys); return vals; @@ -1199,6 +1281,11 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) { kt->attrs |= RAY_ATTR_DICT; return kt; } + if (keys->type == RAY_TABLE || vals->type == RAY_TABLE) { + ray_release(keys); + ray_release(vals); + return ray_error("q: unsupported table dictionary", NULL); + } return ray_dict_new(keys, vals); /* consumes both refs */ } @@ -1373,7 +1460,7 @@ int q_close(int fd) { int q_encode(ray_t *msg, uint8_t **req, int64_t *req_len, char *err, size_t errlen) { - int64_t body_size = q_size_obj(msg); + int64_t body_size = q_size_obj(msg, 0); if (body_size <= 0) { q_set_err(err, errlen, "q: cannot serialize message"); return -1; @@ -1389,7 +1476,7 @@ int q_encode(ray_t *msg, uint8_t **req, int64_t *req_len, char *err, return -1; } int64_t written = q_ser_obj(buf + sizeof(q_header_t), msg); - if (written < 0) { + if (written < 0 || written != body_size) { free(buf); q_set_err(err, errlen, "q: serialization failed"); return -1; @@ -1413,11 +1500,20 @@ int q_exchange(int fd, const uint8_t *req, int64_t req_len, uint8_t **resp, q_set_err(err, errlen, "q: invalid handle"); return -1; } + if (req == NULL || req_len <= 0 || resp == NULL || resp_len == NULL || + compressed == NULL) { + q_set_err(err, errlen, "q: invalid exchange arguments"); + return -1; + } + *resp = NULL; + *resp_len = 0; + *compressed = 0; if (q_send_all(fd, req, (size_t)req_len) < 0) { q_set_err(err, errlen, (errno == EAGAIN || errno == EWOULDBLOCK) ? "q: send timed out" : (errno == EBADF || errno == ENOTSOCK) ? "q: invalid handle" : "q: send failed"); + shutdown(fd, SHUT_RDWR); return -1; } @@ -1427,14 +1523,17 @@ int q_exchange(int fd, const uint8_t *req, int64_t req_len, uint8_t **resp, (errno == EAGAIN || errno == EWOULDBLOCK) ? "q: recv timed out" : "q: recv header failed"); + shutdown(fd, SHUT_RDWR); return -1; } if (header.endianness != Q_LITTLE_ENDIAN) { q_set_err(err, errlen, "q: big-endian peer not supported"); + shutdown(fd, SHUT_RDWR); return -1; } if (header.msgtype != Q_MSG_RESPONSE) { q_set_err(err, errlen, "q: expected response message type"); + shutdown(fd, SHUT_RDWR); return -1; } int64_t body_len = (int64_t)header.size - (int64_t)sizeof header; @@ -1442,11 +1541,13 @@ int q_exchange(int fd, const uint8_t *req, int64_t req_len, uint8_t **resp, q_set_err(err, errlen, body_len <= 0 ? "q: empty response body" : "q: response body too large"); + shutdown(fd, SHUT_RDWR); return -1; } uint8_t *body = (uint8_t *)malloc((size_t)body_len); if (body == NULL) { q_set_err(err, errlen, "q: out of memory"); + shutdown(fd, SHUT_RDWR); return -1; } if (q_recv_all(fd, body, (size_t)body_len) < 0) { @@ -1455,6 +1556,7 @@ int q_exchange(int fd, const uint8_t *req, int64_t req_len, uint8_t **resp, (errno == EAGAIN || errno == EWOULDBLOCK) ? "q: recv timed out" : "q: recv body failed"); + shutdown(fd, SHUT_RDWR); return -1; } *resp = body; @@ -1477,7 +1579,7 @@ ray_t *q_decode(uint8_t *resp, int64_t resp_len, int compressed, char *err, } uint8_t *cursor = decoded; int64_t remaining = decoded_len; - ray_t *result = q_des_obj(&cursor, &remaining); + ray_t *result = q_des_obj(&cursor, &remaining, 0); if (decompressed) free(decompressed); if (result == NULL) @@ -1509,5 +1611,7 @@ ray_t *q_send(int fd, ray_t *msg, char *err, size_t errlen) { ray_t *result = q_decode(resp, resp_len, compressed, err, errlen); free(resp); + if (result == NULL) + shutdown(fd, SHUT_RDWR); return result; } diff --git a/q.h b/q.h index eb43754..9f772d7 100644 --- a/q.h +++ b/q.h @@ -32,8 +32,9 @@ * converts between the wire and rayforce core `ray_t` objects. * * A connection handle is the raw socket file descriptor (>= 0). There is no - * global connection table, so the client is thread-safe and unbounded: each - * fd is independent and owned by the caller. + * global connection table. Independent fds may be used concurrently; calls + * sharing one fd must serialize complete request/response exchanges. Each fd + * remains owned by the caller. */ #include @@ -58,7 +59,10 @@ int q_close(int fd); int q_encode(ray_t *msg, uint8_t **req, int64_t *req_len, char *err, size_t errlen); -/* Send a pre-encoded request on `fd` and read the raw response body. */ +/* Send a pre-encoded request on `fd` and read the raw response body. A + * transport or framing failure shuts down the socket because the request or + * response boundary may be out of sync; the caller still owns and must close + * the fd. */ int q_exchange(int fd, const uint8_t *req, int64_t req_len, uint8_t **resp, int64_t *resp_len, int *compressed, char *err, size_t errlen); diff --git a/q_server.c b/q_server.c index 349dd4e..9bdcff3 100644 --- a/q_server.c +++ b/q_server.c @@ -47,6 +47,8 @@ #include "ops/ops.h" /* ray_is_lazy, ray_lazy_materialize */ #include +#include +#include #include #include #include @@ -70,6 +72,7 @@ typedef struct { #define Q_CAP_MAX 3 /* matches q.c client capability */ #define Q_MAX_BODY ((int64_t)256 << 20) /* reject absurd frames (256 MiB) */ #define Q_MAX_HANDSHAKE 512 /* credential blob upper bound */ +#define Q_CONN_MAGIC UINT32_C(0x7150434e) /* Per-connection rx state. `hdr` carries the current frame's header between * the header and body reads; `cap` accumulates the last byte seen during the @@ -78,6 +81,7 @@ typedef struct { * The sync_* triple is the outbound side's round-trip slot: q_conn_send parks * on the connection and the rx machine deposits the RESPONSE frame here. */ typedef struct { + uint32_t magic; q_header_t hdr; uint8_t cap; int hs_len; @@ -149,7 +153,15 @@ static ray_t *eval_request(ray_t *req) { /* Encode `result` as a Q response and write it. A NULL result encodes as the * identity (::). q_encode stamps the header as SYNC, so flip the message-type * byte to RESPONSE before sending. */ -static void q_send_result(ray_sock_t fd, ray_t *result) { +static int64_t q_send_fn(int64_t fd, uint8_t *buf, int64_t len) { + int flags = 0; +#ifdef MSG_NOSIGNAL + flags |= MSG_NOSIGNAL; +#endif + return (int64_t)send((int)fd, buf, (size_t)len, flags); +} + +static int q_send_result(ray_poll_t *poll, ray_selector_t *sel, ray_t *result) { uint8_t *buf = NULL; int64_t len = 0; char err[128] = {0}; @@ -158,11 +170,53 @@ static void q_send_result(ray_sock_t fd, ray_t *result) { int rc = q_encode(e, &buf, &len, err, sizeof err); q_release_any(e); if (rc < 0) - return; + return -1; } buf[1] = Q_MSG_RESPONSE; - ray_sock_send(fd, buf, (size_t)len); + int64_t queued = 0; + for (ray_poll_buf_t *p = sel->tx.buf; p != NULL; p = p->next) { + int64_t pending = p->size - p->offset; + if (pending < 0 || queued > Q_MAX_BODY - pending) { + free(buf); + return -1; + } + queued += pending; + } + if (len > Q_MAX_BODY - queued) { + free(buf); + return -1; + } + ray_poll_buf_t *out = ray_poll_buf_new(len); + if (out == NULL) { + free(buf); + return -1; + } + memcpy(out->data, buf, (size_t)len); free(buf); + if (sel->tx.buf == NULL) { + sel->tx.buf = out; + } else { + ray_poll_buf_t *tail = sel->tx.buf; + while (tail->next != NULL) + tail = tail->next; + tail->next = out; + } + ray_poll_tx_request(poll, sel); + return 0; +} + +static void q_close_listener(ray_poll_t *poll, ray_selector_t *sel) { + (void)poll; + ray_sock_close((ray_sock_t)sel->fd); +} + +static void q_close_connection_sigpipe(ray_sock_t fd) { +#ifdef SO_NOSIGPIPE + int yes = 1; + setsockopt((int)fd, SOL_SOCKET, SO_NOSIGPIPE, &yes, sizeof yes); +#else + (void)fd; +#endif } static int64_t q_recv_fn(int64_t fd, uint8_t *buf, int64_t len) { @@ -180,17 +234,20 @@ static ray_t *q_accept(ray_poll_t *poll, ray_selector_t *sel) { if (nfd == RAY_INVALID_SOCK) return NULL; ray_sock_set_nonblocking(nfd); + q_close_connection_sigpipe(nfd); q_conn_t *cd = (q_conn_t *)calloc(1, sizeof *cd); if (cd == NULL) { ray_sock_close(nfd); return NULL; } + cd->magic = Q_CONN_MAGIC; ray_poll_reg_t reg = {0}; reg.fd = (int64_t)nfd; reg.type = RAY_SEL_SOCKET; reg.recv_fn = q_recv_fn; + reg.send_fn = q_send_fn; reg.read_fn = q_read_handshake; reg.close_fn = q_on_close; reg.data = cd; @@ -247,8 +304,10 @@ static ray_t *q_read_header(ray_poll_t *poll, ray_selector_t *sel) { return NULL; } if (body == 0) { - if (cd->hdr.msgtype != 0) /* sync -> identity reply */ - q_send_result((ray_sock_t)sel->fd, NULL); + if (cd->hdr.msgtype != Q_MSG_ASYNC && q_send_result(poll, sel, NULL) < 0) { + ray_poll_deregister(poll, sel->id); + return NULL; + } ray_poll_rx_request(poll, sel, (int64_t)sizeof(q_header_t)); return NULL; } @@ -306,8 +365,8 @@ static ray_t *q_read_body(ray_poll_t *poll, ray_selector_t *sel) { if (hdr.msgtype != Q_MSG_ASYNC) { /* sync expects a response, async does not */ ray_selector_t *cur = ray_poll_get(poll, id); /* eval may have closed it */ - if (cur) - q_send_result((ray_sock_t)cur->fd, result); + if (cur && q_send_result(poll, cur, result) < 0) + ray_poll_deregister(poll, id); } else if (result != NULL && RAY_IS_ERR(result)) { /* Async has no reply channel, so an error here would vanish silently — * the one place an operator could learn a push handler is broken. */ @@ -326,6 +385,7 @@ static void q_on_close(ray_poll_t *poll, ray_selector_t *sel) { * mid-round-trip) would otherwise leak. */ if (cd->sync_resp) q_release_any(cd->sync_resp); + cd->magic = 0; free(cd); sel->data = NULL; } @@ -351,6 +411,7 @@ int64_t q_serve(ray_poll_t *poll, int port) { reg.fd = (int64_t)fd; reg.type = RAY_SEL_SOCKET; reg.read_fn = q_accept; + reg.close_fn = q_close_listener; int64_t id = ray_poll_register(poll, ®); if (id < 0) { @@ -398,8 +459,63 @@ static int q_conn_pump(ray_poll_t *poll, int64_t id) { } /* Phase buffer complete — advance the state machine. read_fn may * deregister or re-arm the selector; the loop re-validates either way. */ + void *data = sel->data; sel->rx.read_fn(poll, sel); + sel = ray_poll_get(poll, id); + if (sel == NULL || sel->data != data) + return -1; + if (((q_conn_t *)data)->sync_ready) + return 0; + } +} + +/* Send a complete frame without letting a slow peer park the shared event + * loop forever. The socket is nonblocking; a configured timeout covers both + * readiness waits and short writes. */ +static int q_send_frame_until(int fd, const uint8_t *buf, size_t len, + int bounded, int64_t deadline_ms) { + size_t sent = 0; + while (sent < len) { + int flags = 0; +#ifdef MSG_NOSIGNAL + flags |= MSG_NOSIGNAL; +#endif + ssize_t n = send(fd, buf + sent, len - sent, flags); + if (n > 0) { + sent += (size_t)n; + continue; + } + if (n == 0) { + errno = EPIPE; + return -1; + } + if (errno == EINTR) + continue; + if (errno != EAGAIN && errno != EWOULDBLOCK) + return -1; + + int wait_ms = -1; + if (bounded) { + int64_t left = deadline_ms - ray_time_now_ms(); + if (left <= 0) { + errno = ETIMEDOUT; + return -1; + } + wait_ms = left > INT_MAX ? INT_MAX : (int)left; + } + struct pollfd pfd = {.fd = fd, .events = POLLOUT}; + int ready; + do { + ready = poll(&pfd, 1, wait_ms); + } while (ready < 0 && errno == EINTR); + if (ready == 0) { + errno = ETIMEDOUT; + return -1; + } + if (ready < 0) + return -1; } + return 0; } /* The per-operation timeout q_connect applied to the socket (SO_RCVTIMEO). @@ -419,16 +535,19 @@ int64_t q_conn_attach(ray_poll_t *poll, int fd) { return -1; int timeout_ms = q_sock_recv_timeout_ms(fd); ray_sock_set_nonblocking((ray_sock_t)fd); + q_close_connection_sigpipe((ray_sock_t)fd); q_conn_t *cd = (q_conn_t *)calloc(1, sizeof *cd); if (cd == NULL) return -1; + cd->magic = Q_CONN_MAGIC; cd->sync_timeout_ms = timeout_ms; ray_poll_reg_t reg = {0}; reg.fd = (int64_t)fd; reg.type = RAY_SEL_SOCKET; reg.recv_fn = q_recv_fn; + reg.send_fn = q_send_fn; reg.read_fn = q_read_header; /* q_connect already did the handshake */ reg.close_fn = q_on_close; reg.data = cd; @@ -446,7 +565,8 @@ int64_t q_conn_attach(ray_poll_t *poll, int fd) { ray_t *q_conn_send(ray_poll_t *poll, int64_t id, ray_t *msg) { ray_selector_t *sel = poll ? ray_poll_get(poll, id) : NULL; - if (sel == NULL || sel->data == NULL) + if (sel == NULL || sel->close_fn != q_on_close || sel->data == NULL || + ((q_conn_t *)sel->data)->magic != Q_CONN_MAGIC) return ray_error("handle", "q: not an open connection"); q_conn_t *cd = (q_conn_t *)sel->data; if (cd->sync_waiting) @@ -459,10 +579,16 @@ ray_t *q_conn_send(ray_poll_t *poll, int64_t id, ray_t *msg) { char err[128] = {0}; if (q_encode(msg, &req, &req_len, err, sizeof err) < 0) return ray_error("send", "%s", err[0] ? err : "q: encode failed"); - int64_t sent = ray_sock_send((ray_sock_t)sel->fd, req, (size_t)req_len); + int bounded = cd->sync_timeout_ms > 0; + int64_t deadline_ms = bounded ? ray_time_now_ms() + cd->sync_timeout_ms : 0; + int send_rc = q_send_frame_until((int)sel->fd, req, (size_t)req_len, bounded, + deadline_ms); free(req); - if (sent < 0) - return ray_error("io", "q: send failed"); + if (send_rc < 0) { + const char *code = errno == ETIMEDOUT ? "timeout" : "io"; + ray_poll_deregister(poll, id); + return ray_error(code, "q: send failed"); + } cd->sync_waiting = 1; cd->sync_ready = 0; @@ -470,9 +596,6 @@ ray_t *q_conn_send(ray_poll_t *poll, int64_t id, ray_t *msg) { /* The socket's recv timeout is the whole round-trip's budget, as it is on * the blocking path where the kernel enforces it per recv. */ - int bounded = cd->sync_timeout_ms > 0; - int64_t deadline_ms = bounded ? ray_time_now_ms() + cd->sync_timeout_ms : 0; - for (;;) { /* Process what is already readable, then block for more. The pump can * deregister the selector (peer died), which frees cd — so re-resolve and @@ -515,8 +638,16 @@ ray_t *q_conn_send(ray_poll_t *poll, int64_t id, ray_t *msg) { } } -void q_conn_close(ray_poll_t *poll, int64_t id) { +int q_conn_is_handle(ray_poll_t *poll, int64_t id) { if (poll == NULL) + return 0; + ray_selector_t *sel = ray_poll_get(poll, id); + return sel != NULL && sel->close_fn == q_on_close && sel->data != NULL && + ((q_conn_t *)sel->data)->magic == Q_CONN_MAGIC; +} + +void q_conn_close(ray_poll_t *poll, int64_t id) { + if (!q_conn_is_handle(poll, id)) return; ray_poll_deregister(poll, id); /* q_on_close frees the state and the fd */ } diff --git a/q_server.h b/q_server.h index c9c461a..000e2ca 100644 --- a/q_server.h +++ b/q_server.h @@ -81,6 +81,9 @@ int64_t q_conn_attach(ray_poll_t *poll, int fd); * `timeout` error is returned; the handle is then no longer open. */ ray_t *q_conn_send(ray_poll_t *poll, int64_t id, ray_t *msg); +/* Check whether `id` identifies an open Q connection owned by this module. */ +int q_conn_is_handle(ray_poll_t *poll, int64_t id); + /* Close an attached connection and release its state. */ void q_conn_close(ray_poll_t *poll, int64_t id); diff --git a/test/driver.c b/test/driver.c index 2c2e139..41052a1 100644 --- a/test/driver.c +++ b/test/driver.c @@ -28,14 +28,17 @@ #include "q.h" /* q_decode / q_connect / q_exchange */ #include "q_server.h" /* q_serve */ +#include +#include +#include +#include #include #include #include #include #include -#include -#include #include +#include #include /* Registers `.q.connect` / `.q.send` / `.q.close` */ @@ -373,9 +376,66 @@ static int run_codec_selftest(void) { } release_any(r); + /* An unsupported nested value must fail size calculation before the + * serializer writes a partial type byte past the allocated buffer. */ + r = ray_eval_str("(list (fn [x] x))"); + uint8_t *encoded = NULL; + int64_t encoded_len = 0; + err[0] = '\0'; + if (RAY_IS_ERR(r) || + q_encode(r, &encoded, &encoded_len, err, sizeof err) == 0) { + fprintf(stderr, "codec selftest: nested function was serialized\n"); + failures++; + } + free(encoded); + release_any(r); + + /* A small body with excessive nesting is rejected without exhausting the + * C stack. Every six bytes is a one-item Q general list. */ + const int deep_levels = 5000; + const size_t deep_len = (size_t)deep_levels * 6 + 2; + uint8_t *deep = (uint8_t *)calloc(deep_len, 1); + if (deep == NULL) { + fprintf(stderr, "codec selftest: cannot allocate deep frame\n"); + failures++; + } else { + for (int i = 0; i < deep_levels; i++) { + deep[(size_t)i * 6] = 0; /* general list */ + deep[(size_t)i * 6 + 2] = 1; /* one element */ + } + deep[deep_len - 2] = 101; /* identity */ + r = q_decode(deep, (int64_t)deep_len, 0, err, sizeof err); + if (!RAY_IS_ERR(r)) { + fprintf(stderr, "codec selftest: excessive nesting was not rejected\n"); + failures++; + } + release_any(r); + free(deep); + } + + /* Keyed-table dictionaries keep their Q dictionary marker on re-encode. */ + const uint8_t keyed_table[] = { + 0x63, 0x62, 0x00, 0x63, 0x0b, 0x00, 0x01, 0x00, 0x00, 0x00, 0x6b, + 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x06, 0x00, 0x01, 0x00, + 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x62, 0x00, 0x63, 0x0b, 0x00, + 0x01, 0x00, 0x00, 0x00, 0x76, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, + 0x00, 0x06, 0x00, 0x01, 0x00, 0x00, 0x00, 0x02, 0x00, 0x00, 0x00}; + r = q_decode((uint8_t *)keyed_table, sizeof keyed_table, 0, err, sizeof err); + encoded = NULL; + encoded_len = 0; + if (RAY_IS_ERR(r) || + q_encode(r, &encoded, &encoded_len, err, sizeof err) < 0 || + encoded_len < (int64_t)sizeof(test_q_header_t) + 3 || + encoded[sizeof(test_q_header_t)] != 99) { + fprintf(stderr, "codec selftest: keyed-table marker was not preserved\n"); + failures++; + } + free(encoded); + release_any(r); + err[0] = '\0'; - uint8_t bad_table_marker[] = {98, 0, 0, 11, 0, 0, 0, 0, 0, 0, 0, 0, - 0, 0, 0, 0, 0, 0, 0, 0}; + uint8_t bad_table_marker[] = {98, 0, 0, 11, 0, 0, 0, 0, 0, 0, + 0, 0, 0, 0, 0, 0, 0, 0, 0, 0}; r = q_decode(bad_table_marker, (int64_t)sizeof bad_table_marker, 0, err, sizeof err); ray_t *rs = r ? ray_fmt(r, 0) : NULL; @@ -406,6 +466,25 @@ static int run_codec_selftest(void) { failures++; } + /* A large i64 handle must not truncate to and close an unrelated fd. */ + q_env_register(); + int unrelated_fd = open("/dev/null", O_RDONLY); + if (unrelated_fd < 0) { + perror("codec selftest: open handle sentinel"); + failures++; + } else { + char close_expr[96]; + snprintf(close_expr, sizeof close_expr, "(.q.close %lld)", + (long long)unrelated_fd + (1LL << 32)); + r = ray_eval_str(close_expr); + if (!RAY_IS_ERR(r) || fcntl(unrelated_fd, F_GETFD) < 0) { + fprintf(stderr, "codec selftest: oversized handle closed a valid fd\n"); + failures++; + } + release_any(r); + close(unrelated_fd); + } + /* Unit-converted temporals, raw wire bytes so no q is needed: the q * type byte, then the little-endian payload. */ struct { @@ -492,10 +571,10 @@ static int run_codec_selftest(void) { _exit(0); } close(listener); - int bad_cap = - q_connect("127.0.0.1", ntohs(addr.sin_port), "", "", 1000); + int bad_cap = q_connect("127.0.0.1", ntohs(addr.sin_port), "", "", 1000); if (bad_cap != Q_ERR_HANDSHAKE) { - fprintf(stderr, "codec selftest: invalid handshake capability accepted\n"); + fprintf(stderr, + "codec selftest: invalid handshake capability accepted\n"); if (bad_cap >= 0) q_close(bad_cap); failures++; @@ -514,7 +593,36 @@ handshake_done:; fprintf(stderr, "codec selftest: server accepted out-of-range port\n"); failures++; } + struct sockaddr_in probe_addr = {0}; + int probe_fd = socket(AF_INET, SOCK_STREAM, 0); + int port = 0; + if (probe_fd >= 0) { + probe_addr.sin_family = AF_INET; + probe_addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK); + if (bind(probe_fd, (struct sockaddr *)&probe_addr, sizeof probe_addr) == + 0) { + socklen_t probe_len = sizeof probe_addr; + if (getsockname(probe_fd, (struct sockaddr *)&probe_addr, &probe_len) == + 0) + port = ntohs(probe_addr.sin_port); + } + close(probe_fd); + } + int64_t listener_id = port ? q_serve(poll, port) : -1; + ray_selector_t *listener_sel = + listener_id >= 0 ? ray_poll_get(poll, listener_id) : NULL; + int listener_fd = listener_sel ? (int)listener_sel->fd : -1; + if (listener_id < 0 || listener_fd < 0) { + fprintf(stderr, "codec selftest: cannot create close-test listener\n"); + failures++; + } ray_poll_destroy(poll); + if (listener_fd >= 0 && + (fcntl(listener_fd, F_GETFD) >= 0 || errno != EBADF)) { + fprintf(stderr, "codec selftest: poll destroy left listener open\n"); + failures++; + close(listener_fd); + } } ray_runtime_destroy(rt); @@ -574,14 +682,46 @@ static int run_exchange_selftest(void) { sizeof err); if (rc == 0 || strstr(err, "send") == NULL) { fprintf(stderr, - "exchange selftest: closed peer did not fail cleanly: %s\n", - err); + "exchange selftest: closed peer did not fail cleanly: %s\n", err); failures++; } free(resp); close(sv[0]); } + /* A timeout poisons the blocking stream: a late reply to request A must + * never be accepted as the response to request B. */ + ray_t *r = NULL; + if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) { + perror("exchange selftest: socketpair late-response"); + failures++; + } else { + struct timeval tv = {.tv_sec = 0, .tv_usec = 50 * 1000}; + setsockopt(sv[0], SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof tv); + ray_t *msg = ray_i64(1); + err[0] = '\0'; + r = q_send(sv[0], msg, err, sizeof err); + if (r != NULL || strstr(err, "timed out") == NULL) { + fprintf(stderr, "exchange selftest: silent request did not time out\n"); + failures++; + } + release_any(r); + const uint8_t late_response[] = {1, 2, 0, 0, 13, 0, 0, + 0, 250, 111, 0, 0, 0}; + (void)send(sv[1], late_response, sizeof late_response, MSG_NOSIGNAL); + err[0] = '\0'; + r = q_send(sv[0], msg, err, sizeof err); + if (r != NULL) { + fprintf(stderr, + "exchange selftest: late response was reused by next request\n"); + failures++; + } + release_any(r); + ray_release(msg); + close(sv[0]); + close(sv[1]); + } + /* Poll-attached round-trip against a peer that never answers: the recv * timeout on the socket (what q_connect's timeout_ms sets) must bound * q_conn_send and close the handle, instead of parking the event loop @@ -605,7 +745,7 @@ static int run_exchange_selftest(void) { /* A regression here is an unbounded wait: let SIGALRM fail the run * instead of hanging it. */ alarm(10); - ray_t *r = id >= 0 ? q_conn_send(poll, id, msg) : NULL; + r = id >= 0 ? q_conn_send(poll, id, msg) : NULL; alarm(0); ray_t *es = (r && RAY_IS_ERR(r)) ? ray_fmt(r, 0) : NULL; const char *ep = es ? ray_str_ptr(es) : ""; @@ -629,6 +769,89 @@ static int run_exchange_selftest(void) { close(sv[1]); } + /* Backpressure during a large write is covered by the same round-trip + * deadline as waiting for the response. */ + if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) { + perror("exchange selftest: socketpair send-timeout"); + failures++; + } else { + ray_runtime_t *rt = ray_runtime_create(0, NULL); + ray_poll_t *poll = rt ? ray_poll_create() : NULL; + if (poll == NULL) { + fprintf(stderr, "exchange selftest: send-timeout runtime/poll failed\n"); + failures++; + close(sv[0]); + } else { + int small_sndbuf = 4096; + struct timeval tv = {.tv_sec = 0, .tv_usec = 100 * 1000}; + setsockopt(sv[0], SOL_SOCKET, SO_SNDBUF, &small_sndbuf, + sizeof small_sndbuf); + setsockopt(sv[0], SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof tv); + int64_t id = q_conn_attach(poll, sv[0]); + ray_t *msg = ray_vec_new(RAY_I32, 1000000); + alarm(5); + ray_t *r = id >= 0 ? q_conn_send(poll, id, msg) : NULL; + alarm(0); + ray_t *es = (r && RAY_IS_ERR(r)) ? ray_fmt(r, 0) : NULL; + const char *ep = es ? ray_str_ptr(es) : ""; + if (id < 0 || r == NULL || !RAY_IS_ERR(r) || + strstr(ep, "timeout") == NULL || ray_poll_get(poll, id) != NULL) { + fprintf(stderr, + "exchange selftest: blocked send did not time out: %s\n", ep); + failures++; + } + if (es) + ray_release(es); + release_any(r); + ray_release(msg); + ray_poll_destroy(poll); + ray_runtime_destroy(rt); + close(sv[1]); + } + } + + /* A complete queued response remains available if the peer then sends FIN. */ + if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) { + perror("exchange selftest: socketpair response-fin"); + failures++; + } else { + ray_runtime_t *rt = ray_runtime_create(0, NULL); + ray_poll_t *poll = rt ? ray_poll_create() : NULL; + if (poll == NULL) { + fprintf(stderr, "exchange selftest: response-fin runtime/poll failed\n"); + failures++; + close(sv[0]); + } else { + test_q_header_t h = {.endianness = 1, + .msgtype = 2, + .compressed = 0, + .reserved = 0, + .size = 13}; + uint8_t response[sizeof h + 5]; + int32_t value = 42; + memcpy(response, &h, sizeof h); + response[sizeof h] = 250; /* Q int atom */ + memcpy(response + sizeof h + 1, &value, sizeof value); + int64_t id = q_conn_attach(poll, sv[0]); + if (send(sv[1], response, sizeof response, MSG_NOSIGNAL) != + (ssize_t)sizeof response) + failures++; + shutdown(sv[1], SHUT_WR); + ray_t *msg = ray_i64(1); + ray_t *r = id >= 0 ? q_conn_send(poll, id, msg) : NULL; + if (r == NULL || RAY_IS_ERR(r) || r->type != -RAY_I32 || r->i32 != 42) { + fprintf(stderr, "exchange selftest: response before FIN was lost\n"); + failures++; + } + release_any(r); + release_any(q_conn_send(poll, id, msg)); /* detects and closes the FIN */ + ray_release(msg); + ray_poll_destroy(poll); + ray_runtime_destroy(rt); + close(sv[1]); + } + } + printf("exchange selftest: %s\n", failures ? "FAIL" : "ok"); return failures ? 1 : 0; } From e675994cf56abe0c53137631ac079e30d8dd9734 Mon Sep 17 00:00:00 2001 From: ATH User Date: Sun, 4 Oct 2026 15:34:39 -0400 Subject: [PATCH 2/2] fix(q): distinguish decode failures from Q errors --- q.c | 16 +++++++++++++--- test/driver.c | 31 +++++++++++++++++++++++++------ 2 files changed, 38 insertions(+), 9 deletions(-) diff --git a/q.c b/q.c index 9fd25a5..80c8f53 100644 --- a/q.c +++ b/q.c @@ -1580,16 +1580,26 @@ ray_t *q_decode(uint8_t *resp, int64_t resp_len, int compressed, char *err, uint8_t *cursor = decoded; int64_t remaining = decoded_len; ray_t *result = q_des_obj(&cursor, &remaining, 0); + int is_q_error_frame = + decoded_len >= 2 && decoded[0] == (uint8_t)Q_ERR && + memchr(decoded + 1, '\0', (size_t)(decoded_len - 1)) != NULL; if (decompressed) free(decompressed); if (result == NULL) q_set_err(err, errlen, "q: deserialization returned null"); - else if (RAY_IS_ERR(result)) { - return result; - } else if (remaining != 0) { + else if (remaining != 0) { q_release_any(result); q_set_err(err, errlen, "q: trailing bytes after object"); return NULL; + } else if (RAY_IS_ERR(result)) { + /* A Q error frame is a valid decoded server response. Other errors + * produced by q_des_obj represent malformed or unsupported wire data and + * must use the documented decode-failure channel instead. */ + if (!is_q_error_frame) { + q_release_any(result); + q_set_err(err, errlen, "q: deserialization failed"); + return NULL; + } } return result; } diff --git a/test/driver.c b/test/driver.c index 41052a1..38237db 100644 --- a/test/driver.c +++ b/test/driver.c @@ -366,6 +366,28 @@ static int run_codec_selftest(void) { } release_any(r); + err[0] = '\0'; + uint8_t unsupported_type[] = {127}; + r = q_decode(unsupported_type, (int64_t)sizeof unsupported_type, 0, err, + sizeof err); + if (r != NULL || strstr(err, "deserialization failed") == NULL) { + fprintf(stderr, "codec selftest: unsupported wire type was not reported as " + "decode failure\n"); + failures++; + } + release_any(r); + + err[0] = '\0'; + uint8_t malformed_qerr[] = {128, 'b', 'a', 'd'}; + r = q_decode(malformed_qerr, (int64_t)sizeof malformed_qerr, 0, err, + sizeof err); + if (r != NULL || err[0] == '\0') { + fprintf(stderr, + "codec selftest: malformed Q error frame was not rejected\n"); + failures++; + } + release_any(r); + err[0] = '\0'; uint8_t qidentity[] = {101, 0}; r = q_decode(qidentity, (int64_t)sizeof qidentity, 0, err, sizeof err); @@ -399,13 +421,14 @@ static int run_codec_selftest(void) { fprintf(stderr, "codec selftest: cannot allocate deep frame\n"); failures++; } else { + err[0] = '\0'; for (int i = 0; i < deep_levels; i++) { deep[(size_t)i * 6] = 0; /* general list */ deep[(size_t)i * 6 + 2] = 1; /* one element */ } deep[deep_len - 2] = 101; /* identity */ r = q_decode(deep, (int64_t)deep_len, 0, err, sizeof err); - if (!RAY_IS_ERR(r)) { + if (r != NULL || err[0] == '\0') { fprintf(stderr, "codec selftest: excessive nesting was not rejected\n"); failures++; } @@ -438,14 +461,10 @@ static int run_codec_selftest(void) { 0, 0, 0, 0, 0, 0, 0, 0, 0, 0}; r = q_decode(bad_table_marker, (int64_t)sizeof bad_table_marker, 0, err, sizeof err); - ray_t *rs = r ? ray_fmt(r, 0) : NULL; - const char *rp = rs ? ray_str_ptr(rs) : err; - if (r == NULL || !RAY_IS_ERR(r) || strstr(rp, "q: malf") == NULL) { + if (r != NULL || err[0] == '\0') { fprintf(stderr, "codec selftest: malformed table marker was accepted\n"); failures++; } - if (rs) - ray_release(rs); release_any(r); uint8_t oversized_compressed[4];