Skip to content

Port the Java HA sender to the Rust, Python, and C/C++ clients - #6

Open
javier wants to merge 80 commits into
mainfrom
jv/add_rust_based_clients
Open

javier wants to merge 80 commits into
mainfrom
jv/add_rust_based_clients

Conversation

@javier

@javier javier commented Jul 10, 2026 •

Copy link
Copy Markdown
Collaborator

This branch ports the Java CsvParallelSender — the QWP high-availability sender (failover, store-and-forward, dual ingest/query, and the shared connection options) — to the three clients built on the C/Rust core (the "CRusty" clients): Rust, Python, and C/C++.

Each port keeps the Java design (CSV replay loop, per-worker senders, timestamp / O3 semantics) and mirrors the same --protocol transports, connection options, and flags. Per-language differences are documented in each folder's README.md.

Status — all three complete

  • Rust (rust/) — QWP (WebSocket) with per-worker store-and-forward and multi-host failover, QWP/UDP, and ILP/HTTP; a Reader-based probe (latest ingested timestamp + live serving role via switch status, target=any replica-fallback reads, automatic failover). Validated end-to-end, including a primary-crash failover with a gap-free store-and-forward handoff.
  • Python (python/) — the same three transports plus a pooled query Client, and a pandas/polars ingestion + egress demo (dataframe_demo.py). Built on the client's QWP + egress branch (built from source; the PyPI release is ILP-only). Validated end-to-end, including a primary-crash failover and auto-flush.
  • C/C++ (c/) — the same three transports (C++17, line_sender.hpp / line_reader.hpp), a probe with a handshake-role fallback, and a CMake build via Corrosion. Reads gzipped CSV via zlib. Validated end-to-end, including a primary-crash failover with a gap-free store-and-forward handoff.

Each README.md documents the per-language build and the client differences (e.g. auto-flush present in Java/Python but not Rust/C++; the handshake role exposed everywhere except the Python binding).

javier added 13 commits July 2, 2026 14:34
Emit trade_id = <worker>-<1-based sequence> as a VARCHAR on each row, so completeness/gaps can be checked independent of timestamps and, with DEDUP UPSERT KEYS(timestamp, trade_id) on a single-worker client-timestamped run, store-and-forward replay after a failover is idempotent (drops the at-least-once boundary duplicates). High-cardinality, so a string column, not a symbol.
Add a Row id and dedup section covering the trade_id format, the completeness check, idempotent store-and-forward replay via DEDUP UPSERT KEYS(timestamp, trade_id), and the client-side-timestamp requirement.
Sender:
- Detect the QWP upgrade "HTTP/1.1 400 Bad request" that a missing/invalid
  ILP token produces and append an explicit auth hint wherever it surfaces
  (endpoint-failed, worker error, probe-stopped).
- Probe now reports the LIVE serving role via `switch status` each poll
  instead of the QWP handshake SERVER_INFO, which only refreshes on a
  reconnect and so goes stale after an in-place primary<->replica switch.

Watchdog:
- Monitor the primary's currentRole, not just health: react to a silent
  demotion to REPLICA as well as an unreachable node.
- Bootstrap a fresh cluster: if no node is primary, promote the first
  healthy one in list order.
- Fail over to the first HEALTHY node in preference order; adopt a node
  that is already PRIMARY instead of forcing another switch.
- Add a configurable grace period (QDB_WD_GRACE_PERIOD, default 5s) so a
  manual operator switch can settle before the watchdog acts.
- Never exit on its own; keep retrying when nothing can be promoted.
- Remove the now-unused get_http_code and STARTUP_RETRIES.

Update README and the sample env accordingly.
Build/target:
- pom default questdb.client.version -> 1.3.6-SNAPSHOT (required; QWP and the
  lifecycle switch features are not stable earlier, and it is not on Central).
- Drop the RECONNECT_BUDGET_EXHAUSTED switch case in both senders; that kind was
  removed from SenderConnectionEvent.Kind in 1.3.6, so it no longer compiles.

Diagnostics:
- Replace the brittle "400 Bad request" string-match with the client's typed
  signals: QwpAuthFailedException (401/403), WebSocketUpgradeException.isRoleMismatch()
  (421 = endpoint is a replica, no primary to accept writes), and status-400
  (missing/malformed token). The listener passes the real throwable, not its message.

Timeouts:
- --retry-timeout now also drives QWP (reconnectMaxDurationMillis), not just ILP's
  retryTimeoutMillis, so one flag means the same "keep retrying" budget on both.
- New --connect-timeout-ms (default 3000, 0 = OS default): connectTimeoutMillis on
  the senders and connect_timeout= on the probe, so a black-holed host fails over in
  seconds instead of riding the OS connect timeout. ILP left unbounded/unchanged.
- Watchdog: make the role/health poll timeouts configurable
  (QDB_WD_CHECK_CONNECT_TIMEOUT / QDB_WD_CHECK_MAX_TIME; defaults 2/3, no behaviour change).

Update README and the watchdog sample env accordingly.
The probe fallback printed only "node=(none)" and dropped the role entirely
when `switch status` returned nothing, which was worse than before. Restore the
role/node/zone from getServerInfo() on every line (it tracks connect + failover),
and append the authoritative live role from `switch status` only when present.

When switch status yields no live role, log why at most once per 30s:
an error status+message, a returned batch with no *role* column (names listed),
or no row batch at all -- so the empty case is diagnosable instead of silent.
…dshake

