diff --git a/q.c b/q.c index 6a122f8..cb8c17b 100644 --- a/q.c +++ b/q.c @@ -98,9 +98,11 @@ static inline size_t ray_scalar_elem_size(int8_t type) { #define Q_KT 19 /* time */ #define Q_XT 98 /* table */ #define Q_XD 99 /* dict */ +#define Q_ID 101 /* identity */ #define Q_ERR (-128) #define Q_MSG_SYNC 1 +#define Q_MSG_RESPONSE 2 #define Q_MAX_BODY ((int64_t)256 << 20) typedef struct { @@ -118,6 +120,15 @@ static void q_set_err(char *err, size_t errlen, const char *msg) { snprintf(err, errlen, "%s", msg); } +static void q_suppress_sigpipe(int fd) { +#ifdef SO_NOSIGPIPE + int yes = 1; + setsockopt(fd, SOL_SOCKET, SO_NOSIGPIPE, &yes, sizeof yes); +#else + (void)fd; +#endif +} + /* Build a v2 table from a RAY_SYM-vec of column ids and an array of column * vectors. Retains each column internally; the caller keeps ownership of the * inputs. (Local copy so q has no binding-specific dependencies.) */ @@ -145,6 +156,15 @@ static ray_t *q_build_table(const int64_t *col_ids, ray_t *const *cols, return tbl; } +static void q_release_any(ray_t *obj) { + if (obj == NULL) + return; + if (RAY_IS_ERR(obj)) + ray_error_free(obj); + else + ray_release(obj); +} + static ssize_t q_recv_all(int fd, void *buf, size_t n) { size_t total = 0; uint8_t *p = (uint8_t *)buf; @@ -163,10 +183,15 @@ static ssize_t q_recv_all(int fd, void *buf, size_t n) { } static ssize_t q_send_all(int fd, const void *buf, size_t n) { + q_suppress_sigpipe(fd); size_t total = 0; const uint8_t *p = (const uint8_t *)buf; while (total < n) { +#ifdef MSG_NOSIGNAL + ssize_t r = send(fd, p + total, n - total, MSG_NOSIGNAL); +#else ssize_t r = send(fd, p + total, n - total, 0); +#endif if (r <= 0) { if (r < 0 && errno == EINTR) continue; @@ -325,7 +350,7 @@ static ray_t *q_make_table(ray_t *keys, ray_t *vals) { static int64_t q_size_obj(ray_t *obj) { if (obj == NULL || obj == RAY_NULL_OBJ) - return 1 + 1 + 4; /* type + attrs + len(0) */ + return 1 + 1; /* identity type + primitive code */ int8_t t = obj->type; @@ -403,6 +428,8 @@ static int64_t q_size_obj(ray_t *obj) { cols += q_size_obj(ray_table_get_col_idx(obj, i)); return 3 + names + cols; /* XT + attrs + XD */ } + 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; @@ -431,7 +458,7 @@ static int64_t q_ser_obj(uint8_t *buf, ray_t *obj) { uint8_t *start = buf; if (obj == NULL || obj == RAY_NULL_OBJ) { - *buf++ = 101; /* identity / null */ + *buf++ = Q_ID; /* identity / null */ *buf++ = 0; return buf - start; } @@ -568,9 +595,17 @@ static int64_t q_ser_obj(uint8_t *buf, ray_t *obj) { } return buf - start; } - /* RAY_DICT (type code 99) is never produced by v2; dicts present as - * RAY_LIST + RAY_ATTR_DICT and are serialized through the RAY_LIST path - * above (losing the dict-ness on the wire). */ + if (t == RAY_DICT) { + int64_t r = q_ser_obj(buf, ray_dict_keys(obj)); + if (r < 0) + return -1; + buf += r; + r = q_ser_obj(buf, ray_dict_vals(obj)); + if (r < 0) + return -1; + buf += r; + return buf - start; + } if (t == RAY_ERROR) { const char *msg = ray_err_code(obj); size_t n = msg ? strlen(msg) : 0; @@ -987,6 +1022,16 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) { return ray_error(s, "%s", s); } + case Q_ID: { + Q_NEED(1); + uint8_t primitive = **buf; + *buf += 1; + *len -= 1; + if (primitive == 0) + return RAY_NULL_OBJ; + return ray_error("q: unsupported q primitive", NULL); + } + default: return ray_error("q: unsupported wire type", NULL); } @@ -1187,6 +1232,10 @@ int q_exchange(int fd, const uint8_t *req, int64_t req_len, uint8_t **resp, q_set_err(err, errlen, "q: big-endian peer not supported"); return -1; } + if (header.msgtype != Q_MSG_RESPONSE) { + q_set_err(err, errlen, "q: expected response message type"); + return -1; + } int64_t body_len = (int64_t)header.size - (int64_t)sizeof header; if (body_len <= 0 || body_len > Q_MAX_BODY) { q_set_err(err, errlen, @@ -1233,7 +1282,7 @@ ray_t *q_decode(uint8_t *resp, int64_t resp_len, int compressed, char *err, if (result == NULL) q_set_err(err, errlen, "q: deserialization returned null"); else if (remaining != 0) { - ray_release(result); + q_release_any(result); q_set_err(err, errlen, "q: trailing bytes after object"); return NULL; } diff --git a/q_server.c b/q_server.c index 087d37d..4ea4a8b 100644 --- a/q_server.c +++ b/q_server.c @@ -64,6 +64,15 @@ typedef struct { int hs_len; } q_conn_t; +static void q_release_any(ray_t *obj) { + if (obj == NULL) + return; + if (RAY_IS_ERR(obj)) + ray_error_free(obj); + else + ray_release(obj); +} + /* Turn a decoded request into a result. A char-vector (RAY_STR) is evaluated as * Rayfall. */ static ray_t *eval_request(ray_t *req) { @@ -96,7 +105,7 @@ static void q_send_result(ray_sock_t fd, ray_t *result) { if (q_encode(result, &buf, &len, err, sizeof err) < 0) { ray_t *e = ray_error(err[0] ? err : "q server: encode failed", NULL); int rc = q_encode(e, &buf, &len, err, sizeof err); - ray_release(e); + q_release_any(e); if (rc < 0) return; } @@ -218,7 +227,7 @@ static ray_t *q_read_body(ray_poll_t *poll, ray_selector_t *sel) { : ray_error("q server: malformed request", "%s", err[0] ? err : "q server: decode failed"); if (req) - ray_release(req); + q_release_any(req); if (hdr.msgtype != 0) { /* sync expects a response, async does not */ ray_selector_t *cur = ray_poll_get(poll, id); /* eval may have closed it */ @@ -226,7 +235,7 @@ static ray_t *q_read_body(ray_poll_t *poll, ray_selector_t *sel) { q_send_result((ray_sock_t)cur->fd, result); } if (result) - ray_release(result); + q_release_any(result); return NULL; } diff --git a/test/driver.c b/test/driver.c index 8a68ec0..e2e4913 100644 --- a/test/driver.c +++ b/test/driver.c @@ -25,16 +25,27 @@ #include "core/poll.h" /* ray_poll_create / run / destroy */ #include "core/runtime.h" /* ray_runtime_set_poll */ -#include "q.h" /* q_decode / q_connect */ +#include "q.h" /* q_decode / q_connect / q_exchange */ #include "q_server.h" /* q_serve */ +#include #include #include #include +#include +#include /* Registers `.q.connect` / `.q.send` / `.q.close` */ void q_env_register(void); +typedef struct { + uint8_t endianness; + uint8_t msgtype; + uint8_t compressed; + uint8_t reserved; + uint32_t size; +} test_q_header_t; + /* --serve PORT: build a runtime + poll, start the Q server, run the loop. */ static int run_server(int port) { ray_runtime_t *rt = ray_runtime_create(0, NULL); @@ -312,6 +323,16 @@ static int run_codec_selftest(void) { } release_any(r); + err[0] = '\0'; + uint8_t qidentity[] = {101, 0}; + r = q_decode(qidentity, (int64_t)sizeof qidentity, 0, err, sizeof err); + if (r != RAY_NULL_OBJ) { + fprintf(stderr, "codec selftest: Q identity did not decode as null: %s\n", + err); + failures++; + } + release_any(r); + if (q_connect("127.0.0.1", 70000, "", "", 1) != Q_ERR_SOCKET) { fprintf(stderr, "codec selftest: client accepted out-of-range port\n"); failures++; @@ -334,10 +355,77 @@ static int run_codec_selftest(void) { return failures ? 1 : 0; } +static int run_exchange_selftest(void) { + int sv[2]; + if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) { + perror("exchange selftest: socketpair"); + return 1; + } + + test_q_header_t h = { + .endianness = 1, + .msgtype = 1, + .compressed = 0, + .reserved = 0, + .size = (uint32_t)(sizeof(test_q_header_t) + 1), + }; + uint8_t body = 101; /* identity */ + int failures = 0; + if (send(sv[1], &h, sizeof h, 0) != (ssize_t)sizeof h || + send(sv[1], &body, sizeof body, 0) != (ssize_t)sizeof body) { + perror("exchange selftest: send"); + failures++; + } + + uint8_t req = 0; + uint8_t *resp = NULL; + int64_t resp_len = 0; + int compressed = 0; + char err[128] = {0}; + int rc = q_exchange(sv[0], &req, 1, &resp, &resp_len, &compressed, err, + sizeof err); + if (rc == 0 || strstr(err, "response message type") == NULL) { + fprintf(stderr, + "exchange selftest: non-response frame was not rejected: %s\n", + err); + failures++; + } + free(resp); + close(sv[0]); + close(sv[1]); + + if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) { + perror("exchange selftest: socketpair closed-peer"); + failures++; + } else { + close(sv[1]); + resp = NULL; + resp_len = 0; + compressed = 0; + err[0] = '\0'; + rc = q_exchange(sv[0], &req, 1, &resp, &resp_len, &compressed, err, + sizeof err); + if (rc == 0 || strstr(err, "send") == NULL) { + fprintf(stderr, + "exchange selftest: closed peer did not fail cleanly: %s\n", + err); + failures++; + } + free(resp); + close(sv[0]); + } + + printf("exchange selftest: %s\n", failures ? "FAIL" : "ok"); + return failures ? 1 : 0; +} + int main(int argc, char **argv) { if (argc >= 2 && strcmp(argv[1], "--codec-selftest") == 0) return run_codec_selftest(); + if (argc >= 2 && strcmp(argv[1], "--exchange-selftest") == 0) + return run_exchange_selftest(); + /* Server role: `driver --serve PORT`. */ if (argc >= 3 && strcmp(argv[1], "--serve") == 0) return run_server(atoi(argv[2])); diff --git a/test/rfl/client/08_nulls.rfl b/test/rfl/client/08_nulls.rfl index 26a1fc5..fb7f52d 100644 --- a/test/rfl/client/08_nulls.rfl +++ b/test/rfl/client/08_nulls.rfl @@ -9,6 +9,7 @@ (nil? (.q.send h "0Nh")) -- true (nil? (.q.send h "0Ni")) -- true (nil? (.q.send h "0n")) -- true +(nil? (.q.send h "::")) -- true ;; empty typed vectors (count (.q.send h "`long$()")) -- 0 diff --git a/test/rfl/server/types.rfl b/test/rfl/server/types.rfl index 06e0ce5..c3683f6 100644 --- a/test/rfl/server/types.rfl +++ b/test/rfl/server/types.rfl @@ -41,10 +41,12 @@ ;; table (sym-vec names + list of column vectors, on the wire as XT) (.q.send h "(table [a b] (list [1 2 3] ['x 'y 'z]))") -- (table [a b] (list [1 2 3] ['x 'y 'z])) -;; typed nulls (i64 / i32 / f64 / symbol). Note: rayforce dicts have no wire -;; encode path (q.c only decodes q dicts), so dict-result payloads are not -;; round-tripped here. Rayfall has no `0Ns` literal (it parses as a name), so -;; the null symbol is written as the empty symbol `(as 'sym "")`. +;; dict (native RAY_DICT encoded as XD) +(.q.send h "(dict ['a 'b] [1 2])") -- (dict ['a 'b] [1 2]) + +;; typed nulls (i64 / i32 / f64 / symbol). Rayfall has no `0Ns` literal (it +;; parses as a name), so the null symbol is written as the empty symbol +;; `(as 'sym "")`. (.q.send h "0Nl") -- 0Nl (.q.send h "0Ni") -- 0Ni (.q.send h "0Nf") -- 0Nf diff --git a/test/run.sh b/test/run.sh index 80dd17e..9dfbca8 100755 --- a/test/run.sh +++ b/test/run.sh @@ -88,6 +88,9 @@ trap cleanup EXIT echo "running codec selftest..." "$DRIVER" --codec-selftest +echo "running exchange selftest..." +"$DRIVER" --exchange-selftest + # ---- Leg 1: Rayforce server SERVERPORT="${SERVERPORT:-$(free_port)}" "$DRIVER" --serve "$SERVERPORT" & @@ -159,6 +162,8 @@ T["guid atom"; ("G"$"d49f18a4-1969-49e8-9b8a-6bb9a4832eea")~h"(as 'guid \"d49f T["guid vec"; (enlist "G"$"d49f18a4-1969-49e8-9b8a-6bb9a4832eea")~h"(as 'GUID (list \"d49f18a4-1969-49e8-9b8a-6bb9a4832eea\"))"] / table T["table"; ([] a:1 2 3; b:`x`y`z)~h"(table [a b] (list [1 2 3] ['x 'y 'z]))"] +/ dict +T["dict"; (`a`b!1 2)~h"(dict ['a 'b] [1 2])"] / typed nulls T["null i64"; 0N~h"0Nl"] T["null i32"; 0Ni~h"0Ni"]