Skip to content

feat(qwp): add per-query timeout to the query client - #105

Open
bluestreak01 wants to merge 2 commits into
mainfrom
vi_query_timeout
Open

bluestreak01 wants to merge 2 commits into
mainfrom
vi_query_timeout

Conversation

@bluestreak01

@bluestreak01 bluestreak01 commented Oct 6, 2026 •

Copy link
Copy Markdown
Member

Tandem: questdb/questdb#7768 (server side). It decodes timeout_ms, applies it to the query's circuit breaker, reports timeouts as STATUS_QUERY_TIMEOUT, and advertises CAP_QUERY_TIMEOUT. Against servers without it, this client enforces the timeout on its own (client deadline + CANCEL).

Merge order: this PR first. questdb/questdb#7768 pins the submodule to this branch's head and moves the pointer to the merged commit before it merges. Both sides use the same protocol constants: CAP_QUERY_TIMEOUT = 0x08, QUERY_FLAG_TIMEOUT = 0x02, STATUS_QUERY_TIMEOUT = 0x0E.

What

Queries can now be given a timeout:

  • query_timeout_ms=N in the connection string sets a default for every query (0, the default, means none);
  • Query.timeout(long, TimeUnit) overrides it per handle;
  • on the low-level client, QwpQueryClient.withQueryTimeout(long) sets the default and execute(sql, binds, handler, resetSymbolDict, timeoutMs) overrides it per call.

A timed-out query fails with QueryException.isTimeout() (status 0x0E), and the pooled, authenticated connection stays open for the next query. Completion.await(timeout) keeps its meaning: it only bounds the wait. Its javadoc now points to the query timeout.

How

Wire. When SERVER_INFO advertises CAP_QUERY_TIMEOUT, the client sets QUERY_FLAG_TIMEOUT and appends timeout_ms:varint (the remaining budget) after query_flags. The server applies it to the query's NetworkSqlExecutionCircuitBreaker, the same mechanism as HTTP's Statement-Timeout, which trips inside the SQL engine. It then ends the query with a QUERY_ERROR. That is a per-query error, so the connection is kept. Without the capability, frames are byte-identical to today.

Client deadline (QwpQueryClient.executeOnce), always on:

  1. Before the deadline, nothing changes.
  2. Once the deadline passes, no further batch reaches the handler.
    • Batches are still decoded and released. That keeps the connection's SYMBOL dict in step with the server, and drains the socket so a server stalled on a slow consumer resumes.
    • A server with the capability reports the timeout itself.
    • An older server is sent CANCEL, and its STATUS_CANCELLED is reported as STATUS_QUERY_TIMEOUT, unless the user had cancelled first.
    • A complete END/EXEC_DONE that withheld nothing from the handler still reports success.
  3. With no terminal frame after one grace period, the caller is told about the timeout. The connection keeps draining the aborted query (the server runs one query per connection) for up to one more grace period.
  4. If still nothing arrives, the connection is unresponsive. It is closed (the I/O thread is stopped while its published batches are released, and a terminal failure is latched) and replaced.

The grace period is query_close_timeout_ms on the facade (default 5 s) and withQueryTimeoutGrace on the low-level client.

Failover. One deadline spans every attempt:

  • replays carry the remaining budget;
  • backoff and each reconnect step (TCP connect, upgrade, SERVER_INFO) are bounded by it, so a black-holed endpoint can't hold the caller past the timeout;
  • a reconnect that ends early is retried by the next execute().

Behaviour change: the last point also covers failover reconnects that fail for any non-auth reason. They previously left the client permanently "not connected" (every later execute() threw IllegalStateException), which is now reachable just by a short timeout cutting the reconnect off.

Pool.

  • Query.close() now waits, within query_close_timeout_ms, for a worker that is still draining a timed-out query before returning it to the pool.
  • It discards a worker whose connection failed instead of returning it.
  • A cancel issued while a submission waits behind such a drain used to be dropped, because execute() clears the client's cancel latch on entry. The handle now remembers the cancel and re-applies it before the request is sent.

The per-submit path allocates nothing: the deadline state is primitives, and QwpSpscQueue.take(deadline) uses parkNanos. Java 8 APIs only.

API notes: Query gains an abstract timeout(long, TimeUnit), which breaks compilation only for third-party implementations of the interface. The internal QwpEgressIoThread.submitQuery takes an extra timeoutMs argument.

Test plan

  • New QwpQueryClientQueryTimeoutTest (13 tests, scripted mock server):
    • the wire field: present, absent, and carrying the remaining budget;
    • a server-reported timeout, and the cancel sent to an older server;
    • batches discarded after the deadline, and success just past it;
    • the caller released after the grace period while the connection drains;
    • an unresponsive connection replaced, with failover on and off;
    • a user cancel taking precedence over the timeout;
    • the replay budget, and a timeout while failing over;
    • the reconnect walk bounded by the deadline, plus the lazy reconnect.
  • New QueryTimeoutFacadeTest (5 tests): the configured default and its per-handle override; the connection kept (one handshake); a draining worker returned only once idle; an unresponsive worker replaced; a cancel of a queued submission not lost.
  • Extended QwpSpscQueueTest (timed take), QueryCloseDrainTest (busy worker), QueryImplResetTest, QwpQueryClientConfigHonoredTest.
  • Mutation checks: disabling any one of these fails at least one new test: the discard, the deadline-bounded reconnect, user-cancel precedence, the draining phase, the idle wait, the discard of unhealthy workers, the queued-cancel re-apply.
  • Full core suite on JDK 8: 3499 tests, 0 failures.
  • mvn -P javadoc -DskipTests clean package passes on JDK 8 (with JAVA11_HOME) and JDK 25, with no new javadoc warnings.
  • The new timing-sensitive tests passed 5 runs out of 5 with all cores saturated.