The probe printed getServerInfo()'s role as the headline 'served by role=',
but that is the QWP handshake role, refreshed only on connect/failover. After
an in-place promotion (the read connection never drops) it stayed stuck at
REPLICA even though the node was serving reads as PRIMARY.

Make switch status' current_role the authoritative role shown, append
'(switching -> ROLE)' while a switch is in flight, and demote getServerInfo()
to supplying node/zone only -- its handshake role is now a clearly labelled
fallback used solely when the status query is unavailable. target=any is
unchanged, so replica-fallback reads still work.
The watchdog Configuration section only named a few QDB_WD_* vars inline and
deferred the rest to the env example; enumerate all 14 (required + optional)
with defaults and purpose. Also update the probe example to the current
'served by role=<current_role> node=... zone=...' format, covering the
'(switching -> ROLE)' in-flight case and the labelled handshake fallback.
Add --protocol qwpudp for fire-and-forget UDP datagram ingest to the
UDP port (:9007). It is ingest-only, unauthenticated (auth flags are
warned about and ignored), and single-endpoint with no failover,
store-and-forward, TLS, or query client (the probe is skipped).

The worker flushes datagrams every --batch-size rows and, for a single
worker, stamps rows client-side; multiple workers use server-side
atNow(). Document the transport, including its best-effort nature and
the burst/packet-loss behaviour, with guidance to keep --batch-size
small for UDP.
Port the Java CsvParallelSender to Rust, built on the local questdb-rs
(c-questdb-client) as a path dependency. Mirrors the Java client:

- qwp (WebSocket) with per-worker store-and-forward, transactional
  commit cadence, and multi-host failover.
- qwpudp (UDP) fire-and-forget datagrams: ingest-only, unauthenticated,
  single-endpoint, with upfront single-address validation.
- ilp (HTTP) legacy transport.
- Same connection options (address list, token / basic auth, TLS,
  zone, enterprise durable ack, reconnect budget) and timestamp
  semantics (single-worker client micros vs multi-worker at_now).
- A Reader-based probe (qwp only) polling the latest ingested timestamp
  and the serving node's live role via 'switch status', with target=any
  replica-fallback reads and automatic failover.

Batching is driven by explicit flush() calls since the Rust client has
no auto-flush by design. Validated end-to-end against a local server,
including a primary-crash failover with gap-free store-and-forward
handoff.
Port the Java/Rust CsvParallelSender to Python on the questdb client's
QWP + egress build (the sm_qwp_dataframe_bench branch, built from source;
the PyPI release is ILP-only). csv_parallel_sender.py mirrors the other
ports: qwp (WebSocket) with per-worker store-and-forward and manual flush
cadence, qwpudp, and ilp, plus a probe using the pooled query Client
(select ... limit -1 + switch status role reporting, target=any,
failover). dataframe_demo.py shows pandas ingestion (Sender.dataframe),
polars ingestion (Client.dataframe), and egress to pandas/polars via
Client.query(...).to_pandas()/.to_polars().

Validated against a local server: qwp/ilp/udp ingest, a primary-crash
failover with gap-free store-and-forward handoff, and row-threshold
auto-flush. README documents the build-from-source requirement and the
Python-specific differences.
Port the Java/Rust/Python CsvParallelSender to C++ on the c-questdb-client C++ headers (line_sender.hpp ingestion, line_reader.hpp egress query client). Mirrors the other ports: qwp (WebSocket) with per-worker store-and-forward and manual flush cadence, qwpudp, and ilp, plus a probe (latest ingested timestamp + live serving role via switch status with a handshake-role fallback, target=any, failover). CMake build via Corrosion (cargo FFI), reads gzipped CSV via zlib. Validated against a local server: qwp/ilp/udp ingest and a primary-crash failover with a gap-free store-and-forward handoff.
@javier
javier marked this pull request as ready for review July 13, 2026 09:14
javier added 16 commits July 13, 2026 11:18
Add a Client implementations section explaining the Java implementation is the reference, ported in parity to the Rust, Python, and C/C++ clients, with links to each folder.
Add a --rate flag that targets an aggregate rows/second across all workers, pacing each worker to its share against a deadline schedule so it reaches high targets a per-row --delay-ms cannot. It takes precedence over --delay-ms. Ported to Java, Rust, Python, and C++.

Send designated timestamps at nanosecond resolution everywhere. QuestDB stores at the target column resolution: TIMESTAMP_NS keeps full nanos, micros TIMESTAMP truncates, and a non-existent table is auto-created as TIMESTAMP_NS. Java now uses at long NANOS, Rust uses TimestampNanos with a nanosecond CSV parse, Python gained a nanosecond-preserving ISO parse, and C++ already sent nanos.

Add regenerate_csv.sh and regenerate_csv_fx.sh to export fresh crypto and FX data from the always-on demo box. Document --rate, the nanosecond behavior, and the two scripts in each README.
Stream a table through polars, add a random enriched_rnd SYMBOL column, and
write it back to enriched_<table>_demo. Reads via QueryResult.iter_arrow() one
Arrow batch at a time (peak memory is a single batch, not the whole result),
enriches each chunk, and ingests it via Client.dataframe. Reader and writer use
separate QWP connections. Supports token/basic auth and TLS (qwpwss) for
Enterprise, applied to both the QWP path and the HTTP drop/verify calls.
csv_columnar_sender.py: high-throughput Python ingestion over the columnar QWP
path (Client.dataframe with polars), replacing row-by-row for volume. Streams the
CSV in bounded chunks (flat memory), with --rate pacing, --num-senders, and
enterprise auth/TLS. Row-by-row csv_parallel_sender.py stays for HA/failover demos
(store-and-forward, probe, UDP/ILP), which the columnar path bypasses.

