Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 55 additions & 6 deletions q.c
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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.) */
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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;
}
Expand Down
15 changes: 12 additions & 3 deletions q_server.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -218,15 +227,15 @@ 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 */
if (cur)
q_send_result((ray_sock_t)cur->fd, result);
}
if (result)
ray_release(result);
q_release_any(result);
return NULL;
}

Expand Down
90 changes: 89 additions & 1 deletion test/driver.c
Original file line number Diff line number Diff line change
Expand Up @@ -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 <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/socket.h>
#include <unistd.h>

/* 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);
Expand Down Expand Up @@ -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++;
Expand All @@ -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]));
Expand Down
1 change: 1 addition & 0 deletions test/rfl/client/08_nulls.rfl
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
10 changes: 6 additions & 4 deletions test/rfl/server/types.rfl
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions test/run.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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" &
Expand Down Expand Up @@ -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"]
Expand Down