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
13 changes: 13 additions & 0 deletions docs/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,19 @@

All notable changes to `rayforce-q` are documented here. The format follows [Keep a Changelog](https://keepachangelog.com/), and the project adheres to [Semantic Versioning](https://semver.org/). Bindings pin a tag, so each release is a stable point they can build against.

## [Unreleased]

### Added

- **Requests land in the core's query log.** With `.sys.querylog.enable` on, every request
the server evaluates is one row of `.sys.querylog`, the ring the native IPC server already
feeds: finish time, duration with decode included, status and the source. A string request
is logged verbatim; a list request, the shape a q publisher pushes, is logged as its head
and the frame's size (`(upd ...) 132 B`), because formatting a pushed table back to source
would cost more than the push. A RESPONSE frame is data for a parked `q_conn_send` and is
not logged. An operator can now read what each push costs the poll thread, which is what a
query queued behind it waits for.

## [2.1.1]

### Fixed
Expand Down
52 changes: 52 additions & 0 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/qlog.h" /* ray_qlog_begin/end: the .sys.querylog ring */
#include "lang/eval.h" /* ray_eval, RAY_EVAL_LITERAL_FALLBACK */
#include "ops/ops.h" /* ray_is_lazy, ray_lazy_materialize */

Expand Down Expand Up @@ -160,6 +161,42 @@ static void q_send_result(ray_sock_t fd, ray_t *result) {
free(buf);
}

/* The query text a log row carries. A string request is the text itself.
* Anything else is summarised, not formatted: a push is a list whose head
* names the handler and whose payload can be a hundred megabytes of table,
* and formatting that back to source would cost more than evaluating it.
* The head and the frame's size say what it was. Returns the length written
* into `out`, at most `cap`. */
static size_t q_log_source(ray_t *req, int64_t body, char *out, size_t cap) {
int n;
if (req == NULL || RAY_IS_ERR(req)) {
n = snprintf(out, cap, "malformed request, %lld B", (long long)body);
} else if (req->type == -RAY_STR) {
size_t len = ray_str_len(req);
if (len > cap)
len = cap;
memcpy(out, ray_str_ptr(req), len);
return len;
} else {
ray_t *name = NULL;
if (req->type == RAY_LIST && ray_len(req) > 0) {
ray_t *head = ((ray_t **)ray_data(req))[0];
if (head && head->type == -RAY_SYM)
name = ray_sym_str(head->i64);
}
if (name && !RAY_IS_ERR(name))
n = snprintf(out, cap, "(%.*s ...) %lld B", (int)ray_str_len(name),
ray_str_ptr(name), (long long)body);
else
n = snprintf(out, cap, "(...) %lld B", (long long)body);
if (name)
q_release_any(name);
}
if (n < 0)
return 0;
return (size_t)n < cap ? (size_t)n : cap;
}

static int64_t q_recv_fn(int64_t fd, uint8_t *buf, int64_t len) {
return ray_sock_recv((ray_sock_t)fd, buf, (size_t)len);
}
Expand Down Expand Up @@ -263,6 +300,15 @@ static ray_t *q_read_body(ray_poll_t *poll, ray_selector_t *sel) {
q_header_t hdr = cd->hdr;
int64_t id = sel->id;
char err[128] = {0};

/* The core's query log (`.sys.querylog`): one row per request this server
* evaluates, decode included, so what the q wire costs the poll thread is
* readable beside the native IPC's rows. A RESPONSE is data for a parked
* q_conn_send, not a query, and stays out. A no-op while the log is off. */
ray_qlog_ctx_t qc = {0};
if (hdr.msgtype != Q_MSG_RESPONSE)
ray_qlog_begin(&qc);

ray_t *req =
q_decode(sel->rx.buf->data, body, hdr.compressed, err, sizeof err);

Expand All @@ -285,11 +331,17 @@ static ray_t *q_read_body(ray_poll_t *poll, ray_selector_t *sel) {
return NULL;
}

char qsrc[RAY_QLOG_QUERY_MAX];
size_t qsrc_len = 0;
if (qc.measure.active)
qsrc_len = q_log_source(req, body, qsrc, sizeof qsrc);

ray_t *result = req ? eval_request(req)
: ray_error("q server: malformed request", "%s",
err[0] ? err : "q server: decode failed");
if (req)
q_release_any(req);
ray_qlog_end(&qc, qsrc, qsrc_len, result);

if (hdr.msgtype !=
Q_MSG_ASYNC) { /* sync expects a response, async does not */
Expand Down
4 changes: 3 additions & 1 deletion q_server.h
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,9 @@
* The listener is non-blocking and runs entirely on the rayforce poll you
* hand it — the same event loop the REPL and native IPC use — so it never
* spawns a thread. Every request is evaluated on the poll thread, serialized
* with the rest of the runtime.
* with the rest of the runtime. With the core's query log on
* (`.sys.querylog.enable`), every request evaluated here lands in it beside
* the native IPC's rows, decode included, with its duration and status.
*/

#include <rayforce.h>
Expand Down
33 changes: 33 additions & 0 deletions test/rfl/server/querylog.rfl
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
;; Every request the server evaluates lands in the core's query log while the
;; log is on (`.sys.querylog.enable`), the ring the native IPC server feeds, so
;; an operator reads what the q wire costs the poll thread: duration, status
;; and the source. A string request is logged verbatim; a list request, the
;; shape a q publisher pushes, is logged as its head and the frame's size,
;; because formatting a pushed table back to source would cost more than the
;; push. The `enable` call itself is not a row: the log was off when it began.

(set h (.q.connect qhost qport))
(>= h 0) -- true

(.q.send h "(.sys.querylog.enable 1)") -- 1
(.q.send h "(+ 1 2)") -- 3
(.q.send h "(+ 1 2") !- parse

;; Two rows so far: the sum and the parse error. The count runs inside a
;; third request, whose row lands after it answers.
(.q.send h "(count (.sys.querylog))") -- 2
(.q.send h "(at (get (.sys.querylog) 'query) 0)") -- "(+ 1 2)"
(.q.send h "(at (get (.sys.querylog) 'status) 0)") -- 'ok
(.q.send h "(at (get (.sys.querylog) 'query) 1)") -- "(+ 1 2"
(.q.send h "(at (get (.sys.querylog) 'status) 1)") -- 'parse
(.q.send h "(>= (at (get (.sys.querylog) 'duration-ms) 0) 0.0)") -- true

;; A list request calls its head, as a publisher's (`upd;packet) does, and is
;; logged by that head and its size on the wire. `set` answers with the lambda,
;; which the q codec cannot encode, so the definition answers 0 instead.
(.q.send h "(do (set upd (fn [x] (* 2 x))) 0)") -- 0
(.q.send h (list 'upd 21)) -- 42
(.q.send h "(like (last (get (.sys.querylog) 'query)) \"(upd ...) * B\")") -- true

(.q.send h "(.sys.querylog.enable 0)") -- 0
(.q.close h)