read_bench.py: reads the last N rows as fast as possible via streaming iter_arrow()
and reports rows/s plus MB/s and Gb/s (decoded Arrow payload). --readers splits the
scan across N parallel connections by timestamp to get past the per-connection
socket-buffer cap.

boost_tcp.sh: raises net.core.wmem_max/rmem_max so the client's 4 MiB socket buffer
is not clamped to ~416 KB (the per-connection throughput cap on real networks). Safe
to source or exec.

README: 'which script to use' guidance (columnar vs row-by-row) and network tuning.
An empty --token (typically --token "$ILP_TOKEN" with the env var unset) silently
disabled auth+TLS, so the client connected plaintext to a secured server and hung
retrying the handshake at sent=0, uninterruptible. Now: an empty --token (or a
--username with empty --password) exits with a clear error instead of hanging.

Worker threads are now daemon and joined via a polling loop, so KeyboardInterrupt
is honored even while a worker is blocked in a native connect/flush.
Before starting workers, probe the server over the full path (TCP+TLS+auth+
'select 1') in a daemon thread with a timeout. A down/unreachable/misconfigured
server now fails loudly with a clear error instead of the ingest client retrying
forever at sent=0. Default 10s; --connect-timeout 0 skips the check.
The progress counter only tracked rows appended to the client buffer, which with
QWP store-and-forward runs far ahead of what the server has committed - so the
per-second rate looked inflated and the end-of-run showed a long '0 rows/s' tail
while the client drained the backlog.

Now the reporter prints two counters: submitted (client-side, as before) and
acknowledged (rows the server has actually committed). The acked count comes from
the QWP ack watermark: each flush records its sequence via flushAndGetSequence(),
and getAckedFsn() is mapped back to committed rows through a per-worker fsn->rows
map - no extra query round-trips. A bounded awaitAckedFsn loop drains at the end
while the acked counter climbs to the full count, so the tail shows real progress
instead of zeros. The summary throughput is now labeled (acknowledged, end-to-end)
and split into submit phase vs commit drain.

Validated against a local QWP server: acknowledged trails submitted and converges
to 100%. ILP/UDP fall back to the plain flush path.
Once all rows are submitted, silence the per-second progress and probe lines so the
commit-drain tail no longer prints repeated '+0/s' / probe lines. The single final
summary line (acknowledged, end-to-end + submit/drain split) is the last output.
Connection/failover probe events still pass through.
qwpws/qwpwss are deprecated aliases; the client maps ws->QwpWs and wss->QwpWss.
Switch the Rust and C++ ingest conf strings to ws/wss, matching their readers
(which already use them) and the Java query client. Verified: both build against
their client checkouts and ingest cleanly over ws against a local server.

Python is left on qwpws/qwpwss - its binding's Protocol enum only exposes
QwpWs/QwpWss and rejects ws/wss (ValueError: Invalid value for Protocol), even
though the underlying Rust client accepts the aliases. Docs updated accordingly.
backfill.py spreads rows evenly across a historical window ([--start,--end),
default yesterday 00:00 UTC .. today 13:00 UTC) at --rate-per-day density
(default 500M/day -> ~771M rows over the window), ingesting them with those
historical timestamps. Streams in bounded chunks (flat memory), splits the
window across --num-senders workers (each a contiguous time block), columnar
QWP path, with the same auth/TLS + connect-timeout preflight as the other
scripts. Verified at small scale: timestamps land spread across the window.

Cosmetic: the [conf] line now prints wss/ws instead of qwpwss/qwpws in
csv_columnar_sender / read_bench / enrich_polars_demo / backfill. The actual
connect string still uses qwpwss/qwpws (the Python binding requires them).
Target table is now configurable instead of hardcoded; auto-created if absent.
Echoed in the startup line. Verified ingesting into a custom table.
Add --sample N (default 5): reader 0 wraps its first received Arrow batch as a
polars DataFrame (zero-copy) and keeps its head - reusing data already streamed,
no extra query and no re-scan. Rendered after the timing, so throughput is
unaffected. Only one batch is ever held, so memory stays bounded.
blotter.py polls a table (or live view) at --rate Hz (default 5, max 20) and redraws
a polars table in place (ANSI, no flicker). Builds 'select * from TABLE [WHERE ...]
limit N' from three params: table, --where (WHERE auto-prepended if absent; trailing
clauses like ORDER BY ride along; not sanitised - demo use), and --limit (default -10,
magnitude clamped to 100). Query errors render in place instead of crashing, so the
SELECT can be tweaked live. --once renders a single plain frame for scripting.

Enterprise auth/TLS via --token/--username/--password (+ --tls-verify), plus
--token-file/--token-label to read a bearer token from a file by label (keeps the
secret off the command line). Verified against a remote core_price_lv live view.

