Skip to content
Merged
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
7 changes: 5 additions & 2 deletions embed/rayforce_q.c
Original file line number Diff line number Diff line change
Expand Up @@ -164,8 +164,11 @@ static ray_t *qb_send(ray_t *handle, ray_t *msg) {
char err[128] = {0};
ray_t *res = q_send((int)fd, msg, err, sizeof err);
if (res == NULL) {
/* Distinguish a closed/invalid fd from a transport failure. */
const char *code = strstr(err, "handle") ? "handle" : "send";
/* Distinguish a closed/invalid fd and an expired .q.connect timeout
* from a transport failure. */
const char *code = strstr(err, "handle") ? "handle"
: strstr(err, "timed out") ? "timeout"
: "send";
return ray_error(code, "%s", err[0] ? err : ".q.send: send failed");
}
return res;
Expand Down
42 changes: 39 additions & 3 deletions q_server.c
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
#include "q.h" /* pulls in <rayforce.h>: ray_eval_str, q_encode/q_decode */

#include "core/sock.h" /* ray_sock_listen/accept/recv/send/close */
#include "core/timer.h" /* ray_time_now_ms */
#include "lang/eval.h" /* ray_eval, RAY_EVAL_LITERAL_FALLBACK */
#include "ops/ops.h" /* ray_is_lazy, ray_lazy_materialize */

Expand All @@ -50,6 +51,8 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/socket.h>
#include <sys/time.h>

/* Wire header — must match q.c byte-for-byte */
typedef struct {
Expand Down Expand Up @@ -80,6 +83,7 @@ typedef struct {
uint8_t sync_waiting; /* a q_conn_send is parked on this connection */
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 */
} q_conn_t;

static void q_release_any(ray_t *obj) {
Expand Down Expand Up @@ -390,14 +394,28 @@ static int q_conn_pump(ray_poll_t *poll, int64_t id) {
}
}

/* The per-operation timeout q_connect applied to the socket (SO_RCVTIMEO).
* Once the fd is non-blocking under a poll, the kernel no longer enforces
* it, so q_conn_send has to — read it here before the switch. 0 = none. */
static int q_sock_recv_timeout_ms(int fd) {
struct timeval tv = {0, 0};
socklen_t n = sizeof tv;
if (getsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, &n) < 0)
return 0;
int64_t ms = (int64_t)tv.tv_sec * 1000 + tv.tv_usec / 1000;
return ms > INT32_MAX ? INT32_MAX : (int)ms;
}

int64_t q_conn_attach(ray_poll_t *poll, int fd) {
if (poll == NULL || fd < 0)
return -1;
int timeout_ms = q_sock_recv_timeout_ms(fd);
ray_sock_set_nonblocking((ray_sock_t)fd);

q_conn_t *cd = (q_conn_t *)calloc(1, sizeof *cd);
if (cd == NULL)
return -1;
cd->sync_timeout_ms = timeout_ms;

ray_poll_reg_t reg = {0};
reg.fd = (int64_t)fd;
Expand Down Expand Up @@ -442,6 +460,11 @@ ray_t *q_conn_send(ray_poll_t *poll, int64_t id, ray_t *msg) {
cd->sync_ready = 0;
cd->sync_resp = NULL;

/* 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
Expand All @@ -461,9 +484,22 @@ ray_t *q_conn_send(ray_poll_t *poll, int64_t id, ray_t *msg) {
}
if (rc < 0)
return ray_error("io", "q: connection closed");
int w = ray_sock_wait_readable_intr((ray_sock_t)sel->fd, -1);
if (w == -2)
continue; /* interrupted by a signal — keep waiting */
int wait_ms = -1;
if (bounded) {
int64_t left = deadline_ms - ray_time_now_ms();
if (left <= 0) {
/* The request is in flight and its RESPONSE may still land; left
* open, the next sync send on this handle would claim it as its
* own reply. Tear the connection down and let the caller
* reconnect. q_on_close frees cd. */
ray_poll_deregister(poll, id);
return ray_error("timeout", "q: response timed out");
}
wait_ms = left > INT32_MAX ? INT32_MAX : (int)left;
}
int w = ray_sock_wait_readable_intr((ray_sock_t)sel->fd, wait_ms);
if (w == -2 || w == 0)
continue; /* signal, or the slice elapsed — the deadline check decides */
if (w < 0) {
cd->sync_waiting = 0;
return ray_error("io", "q: recv failed");
Expand Down
7 changes: 6 additions & 1 deletion q_server.h
Original file line number Diff line number Diff line change
Expand Up @@ -66,14 +66,19 @@ int64_t q_serve(ray_poll_t *poll, int port);
/* Put an already-connected, already-handshaken q_connect fd under `poll`'s
* rx machine. Returns a poll selector id (>= 0) — the handle for the calls
* below — or -1. The fd is owned by the poll from here on; do not q_close it.
* The socket's recv timeout (what q_connect's timeout_ms set) is read here
* and becomes q_conn_send's round-trip budget; none = wait indefinitely.
*/
int64_t q_conn_attach(ray_poll_t *poll, int fd);

/* Sync round-trip on an attached connection: write a SYNC frame, then pump
* this connection until its RESPONSE arrives — dispatching, not swallowing,
* any frame that arrives in between (a pushed async runs its handler, an
* inbound sync request is answered). Returns a freshly-owned object, which
* may be a RAY_ERROR from the peer or a local io/handle error. */
* may be a RAY_ERROR from the peer or a local io/handle error. When the
* budget inherited at attach time runs out first, the connection is closed
* (a late RESPONSE must not be taken for the next request's reply) and a
* `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);

/* Close an attached connection and release its state. */
Expand Down
48 changes: 48 additions & 0 deletions test/driver.c
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
#include <stdlib.h>
#include <string.h>
#include <sys/socket.h>
#include <sys/time.h>
#include <unistd.h>

/* Registers `.q.connect` / `.q.send` / `.q.close` */
Expand Down Expand Up @@ -415,6 +416,53 @@ static int run_exchange_selftest(void) {
close(sv[0]);
}

/* 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
* indefinitely. */
if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) {
perror("exchange selftest: socketpair silent-peer");
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: failed to create runtime/poll\n");
failures++;
} else {
struct timeval tv = {.tv_sec = 0, .tv_usec = 200 * 1000};
setsockopt(sv[0], SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof tv);
int64_t id = q_conn_attach(poll, sv[0]);
if (id < 0)
close(sv[0]);
ray_t *msg = ray_str("1+1", 3);
/* 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;
alarm(0);
ray_t *es = (r && RAY_IS_ERR(r)) ? ray_fmt(r, 0) : NULL;
const char *ep = es ? ray_str_ptr(es) : "";
if (r == NULL || !RAY_IS_ERR(r) || strstr(ep, "timeout") == NULL) {
fprintf(stderr, "exchange selftest: silent peer did not time out: %s\n",
ep);
failures++;
}
if (ray_poll_get(poll, id) != NULL) {
fprintf(stderr, "exchange selftest: timed-out handle was left open\n");
failures++;
}
if (es)
ray_release(es);
release_any(r);
ray_release(msg);
ray_poll_destroy(poll);
}
if (rt)
ray_runtime_destroy(rt);
close(sv[1]);
}

printf("exchange selftest: %s\n", failures ? "FAIL" : "ok");
return failures ? 1 : 0;
}
Expand Down
13 changes: 13 additions & 0 deletions test/rfl/client/06_errors.rfl
Original file line number Diff line number Diff line change
Expand Up @@ -20,3 +20,16 @@
(.q.send h "2+2") -- 4

(.q.close h)

;; The .q.connect timeout bounds every sync round-trip, not just the
;; connect: a peer that takes longer than the budget to answer surfaces
;; as a `timeout` error instead of a wait with no upper bound (the
;; event-loop path used to block forever). q itself is the slow peer here.
(set slow (.q.connect qhost qport "" "" 300))
(.q.send slow "system \"sleep 2\"") !- timeout
(.q.close slow)

;; a fresh connection with a budget the peer meets works as before
(set fast (.q.connect qhost qport "" "" 5000))
(.q.send fast "3+4") -- 7
(.q.close fast)
Loading