Repository navigation
feat(qwp): align ACK waits, reconnect and startup with other QuestDB clients - #71
Open
glasstiger wants to merge 51 commits into
Open
glasstiger wants to merge 51 commits into
glasstiger wants to merge 51 commits into
Conversation
QwpSenderSession and QwpSenderSessionFactory were exported from both package roots, which made the seam between QwpSender and its transports a public extension point. No other QuestDB client exposes a pluggable sender session: Java, Rust and Python keep that layer internal and hand out senders from their factories. Move both types to an internal module so neither package exports them, and mark the QwpSender constructor @internal. Applications keep obtaining senders from the runtime factories or a pooled client. The public type contracts now also assert that neither root re-exports the two names. Remove the QWP.md paragraph that documented custom session implementations, and regenerate the API reference.
Bring the QWP API closer to the Java, Rust and Python clients. ACK waiting: - Replace the awaitServerAck/awaitDurableAck sender options with flushAndWait(timeoutMs?), which resolves false when acknowledgements stop progressing. waitForAcknowledged() resolves a boolean too, and QwpIngressAckTimeoutError is removed. - Replace sendFrameWithPublication()/QwpIngressSendResult with publishFrame(). - Add QwpIngressSessionCloseTimeoutError. Reconnect: - Split QwpReconnectOptions into QwpIngressReconnectOptions (reconnect* names) and QwpEgressReconnectOptions (failover* names). A connected ingress session retries every outage until close(); reconnectMaxDurationMs bounds only a synchronous startup. Startup (initial_connect_retry, lazy_connect): - Accept lazy_connect=on|off|true|false in any case, and read initial_connect_retry case-insensitively, as the Java and Rust parsers do. - A standalone Sender ignores lazy_connect and warns, as the other clients' standalone senders do; it stays a pooled-client flag. - Make QwpIngressSessionOptions.initialConnectMode public and validate it. "async" now starts in the background in browsers as well, and the Node adapter honours the session spelling with store-and-forward and rejects two spellings that disagree. - Stop documenting initial_connect_retry as requiring sf_dir, and note that the Rust and Python pools do not promote startup to sync from reconnect_* keys.
glasstiger
force-pushed
the
ia_qwp_fixes
branch
from
October 5, 2026 17:24
9f6f1b2 to
f5b5316
Compare
QwpIngressSession was exported from both package roots, along with connectQwpNodeIngress() and connectQwpBrowserIngress(), which hand one out. That made the layer below QwpSender a second, parallel ingress API. The other QuestDB clients keep that layer internal: Java, Go and Python publish through their senders, and the batches Rust and C accept -- a row Buffer, or a column Chunk for zero-copy Arrow and Polars ingest -- are payloads handed to a sender, not a session. Stop exporting QwpIngressSession, the two factories, and QwpIngressSessionCloseTimeoutError, which only the session's own close() throws. The session options, metrics, notification and error types that the sender surfaces stay public. Both roots now re-export their runtime adapter by name instead of with `export *`, because the adapter module also exports the factory its senders use. The browser adapter therefore moves from index.ts to qwp.ts, verbatim apart from its header (git diff -C shows the move), matching the Node package's layout. The runtime contract asserts that the four names stay out of all three barrels, and the type contracts pin the same for the source and the built declarations. The browser bundle tests read the negotiated handshake from createQwpBrowserConnectionFactory() instead, which exercises the same ingress negotiation. QWP.md drops the low-level session guide and notes where the ordinary OK remains observable once durable ACK tracking moves the watermark.
Regenerated with pnpm run docs after the ingress session became internal. The pages for QwpIngressSession, QwpIngressSessionCloseTimeoutError, connectQwpNodeIngress and connectQwpBrowserIngress are removed, and source links now point at the browser adapter's new qwp.ts.
QwpNodeOrphanDrainer was exported from the Node root together with its options and session types and scanQwpNodeOrphanSlots(), so applications could construct and drive the drainer themselves. The other QuestDB clients configure orphan recovery instead -- drain_orphans and max_background_drainers in Java, Go, Rust and Python -- and expose at most a drainer listener and statistics. QwpNodeUdpSession was likewise a session-level UDP API beside connectQwpNodeUdpSender(), while Java publishes UDP rows through its Sender. Stop exporting QwpNodeOrphanDrainer, QwpNodeOrphanDrainerOptions, QwpNodeOrphanDrainSession, scanQwpNodeOrphanSlots, QwpNodeUdpSession and QwpNodeUdpMetrics, and remove connectQwpNodeUdp(), which only returned the session: createQwpNodeUdpSender() connects it directly. The drainer's notification surface -- onOrphanDrainEvent with its event, kind and metrics types -- stays public, as do retryQwpNodeOrphanSlot() and QWP_ORPHAN_FAILED_SENTINEL for re-enabling a quarantined slot. The UDP datagram counters were reachable only through the session; UDP senders keep their own metrics and onError. The runtime contract asserts that the four runtime names stay out of the Node root, and the type contracts pin all seven for the source and the built declarations. Tests reach the internals through their qwp-node/ modules, and config-docs no longer counts connectQwpNodeUdp() as a high-level entry point. QWP.md drops connectQwpNodeUdp() and scanQwpNodeOrphanSlots() and says how to configure and observe orphan recovery instead.
Regenerated with pnpm run docs after the orphan drainer and UDP session became internal. The pages for QwpNodeOrphanDrainer, QwpNodeOrphanDrainerOptions, QwpNodeOrphanDrainSession, scanQwpNodeOrphanSlots, QwpNodeUdpSession, QwpNodeUdpMetrics and connectQwpNodeUdp are removed, and the embedded QWP.md matches the updated guide.
The public option interfaces mirrored internal layers more than caller concerns, and several carried fields that belonged elsewhere. None of this API has shipped yet, so trim it before 5.0.0 does. The browser client accepted a split form with complete ingress and egress trees beside the cluster form. Both arrived in the same commit, so the split form was compatible with nothing published. Drop QwpBrowserSplitClientOptions and QwpBrowserUnifiedClientOptions: QwpBrowserClientOptions is now the cluster form, and a client without a cluster is rejected by name. Callers that must connect the two sides differently can give QwpClient their own factories. requestDurableAck sat on the WebSocket options both sides extend, so egress inherited it and two places had to strip it again: sent on /read/v1 it failed every query session. It now lives on the ingress types only -- QwpNodeIngressOptions and the new QwpBrowserIngressOptions, which also takes ingressNegotiationTimeoutMs -- and the Node egress path passes false to the endpoint connector instead of relying on a strip. The pooled client's typed override moves from webSocket.requestDurableAck to a new ingress section, and Sender.fromConfig() routes qwp.webSocket.requestDurableAck there. Without the ingress fields, QwpBrowserClusterOptions was identical to QwpBrowserWebSocketOptions and is removed. QwpEgressRoutingOptions becomes QwpRoutingOptions, since ingress honours target and zone too, and QwpSenderOptions takes gorilla and symbolDictionary directly instead of an encode object, which removes QwpSenderEncodeOptions. Five @internal fields on QwpIngressSessionOptions -- replayStore and the background and orphan store-and-forward handoffs -- were published in both packages' declarations, because nothing strips @internal. They move to QwpIngressSessionInternalOptions, which no root exports, together with the pre-session error delivery count, which had been attached under a symbol to avoid exactly that leak. Lazy startup no longer sets backgroundStoreAndForward on the public ingressSession options: the async initial connect mode already implies it. The type contracts pin the renamed and removed names for the source and the built declarations, and fail if an internal session field or an ingress-only field returns to a shared or egress type. New tests cover the Sender's requestDurableAck routing, an egress upgrade given requestDurableAck anyway, and a browser client without a cluster. QWP.md drops the split form and the encode object.
Regenerated with pnpm run docs after the option types were reduced. The pages for QwpBrowserClusterOptions, QwpBrowserSplitClientOptions, QwpBrowserUnifiedClientOptions, QwpEgressRoutingOptions and QwpSenderEncodeOptions are removed, QwpBrowserClientOptions is now an interface page rather than a type alias, and QwpBrowserIngressOptions and QwpRoutingOptions are added.
Each QWP entry point took its options in the shape of the layers behind it: a sender took connection, sender and session options as three arguments, an egress connection took connection and session options as two, and the pooled client took five sections split the same way. With the ingress session internal, two of those types named objects no caller can reach, and some settings lived in two places -- initialConnectMode on the session and on storeAndForward, maxBatchRows on the egress connection and on its session -- with code to reconcile them. The ingress and egress options of each runtime now absorb the sender's buffering options and the session's delivery options. connectQwpNodeSender(), connectQwpBrowserSender() and their create forms take one QwpNodeIngressOptions or QwpBrowserIngressOptions, the UDP senders one QwpNodeUdpOptions, and connectQwpNodeEgress() and connectQwpBrowserEgress() one egress options object plus an optional signal. QwpSenderOptions and QwpIngressSessionOptions become internal building blocks, which the published declarations still carry for the types that extend them; QwpEgressSessionOptions stays public for QwpEgressSession.connect(). storeAndForward loses its own initialConnectMode, so assertConsistentQwpInitialConnectMode() and the conflict it reported are gone, and one maxBatchRows is both the requested row cap and the decoder bound. The Node pooled client takes the cluster form the browser client already had -- cluster, ingress, egress, pool and lazyConnect -- with the endpoints and authorization owned by cluster and each side deriving its route from the cluster URL through a helper both runtimes now share. Overrides for a connect string take the same sections, all optional, so QwpNodeClientConfigOptions is no longer a vocabulary of its own, and lazyConnect gains a typed override. On the browser side QwpBrowserClientIngressOptions and QwpBrowserClientEgressOptions give way to the merged side options, minus what the cluster owns. Sender.fromConfig() keeps two qwp sections, webSocket and udp, each taking that sender's complete typed options; session and sender are folded into them. A section from the old shapes would otherwise be skipped with every setting inside it -- two tests here were passing qwp.session settings that no longer applied -- so the qwp sections, the pooled client's options and its overrides now reject an unknown section by name. Public option types fall from 19 to 17 on the Node root and from 18 to 14 on the browser root. The contracts pin the new signatures and shapes, and fail if a side accepts a cluster-owned field, if storeAndForward regains a startup mode, or if the merged ingress options regain an adapter-only session field. QWP.md, the READMEs and the example use the one-object forms.
Regenerated with pnpm run docs after each QWP entry point took one options object. The pages for QwpSenderOptions, QwpIngressSessionOptions, QwpBrowserClientIngressOptions and QwpBrowserClientEgressOptions are removed; the ingress and egress option pages now list the buffering, delivery and session members they include.
The connect-string overrides gained lazyConnect when they took the pooled client's own sections. Pin that a typed value wins over lazy_connect in both directions, and that enabling it applies the same async startup and cold query pool the key does.
Both packages exported a raw WebSocket connector and a stateful endpoint walker -- connectQwpNodeWebSocket() and createQwpNodeConnectionFactory(), connectQwpBrowserWebSocket() and createQwpBrowserConnectionFactory() -- each returning a bare QwpBinaryConnection with no session over it. With the ingress session internal, their only public consumer was QwpEgressSession.connect(), which connectQwpNodeEgress() and connectQwpBrowserEgress() already drive with factories of their own, and every sender, query session and pooled client opens and reconnects its own connections from the same endpoint options. Stop exporting all four. connectQwpNodeWebSocket() and connectQwpBrowserWebSocket() stay in their adapters as test seams, and createQwpBrowserConnectionFactory() as the factory behind the browser sender. The Node factory had no caller besides its own connector, so its public wrapper is removed and the internal walker takes its name. The same-origin and session-cookie guidance on the browser connector moves to QwpBrowserWebSocketOptions, so the API reference keeps it. The runtime contract asserts that the four names stay out of all three roots, and the type contracts pin them for the source and the built declarations. Three browser e2e tests drove the raw factory from the built bundle to read its handshake; they now connect a public sender and read its metrics, where effectiveAutoFlushBytes reflects the batch cap the SERVER_INFO frame advertised. QWP.md says the raw connectors are internal.
QwpWebSocketConnectOptions holds the endpoints and WebSocket deadlines both runtime adapters read, and both packages exported it. Nothing takes it on its own: QwpNodeWebSocketOptions and QwpBrowserWebSocketOptions extend it, and since the raw connection helpers became internal no public signature names it either. Exported, it was one more option type for a caller to tell apart from the two that are used. Move it from transport.ts, which the shared barrel re-exports whole, to _internal/websocket-connection.ts beside validateQwpWebSocketTimeouts(), which checks the same deadlines. The published declarations still carry it as the base of the two runtime options, as they do QwpSenderOptions, so its fields keep their documentation on those pages. The authTimeoutMs note now names connectTimeoutMs in prose rather than linking through the unexported base. The type contracts assert that neither root exports it, for the source and the built declarations, and fail if either runtime's WebSocket options lose one of its fields.
QwpNodeFileReplayStore was exported with its options and metrics, and both packages exported QwpIngressReplayStore with its record and reference types. The only way to hand a sender a replay store was the replayStore session option, which is internal, so applications could construct the journal but not give it to a sender: storeAndForward builds it, and QwpNodeStoreAndForwardOptions merely extended its options. Stop exporting QwpNodeFileReplayStore, QwpNodeFileReplayStoreOptions and QwpNodeFileReplayStoreMetrics, and move QwpIngressReplayStore, QwpIngressReplayRecord and QwpIngressReplayReference from transport.ts to a new _internal/replay-store.ts, which no root re-exports; the browser package carried the contract for a store it can never be given. QwpNodeStoreAndForwardOptions now declares the journal settings itself -- directory, maxBytes, maxSegmentBytes, durability, checkpointIntervalMs, backpressurePolicy, appendDeadlineMs and onRecoveryDataLoss -- and the store's internal options are a Pick of it, so the two cannot drift. The policies, the replay-store errors and the data-loss report stay public, since a sender surfaces them. The store's metrics were reachable only through a store a caller built; QWP.md now points at the backlog and watermarks a sender's metrics.ingress reports. Two dist tests built the store from the published bundle: the ESM worker check now opens a journal through an async sender, and the multi-process child drives its slot through a sender whose endpoint never answers, one journalled frame per flushed row, which keeps all five cross-process lock and recovery contracts. The runtime contract asserts that the store stays out of the Node root, the type contracts pin all six names for the source and the built declarations, and they fail if the store-and-forward options lose one of the journal settings. Unit tests import the store from its qwp-node/ module.
The typed storeAndForward options and the sf_* connect-string keys configure the same journal, yet they disagreed on four of its defaults: durability (append typed, memory in a connect string), the size cap (1 GiB against 10 GiB), a full journal (failing against waiting), and the journal's location (the directory itself against its `default` slot). Moving between Sender.fromConfig() and the typed options -- a change that looks like pure configuration style -- changed how durable and how bounded the journal was and pointed it at a different directory, stranding any backlog in the old one. The typed defaults were the journal's historical ones, the backpressure policy labelled backwards compatible, but no release exists to be compatible with; the connect string's follow the Java client. Take the connect string's everywhere. The journal's defaults come from one QWP_SF_DEFAULTS object, which the parser now resolves the sf_* keys from as well, so the two spellings cannot drift apart again. senderId defaults to `default`, as sender_id does, so storeAndForward.directory is always the slot root: a sender journals into <directory>/<senderId>, and pooled senders into <directory>/<senderId>-<slot>, `default-N` rather than `sender-N`. An adopted orphan still opens the slot the scanner found. With every journal a named slot, a standalone drainer always scans the configured directory, so the branch that ignored drainOrphans for an unnamed journal is gone. The warning that one construction style had left a journal where the other would not look becomes a warning that the slot root itself holds segments, which is what a directory naming a slot rather than its root leaves behind. Store tests that relied on the old defaults now ask for append durability or the error policy explicitly, and the transport tests use the default slot names. New tests pin that a typed storeAndForward object gets the parser's journal defaults, that a typed sender journals into its slot rather than the root, and the slot-root warning. QWP.md and the READMEs describe one set of defaults and one layout.
Regenerated with pnpm run docs after the four follow-ups. The pages for connectQwpNodeWebSocket, createQwpNodeConnectionFactory, connectQwpBrowserWebSocket, createQwpBrowserConnectionFactory, QwpWebSocketConnectOptions, QwpNodeFileReplayStore, QwpNodeFileReplayStoreOptions, QwpNodeFileReplayStoreMetrics, QwpIngressReplayStore, QwpIngressReplayRecord and QwpIngressReplayReference are removed. The WebSocket option pages still list the members of their unexported base, QwpNodeStoreAndForwardOptions lists the journal settings with the connect-string defaults, and the embedded QWP.md matches the updated guide.
QwpEgressSession is the query API, which connectQwpNodeEgress(), connectQwpBrowserEgress() and the pooled clients return, so the class has to stay public. Constructing one did not: the static connect() took a QwpConnectionFactory and the constructor a QwpBinaryConnection, and since the raw connection helpers became internal, nothing public produced either. A caller would have had to reimplement the QWP upgrade to use them; the documented ways to connect differently -- the webSocketFactory option and custom QwpClient factories -- go through the public egress functions instead. They were the egress counterpart of the ingress session made internal earlier, left public only because the class itself had to stay. The constructor now takes a token only this module holds, as QwpTableWriter's does, and the static connect() becomes connectQwpEgressSession(), with createQwpEgressSession() for one fixed connection. Neither root exports them, so the shared barrel lists the egress session's exports by name, as it already did for the ingress session and the sender. Opening a reconnecting connection used three of the session's private replay methods from inside the class; the constructor now fills in an internal hooks object with them instead, which the reconnecting connection calls through. With no public signature left taking them, QwpBinaryConnection, QwpConnectionFactory and the two transport-metrics types only that connection exposes move from transport.ts to a new _internal/binary-connection.ts. QwpEgressSessionOptions becomes a non-exported base of the runtime egress options, as QwpIngressSessionOptions did. QwpHandshakeMetadata stays public for the session's handshake getter. The runtime contract asserts that the factories stay internal, that the class has no static connect(), and that its constructor rejects a call without the token. The type contracts pin all five type names for the source and the built declarations, and fail if connect() returns, if the constructor stops leading with its token, or if the runtime egress options lose a session field. Tests construct sessions through the internal factories. QWP.md says how query sessions are opened.
QwpSender's constructor was documented as internal but callable: it took a QwpSenderSessionFactory, a type neither root exports since the sender session interface became internal, so the published declarations offered a constructor no caller could satisfy without reaching into unexported types. The runtime factories and the pooled clients' borrowSender() are how applications obtain a sender. The constructor now takes a token only the sender module holds, as QwpTableWriter's and now QwpEgressSession's do, and createQwpSender() builds senders for the runtime adapters. The shared barrel already lists the sender module's exports by name, so the factory stays internal. The runtime contract asserts that createQwpSender stays out of all three roots and that the constructor rejects a call without the token, and the type contracts fail, for the source and the built declarations, if the constructor stops leading with its token. Tests and the sender benchmark build senders through createQwpSender(). QWP.md says where a QwpSender comes from.
Regenerated with pnpm run docs after sender and query-session construction became internal. The pages for QwpBinaryConnection, QwpConnectionFactory, QwpEgressSessionOptions, QwpEgressTransportMetrics and QwpIngressTransportMetrics are removed from both packages. QwpEgressSession no longer lists a static connect(), its constructor and QwpSender's are marked internal, the egress option pages list the session members they include, and the embedded QWP.md matches the updated guide.
QwpSender, QwpEgressSession and QwpTableWriter are public classes whose constructors take a token only the package holds, so a caller can never use them. TypeDoc still rendered each one, token parameter and all: excludeInternal does not remove a constructor tagged @internal. Tag them @hidden as well, so the reference pages list only what a caller can use. The comments keep @internal and say where instances come from, for readers of the source and the published typings, which still declare the constructors.
Regenerated with pnpm run docs after the token-guarded constructors were hidden. The QwpSender, QwpEgressSession and QwpTableWriter pages in both packages no longer list a constructor; QwpClient's public constructor is still documented, and no page is added or removed.
…ested Without requestDurableAck, the Node client took durableAckKeepaliveMs (durable_ack_keepalive_interval_millis) as a request for durable ACKs and rejected it next to requestDurableAck=false, while the browser client rejected it outright. The Java and Rust clients accept the keepalive in both cases and ignore it: it only paces the durable-ACK poll once durable ACK has been requested. Promoting it also changed the upgrade Node sent, so a connect string those clients use against a server without durable ACK failed to connect from Node with QwpDurableAckUnavailableError. Both adapters now hand the session a keepalive only when durable ACK was requested -- the configured interval, or the 200ms default -- through resolveQwpDurableAckKeepaliveMs(), which the Node orphan drainer uses as well. The session reads a defined interval as the request to track durable progress, so dropping it is what keeps an ignored keepalive inert. The value is still validated when it is ignored. The connect-string parser no longer sets request_durable_ack from the keepalive or rejects it next to request_durable_ack=off. The Node and browser tests connect to a server that offers durable ACK unasked and check that the ordinary OK still advances the watermark, with no durable-ACK request and no polls. The Node test that polls with the keepalive now requests durable ACK explicitly. QWP.md documents the 200ms default, which config-docs pins to the code, and says the key is ignored without request_durable_ack=on.
Regenerated with pnpm run docs after the durable-ACK keepalive change. The QwpNodeIngressOptions and QwpBrowserIngressOptions pages say durableAckKeepaliveMs takes effect only alongside requestDurableAck and is otherwise ignored, the embedded README and QWP.md match the updated guides, and no page is added or removed.
QwpIngressSession.close() gave frames still in the in-memory replay queue up to 5 seconds to reach the socket, and rejected with QwpIngressSessionCloseTimeoutError when some could not. QwpSender never used it: its close() already flushes, waits up to closeFlushTimeoutMs for the committed-frame ACK watermark -- or, in a fast close, for those frames to reach the socket -- and rejects when that drain cannot finish, then closed the session through closeWithoutDrain() so the two budgets could not stack. Since the ingress session became internal, its only other caller is the orphan drainer, whose sessions journal to disk, where the drain did nothing. As with the Java client's close() and the Rust and Python clients' close_drain(), the drain belongs to the sender, which is what applications call. The session's close() now closes without a drain, as closeWithoutDrain() did, and closeWithoutDrain() goes from the session, the sender session interface and the sender. So do QwpIngressSessionCloseTimeoutError and the connection's stopPublishing() and unsentFrameCount, which only that drain used. The unsent-frame count stays behind waitForPendingSends(), which the sender's fast close still waits on. The three tests of the session drain go, the five-minute RAM replay test no longer closes into a timeout, and the public API contracts stop listing the removed error.
QwpSender.commit() only called flush(), so it added a second name for the same operation: an explicit flush() already publishes the group-closing frame of a transactional sender, rows auto-flushed earlier and pending rows alike. No other QuestDB client has the alias. The Java client, whose transactional mode this one follows, commits on an explicit flush() and has no commit() method; the Rust, Python, Go and .NET clients have no QWP transactional mode at all. flush() now documents that it is the commit point of a transactional sender, and so does the transactional option. Tests that committed through the alias call flush() instead; the closed-sender check no longer repeats its flush() case for it. README.md, QWP.md and the browser README commit with flush(), and QWP.md maps the Java client's publish/commit concept to flush() alone.
The Java client's QuestDB facade reads query_timeout_ms as the default timeout of every query, with 0 meaning none. This parser rejected the key as unknown, so a connect string written for that client failed before a single query ran, and a connect string had no way to set the default at all. The key now sets queryTimeoutMs on the egress session options. It accepts 0 and, like every millisecond key that arms a timer, is capped at the host timer ceiling. A typed egress.queryTimeoutMs still wins over it, and Sender.fromConfig() names it among the query-side keys it ignores. QWP.md lists it with the egress keys.
A query timeout rejected the caller the moment it expired, then cancelled the query. A statement that completed just past the deadline had taken effect but was reported as timed out, which invites a retry that applies it twice. The Java client's per-query timeout (questdb/java-questdb-client#105) waits for the server instead, and this client now does the same: - Once the timeout expires no further batch reaches the consumer and the query is cancelled, but the outcome waits up to one grace period, cancelDrainTimeoutMs, for the server's terminal response. EXEC_DONE, and a RESULT_END after which nothing was withheld, still succeed. A CANCELLED reply, or a result that lost a batch to the timeout, rejects with QwpEgressQueryTimeoutError. A server still silent after the grace period releases the caller with the timeout, and the session gives the connection up one period later. - The timeout runs from the query() call, so waiting for SERVER_INFO, for a previous query to drain or for a reconnect counts against it. A request whose timeout expires while a reconnect holds it is released at the deadline and never sent. - On a lost connection, a query past its timeout or whose consumer has retired is settled at once instead of being replayed only to be cancelled, and the cancel drain no longer runs without a connection. A reconnect longer than cancelDrainTimeoutMs used to fail the session. - A cancel the application issued before the timeout is reported as a cancellation, a query sends at most one CANCEL, and query() waits for a released query to finish draining instead of throwing "a QWP query is already active", which it did even right after breaking out of for await. The reconnecting egress connection now tells the session when it loses its connection, lets it drop a query from the replay and withdraw a request a reconnect still holds, and reports whether a terminal response is still queued: one that arrived before the connection went still decides the outcome. The session tests cover each outcome past the deadline, both grace periods, withheld materialized batches and views, an application cancel, the single CANCEL, the clock starting at query() and the wait behind a drain. The reconnect tests cover a timeout during an outage, a connection lost during the grace period, a draining query, a replayed application cancel, a terminal queued before the loss and a withdrawn request. README.md and QWP.md describe the behaviour, and query_close_timeout_ms is documented as the grace period too. The timeout is still enforced by the client alone; sending it to the server waits for questdb/questdb#7768. BREAKING CHANGE: a timed-out query is reported when the server ends it or after cancelDrainTimeoutMs, rather than at the deadline; a statement that completes past its timeout succeeds; the timeout includes the wait before the request is sent; and query() waits for a query that is still draining instead of throwing.
Regenerated with pnpm run docs after the query-timeout change. The QwpEgressQueryOptions, QwpNodeEgressOptions, QwpBrowserEgressOptions, QwpEgressQueryTimeoutError and QwpEgressQueryCancelTimeoutError pages describe the timeout measured from query() and its grace period, the embedded README and QWP.md match the updated guides -- the QWP.md copy had fallen behind since the commit() alias was removed, which test/docs-reference.test.ts reported -- and no page is added or removed.
A browser sender that requested durable ACKs polls for durable progress with table-less QWP frames, and those pass through the in-memory replay queue like data. Since RAM publication completes at enqueue, a sender keeps accepting flushes during an outage until that queue is full. The next poll then waited out memoryReplayAppendDeadlineMs for trimming that only durable progress -- which only a poll can report -- could bring, and its QwpMemoryReplayAppendTimeoutError went to the poll's failure handler, which ended the session. Every row in the queue was lost, even though the same full queue only rejects an ordinary publication and the sender survives the outage without durable ACK. A queue filled by frames awaiting durable upload failed the same way while connected. The memory replay queue now admits a durable-ACK poll above its target, as it already admits a transaction-closing frame: it is a fixed-size control frame, and the progress it asks for is what trims the queue. To keep that from growing during an outage, a browser poll is published only when no frame is waiting for the socket; the frame or poll already queued reaches the server first, and the next keepalive tick tries again. Node senders poll with WebSocket PING and are unaffected. A new reconnect test fills the queue with frames awaiting durability, drops the connection and checks that one poll is queued, that a data publication still meets non-terminal backpressure, and that the session recovers on the replacement connection. It fails without either change. QWP.md and the README describe the exception to the memory cap.
The regression test for #70 waited with a 1-second timeout and asserted only that the wait resolved true. The wait's own deadline re-checks the ACK watermark before it expires, so the test also passed with the listener that wakes waiters for an OK the reconnecting connection does not forward removed: the wait then resolved a second late instead of when the OK arrived, and the whole suite stayed green. The wait now uses a timeout longer than the test's, so only the OK itself can settle it. Without the listener the test times out.
Four passages still described behaviour this branch has changed: - The E2E benchmark's local-publication arm was documented, and labelled, as flush() returning once the WebSocket accepted the frame. Without a journal, flush() now completes when the frame enters the in-memory replay queue and a background drainer hands it to the socket, so the arm measures that boundary. - README.md and the QWP.md migration notes called the store-and-forward publication boundary a durable journal append. The journal's default durability is now `memory` for typed options as for connect strings, which relies on the page cache and makes no power-loss promise; they now say so and name `durability: "append"`. - QWP_INITIAL_CONNECT_MODE said browser and memory-only senders resolve their startup mode internally and that only Node store-and-forward exposes the three modes. initialConnectMode is an ingress option of every sender, with or without a journal. - The QWP.md pool section pointed at an exception noted in its table, but the only such note was removed with the egressSession section.
Regenerated with pnpm run docs after the durable-ACK poll fix and the documentation corrections. The QwpNodeIngressOptions and QwpBrowserIngressOptions pages say a durable-ACK poll is admitted above memoryReplayMaxBytes, the QWP_INITIAL_CONNECT_MODE page describes the initialConnectMode option every sender accepts, the embedded README and QWP.md match the updated guides, and no page is added or removed.
QuestDB accepts writes on the primary alone. A replica, or a primary still catching up, answers the /write/v4 upgrade with 421 and the role it holds, and the endpoint sweep moves on, so ingress reaches the primary whatever it asks for. The Java sender therefore accepts target and zone and ignores them, and the shared connect-string reference scopes both keys to egress. The Node client applied them to ingress as well: as a role filter on the X-QuestDB-Role a completed write upgrade advertises, and as zone affinity in the endpoint ranking. That could not route a write anywhere, and target=replica -- the usual read-scaling setting -- made the sweep refuse the primary it had found. A pooled client built from a cluster string carrying it could neither prewarm nor borrow a sender (QwpPoolResourceError), and Sender.fromConfig() never connected. The server sends no zone on the ingress upgrade, so zone had no input to rank by either. Ingress now passes neither key to its endpoint walker or its health trackers, and the connect-string parser puts them in the egress section only: - QwpNodeIngressOptions no longer includes QwpRoutingOptions. Setting target or zone on a sender's own options, on a pooled client's ingress or cluster section, or in a Sender's qwp.webSocket section throws a TypeError naming the option, rather than silently dropping it for a JavaScript caller. - Sender.fromConfig() lists target and zone among the keys it warns it ignores, with the other egress and pool keys. - The endpoint walker no longer lets an endpoint that declares no role satisfy every target. That rule existed for ingress, which reads the role from an upgrade header; egress always learns one from SERVER_INFO. The test whose mock replica completed the write upgrade, which a real replica never does, is replaced by one where the replica answers 421: a raw ingress session, and a pooled client built from a target=replica cluster string, both reach the primary. The Sender.fromConfig() test that sets target=replica now has its mock advertise PRIMARY, the configuration, typed-override and warning tests pin the new behaviour, and the public API contract checks that no ingress options type accepts either key. QWP.md scopes both keys to egress.
Regenerated with pnpm run docs after the ingress routing fix. The QwpNodeIngressOptions page no longer lists target and zone and points to QwpNodeEgressOptions for them, the QwpRoutingOptions and QwpTarget pages in both packages scope the routing options to query sessions, the embedded QWP.md matches the updated guide, and no page is added or removed.
The Java, Rust, Go and Python clients deliver the same seven connection events to their listeners: CONNECTED, DISCONNECTED, RECONNECTED, FAILED_OVER, ENDPOINT_ATTEMPT_FAILED, ALL_ENDPOINTS_UNREACHABLE and AUTH_FAILED (Java's SenderConnectionEvent.Kind). reconnect.onEvent called the outage `reconnecting`, had no per-endpoint or authentication event, and folded every failed attempt into `attempt-failed`. A listener learned which endpoints had failed only once the whole sweep had, from the QwpFailoverError attempts, and could tell a rejected credential from an unreachable cluster only by inspecting the cause. A listener ported from another client, or a dashboard keyed on these kinds, found three of them missing and one renamed. QWP_RECONNECT_EVENT_KIND now carries all seven, on ingress and egress sessions alike: - `reconnecting` is renamed `disconnected`. It still fires once per outage, before the first retry. - `endpoint-attempt-failed` fires as each endpoint fails, while the sweep is still running. The endpoint walker reports each failure through a new optional parameter of the internal connection factory, and skips the attempts an abort ended. - `attempt-failed` is replaced by `all-endpoints-unreachable`, which fires once per failed sweep and names the last endpoint tried. - `auth-failed` reports a 401 or 403. The rejection ends the sweep, so it takes the place of both events above for that endpoint. An endpoint that opens and then fails before the session can use it, as when it refuses the replayed frames, is reported as endpoint-attempt-failed, because this client raises its success events only once the replay has finished. A sweep that close() or the startup deadline cuts short reports no failure: no endpoint was judged, and the caller receives the error. The ingress loop used to emit attempt-failed for both. QWP.md lists the kinds. The new tests pin the events of a failed sweep, of an authentication rejection from a walking and from a non-walking factory, of an endpoint that fails its replay, of a sweep that close() interrupts and of one the startup deadline cuts short, and the egress equivalents of the first two. They also pin the walker reporting each failed endpoint before it tries the next, and staying silent for an aborted sweep. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016VDgngcJ9c6YZB5wtQt3N1
Regenerated with pnpm run docs after the connection-event fix. The QWP_RECONNECT_EVENT_KIND pages in both packages describe each kind and list disconnected, endpoint-attempt-failed, all-endpoints-unreachable and auth-failed in place of reconnecting and attempt-failed, the QwpReconnectEvent pages say which endpoint an all-endpoints-unreachable event names, the embedded QWP.md matches the updated guide, and no page is added or removed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016VDgngcJ9c6YZB5wtQt3N1
The Java, Rust and Python clients escalate a poison frame once it has been rejected maxFrameRejections times over at least 5 seconds of connected time (DEFAULT_POISON_MIN_ESCALATION_WINDOW_MILLIS in Java, QWP_WS_DEFAULT_POISON_MIN_ESCALATION_WINDOW in Rust, which the Python client uses). This client waited 5 minutes, so with the same connect string it kept retrying a frame the other clients had already given up on. poisonMinEscalationWindowMs, and so poison_min_escalation_window_millis, now defaults to 5_000. Escalation still requires maxFrameRejections strikes as well, and outage time still does not count toward the window. A retriable rejection that outlasts 5 seconds -- a concurrent DDL, a checkpoint, a briefly full volume -- can now end a producer the old default kept alive, so a deployment that sees such faults should raise the window. The comment that argued for the 5-minute default now says that instead, and QWP.md and the option's JSDoc state the new value. config-docs.test.ts now checks the documented default against the constant, as it already did for the reconnect_* keys. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016VDgngcJ9c6YZB5wtQt3N1
The other QuestDB clients name their ACK wait and watermark with "acked": awaitAckedFsn() and getAckedFsn() in Java, await_acked_fsn() and acked_fsn() in Python, AwaitAckedFsn() and AckedFsn() in Go, and acked_fsn() in Rust. This client spelled them out as waitForAcknowledged() and acknowledgedSequence. They are now: - waitForAck(), on QwpSender and the Node.js Sender. The sender session interface, and the ingress and UDP sessions that implement it, rename the method the sender waits through as well. - ackedSequence, on QwpSender and Sender, and on the two public types that report the same watermark: QwpIngressMetrics, which a sender's metrics.ingress and the ingress progress and error events carry, and QwpSenderCloseTimeoutError, whose message names the field ackedSequence too. Neither name has shipped in a release, so the old ones go without deprecated aliases. Other internal names, such as the ingress session's acknowledgedFrameSequence, keep the long form. QWP.md, the tests, the benchmark's encoding session and the public API contracts use the new names. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016VDgngcJ9c6YZB5wtQt3N1
Regenerated with pnpm run docs after the poison escalation window fix and the ACK renames. The QwpSender, QwpSenderCloseTimeoutError and QwpIngressMetrics pages in both packages, and the Node.js Sender page, list waitForAck() and ackedSequence in place of waitForAcknowledged() and acknowledgedSequence, the QwpIngressReconnectOptions pages give the 5-second default of the poison escalation window, the QwpNodeIngressOptions and QwpBrowserIngressOptions pages name waitForAck() in their ackTimeoutMs and onError descriptions, the embedded QWP.md matches the updated guide, and no page is added or removed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016VDgngcJ9c6YZB5wtQt3N1
The other QuestDB clients call this callback a failover reset: onFailoverReset() on Java's QwpColumnBatchHandler, on_failover_reset() and FailoverResetEvent in Rust, qwp_reader_query_on_failover_reset() in C, and QwpFailoverReset in Go. This client called it onReplayReset, although "replay" here names the ingress store-and-forward journal (QwpReplayStore*, QwpNodeReplayRecoveryEvent), and the egress policy the callback belongs to is already configured through failover* options. - onFailoverReset replaces onReplayReset on QwpEgressSessionOptions and QwpEgressQueryOptions, and QwpEgressFailoverResetEvent replaces QwpEgressReplayResetEvent. The event type now says the callback runs before every replay, whether the reconnect reached another endpoint or the same one, as Java's does, and QWP.md says so too. - The internal handler, notifier and callback-failure error follow, so the error a throwing callback ends the connection with now reads "QWP egress failover reset callback failed". - QwpEgressReplayRequiredError is gone. It was deprecated from the start, nothing threw it, the release it was kept for never shipped, and its message named onReplayReset. Neither name has shipped in a release, so the old ones go without deprecated aliases. The replay mechanism keeps its internal names, such as QwpEgressReplayHooks. QWP.md, the browser README, the tests and the public API contracts use the new names. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016VDgngcJ9c6YZB5wtQt3N1
Regenerated with pnpm run docs after the onFailoverReset rename. The QwpEgressQueryOptions pages in both packages, and the QwpNodeEgressOptions and QwpBrowserEgressOptions pages, list onFailoverReset in place of onReplayReset, QwpEgressFailoverResetEvent pages replace the QwpEgressReplayResetEvent ones and say the callback runs before every replay, whether the reconnect reached another endpoint or the same one, the QwpEgressReplayRequiredError pages are removed, the package index pages, navigation and search index follow, and the embedded QWP.md matches the updated guide. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016VDgngcJ9c6YZB5wtQt3N1
A fast close (closeFlushTimeoutMs <= 0) still bounds publication and the RAM-frame send drain by the 5 s default, but its QwpSenderCloseTimeoutError reported the configured value, so a close that waited 5 s during an outage said "timed out after 0ms". The deadline and the error now share one effective timeout. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… fails With durable ACK, a DURABLE_ACK could trim the journal and then fail to retire a recovered uncommitted tail. The store's watermark advanced but the session's did not, and with nothing left to replay no later response published it: waitForAck() and flushAndWait() returned false and close() reported "pending data may be lost" for durable data. The DURABLE_ACK path now publishes the watermark before rethrowing, as the OK path does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Only source links and line anchors change; no page's content changes and no page is added or removed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A typed pooled client without senderId used to name its store-and-forward slots sender-N. It now journals into default-N, and its drainer adopts only those names unless drainOrphans is set, so a backlog left under sender-N was neither replayed nor reported. The pool now warns on start when such a slot still holds journal segments, naming each one, and QWP.md describes the migration: rename the slots to default-N, or drain them once with drainOrphans. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The test for a durable trim whose tail retirement fails waited only after the watermark was published, so it passed without the wake-up the fix also added. The waiter is now registered before the DURABLE_ACK arrives, with a timeout longer than the test's, so a missed wake-up fails the test. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The embedded QWP.md page gains the migration note for a typed pool's earlier sender-N slots; every other page changes only its source links and line anchors, and no page is added or removed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Aligns the QWP client with the Java, Rust and Python clients. ACK waits, reconnect, startup and store-and-forward defaults now behave the same way. Applications reach QWP only through senders, query sessions and pooled clients: the layers below them are internal, and each entry point takes one options object. Along the way, the PR fixes ACK-wait, in-memory replay, pool teardown and query-timeout bugs. Fixes #70.
QWP has not shipped in a release yet: npm
latestis 4.2.0, and@questdb/browser-clientis unpublished. The API changes below therefore break only code written againstmain. For the released ILP transports, the only change is the newSender.flushAndWait().Fixes
Replay ACK waiters (waitForAcknowledged() never resolves for frames replayed from the store-and-forward journal #70).
waitForAck()now wakes when a store-and-forward replay advances the ACK watermark without forwarding the replayed OK to the session. The session-visible watermark advances only after the response carrying it has been queued, so a woken waiter orclose()cannot overtake that response'sonResponseand progress notifications.In-memory replay delivery. Without a journal,
flush()completes once frames enter the bounded in-memory replay queue. A background drainer then sends them in order, so a sender keeps accepting flushes during an outage until the queue applies backpressure. To keep those frames from being dropped silently:closeFlushTimeoutMs <= 0) still waits up to 5 s for them to reach the socket, and otherwise rejects withQwpSenderCloseTimeoutError;requestDurableAck) are admitted above the queue's cap, and none is queued while frames still wait for the socket. Before this, a poll that met a full queue during an outage, or a queue full of frames awaiting durable upload, ended the session and lost every queued row.Pooled sender teardown is best effort, as in the other clients' pools. A failed close logs a warning, and the warning notes when frames remain in the store-and-forward journal. Neither
QwpClient.close()nor a lease'sclose()rejects.Query timeouts settle on the server's answer (commit
f8ad871c, markedBREAKING CHANGEagainstmain). When the timeout expires:cancelDrainTimeoutMs(query_close_timeout_ms) for the server's terminal response;CANCELLEDreply, or a result that lost a batch to the timeout, rejects withQwpEgressQueryTimeoutError.The timeout now runs from the
query()call, so waiting for a reconnect or for a previous query to drain counts against it.query()waits for a released query that is still draining instead of throwing "a QWP query is already active". The timeout is still enforced only by the client.Egress: resetting a query schema no longer throws on skipped batch slots.
QwpEgressQueryOptions.onReplayResetadds a per-query replay callback, which pooled leases need because request IDs are not unique across pooled sessions.Durable-ACK keepalive.
durableAckKeepaliveMs(durable_ack_keepalive_interval_millis) takes effect only together withrequestDurableAck, as in the Java and Rust clients. Node no longer treats it as a request for durable ACK. The browser client no longer rejects it.API changes
ACK waiting
awaitServerAck,awaitDurableAckanddurableAckTimeoutMssender options are removed. Instead,flushAndWait(timeoutMs?)onQwpSenderand on the NodeSenderflushes and then:trueonce every published frame is acknowledged;falsewhen the ACK watermark makes no progress fortimeoutMs. That timeout defaults toackTimeoutMs(15 s) and restarts whenever the watermark advances.trueonce the data is sent.waitForAcknowledged()is renamedwaitForAck()and resolves a boolean with the same timeout semantics.QwpIngressAckTimeoutErroris removed, and a timed-out wait is reported only to its caller, not toonError.requestDurableAck: true. The watermark then follows durable progress, soflushAndWait(),waitForAck()and the close drain all wait for durability.QwpSender.commit()is removed.flush()is the commit point of a transactional sender, as in the Java client.One options object per entry point
connectQwpNodeSender()/createQwpNodeSender()take oneQwpNodeIngressOptions, which holds connection, buffering and delivery settings.connectQwpBrowserSender()/createQwpBrowserSender()take oneQwpBrowserIngressOptions.QwpNodeUdpOptions.connectQwpNodeEgress()/connectQwpBrowserEgress()take one egress options object plus an optional signal.QwpSenderOptions,QwpIngressSessionOptions,QwpEgressSessionOptionsandQwpWebSocketConnectOptionsare now non-exported bases of those types.gorillaandsymbolDictionaryreplace theencodeobject.requestDurableAckis an ingress option only; the egress upgrade never requests durable ACK.QwpEgressRoutingOptionsbecomesQwpRoutingOptions.{ cluster, ingress, egress, pool, lazyConnect }in Node and{ cluster, ingress, egress, pool }in browsers.clusterowns the endpoints and authorization, and each side derives its route from the cluster URL. The browser split form is removed.Sender.fromConfig()keeps twoqwpsections,webSocketandudp. Each takes that sender's complete typed options.Internal layers
No package root exports the following any more. Applications use senders,
connectQwp*Egress()and the pooled clients.QwpIngressSession,connectQwpNodeIngress()andconnectQwpBrowserIngress();QwpSenderSessionandQwpSenderSessionFactory;connectQwpNodeWebSocket(),createQwpNodeConnectionFactory(),connectQwpBrowserWebSocket()andcreateQwpBrowserConnectionFactory(), together withQwpBinaryConnection,QwpConnectionFactoryand the transport metrics types;QwpNodeFileReplayStoreand its options and metrics, plusQwpIngressReplayStore,QwpIngressReplayRecordandQwpIngressReplayReference;QwpNodeOrphanDrainer, its options and session types, andscanQwpNodeOrphanSlots(). Orphan recovery is configured throughstoreAndForwardand observed throughonOrphanDrainEvent.QwpNodeUdpSession,QwpNodeUdpMetricsandconnectQwpNodeUdp();QwpEgressSession.connect(). TheQwpSenderandQwpEgressSessionconstructors now require a token that only the package holds.Each package root therefore lists its runtime adapter's exports by name. The browser adapter moved from
src/index.tstosrc/qwp.ts, matching the Node package's layout.Reconnect
QwpReconnectOptionsis split in two:QwpIngressReconnectOptions(reconnect*names) andQwpEgressReconnectOptions(failover*names), mirroring the connect-string keys.close()or a terminal error.reconnectMaxDurationMsbounds only a synchronous startup, and0now allows one attempt without retries. Egress failover episodes remain bounded.reconnect*tuning settings promote an unset startup mode to"sync".Startup
initialConnectModeis an ingress option of every sender, with or without a journal."async"now also works in browsers.initial_connect_retryis case-insensitive.lazy_connectacceptson,off,trueandfalsein any case. Only the pooled client applies it:Sender.fromConfig()validates it, ignores it and logs a warning, as the other clients' standalone senders do.Store-and-forward defaults and layout
The typed
storeAndForwardoptions now default exactly as thesf_*keys do:durability: "memory"(was"append");maxBytes10 GiB (was 1 GiB);backpressurePolicy: "wait"(was"error");appendDeadlineMs30 s.Set
durability: "append"if every append must survive power loss.Each journal is a slot below the configured directory:
<directory>/<senderId>, wheresenderIddefaults todefault;<directory>/<senderId>-<slot>, i.e.default-Nwhere they usedsender-Nbefore.A typed sender built from
mainjournalled directly into<directory>. Move such a backlog into the slot directory before upgrading. A typed pool started withoutsenderIdadopts onlydefault-Nslots and warns when an earliersender-Nslot still holds journal segments; rename those todefault-N, or start the pool once withdrainOrphans: true. The client warns when the slot root itself holds journal segments.Connect-string keys
query_timeout_mssets the default query timeout;0disables it.Documentation
QWP.md,CLAUDE.md, the root, browser-client and nodejs-client READMEs and the benchmark README are updated to match. The regenerated TypeDoc reference underdocs/accounts for most of the changed files.Validation
Local runs on the head of this branch:
pnpm test: 1,200 tests pass.pnpm test:dist,pnpm check:packagesandpnpm test:qwp-browser(8 real-browser tests) pass.GitHub Actions, including the Enterprise durable-ACK E2E, passed on the previous head
dfaa95dand runs again on this one.