run_blotter.sh launches it with a default cluster address and $ILP_TOKEN, passing
table/--where/--limit/--rate through; ADDR and PYTHON overridable via env.
--table replaces the positional arg. New --query runs full SQL verbatim (CTEs,
window functions, ...); when set, --table/--where/--limit are ignored and a trailing
';' is stripped. Provide one of --table/--query. Multi-line SQL is collapsed to one
line in the header; display is capped at 100 rows. Verified against a live view with
a CTE + window query. run_blotter.sh uses --table with a commented --query example.
…resh ceiling

The full-screen redraw each tick was terminal/SSH-output bound (froze past ~4 Hz over
SSH), not query bound. draw() now diffs the new frame against the last and rewrites
only the lines that changed - positioning each absolutely, writing it, and clearing to
end-of-line - so static borders/header/rows are never re-sent. Cuts terminal output
sharply and lets the refresh rate go well past the old ceiling.
A small http.server that queries QuestDB (QWP, token held server-side) and serves a
self-contained dark-themed page that polls /data and renders a live table with cells
flashing green/red on change. Rendering is in the browser, so it avoids the terminal/
SSH redraw bottleneck; only small JSON crosses the wire. Same query/auth params as
blotter.py (--table/--query/--where/--limit, --token[-file/-label], --tls-verify).
--host controls the bind address (default 127.0.0.1; --host 0.0.0.0 exposes it - the
web server has no auth, so firewall the port).
All of this came from running a slow, small-batch config (--delay-ms 1,
--batches-per-transaction 1) that the earlier high-throughput testing never covered.
That turned out to be where the reporting was weakest, and it is the configuration
closest to an actual demo.

Banner the three headline events. A write path moving to another node, reads moving
to a replica, and a spilled backlog replaying are the whole point of an HA demo, and
as ordinary single lines they scrolled past unnoticed between per-second progress
lines. Each is now wrapped in ### rules:

    ############################################################################
    ###  READS FAILED OVER -- now served by role=REPLICA
    ###  node=... zone=eu-west-1
    ###  queries keep being answered; only writes need a primary
    ############################################################################

Only RECONNECTED and FAILED_OVER are promoted. The per-attempt
"endpoint ... failed, trying next" stays an ordinary line, because during an outage
it repeats every few seconds and bannering it would bury the real transitions. Each
banner is built into one StringBuilder and emitted with a single write, so the
reporter and probe threads cannot interleave a line into the middle of it.

Name the flush threshold. "frames acked=0/pending" was opaque: a flush happens every
batch-size x batches-per-transaction rows, so with 10000 x 1 at ~944 rows/s the first
frame is 10.6s away and the field looked broken. It now reads
"none yet, first flush at 10,000 rows".

Qualify the spill figure. sfDirBytes reports ALLOCATED bytes, and segments are
memory-mapped and pre-allocated to sf_max_segment_bytes (4 MiB default), so a run
holding under a megabyte displayed "8.0 MiB". That is ~10x overstated at demo scale
and invisible at GB scale, which is why earlier runs looked fine. Now
"spill: 8.0 MiB allocated in 5 segment file(s)", from a single directory walk rather
than two that could disagree while segments rotate.

Clamp the second -1 sentinel. getAckedFsn() returns -1 for "nothing acked yet",
exactly as flushAndGetSequence() returns -1 for "nothing to flush". Guarding only the
latter left "frames acked=-1/11" on the first flush of a slow feed. Both accessors are
clamped at every assignment now, not at the point of display, and the convention is
recorded next to the fields.

Omit frames for non-QWP protocols. UDP is fire-and-forget and ILP auto-flushes over
HTTP; neither carries frame sequence numbers, so --protocol qwpudp/ilp would have
printed "first flush at 0 rows" indefinitely. Found by review, not by testing, since
the demo path is QWP.

Verified against a live cluster: the startup-replay banner, the "none yet" wording,
the spill wording, and zero sentinel leaks across --batch-size 10000 and 1000.

NOT verified: the read-failover and write-failover banners need a writable cluster to
trigger, and the cluster had no primary at the time (enterprise-primary reporting
role=REPLICA, no writable node among all three addresses). The code path is the same
listener that already produced the unbannered versions of both messages earlier, but
the rendering itself is unobserved.
A single ### rule above and below still got lost in the scroll of per-second progress
and probe lines during a live demo. Three each side, plus a blank line outside the
block so it does not butt up against the surrounding output.

Rendered from identical code to confirm the spacing:

    [progress] submitted=8,332 (+944/s) | frames acked=31/33

    ############################################################################
    ############################################################################
    ############################################################################
    ###  INGESTION failed over 172.31.42.41:9000 -> 172.31.41.35:9000 -- REPLAYING ...
    ###  sender=ha_sender-0
    ############################################################################
    ############################################################################
    ############################################################################

    [progress] submitted=9,276 (+944/s) | frames acked=31/35

Known cosmetic wrinkle: the INGESTION line runs to ~113 characters against a
76-character rule, so it wraps on an 80-column terminal and the block looks ragged
there. Left alone rather than guessed at, since the fix depends on the demo terminal
width: either widen BANNER_RULE or split the hosts and the backlog onto separate
### lines.
Port python/read_bench.py to Java (com.example.reader.ReadBench). Streams records off QWP column batches without materialising them, reports rows/s and payload throughput live, and prints the first and last rows of the range at the end.

--addr accepts a comma-separated host list for failover. Chunks re-run after a failover or retried after an error roll back their counts first, so rows are never double-counted.

