diff --git a/embed/rayforce_q.c b/embed/rayforce_q.c index ec886da..6e1b326 100644 --- a/embed/rayforce_q.c +++ b/embed/rayforce_q.c @@ -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; diff --git a/q_server.c b/q_server.c index 4e19cbc..f5be715 100644 --- a/q_server.c +++ b/q_server.c @@ -42,6 +42,7 @@ #include "q.h" /* pulls in : 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 */ @@ -50,6 +51,8 @@ #include #include #include +#include +#include /* Wire header — must match q.c byte-for-byte */ typedef struct { @@ -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) { @@ -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; @@ -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 @@ -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"); diff --git a/q_server.h b/q_server.h index d81c094..c9c461a 100644 --- a/q_server.h +++ b/q_server.h @@ -66,6 +66,8 @@ 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); @@ -73,7 +75,10 @@ int64_t q_conn_attach(ray_poll_t *poll, int fd); * 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. */ diff --git a/test/driver.c b/test/driver.c index 8bdf204..0023ced 100644 --- a/test/driver.c +++ b/test/driver.c @@ -33,6 +33,7 @@ #include #include #include +#include #include /* Registers `.q.connect` / `.q.send` / `.q.close` */ @@ -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; } diff --git a/test/rfl/client/06_errors.rfl b/test/rfl/client/06_errors.rfl index 903a32d..eb762f1 100644 --- a/test/rfl/client/06_errors.rfl +++ b/test/rfl/client/06_errors.rfl @@ -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)