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
2 changes: 0 additions & 2 deletions q.c
Original file line number Diff line number Diff line change
Expand Up @@ -938,7 +938,6 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) {
case -Q_KJ:
return q_des_atom_i(buf, len, RAY_I64, 8);
case -Q_KP:
case -Q_KN:
return q_des_atom_i(buf, len, RAY_TIMESTAMP, 8);
case -Q_KD:
return q_des_atom_i(buf, len, RAY_DATE, 4);
Expand Down Expand Up @@ -1019,7 +1018,6 @@ static ray_t *q_des_obj(uint8_t **buf, int64_t *len) {
case Q_KJ:
return q_des_vec_i(buf, len, RAY_I64, 8);
case Q_KP:
case Q_KN:
return q_des_vec_i(buf, len, RAY_TIMESTAMP, 8);
case Q_KD:
return q_des_vec_i(buf, len, RAY_DATE, 4);
Expand Down
64 changes: 59 additions & 5 deletions q_server.c
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ typedef struct {
#define Q_MSG_RESPONSE 2
#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_RX_CHUNK_SIZE ((int64_t)64 << 10)
#define Q_MAX_HANDSHAKE 512 /* credential blob upper bound */

/* Per-connection rx state. `hdr` carries the current frame's header between
Expand All @@ -85,6 +86,10 @@ typedef struct {
uint8_t sync_ready; /* its RESPONSE has been deposited below */
ray_t *sync_resp;
int sync_timeout_ms; /* q_conn_send round-trip budget; <= 0 blocks */
uint8_t *rx_body;
int64_t rx_body_len;
int64_t rx_body_cap;
int64_t rx_body_total;
} q_conn_t;

static void q_release_any(ray_t *obj) {
Expand Down Expand Up @@ -246,30 +251,78 @@ static ray_t *q_read_header(ray_poll_t *poll, ray_selector_t *sel) {
ray_poll_deregister(poll, sel->id);
return NULL;
}
free(cd->rx_body);
cd->rx_body = NULL;
cd->rx_body_len = 0;
cd->rx_body_cap = 0;
cd->rx_body_total = body;
if (body == 0) {
if (cd->hdr.msgtype != 0) /* sync -> identity reply */
q_send_result((ray_sock_t)sel->fd, NULL);
ray_poll_rx_request(poll, sel, (int64_t)sizeof(q_header_t));
return NULL;
}
sel->rx.read_fn = q_read_body;
ray_poll_rx_request(poll, sel, body);
ray_poll_rx_request(poll, sel,
body < Q_RX_CHUNK_SIZE ? body : Q_RX_CHUNK_SIZE);
return NULL;
}

static ray_t *q_read_body(ray_poll_t *poll, ray_selector_t *sel) {
q_conn_t *cd = (q_conn_t *)sel->data;
int64_t body = (int64_t)cd->hdr.size - (int64_t)sizeof(q_header_t);
if (!sel->rx.buf || sel->rx.buf->offset < body)
if (!sel->rx.buf || sel->rx.buf->offset <= 0)
return NULL;

int64_t received = sel->rx.buf->offset;
int64_t needed = cd->rx_body_len + received;
if (needed < cd->rx_body_len || needed > cd->rx_body_total) {
ray_poll_deregister(poll, sel->id);
return NULL;
}
if (needed > cd->rx_body_cap) {
int64_t cap = cd->rx_body_cap
? cd->rx_body_cap
: (cd->rx_body_total < Q_RX_CHUNK_SIZE
? cd->rx_body_total
: Q_RX_CHUNK_SIZE);
while (cap < needed) {
if (cap > cd->rx_body_total / 2) {
cap = cd->rx_body_total;
break;
}
cap *= 2;
}
uint8_t *body = (uint8_t *)realloc(cd->rx_body, (size_t)cap);
if (body == NULL) {
ray_poll_deregister(poll, sel->id);
return NULL;
}
cd->rx_body = body;
cd->rx_body_cap = cap;
}
memcpy(cd->rx_body + cd->rx_body_len, sel->rx.buf->data,
(size_t)received);
cd->rx_body_len = needed;

if (cd->rx_body_len < cd->rx_body_total) {
int64_t left = cd->rx_body_total - cd->rx_body_len;
ray_poll_rx_request(poll, sel,
left < Q_RX_CHUNK_SIZE ? left : Q_RX_CHUNK_SIZE);
return NULL;
}

/* q_decode fully materializes the request into ray_t objects, so the rx
* buffer is free to reuse the moment it returns. */
q_header_t hdr = cd->hdr;
int64_t id = sel->id;
char err[128] = {0};
ray_t *req =
q_decode(sel->rx.buf->data, body, hdr.compressed, err, sizeof err);
ray_t *req = q_decode(cd->rx_body, cd->rx_body_len, hdr.compressed, err,
sizeof err);
free(cd->rx_body);
cd->rx_body = NULL;
cd->rx_body_len = 0;
cd->rx_body_cap = 0;
cd->rx_body_total = 0;

sel->rx.read_fn = q_read_header;
ray_poll_rx_request(poll, sel, (int64_t)sizeof(q_header_t));
Expand Down Expand Up @@ -326,6 +379,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);
free(cd->rx_body);
free(cd);
sel->data = NULL;
}
Expand Down
29 changes: 28 additions & 1 deletion test/driver.c
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,9 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/socket.h>
#include <arpa/inet.h>
#include <netinet/in.h>
#include <sys/socket.h>
#include <sys/wait.h>
#include <sys/time.h>
#include <unistd.h>
Expand Down Expand Up @@ -450,6 +451,32 @@ static int run_codec_selftest(void) {
release_any(r);
}

/* Q timespans (KN) are durations, not timestamps. Rayforce has no
* duration type, so reject them rather than silently changing meaning. */
const struct {
const char *name;
uint8_t body[6];
int64_t body_len;
} unsupported_temporals[] = {
{"timespan atom", {0xF0, 0, 0, 0, 0, 0}, 1},
{"timespan vector", {0x10, 0, 0, 0, 0, 0}, 6},
};
for (size_t i = 0;
i < sizeof unsupported_temporals / sizeof unsupported_temporals[0];
i++) {
err[0] = '\0';
r = q_decode((uint8_t *)unsupported_temporals[i].body,
unsupported_temporals[i].body_len, 0, err, sizeof err);
if (r == NULL || !RAY_IS_ERR(r) ||
strcmp(ray_err_code(r), "q: unsu") != 0) {
fprintf(stderr, "codec selftest: %s was not rejected: %s\n",
unsupported_temporals[i].name,
r && RAY_IS_ERR(r) ? ray_err_code(r) : err);
failures++;
}
release_any(r);
}

int listener = socket(AF_INET, SOCK_STREAM, 0);
if (listener < 0) {
perror("codec selftest: handshake socket");
Expand Down
73 changes: 73 additions & 0 deletions test/run.sh
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,79 @@ with socket.create_connection((host, port), 1.0) as s:
raise SystemExit("unknown message type was executed or answered")
PY

echo "checking multi-chunk Q frame decoding..."
python3 - "$HOST" "$SERVERPORT" <<'PY'
import socket
import struct
import sys

host, port = sys.argv[1], int(sys.argv[2])
count = 40_000
values = list(range(count))
request = bytes([6, 0]) + struct.pack("<I", count) + struct.pack(
f"<{count}i", *values
)

def recv_exact(sock, size):
chunks = bytearray()
while len(chunks) < size:
part = sock.recv(size - len(chunks))
if not part:
raise SystemExit("multi-chunk test: server closed early")
chunks.extend(part)
return bytes(chunks)

with socket.create_connection((host, port), 1.0) as s:
s.sendall(bytes([3, 0]))
if recv_exact(s, 1) != bytes([3]):
raise SystemExit("bad Q handshake response")
s.sendall(struct.pack("<BBBBI", 1, 1, 0, 0, len(request) + 8))
s.sendall(request)
header = recv_exact(s, 8)
if header[:4] != bytes([1, 2, 0, 0]):
raise SystemExit(f"multi-chunk test: unexpected response header {header.hex()}")
body = recv_exact(s, struct.unpack_from("<I", header, 4)[0] - 8)
if body != request:
raise SystemExit("multi-chunk test: vector changed during round-trip")
PY

if [[ -r "/proc/${PIDS[0]}/status" ]]; then
echo "checking partial large-frame memory use..."
python3 - "$HOST" "$SERVERPORT" "${PIDS[0]}" <<'PY'
import socket
import struct
import sys
import time

host, port, pid = sys.argv[1], int(sys.argv[2]), sys.argv[3]
status_path = f"/proc/{pid}/status"

def vm_size_kib():
with open(status_path, encoding="ascii") as f:
for line in f:
if line.startswith("VmSize:"):
return int(line.split()[1])
raise RuntimeError("VmSize missing from server status")

before = vm_size_kib()
with socket.create_connection((host, port), 1.0) as s:
s.sendall(bytes([3, 0]))
if s.recv(1) != bytes([3]):
raise SystemExit("bad Q handshake response")
# Advertise a 128 MiB body, then hold the connection without sending it.
total_size = 128 * 1024 * 1024
s.sendall(struct.pack("<BBBBI", 1, 1, 0, 0, total_size))
time.sleep(0.25)
delta_kib = vm_size_kib() - before
if delta_kib > 32 * 1024:
raise SystemExit(
f"partial-frame test: server reserved {delta_kib} KiB before body arrived"
)

print(f"partial-frame memory check ok ({delta_kib} KiB growth)")
PY
fi

# ---- Leg 2: real-q interop against Rayforce server
find_q
if [[ -n "$QBIN" ]]; then
Expand Down
Loading