Add query_table.sh to build and run it; with no arguments it runs the demo read.
Was "questdb:9000", a single host that exercises no failover. Default to the
cluster's three nodes so a bare run without --addrs already covers the HA path:

    172.31.42.41:9000,172.31.41.35:9000,10.0.0.8:9000

These are VPC-internal addresses, reachable from a sender host inside the VPC and not
from outside it.

The target table is unchanged and remains hardcoded to "trades", both in
sender.table(...) and in PROBE_QUERY. There is no --table flag to default.
bash ./query_table.sh with no arguments ran against a single host and so could never
fail over, which read as "works, but no failover".

Two defaults were single-host:
- query_table.sh handed ReadBench "${QDB_ADDR:-172.31.42.41:9000}"
- ReadBench's own --addr default was "localhost:9000"

Both now carry 172.31.42.41:9000,172.31.41.35:9000,10.0.0.8:9000, matching
CsvParallelSender.DEFAULT_ADDRS. ReadBench already passed the value straight into
addr= and the client parses a comma-separated list there, so no other change was
needed. The header comment in query_table.sh was also claiming a single-host default.

Overrides still work: QDB_ADDR=host:9000 ./query_table.sh for one node, or pass
--addr explicitly.

This should have landed with the sender's default in f520543. That change was made
three minutes after the reader was committed, so both files were present, but the
search was scoped to CsvParallelSender.java alone and would not have matched either
of them in any case: ReadBench uses --addr, not --addrs, and query_table.sh is not a
Java file. A repo-wide grep for :9000 finds all three at once.

Not exercised against the cluster: build only.
Four panels over QWP, redrawn once a second, each a rolling window so cost and memory
stay flat however long it runs. Verified against a live demo instance.

The layout is organised around the DELAY each measurement needs. Anything looking
forward in time cannot be evaluated until that window has elapsed, so each query
shifts back by exactly its own forward horizon:

    panel            window                forward   delay
    slippage         $now-1m .. $now       none      live
    markout +-1m     $now-2m .. $now-1m    +1m       1m
    markout past     $now-1m .. $now       none      live
    min/max +-10s    $now-70s .. $now-10s  +10s      10s

The two markout panels side by side are the point: the delayed one shows the complete
curve either side of the fill, the live one only what has already happened. Same
calculation, different trade between completeness and latency.

Each markout panel issues two statements, and --view selects which run, so a pivot-only
run never pays for the raw one. The raw statement is the query as authored. The pivot
is the same calculation grouped by ecn and horizon alone, ~52 rows against ~1,700, so
the displayed curve is a pure reshape of a server-side aggregate. Folding the wider
result down in polars instead would need a weighted mean, sum(avg * n) / sum(n), to
stop a small group counting the same as a large one; asking the server for the
granularity actually displayed removes that arithmetic entirely.

counterparty is not a grouping key. With 1,061 distinct counterparties and ~1,600 fills
a minute it put 1 fill in nearly every cell (measured busiest: n=2), so avg() and sum()
were aggregating a single observation.

Raw rows are a fresh random sample each tick. The authored ORDER BY is by key, so a
head() always showed the alphabetically first rows and looked frozen while the data
moved underneath.

Every displayed column is pre-formatted to a fixed width. polars sizes each column to
its widest CURRENT value, so a window without a Currenex fill, or with smaller numbers,
redrew narrower and every border shifted each second. There is no setting for this:
tbl_width_chars only truncates, it never pads. Widths come from measuring the real
maxima and adding headroom, since pinning at exactly the observed maximum lets one
larger value grow the column again, and over-padding pushes the table past the terminal
so polars truncates every cell.

Flags: --symbol (applied in SQL), --interval-ms, --view, --layout, --panels, --precision,
--raw-sort, --raw-rows, --rows, --minmax-rows, --left-width, --width, plus --addr/--conf
and the usual auth.

Not covered: the derived slippage summary is a client-side mean, correct because every
row there is exactly one fill, but it is not query output and is labelled as such.
…cript

Ctrl+C now works in any state. The polling loop moved to a daemon thread and the main
thread only ever waits in worker.join(0.2). Python can raise KeyboardInterrupt between
bytecodes but never inside the client's blocking native calls, so a loop that queries on
the main thread is unkillable for as long as one blocks, and during an outage every tick
blocks on a connect, which is exactly when you want to stop it. Tuning the sleep does not
fix this; keeping the main thread out of native code does. Exit uses os._exit so shutdown
cannot block in the same native code it is escaping.

Verified through a pty, sending a real ^C: "interrupted", exit code 130. Three earlier
harnesses gave false verdicts before that, so the final one was validated against a
trivial sleep loop first: a background-launched process inherits SIGINT set to SIG_IGN
and Python preserves it, and a harness that blocks in os.read never reports. Also
confirmed the client leaves the SIGINT handler alone (still default_int_handler after
connect), so nothing else was interfering.

The reader lease is rebuilt after any failure rather than held for the whole run. Reader
failover is on by default and to_polars() replays a mid-query failover transparently, so
a node dying mid-result is invisible. What is not automatic is the lease: a terminal
rejection is latched by the connection and raised by the next call on it, so a lease held
for the process lifetime can be poisoned permanently. Verified by killing and restarting a
TCP proxy under a running dashboard: it reported the reason each tick and resumed on its
own. Note that transparent replay depends on to_polars(); the iter_* consumers raise
FailoverWouldDuplicate instead, so switching a panel to streaming would need that handled.