Adds a query timeout to QwpQueryClient and the QuestDB facade:
query_timeout_ms (connection-string default), Query.timeout(...) per
handle, and an execute() overload with an explicit timeout. A timed-out
query ends with the new STATUS_QUERY_TIMEOUT (QueryException.isTimeout())
on a connection that stays open and authenticated for the next query.

- A server advertising the new CAP_QUERY_TIMEOUT (0x08) capability gets
  the remaining budget as a timeout_ms field after query_flags
  (QUERY_FLAG_TIMEOUT, 0x02) and ends the query itself with
  STATUS_QUERY_TIMEOUT (0x0E). Older servers are cancelled by the client
  at the deadline; the cancellation is reported as the same status.
- Once the deadline passes no further result batch reaches the handler.
  Discarded batches are still decoded, keeping the connection-scoped
  SYMBOL dict in step with the server.
- If the server does not end the query within the grace period
  (query_close_timeout_ms on the facade), the caller is released anyway
  and the connection drains the aborted query in the background. Only a
  connection that stays silent through a second grace period is closed.
- The deadline spans failover: replays carry the remaining budget, the
  backoff and every reconnect step are bounded by it, and a reconnect it
  cut short is retried by the next execute() instead of leaving the
  client disconnected.

Pool: Query.close() waits for a worker still draining a timed-out query
before returning it, discards a worker whose connection failed, and a
cancel issued while a submission waits behind such a drain is no longer
lost.

The server side (decoding timeout_ms, the breaker timeout, the status
mapping and the capability advert) lands in a tandem questdb PR.
A query that ended early on the client side left its remaining frames
on the connection, and the next query read them as its own result.
This fixes the three ways in.

Handler throws on a query timeout. At the end of the grace period
execute() reports the timeout through onError while the aborted query
still runs on the connection. If onError threw, execute() skipped the
drain that follows. FailoverProbeHandler now holds the exception,
executeOnce keeps draining as it does for a handler that returns, and
executeImpl rethrows the original throwable once the drain ends or the
connection is given up.

The facade worker reports an exception that escapes execute() as the
outcome of the current submission. Rethrown after the drain, it could
land after the caller, already released by the timeout, had submitted
again on the same Query, and complete that newer submission with the
old exception. QueryImpl now numbers submissions, and runOn() reports
an escaping exception only while its own submission is current. The
check and the report share doneLock with submit(). This also closes
the same race, until now microseconds wide, for a handler that throws
from onEnd or onExecDone.

Handler throws from onBatch, or the thread is interrupted. The I/O
thread now stamps every event of a query with that query's request id;
connection-level events carry QueryEvent.ANY_REQUEST and still reach
whichever query waits. executeOnce skips the events of other queries,
so the next query no longer receives the leftovers of an abandoned
one, and it handles its own deadline even while leftovers keep coming.
An attempt that ends before its query's last event now cancels that
query, so the leftovers stop coming.

While the I/O thread still works through an abandoned query, the next
query can be queued and cancelled. The I/O thread sent that CANCEL at
once, ahead of the query itself, and the server drops a cancel of a
query it does not know, so the cancel was lost. drainPendingCancel now
sends only a cancel of the query it serves, holds one of a later
query, and drops one of a query that has already ended.

Interrupt before the I/O thread encodes the request. execute()
returned while the I/O thread could still read the request holder,
the bind buffer and the caller's SQL text, and the next execute()
overwrote them: the interrupted query was never sent and the next one
was sent twice, possibly mixing the two queries' SQL and binds. With
leftovers of an abandoned query on the connection, the next execute()
could also block forever on the single request slot, while the I/O
thread waited for it to release a leftover batch. An interrupted
execute() now withdraws a request the I/O thread has not picked up
yet, or waits for the encoding the I/O thread has started, before it
reports the interrupt. An I/O thread still not done after
shutdownJoinMs gives up the connection, as on a query timeout.

Behavior visible to callers: an exception from onError at the end of
the grace period now surfaces when the drain ends, up to one more
grace period later, which is when execute() already returns for a
handler that does not throw. A handler that throws from onBatch, or
an interrupt, now cancels its query, and the next query skips that
query's leftovers instead of returning them. A handler that does not
throw sees no change, and the per-batch path gains only one field
store on the I/O thread and one comparison on the reading thread.

New tests drive each path against the scripted mock server, through
both QwpQueryClient and the QuestDB facade.
@mtopolnik

Copy link
Copy Markdown
Contributor

[PR Coverage check]

😍 pass : 306 / 356 (85.96%)

file detail

path covered line new line coverage
🔵 io/questdb/client/cutlass/qwp/client/QwpQueryClient.java 196 240 81.67%
🔵 io/questdb/client/impl/QueryWorker.java 17 20 85.00%
🔵 io/questdb/client/impl/QueryImpl.java 43 46 93.48%
🔵 io/questdb/client/cutlass/qwp/client/QwpSpscQueue.java 20 20 100.00%
🔵 io/questdb/client/QueryException.java 1 1 100.00%
🔵 io/questdb/client/impl/ConfigSchema.java 1 1 100.00%
🔵 io/questdb/client/cutlass/qwp/client/QueryEvent.java 10 10 100.00%
🔵 io/questdb/client/cutlass/qwp/client/QwpEgressIoThread.java 16 16 100.00%
🔵 io/questdb/client/impl/QueryLease.java 2 2 100.00%

This branch has not been deployed

No deployments
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.

2 participants