sender_pool_min=0, because this tool only reads and connect() otherwise pre-opens a
store-and-forward sender it never uses. The reader pool is still opened eagerly, so an
unreachable server at startup fails fast rather than launching a dashboard that can show
nothing.

README: nine scripts were undocumented, not just this one. Added rows for tca_live,
blotter, web_blotter, backfill and the four slippage variants, plus a section for the
dashboard covering the window/delay table, the two-statements-per-panel design, the flags
and the terminal-width caveat. run_blotter.sh and start_sending.sh are still uncovered;
the table is scoped to Python entry points.
A Node server issues the same four TCA queries as tca_live.py over QWP and
serves the results to a page that renders them as a grid of panels. The
browser cannot open the QWP socket itself yet: @questdb/browser-client is
unpublished, and the WebSocket upgrade requires Origin == Host, which a
separate dev-server port cannot satisfy. data.mjs is therefore the only
module that knows QuestDB exists, so swapping in the browser client later
touches one file.

- symbol dropdown from a LATEST ON query, with ALL as the default, filtering
  every panel in SQL rather than client side
- per-panel SQL view, fetched from the same sqlFor() the query path uses so
  what is shown cannot drift from what ran
- markout P&L curves as inline SVG, one line per ECN, with a crosshair and a
  tooltip ranking venues at the hovered horizon; the hovered horizon is held
  outside the chart because each tick rebuilds it
- refresh dropdown, 250ms to 2000ms plus Off, which holds the last frame at
  full opacity instead of dimming it
The shutdown handler awaited data.close() before process.exit(). Installing a
SIGINT handler replaces node's default terminate behaviour, so a QWP close that
stalls left a process that ignored Ctrl+C outright and had to be killed by PID.

The graceful path is now raced against a 2s deadline, and a second signal exits
at once. Keep-alive sockets are closed explicitly: an open dashboard tab holds
the event loop well past its last request, so server.close() alone would not
have been enough.

Measured with a browser polling /api/tick once a second: SIGINT to exit 0 in
0.52s, without reaching the deadline.
Adds a second tab that streams a scan of the last N rows of a table and reports
only how far it has got, so neither the server nor the page holds the result.
200,000,000 rows of core_price complete in 242s at 826k rows/s against the
remote cluster, with the server at 147MB RSS: memory tracks reader count, not
rows read.

Reading is queryViews(), which hands back a reusable batch view and never
materialises rows. The whole per-batch hot path is two counters and a row count
the batch already carries; the ten rows the page displays are copied out only
when a progress frame is due, ~10 times a second.

Three limits had to be worked around to make a scan this size finish at all:

- QuestDB applies query.timeout (60s) server-side, independently of the client
  deadline. The scan is therefore split into LIMIT -m, -n slices that tile the
  range with no gap or overlap, as python/read_bench.py does.
- Slice count is derived from a bounded rows-per-query, not from a fixed number
  of slices. 64 slices of 200M rows is 3.1M rows each, which times out and
  completes nothing, while the same 64 slices of 4M rows finish comfortably.
- query_pool_max defaults to 4, so more than four parallel readers failed with
  "timed out waiting for a QWP query from the pool".

Stopping a scan mid-flight and relaunching it hit the same pool limit, because
cancelling a query is not instant: Promise.all returned while sibling readers
were still draining, and the server did not wait for the previous scan to
release. Readers are now awaited with allSettled, a new scan waits for the
previous one and says so, and the single-flight slot is claimed synchronously.

Connection details are configuration rather than constants: --addrs takes the
comma-separated list the client resolves with ordered failover, --token-file
keeps the bearer token out of argv and out of the startup log, and tls_verify
defaults to unsafe_off for the self-signed cluster certificates with a warning
at startup. Over the WAN link, zstd is worth 1.9x and parallel readers 9x.

Timestamps now infer their unit from magnitude instead of assuming nanoseconds,
which was rendering microsecond columns as 1970, and SELECT * reaches UUID,
LONG256 and DECIMAL values that previously rendered as [object Object].
Profiling a scan put ~60% of active client CPU in readTimestampView and its
bit reader, with the BigInt allocation behind most of the GC on top. The cause
is the designated timestamp being GORILLA-encoded on the wire: delta-of-delta
bit packing cannot be handed over as a view, so it is unpacked value by value
with BigInt arithmetic, eagerly, before a batch is delivered. Every other
fixed-width column is a zero-copy view costing almost nothing. So the scan was
never decoding rows, but it was always decoding one column per row, and that
column was the only expensive one.

The scan's projection is now selectable. Casting the timestamp to LONG changes
the wire type, which routes it through readFixedView instead of the Gorilla
path while preserving the exact value, and is the new default:

  all           2,476,649 rows/s   timestamp, Gorilla-decoded
  epoch-long    6,293,394 rows/s   timestamp, as an epoch integer
  no-timestamp  8,714,877 rows/s   no timestamp

Measured at 8 readers against a local instance. The cast costs 8 uncompressed
bytes per row on the wire, which is the right trade on a LAN and the wrong one
on a thin link. The page needs no change to display it: format.js already
infers a timestamp's unit from its magnitude.

Empty batches no longer overwrite the displayed row tail. Slicing a small table
into many ranges yields plenty of zero-row batches, and whichever landed on the
final progress tick left the page reporting "last 0 rows received".

The third tab draws candles with a VWAP overlay and a volume histogram, built
with TradingView Lightweight Charts, vendored from node_modules rather than a
CDN so the demo does not depend on the network. SAMPLE BY aggregates the bars
server-side, so a window covering hours of trading arrives as a few thousand
bars rather than the millions of trades behind them. The window is anchored to
the newest row rather than to $now, so the chart still works on an instance
whose writer has stopped. The running VWAP is accumulated from the per-bar vwap
and volume already returned, so it costs no extra query.

Absolute timestamp windows use explicit >= and <= comparisons: the IN 'a..b'
interval form parses the $now-relative literals the TCA panels use, but rejects
an absolute ISO pair outright.
…ible

The tab opened on whichever symbol sorted first alphabetically, which was
AUDCAD at 687 trades in the window: a chart full of gaps that says nothing
about either the data or the database. The symbol query now counts trades and
orders by activity, so the default is the most active instrument (EURUSD, 5,727
trades over the same window). The picker itself stays alphabetical, because
that is how someone looks an instrument up.

Wheel zoom had no way back short of reloading. Double-clicking the plot, or the
new Fit button, restores the full window.
Three reasons the chart felt static, none of them visible from the page.

The live tail updated the candles and the volume but never the VWAP, so the one
line whose entire job is to move stopped at load time. It now carries the
running numerator and denominator past the loaded window and extends the line
from them.

New bars did not pull the viewport, so anything arriving landed off-screen:
the time scale now shifts on a new bar and keeps a small right offset.

Live was an unticked checkbox, which made a realtime chart static by default
unless the box was found. It ships on.

Zoom had no affordance at all: wheel and pinch worked, but nothing said so.
There are now + and - buttons driving the same visible logical range, alongside
Fit and double-click-to-reset, and a line under the chart naming the gestures.

A freshness tile reports how far behind the newest row is, because "the chart
looks static" has two unrelated causes - the page is not tailing, or nothing is
being written - and they were indistinguishable. It counts whether or not bars
arrive, so a stalled writer reads as a rising number rather than as a frozen
page.
Default to 1s bars over a 5m window, and make the live poll rate a control
(100-1000ms, default 250ms) instead of a hardcoded second. The bar width caps
how often a new candle appears; the poll caps how often the in-progress one is
redrawn, and both were too slow to look alive. Tail queries now cover 1m rather
than 5m, so the faster poll is cheaper. Adds 1m and 2m windows, and a legend
naming what the candles, the orange line and the volume bars are.
…query

The chart only auto-followed while the viewport sat at the last bar, so zooming
in looked like it had stopped: bars kept arriving, just outside the visible
range. It now tracks whether the right edge is in view and pans to the newest
bar without changing the zoom, and the hint line says which mode it is in. Fit
and double-click resume following.

Each live poll ran three statements - anchor, symbol list, bars - of which two
are useless for a tail. /api/ohlc-tail runs one, anchored on $now. Locally that
is 17ms against 7ms, but each statement is a separate round trip, so on a remote
cluster it is the difference between one latency and three per poll.
A 1s candle closes once a second however fast the page polls, so candles alone
look like a slideshow and no poll rate fixes that. The newest quote does change
with every tick, so it is now polled on its own: bid and ask are drawn as price
lines with axis labels and a bid/ask/mid/spread readout, and they move on every
poll while the candles close on their own cadence.

Quotes come from a single-row LATEST ON against core_price, so polling at 100ms
is cheap; bars are refetched about once a bar instead of once a poll, which they
were, needlessly.

The freshness tile now also reports measured polls per second, the round trip,
and how long since a new candle. 'It feels slow' has three causes - the browser
throttling the timer, a slow round trip, or the bar width - and they cannot be
told apart by eye.
…rs()

polars now deprecates casting UInt32 into a Categorical, which the QWP client
does internally for symbol columns, once per query. tca_live reprints every
tick, so the dashboard was buried in it. The cast is inside the client, so the
scripts can only filter the warning; matched on its message so other
deprecations still surface.
The scan now reports decoded bytes alongside rows, as GiB total plus MB/s and
Gb/s, measured the way python/read_bench.py measures it so the two are
comparable: valuesBytes() and nullBitmapBytes() are views over buffers the
client has already decoded, so it copies nothing and costs one pass over the
columns per batch rather than the rows. It is the decoded payload, not the
compressed wire bytes, which the client does not expose. Cross-checks at 68.0
bytes a row for core_price, matching both the schema and the Python figure.

--help lists every flag and endpoint and exits before connecting, so it answers
with no database reachable.
How alive a candle chart looks is governed by how WIDE each bar is, not by how
fast the data arrives. The default was 1s bars over 5m: 300 candles a few pixels
across, where a new one per second is invisible. The QuestDB console panel this
was compared against fits 20 bars into roughly 650px, about 30px each; this
chart is roughly 1500px, so about 50 bars matches it. 1m gives 53.

Separately, and measured: new rows become visible once per second on both
tables - 43 reads in 6s returned 7 distinct values, stepping in ~1.0s gaps - so
polling faster than that re-reads identical data. That is a ceiling on candle
cadence, but it was not what made the chart look static.
Three things pinned the newest candle to one update a second regardless of how
fresh the data was.

Bars were fetched on setInterval(1000), on the reasoning that a bar can only
change once per bar width. Only CLOSED bars are immutable; the one being formed
changes with every trade that lands inside it, so it has to be redrawn at the
poll rate.

Even corrected, a timer is throttled to 1Hz by the browser in a background tab.
Measured in one page at one moment: bars 1.0/s on a timer against quotes 5.1/s
on a stream. Bars now ride the same stream, which is not throttled and costs no
round trip per update.

The stream then sent no bars at all, because the change-detection key used
JSON.stringify without the BigInt replacer: result rows carry BigInt, the throw
was swallowed into a 'failed' event, and the symptom was 20 quote frames, zero
bar frames and 20 failures.
…used it

The old text asserted a 1Hz ceiling on visible data as if it were a property of
QuestDB. It is a property of the WRITER: the same measurement that gave 7
distinct values in 43 reads on a sender flushing once a second gave 43 distinct
values in 43 reads once ingestion moved to QWP flushing every 50ms.

Adds the three findings that actually made the chart look slow - browsers
throttle timers but not streams, only closed bars are immutable, and perceived
liveness is mostly bar width.
Start of a worker_threads scan that was never wired up, and swept into an
earlier commit by an over-broad git add. Nothing imports it. The decision was
not to pursue it: 8 vCPUs is about 4 physical cores, so threading buys perhaps
4-5x on a single-thread figure of 3.1M rows/s, and the link saturates around
12 Gb/s, which is where Python already tops out at 22.2M rows/s. The cheap half
of that win is already shipped as the epoch-long projection.
One candle advancing once a second reads as static next to thirty rows whose
prices, spreads and volumes all move at once, which is what makes the console's
FX top of book feel alive. Same query, pushed on the same stream as the candles
and quotes.

Rows are built once and mutated in place: rebuilding the table would restart
every animation and throw the flash away. A cell that moves twice quickly would
only flash once, because re-adding a class does not restart a CSS animation, so
the class is removed and a reflow forced first. The arrow outlives the flash so
direction stays readable, and rows stay alphabetical so none jumps position when
its price changes.
The scan never reads a value on its hot path, so swapping client builds only
exercises the decoder and never the accessor. This reads every row's timestamp
three ways - not at all, via getLong(), via getLongNumber() - with a checksum so
nothing is optimised away, and takes the best of two runs per mode.

It must run next to the database. Over a WAN it reports the link: on a slow link
reading nothing came out slower than reading every row, which cannot be true if
decoding is the bottleneck, so the script now warns when that happens rather
than letting the comparison be believed.
Both knobs were sized against a slow WAN, where a statement had to finish inside
the 60s query.timeout and 500k-row slices were the safe answer. Next to the
database that reasoning inverts: a 500k slice finishes in tens of milliseconds,
so per-statement setup stops being a rounding error, and 200M rows becomes 400
statements. The accessor benchmark on the same box read 2M rows in one statement
at ~8M rows/s while the sliced scan managed 3.1M, which is the gap this is meant
to explain.

Sweeps slice size against readers and prints the matrix, using the same
scanChunks() the app uses so the slicing under test is the slicing that ships,
and reporting how long a slice took at the winning setting so the timeout that
motivated the original sizing stays visible.
The sweep runs slice sizes in the outer loop, so the first cell paid the whole
V8 JIT cost for the decode loops and the smallest slice size looked slow for a
reason unrelated to slicing - which happens to be the result the benchmark was
written expecting, so the bias pointed the wrong way. Two discarded scans run
first, at both ends of the slice range.
Defaults were 4 readers and 500k slices, both chosen against a slow WAN. On a
same-AZ cluster the measurements say 8 readers and 2M slices: 15.31M rows/s with
epoch-long at 8 readers against 11.40M at 2.

Reader scaling turns out to depend on the projection, which is not obvious and
is now written down. With Gorilla the decode is a serial bottleneck on one
JavaScript thread, so throughput peaks at TWO readers and 4/8/16 are all
slightly worse; with the timestamp cast to long the decode is nearly free, the
limit moves to the network, and 8 readers beat 2 by 34%. Tuning readers before
fixing the projection measures the wrong thing.

Slice size makes no difference at all next to the database - 40 statements and 1
statement landed within noise across three repeated runs - so chunk_rows is now
documented as timeout headroom rather than a throughput knob.

The benchmark's warm-up used the full row count, which cost half a minute on a
200M-row sweep; it now uses a fixed small sample.
pnpm materialises a file: dependency as a hard-linked copy under .pnpm, so
rebuilding the client repo does not reach the app: the rebuild writes new files,
breaking the link, and pnpm install --force will not re-resolve a dependency
whose version has not changed. npm symlinks instead, so the same instructions
behave differently on different machines.

An A/B run that way silently compares a build against itself, which is what
happened: an apparent 4.8x improvement was run-to-run variance on identical
code, and the stale copy was only caught by an md5sum. --client imports the
built bundle directly, and the chosen path is printed on every run so a result
can no longer be misattributed to the wrong build.
The install steps hardcoded a relative hop between two repos, which only works
if they sit where mine do; following them on a flat $HOME linked a path that
did not exist and broke the install. The client's location is now captured as
$CLIENT at clone time and never assumed.

Adds the part that actually cost time: npm installs a file: dependency as a
symlink, so a rebuild is picked up, while pnpm makes a hard-linked copy that a
rebuild silently orphans - and pnpm install --force will not refresh it, because
the version has not changed. Comparing two client builds that way compares one
build against itself. The md5sum check that catches it is written down.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant