Skip to content

feat(qwp): recycle the symbol dictionary at a threshold to lift the cardinality cap - #91

Open
jovfer wants to merge 129 commits into
mainfrom
qwp-dict-recycle-r1
Open

jovfer wants to merge 129 commits into
mainfrom
qwp-dict-recycle-r1

Conversation

@jovfer

@jovfer jovfer commented Aug 19, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #80

Tandem:

Summary

A QWP WebSocket sender keeps a dictionary of every distinct symbol value it has sent, and that dictionary only ever grew. A long-lived sender that saw enough distinct values reached a hard cap, and from then on every new symbol value failed. A pooled sender made this worse: it went back to the pool still full, so the next borrower inherited the failure (#80).

The sender now recycles its dictionary. When the dictionary reaches a threshold (100,000 distinct symbols by default), the sender waits for a safe moment, starts a fresh dictionary on a new connection and carries on. No rows are lost, application code does not change, and the feature is on by default.

The change is client-only, with no wire-protocol change and no on-disk format change. It relies on the server clearing its per-connection dictionary state on disconnect, which QuestDB 10.0.0 and later already do. The tandem OSS PR adds a server test that keeps it that way.

How it works

  1. After each flush the sender checks the dictionary size. At or above the threshold, it arms a recycle.
  2. The recycle runs at the next table() call that finds every published frame acknowledged, no row in progress and the connection up.
  3. The sender then closes its send loop, starts a new dictionary and reconnects in the background. In store-and-forward mode it reuses the same slot, which is empty at that point. Rows written in the meantime are buffered as usual.

Because the swap happens only when nothing is unacknowledged, no frame built against the old dictionary can be replayed against the new one.

A producer that never pauses seldom has everything acknowledged at the start of a row. For that case, once a recycle has been armed for symbol_dict_reset_max_wait_millis (2 s by default), one table() call pauses for up to that long to let the backlog drain. This happens at most once per armed recycle, and on a healthy connection it takes about one acknowledgement round trip. If the backlog does not drain in time, the call returns normally and the recycle stays armed until a quieter moment.

Three more rules keep the recycle out of the way:

  • It never runs during an outage. While the sender is reconnecting, the recycle stays armed and runs once the connection is back. A recycle never surfaces a transport error or a reconnect timeout to the producer.
  • The bar rises after each recycle. If the live symbol set is larger than the threshold, recycling at the threshold would repeat forever. After each recycle the bar rises to twice the dictionary size at that recycle, up to 1,000,000. A static set below 1,000,000 therefore stops recycling after about log2(set size / threshold) recycles: four for 900,000 symbols at the default threshold. A set of 1,000,000 or more, or one whose values keep changing, recycles again at each 1,000,000 symbols the sender registers, and each of those recycles re-sends the symbol strings on a new connection. A set that stays below 2,000,000 avoids them with symbol_dict_reset=off. A manual resetSymbolDictionary() ignores the bar for that one recycle.
  • It never hides a terminal error. If the connection reports a terminal error, such as a security rejection, while a recycle is stopping it, the recycle does not run and the sender throws the error from that table() call or the next call, and from every later one, as on a sender that never recycled.

Configuration

All three settings apply to the WebSocket transport only.

Connect-string key Builder method Default Meaning
symbol_dict_reset symbolDictReset(boolean) on Enables recycling. off also turns resetSymbolDictionary() into a no-op.
symbol_dict_reset_threshold symbolDictResetThreshold(int) 100000 Distinct symbols that arm a recycle. Accepts 1 to 1,000,000.
symbol_dict_reset_max_wait_millis symbolDictResetMaxWaitMillis(long) 2000 Longest pause of one table() call for a recycle that has not found a quiet moment. 0 never pauses.

New API

  • Sender.resetSymbolDictionary() asks for a recycle regardless of dictionary size. The request is advisory: the recycle still runs at the next safe table() call. It is a no-op on transports without a symbol dictionary, and PooledSender forwards it to the sender it wraps. Call it on the producer thread, like every other Sender method.
  • QwpWebSocketSender gains monitoring getters: getSymbolDictEpoch() (recycles completed), isResetArmed() and getSymbolDictResetStarvationTimeouts() (pauses that ended without a recycle).
  • QwpWebSocketSender.setEngineRebuildFactory(...) lets a sender created through the low-level connect(...) overloads recycle. Senders from Sender.builder(), Sender.fromConfig() and SenderPool are set up automatically.
  • QwpWebSocketSender.connect(...), QwpWebSocketSender.connectWithCredentialSupplier(...) and the CursorWebSocketSendLoop constructor each gain an overload that carries the recycle settings. The two new sender overloads reject a threshold or a wait outside the builder's ranges. Every existing overload is kept. CursorWebSocketSendLoop.markEverConnected() is new.

Release notes

  • A sender no longer runs out of symbol capacity over its lifetime. The dictionary is recycled automatically at 100,000 distinct symbols; see Configuration to tune or disable it.
  • QuestDB server 10.0.0 or later is required. Earlier servers cap the dictionary at 1,000,000 entries, and the protocol does not negotiate the limit.
  • The client-side dictionary cap is now 2,000,000, up from 1,000,000 and equal to the server's. The error raised at the cap now points to the reset settings.
  • Frame sequence numbers keep increasing across recycles. flushAndGetSequence(), getAckedFsn(), awaitAckedFsn(), progress callbacks and SenderError ranges stay on one continuous scale for the life of a sender instance, and getAckedFsn() does not fall back to −1 during a recycle. The scale is not persisted: a process restarted on the same store-and-forward slot numbers from the slot's on-disk state.
  • Sender counters are lifetime totals. getTotalFramesSent(), getTotalAcks(), getTotalReconnectAttempts() and the related counters carry across recycles and remain readable after close(). A recycle counts as one reconnect attempt and one success.
  • A recycle's reconnect is reported without a disconnect. SenderConnectionListener receives RECONNECTED (or FAILED_OVER) with no preceding DISCONNECTED.
  • SenderErrorHandler can be called on the producer thread. If the sender finds a damaged store-and-forward slot while rebuilding during a recycle, it reports it synchronously from the table(), flush() or drain() call that ran the rebuild, as it already does at build time.

Limits

  • A recycle needs a quiet moment. A producer that continuously outruns the connection may never give it one, and the dictionary then keeps growing toward the 2,000,000 cap. The sign is isResetArmed() staying true while getSymbolDictEpoch() does not advance; the same reading during an outage only means the sender is between connections. Raise symbol_dict_reset_max_wait_millis, or call drain(...) at a quiet point. Moving the swap to the I/O thread, so that it no longer depends on a quiet moment, is the planned follow-up.
  • A long outage defers the recycle. If the dictionary reaches 2,000,000 entries before the connection returns, symbol() fails with the dictionary-full error.
  • An interrupted recycle resumes. If a recycle cannot finish for a transient reason, such as a thread interrupt, a stalled disk or a failed file open, the next table(), flush() or drain() call completes it; an at() or atNow() call made first fails and rolls back its row. While the disk under the slot is stalled mid-recycle, one call can block for up to 30 s and rows are refused until the disk recovers. The sender becomes unusable in one case only: the slot it has just emptied holds unacknowledged frames when reopened, which means the store-and-forward guarantee was broken.
  • Low-level connect(...) senders need two things to recycle: a rebuild factory, and in store-and-forward mode a slot directory whose parent can hold the slot's lock files. If the sender cannot take the slot's lock when a recycle starts, it logs one warning, stops recycling and keeps accepting rows.
  • Pooled senders hand an armed recycle to the next borrower. On a healthy connection the borrower's first table() call recycles without waiting. On an unhealthy one it may pay the pause while holding its lease, which is why the default wait is well below the pool's 5 s acquire timeout. The pooled path has no end-to-end test yet.

Design decisions

  • Recycle by reconnecting, not by a new protocol message. The server already drops its dictionary state when a connection closes, so the feature needs no server change.
  • Swap only when everything is acknowledged. This is what makes the recycle lossless without translating symbol ids on replay: nothing in flight refers to the old dictionary.
  • table() is the only trigger. The check sits at the start of a row, so the per-row hot path has no new branch. A caller that never starts another row never recycles.
  • The default wait is 2 s. Zero would leave the recycle unreachable for continuous producers, the very senders that reach the threshold. A long wait would let a pooled borrower hold its lease past the pool's acquire timeout and starve other borrowers.
  • The sender keeps its store-and-forward slot locked across the swap. Another sender started with the same sf_dir and sender_id fails fast, as it does at any other time, instead of taking the slot mid-recycle.
  • The sequence-number offset stays in memory. Persisting it would add to the on-disk format, which this change avoids.
  • Sender stays single-threaded. Calling close() from another thread during a recycle is outside the documented contract.

Tests

  • Client: new suites cover arming and thresholds, the swap in store-and-forward and memory modes, sequence-number continuity, the bounded wait, outages and reconnect races, a terminal error arriving while a recycle stops the connection, interrupted and resumed recycles, slot ownership, and backward compatibility of the public API.
  • OSS (test(qwp): pin the disconnect-clear contract, e2e and fuzz coverage for client symbol-dict recycle questdb#7525): a server test for the disconnect-clears-dictionary contract, real-client end-to-end tests in both modes, and a seeded fuzz test that interleaves recycles with server restarts.
  • Enterprise (questdb/questdb-enterprise#1171): continuous ingestion through a primary-to-replica failover with recycles in flight, asserting that no row is lost.

🤖 Generated with Claude Code

jovfer and others added 23 commits August 17, 2026 16:13
Pull the lock/quarantine engine-construction block out of build() into
LineSenderBuilder.constructEngineOnSlotLocked()/constructEngineOnSlot(),
and expose a QwpWebSocketSender.EngineRebuildFactory seam that build()
installs on the connected sender once connect() succeeds. Pure refactor,
zero behavior change: build() keeps its wide logical-lock scope spanning
the connect loop; the standalone constructEngineOnSlot() acquires the
lock only around construction, for a later symbol-dictionary epoch
rebuild to reuse the identical construct/quarantine code path.
constructEngineOnSlotLocked() now returns a small ConstructedEngine
result (engine + whether construction itself quarantined the slot)
instead of a bare CursorSendEngine. build() seeds its own quarantined
local from that verdict, restoring the pre-refactor invariant that a
construction-time quarantine counts toward the one quarantine per
build() attempt the connect loop's retry guard allows. Without this,
a construction-time quarantine followed by an UnreplayableSlotException
from connect() would take a second quarantineTornSlot pass instead of
the original close-and-rethrow.

constructEngineOnSlot(), the public factory entry Task 5 consumes,
keeps its 8-arg signature and CursorSendEngine return type: it just
unwraps ConstructedEngine.engine and discards the quarantined verdict,
since the recycle path latches terminal on connect failure rather than
quarantining.
Adds three connect-string keys for the upcoming QWP symbol-dictionary
recycle feature: symbol_dict_reset (on/off, default on),
symbol_dict_reset_threshold (distinct-symbol count that triggers a
recycle, default 100_000, bounded by QwpConstants.MAX_SYMBOL_DICTIONARY_SIZE),
and symbol_dict_reset_max_wait_millis (upper bound on how long a
triggered recycle waits for an opportunistic window before forcing,
default 30_000, 0 means opportunistic-only).

This is config plumbing only: the three values land on QwpWebSocketSender
as resetEnabled/resetThresholdSymbols/resetMaxWaitMillis fields with
@testonly getters, following the catch_up_cap_gap_min_escalation_window_millis
knob end-to-end (ConfigSchema registry, builder methods with WS-transport
guards, both connect-string parse paths, wsConfigSnapshotForTest). Actual
recycle behavior is a follow-up task.
Adds Sender.resetSymbolDictionary(), an advisory request to start a fresh
symbol-dictionary epoch (default no-op; QwpWebSocketSender overrides it).
QwpWebSocketSender.armIfEligible() re-evaluates arming at the tail of
resetTableBuffersAfterFlush() -- the shared exit point for the plain flush,
split flush, and close-path callers -- so the recycle arms once
symbol_dict_reset is enabled and either the global dictionary reaches
symbol_dict_reset_threshold distinct entries or a caller requested a manual
reset. Arming deliberately ignores deltaDictEnabled: a sender degraded to
full self-sufficient frames still benefits from bounding dictionary growth,
and a manual request is honoured regardless of mode.

Covered by SymbolDictRecycleArmingTest: threshold crossing, symbol_dict_reset=off
never arming, the manual advisory API (both immediate and mid-batch-deferred
arming), the split-flush path sharing the same arming tail, arming in
full-dict (degraded) mode, and the no-op default on a non-WebSocket sender.
testArmsInFullDictMode previously substituted the manual
resetSymbolDictionary() advisory request for crossing
symbol_dict_reset_threshold, which never touches globalSymbolDictionary and
so left a future deltaDictEnabled-conditioned regression in threshold-based
arming undetected. Rewrite it to construct the sender through the widest
connect(List<Endpoint>, ...) overload, which accepts both a custom
symbolDictResetThresholdSymbols and the fault-injecting CursorSendEngine, and
genuinely cross the threshold (registering a, b, c) while the sender stays
degraded to full self-sufficient frames. No manual reset call remains in the
test.
- rollFsnEpochBase now throws IllegalStateException when cursorSendLoop is
  non-null: the loop's externalFsnBase is a construction-time snapshot,
  never updated on a live loop, so rolling with a loop attached would
  silently desync sender-level FSN accessors from loop-level FSN emission.
- testPreRollTargetAnswersTrueAfterRoll rolled the same already-connected
  sender that produced fsn1, so its raw engine watermark never reset and
  the pre-fix comparison (ackedFsn() == fsn1 >= fsn1) was also true --
  the test could not fail. Rebuilt around a fresh rolled sender/engine via
  the existing createRolledSender helper, matching tests 3/5/6. Tests 2, 4,
  and 7 also rolled an already-connected sender (now rejected by the new
  guard) and are rewritten the same way.
Implement the table() barrier hook and recycleForDictReset() swap: once
the symbol-dictionary recycle is armed (Task 3) and the ring is proven
drained, table() tears the cursor I/O loop and engine down, rolls the
FSN epoch base (Task 4), replaces the producer's symbol dictionary, and
rebuilds the engine via the Task 1 EngineRebuildFactory before
reconnecting -- all synchronously inside a single call.

A failed rebuild latches recycleFailure as a terminal state: every frame
existed before the swap was already proven acked, so no data is at
risk, but the sender that observed the torn-down engine/loop refuses
further use. checkRecycleFailure() covers table() and the flush-family
entry points (flush, flushAndGetSequence, drain, awaitAckedFsn) --
deliberately not close(), which must still be able to tear down a
latched sender.
Fix round 1 from review:

- maybeRecycleForDictReset() now refuses before any teardown when
  engineRebuildFactory is null (every public connect() overload leaves
  it unset -- only Sender.build() installs one) or the cursor engine
  isn't owned by this sender (setCursorEngine(engine, false)'s
  contract). Without this, a connect()-built sender with the
  default-on recycle feature armed (via resetSymbolDictionary() or a
  threshold crossing) would reach step 6, NPE against the null
  factory, and latch itself terminal for no reason.
- Step 5 also resets lastCommitBoundaryFsn -- it held a raw old-epoch
  FSN that does not survive the roll.
- checkRecycleFailure() now also guards sendRow(), closing the
  fluent-chain corner where a caller continues .symbol(...).atNow()
  against a currentTableBuffer selected before the latch, without an
  intervening table() call.
- Step 6 rewires cursorEngine.setSlotLockReleaseListener(...) on the
  rebuilt engine, restoring the pool early-wakeup notification that
  bypassing setCursorEngine had dropped.
- testPostRecycleSlotContents now asserts the post-recycle slot
  directory listing against the exact expected fresh-state file set,
  replacing a bare Files.exists(...) check that proved nothing (the
  outgoing engine had a same-named file too).
- New testConnectBuiltSenderNeverRecyclesWithoutFactory covers both
  ways such a sender can arm (manual request, threshold crossing):
  neither may recycle, throw, or stop the sender from working.
Adds SymbolDictRecycleMemoryModeTest, the sf_dir-omitted counterpart of
SymbolDictRecycleTest: threshold-triggered recycle at an empty backlog,
a content oracle proving the epoch boundary loses (and duplicates)
nothing acked, and the same recycle under initial_connect_retry=async.

All three pass unmodified against the existing recycle swap -- zero
production changes -- confirming the factory's slotPath == null arm,
CursorSendEngine's file-less close, and the table() barrier are already
mode-agnostic.
Fills the maybeBlockForStarvedReset() stub: when a symbol-dict recycle
is armed but the ring is not yet drained, opportunistically waits
(parked, awaitAckedFsn-shaped) up to symbol_dict_reset_max_wait_millis
for the outstanding acks before giving up for this armed window.
resetMaxWaitMillis<=0 disables the wait entirely; at most one blocking
wait runs per armed window (starvationWaitDoneThisArm); an open
deferred-commit group is never waited on, since the server withholds
its acks by design until the closing commit lands and this producer
thread is the only one that could ever send that commit -- blocking
there would just run out the clock every time. A timeout increments
the new symbolDictResetStarvationTimeouts counter and leaves the
recycle armed so a later drained table() call can still fire it.

Verified the deferred-commit guard is load-bearing by temporarily
removing it and confirming the test fails (blocks the full deadline
instead of returning immediately) before restoring it.
Javadoc for symbol_dict_reset_max_wait_millis (Sender.java builder
method, and the DEFAULT_.../field comments in QwpWebSocketSender.java)
said the knob controls how long the recycle waits "before forcing the
rebuild" and that 0 means "never forces". That is backwards: nothing
is ever forced through. On timeout the wait gives up, increments the
starvation counter, logs a warning, and stays armed -- the dictionary
threshold is the only actual backstop. Rewrite all three to state the
real policy: once the armed window exceeds the knob, the NEXT
table(...) call may block the calling thread for up to the knob's
value waiting for the backlog to drain; 0 disables blocking entirely.

WARN log wording "re-arming opportunistically" -> "staying armed":
the arm is never consumed on a timeout, so nothing re-arms.

Test fixes in SymbolDictRecycleStarvationTest:
- testTimeoutLogsAndReArms: the "second table() must not re-block"
  probe was placed after a row had been queued (pendingRowCount==1),
  so it short-circuited on maybeRecycleForDictReset()'s precondition
  guard before ever reaching the wait -- vacuously true regardless of
  starvationWaitDoneThisArm. Moved the probe earlier, to a table()
  call with pendingRowCount==0, so it actually exercises the "at most
  one blocking wait per armed window" guard.
- testBlocksThenRecyclesWhenAcksArrive: raised maxWaitMillis from 400
  to 700ms (keeping the 150ms release delay) so the full recycle path
  (I/O loop join, engine close, rebuild, fresh handshake) has real
  headroom under the elapsedMs<maxWaitMillis assertion instead of ~250ms.
- Both testBlocksThenRecyclesWhenAcksArrive and
  testLatchedErrorDuringWaitThrows now join their helper thread
  (releaser/poisoner) in a finally block, so an assertion failure
  can no longer leak a non-daemon thread.

Re-ran SymbolDictRecycleStarvationTest (5/5) and the full
SymbolDictRecycle* battery (27/27), all green.
Adds SymbolDictRecycleOutageTest covering two interleavings between the
symbol-dictionary recycle swap and events outside the producer's own
control: a real connection outage on its own stream (recycle triggers
while the pre-recycle I/O thread is mid-reconnect against a killed
server; step 2's close() joins it, step 7's fresh connect recovers once
the endpoint accepts again), and a sibling orphan drainer mid-drain
(the recycle only tears down the foreground sender's own cursor
engine/I/O loop, leaving a concurrently-gated BackgroundDrainer
untouched and able to complete afterward).

The third scenario from the task brief -- OrphanScanner.isCandidateOrphan
rejecting an empty slot directory -- turned out to already be pinned by
OrphanScannerTest#testIsCandidateOrphanDirect and
#testEmptySlotDirIsNotAnOrphan, so no new test was added for it.

Test-only change; no production code touched.
SymbolDictRecycleCatchUpSkipTest pins that recycleForDictReset()'s step 7
reconnect never pays for a delta-dictionary catch-up frame: the rebuilt
engine sits on a freshly-emptied slot, so the new loop's sentDictCount
mirror seeds from PersistedSymbolDict.recoveredSize() == 0 and
setWireBaselineWithCatchUp's gate stays false for the whole first
post-recycle connection.

The core scenario chains the negative and positive observations in one
test so the zero count is provably a property, not a handler blind spot:
after the recycle sends zero zero-table frames and tiles ids from 0, the
handler force-drops the connection, and the resulting UNPLANNED reconnect
does catch up -- bounded to exactly the new epoch's symbols, never
replaying the retired epoch's. A follow-up symbol then ships with a delta
start above 0, pinning that resetSymbolDictStateForNewConnection() on the
plain-reconnect path preserves sentMaxSymbolId rather than folding the
recycle's baseline reset into itself.

A second test repeats the zero-catch-up observation under
initial_connect_retry=async, where step 7's reconnect funnels through
ensureConnected()'s ASYNC arm and the handshake completes on the I/O
thread.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
rev-t10 traced resetSymbolDictStateForNewConnection() to its single call
site in ensureConnected(), which the unplanned I/O-thread reconnect
(swapClient) never reaches -- so the old comment credited a function that
does not even run on this path. sentMaxSymbolId survives the plain
reconnect because nothing touches it there; only recycleForDictReset()'s
step 5 ever zeroes the baseline. The class javadoc's failure-mode
attribution carried the same imprecision and is reworded to match. The
assertion itself was traced sound (symbolDeltaBaseline() ->
encoder.beginMessage -> deltaStart) and is unchanged; comment-only diff.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
Adds SymbolDictRecycleCrashWindowsTest, pinning what a restarted sender
recovers if the process crashes at each of four points around
QwpWebSocketSender.recycleForDictReset()'s 8-step symbol-dictionary
recycle swap:

- (a) before step 2 (the barrier that starts the swap): the pre-recycle
  epoch's slot holds a fully-acked batch on disk. Recovery must find
  that residue, recognize it as already acked (nothing to replay) and
  resume the SAME dictionary rather than starting fresh. Constructed by
  closing fast against a server that never acks (so the fully-drained
  unlink never fires) and then stamping the ack watermark directly to
  declare the batch acked retroactively, mirroring
  DeltaDictRecoveryTest#writeAckWatermark.

- (b) between step 3 (fully-drained close of the old engine) and step 6
  (rebuild): the slot is empty. Constructed by driving a real recycle to
  completion and then closing immediately, before any flush touches the
  freshly-rebuilt engine -- finishClose treats "nothing published yet"
  as fully drained too, so this unlinks everything step 6 just created,
  leaving the same empty state step 3 alone would have left.

- (c) after step 7 (reconnect), before the new epoch's first flush: the
  slot holds a freshly-rebuilt engine's own state files but no data.
  Constructed by snapshotting the rebuilt slot's bytes before closing
  (there is no supported way to release just the slot's OS flock without
  running finishClose's unlink), closing for real so nothing leaks, then
  restoring the snapshot on top of the vacated directory.

- (d) an ordinary mid-operation crash one epoch into the post-recycle
  steady state, to prove the epoch swap does not corrupt normal backlog
  recovery. Uses the same close-fast-against-a-non-acking-server idiom
  as RecoveryReplayTest, but only after the recycle's fresh connection
  is established.

Arms (b) and (c) both replay nothing and both look "empty" at first
glance, but they are not the same recoverable state and a restarted
engine can tell them apart: (c)'s slot carries a manifest with collapsed
boundaries alongside a same-based, zero-frame active segment, which
SegmentRing.recover()'s chain-building accepts as a RECOVERED (if empty)
chain, while (b)'s slot carries no engine state at all and recovers as
EMPTY. wasRecoveredFromDisk() is the pinned, distinguishing observable
between the two, asserted explicitly instead of writing two
assertion-for-assertion duplicate tests.

Every arm's oracle: the recovered sender keeps ingesting; the symbols it
and its predecessor registered are exactly and correctly reconstructable
from the wire (each fresh server handler rebuilds the per-connection
delta dictionary); and no data frame is delivered more than the
at-least-once contract allows (each handler counts data frames so a
spurious re-send would show up as an unexpected count).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
Applies review findings 1, 2, 3, 4, 6, 7, 8 and 10 from the task-11
review (findings 5, 9, 11 are deferred to the whole-branch review):

- Arm (b) never actually pinned the "empty slot" disk image it exists
  to test -- it inferred emptiness from wasRecoveredFromDisk()==false
  on the successor instead of asserting the crashed sender's own slot
  dir. Adds an explicit listDir(slot) == [".lock", ".lock.pid"]
  assertion right after the crashed sender closes (symmetric to arm
  (c)'s pre-existing file-set assertion), plus a second one after the
  successor's own fully-drained close, proving the empty-slot state is
  stable rather than a one-shot coincidence. This also replaces arm
  (b)'s closing assertion, whose message previously credited the
  successor's "recovery" for cleaning up a stale manifest that the
  crashed sender's own close had already removed.

- Arm (c) snapshotted the freshly-rebuilt slot before waiting for the
  manager worker's asynchronous hot-spare provisioning to settle, so a
  mid-provision snapshot could race-capture a zero-magic spare that
  recovery would hard-fail on. Adds a bounded poll
  (awaitExactFileSet) before the snapshot, reusing the same expected
  file list (now a shared FRESH_REBUILD_FILES constant) as the
  post-restore assertion so the two can never drift apart.

- Arms (a) and (d) loosely asserted recoveredMaxSymbolId() >= 1, which
  a leaked epoch-0-plus-epoch-1 dictionary would also satisfy. Tightens
  both to the deterministic exact value (1L), keeping their explanatory
  messages as-is now that the assertion actually proves what the
  message claims.

- Widens QwpWireTestUtils.tableCount to public (its sibling frame
  helpers already are) and deletes this suite's local copy plus its
  now-unnecessary justification javadoc.

- Inlines the AckAllHandler bindings in arms (b) and (c) that were
  never read (dict()/dataFrameCount() were only meaningful on the
  "fresh" server's handler, not the "crashed" one).

- Class javadoc: notes that arm (c)'s snapshot/restore does not cover
  the logical slot lock (lives outside the slot dir; sender.close()
  reclaims it, acquireLogical recreates it -- benign), and corrects
  three mechanism imprecisions the review traced: the "never published"
  fully-drained check lives in close(boolean), not finishClose; a
  fully-drained close does not remove .lock/.lock.pid; and arm (b)'s
  SegmentRing.recover() counterpart is the no-manifest fall-through to
  Recovery.empty(), not the manifest-present collapse branch.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
A sender that degrades to full self-sufficient frames after a symbol-dict
persistence fault (disableDeltaDict) used to stay degraded for the rest of
its life -- deltaDictEnabled was set once at engine construction and never
re-evaluated. The symbol-dictionary recycle already rebuilds the cursor
engine from scratch on every swap; make that rebuild also re-derive
deltaDictEnabled from the fresh engine instead of carrying the old one's
verdict forward. If the underlying fault has cleared, the next recycle
heals the sender back into delta mode; if it has not, the fresh engine
degrades again on its own first append, the same ordinary catchable
LineSenderException as any other persistence fault -- a degrade, never a
latched recycleFailure terminal state.

Promote the epoch and starvation-timeout counters from @testonly
accessors to permanent public API: getSymbolDictEpochForTest() becomes
getSymbolDictEpoch(), and getSymbolDictResetStarvationTimeoutsForTest()
becomes getSymbolDictResetStarvationTimeouts(). Both counters already
existed; this only changes their visibility and documents their
thread-safety contract (producer-thread-written, so a read from another
thread is an eventually-consistent snapshot). Add a third counter,
getSymbolDictResetsPerformed(), incremented alongside the epoch inside
recycleForDictReset() -- the two move together today but are defined and
incremented independently, since a future change could roll the epoch by
some path other than a completed recycle swap.

Migrate every existing call site of the two renamed getters (7 test
files, 48 + 15 occurrences) to the new public names.

SymbolDictRecycleHealingTest covers all of this: a fault-then-heal
recycle proves deltaDictEnabled and wire framing both return to proper
delta encoding (the second post-recycle frame's delta starts where the
first left off and carries only the newly-added symbol, the shape only
delta mode produces); a fault-persists variant proves the fresh engine
degrades again without escaping as a raw error or latching the sender
terminal; and a two-recycle run asserts the three metrics getters track
correctly in lockstep while the starvation counter stays untouched.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
Applies task-12-review.md findings Q1-Q6.

Q1: the three symbol-dictionary-recycle metrics fields (symbolDictEpoch,
symbolDictResetsPerformed, symbolDictResetStarvationTimeouts) become
volatile -- they are permanent public API now, and their obvious reader
is a monitoring thread on some other thread, unlike every other
symbol-dictionary-recycle field they were modelled on. A plain,
non-volatile long gives a cross-thread reader no visibility guarantee at
all under the JMM (a polling loop can legally observe 0 forever), whereas
the javadoc claimed an "eventually-consistent snapshot". Reworded the
three javadocs to state the guarantee volatile actually gives: each read
sees the latest write the producer thread completed, with no atomicity
across the three counters -- a concurrent reader can see the epoch
already advanced while resets-performed still reflects the prior value,
even though the producer thread writes them on adjacent lines.

Q3: getSymbolDictResetStarvationTimeouts()'s javadoc named
recycleForDictReset() as the writer by importing getSymbolDictEpoch()'s
caveat sentence verbatim; the starvation counter is actually written in
maybeBlockForStarvedReset(). Named that method directly instead.

Q4: getSymbolDictResetsPerformed() said swaps this sender has
"completed" without stating that the counter advances at step 5, before
the engine rebuild (step 6) and reconnect (step 7) -- so a recycle that
later latches recycleFailure at step 6/7 still counts. Made that
explicit, mirroring the precision getSymbolDictEpoch() already had.

Q2: SymbolDictRecycleHealingTest's healing test asserted
isDeltaDictEnabledForTest() right after the recycle and attributed the
result to the healed facade, but the same assertion holds unconditionally
-- a fresh engine's construction never touches mmap, so it reports true
whether or not the facade was healed (the persistent-fault sibling test
proves this directly). Reworded the message to state what that assertion
actually pins, and added the discriminating check: re-assert after the
first post-recycle flush, the fresh engine's first real append -- a
still-armed facade would degrade it there, so staying true is real
evidence of healing.

Q5: AckAllHandler never reset its ack sequence per connection, unlike
CapturingAckHandler right below it in the same file. A rebuilt engine
restarts its raw FSNs at 0 after a recycle, so the stale, unreset
sequence could satisfy awaitAckedFsn with an ack for a frame that was
never published on the new connection -- a weaker gate than it looks,
even though the specific sequences in these tests happened not to
collide. Reset nextSeq on every new connection, matching the sibling
handler and existing repo precedent for connection-aware ack sequencing.

Q6: both fault-injection tests caught LineSenderException with a comment
claiming parity with MmapFaultDegradesTest's guard, but never actually
asserted the message MmapFaultDegradesTest checks
("failed to persist symbol dictionary before publish") -- so the catch
could equally have swallowed an unrelated connection-level exception.
Added the same message assertion to all three catch sites (two in the
persistent-fault test, one in the healing test) so the comment's claim
now holds.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
The symbol-dictionary cap error told callers to close the sender and
build a new one, but gave no way to avoid hitting the cap in the first
place. Append a sentence pointing at the automatic dictionary reset
knobs (symbol_dict_reset, symbol_dict_reset_threshold) and the manual
Sender.resetSymbolDictionary() escape hatch, so the error message
matches the recycle feature this client now ships.

Raise QwpConstants.MAX_SYMBOL_DICTIONARY_SIZE from 1_000_000 to
2_000_000 to mirror the server-side constant of the same name, which
moved to 2_000_000 in questdb OSS commit 306062e243 (#7468). That
commit is contained in release tag 10.0.0, the support floor for this
client, so client <= server holds across the whole supported fleet.

Update DeltaDictCeilingTest and GlobalSymbolDictionaryTest, which
pinned the old 1,000,000 value and message text, to the new cap and
add a case that drives the dictionary to the cap with automatic reset
disabled (symbol_dict_reset=off) and a threshold configured at the
cap, confirming the refusal still fires and still names the reset
valve even on a sender that has it switched off.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
The new cap-reached-while-armed test never calls sender.flush(): the
2M fill goes through the raw GlobalSymbolDictionary test accessor, and
the one Sender-routed call throws inside symbol() before a row
completes. armIfEligible() only runs from the tail of a completed
flush(), so it never had a code path available to flip isResetArmed()
to true regardless of whether symbol_dict_reset was on or off -- the
two assertFalse(ws.isResetArmed()) checks passed identically either
way and proved nothing about the knob under test.

Drop both assertions and note in the test's javadoc that arming
semantics are out of scope here and are pinned instead by
SymbolDictRecycleArmingTest.testArmsAtThreshold, which drives real
rows through flush() and is the idiom that actually exercises
armIfEligible(). The test's brief-mandated assertion -- the cap
refusal still fires, with the new reset-valve message, when reset is
disabled -- is untouched and remains the load-bearing check here.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
recycleForDictReset() no longer latches the sender terminal when the
recycle's reconnect fails. Step 7 moves out of the latching try block:
steps 1-6 still latch recycleFailure (a half-swapped sender genuinely
cannot make progress), but a failed ensureConnected() logs a warning and
rethrows to the triggering caller without latching. By step 7 the swap has
committed, so the sender is coherent - connected == false, loop and client
already closed and nulled by ensureConnected's own catch, the fresh engine
attached, the step-5 epoch and swap counters correctly left incremented -
and the ordinary sendRow() -> ensureConnected() path retries the connect,
and only the connect, on the next send. Nothing re-runs a teardown step,
and nothing can fire a second swap meanwhile: the fresh dictionary sits
below the threshold, manualResetRequested was consumed at step 5, and
maybeRecycleForDictReset requires connected.

This removes a default-configuration brick. A sender built with no
reconnect_* knob resolves initialConnectMode to OFF, so step 7 is a
single-shot connect; a server restart or an LB blip across a drained,
armed sender therefore latched every later table()/flush() call forever,
including seconds later once the endpoint was back.

Deferring the connect to the next send exposed a second defect, which this
commit also fixes. ensureConnected() calls
resetSymbolDictStateForNewConnection(), which cleared
currentBatchMaxSymbolId unconditionally. That watermark is batch-scoped,
not connection-scoped - a flush ships exactly [sentMaxSymbolId+1 ..
currentBatchMaxSymbolId] - and clearing it was harmless only because
build() connects before the application can register a symbol. On the
deferred path symbol() runs first, so the clear made the next flush ship
an empty delta while its rows referenced symbol id 0: rows on the wire
pointing at ids the server never received. The reset now runs only from
the drained state the old code assumed (no pending rows, no row in
progress).

testDefaultConfigRecycleSurvivesFailedReconnect pins both. It drives a
default-config sender (no reconnect knobs), kills the listener at a
drained instant, asserts the recycle's throw reaches the caller, asserts
recycleFailure is not latched by ingesting successfully once the endpoint
returns on the same port, and asserts the recovered stream defines every
symbol its rows reference. Against the pre-fix code it fails on the latch;
with the latch fixed but the watermark clear restored it fails on the
empty dictionary.

Documentation and hygiene alongside:

- getTotalFramesReplayed/getTotalFramesSent/getTotalReconnectAttempts/
  getTotalReconnectsSucceeded/getTotalServerErrors now state that they
  read the live send loop and therefore restart at 0 on every recycle
  ("since the last recycle"), and point at the lifetime-scoped
  getSymbolDictEpoch/getSymbolDictResetsPerformed for correlation.
- Sender.resetSymbolDictionary(), its QwpWebSocketSender override and
  LineSenderBuilder.symbolDictReset() now say that the manual valve is a
  permanent no-op while symbol_dict_reset is off, because armIfEligible
  gates on that knob.
- armIfEligible's javadoc names both call sites; resetSymbolDictionary()
  calls it too, and the stale single-call-site claim invited a wrong
  inlining refactor.
- The builder's symbolDictReset default references
  DEFAULT_SYMBOL_DICT_RESET_ENABLED instead of a hardcoded true.
- SymbolDictRecycleOutageTest joins its trigger thread in a finally so an
  assert failure cannot leave a thread inside the sender, and records why
  the revived server's handshake count stays >= 1 rather than == 1.
- SymbolDictRecycleCrashWindowsTest inlines an unread handler binding, and
  SymbolDictRecycleHealingTest's AckAllHandler carries a warning that its
  per-connection ack reset assumes every connection change is a recycle
  and must not be copied into a plain-reconnect test.

Client suite: 3045 run, 0 failures, 0 errors, 2 skipped.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01K2Hw6TfgX9sBB2L8oSdaWn
jovfer and others added 4 commits August 19, 2026 13:46
Review C1: recycleForDictReset step 3 re-acquired the slot flock without
waiting for a deferred close to release it. When the SF worker is wedged
in a syscall past SegmentManager's bounded join, CursorSendEngine.close()
returns with the flock retained and isCloseCompleted() false, releasing
both from the worker's exit path; the step-6 rebuild then threw
SlotLockContentionException on the retained flock and latched the sender
permanently terminal -- converting exactly the transient disk stall the
deferred-close machinery exists to survive into a hard sender death on
the default (SF + recycle-on) path.

recycleForDictReset now mirrors close()'s deferred-close discipline:
awaitDeferredEngineClose parks (awaitAckedFsn-shaped) until the deferred
cleanup confirms the flock release, re-arming the shared flock-release
retry driver each pass like isSlotLockReleased() does. Only exhausting
the 30 s budget -- a genuinely dead worker -- latches terminal, and that
path hands the still-locked engine to retainedEngine so a pool re-probe
recovers the slot's capacity if the worker ever exits; close() no longer
clobbers slotLockReleased to true while such an engine is pending.
SymbolDictRecycleDeferredCloseTest pins both branches with a wedged-
worker harness (red-proofed: without the await, the survival test dies
with the exact SlotLockContentionException the review traced).

Review Mo1: getAckedFsn and awaitAckedFsn snapshot cursorEngine into a
local (the recycle transitions it non-null -> null -> non-null), and
while it is null they report the durable watermark the recycle barrier
proved (new lastRecycleDurableFsn) instead of collapsing to -1.

Review Mo2: PooledSender forwards resetSymbolDictionary() to the live
delegate instead of inheriting the interface's no-op default; pinned by
a pooled-path test.

Review Mo5 (verified: released 9.4.x servers cap the dictionary at 1M):
document the pre-10.0.0 compatibility constraint on the 2M client cap at
QwpConstants.MAX_SYMBOL_DICTIONARY_SIZE and symbolDictReset(boolean).

Review mi7/mi8: step 6 refuses a rebuilt engine that recovered from disk
(the empties-the-slot contract was breached) and resets slotLockReleased
after a successful swap. Review mi3: recycle javadoc now says seven
steps / steps 2-6, the resetArmed and resetSymbolDictionary docs match
the code, and EngineRebuildFactory moves out of the field block.

Review Mo4: SymbolDictRecycleFsnContinuityTest wraps every test in
assertMemoryLeak, and LineSenderBuilderWebSocketTest pins the config
boundaries: threshold == 2M accepted on both config paths, the
symbol_dict_reset=on parse branch, the invalid-value message, and the
three fluent-setter transport guards.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
getAckedFsn()/awaitAckedFsn() read cursorEngine, fsnEpochBase and
lastRecycleDurableFsn from monitoring threads while the recycle
reassigns all three on the producer thread. The sibling observability
counters went volatile in an earlier round for exactly this reader;
these three carry the same contract, so a monitor could observe a
fresh engine with a stale epoch base and report an FSN dip to -1.
Also document resetSymbolDictionary() as producer-thread-only: it
mutates unsynchronized producer state, and the fire-and-forget javadoc
framing invited cross-thread calls (notably via PooledSender).
Red first: the outage tests now assert the store-and-forward contract
(producer never sees a transport error nor a reconnect budget after the
initial connect) instead of pinning the step-7 foreground connect they
used to. The production change lands in the next commit.
Step 7 of the symbol-dict recycle re-ran the initial-connect policy on
the producer thread: OFF senders got a single-shot connect whose failure
made every subsequent send throw before buffering until the endpoint
returned, and SYNC senders blocked the producer up to
reconnect_max_duration_millis. Both violate the store-and-forward
contract (post-init, the client never exposes transport problems and
never imposes a reconnect budget on the producer).

ensureConnected() now latches hasConnectedOnce on its first completion
and routes every later entry through the existing ASYNC (deferred)
branch: the loop is built and started synchronously, the socket connect
happens on the I/O thread with indefinite retry, and the producer keeps
buffering into the fresh epoch's slot. This is the same path ASYNC-mode
senders already took at step 7. The server clears its per-connection
dictionary on disconnect, so the buffered fresh-epoch frames
(deltaStart=0) replay correctly on reconnect.
jovfer and others added 9 commits September 25, 2026 03:58
The recycle's step-3 fully-drained close released the slot's directory
flock and unlinked its parent-anchored logical lock, and step 4 took both
again, so a colliding build() with the same sf_dir and sender_id could
land in between. The incumbent's rebuild then failed with contention on
every later row start for as long as the intruder lived, and latched
terminal if the intruder exited with unacknowledged frames -- against the
senderId contract that the second sender fails fast and the first is
untouched.

Take the logical lock before tearing the outgoing stack down, close the
outgoing and any healing engine without reclaiming it, rebuild through
constructEngineOnSlotLocked under it, and release it once the fresh engine
holds the directory flock -- also on the breach latch and in close(), so
an abandoned swap never keeps the slot from a successor.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…f-locking

The recycle previously delegated taking the parent-anchored logical lock
to a new pair of EngineRebuildFactory methods (acquireLogicalSlotLock /
rebuildLocked). That broke SymbolDictRecycleDeferredCloseTest's mid-recycle
factory swap: its stand-in lambda fell back to the real factory's
rebuild(), which self-locked through constructEngineOnSlot and contended
with the lock the recycle already held.

The sender now takes the logical lock itself in recycleForDictReset's
step 0, from the live engine's slot directory (CursorSendEngine.sfDir()),
instead of asking the factory for it. Every factory's rebuild(...)
constructs under that lock without taking it again: Sender.build()'s
factory now calls constructEngineOnSlotLocked directly, so the unlocked
constructEngineOnSlot entry point has no callers left and is deleted.
EngineRebuildFactory keeps its original two methods; only its javadoc
documents the lock contract.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
A colliding build at the rebuild boundary, while a rebuild is pending,
after the breach latch and after close(); and a contended step-0 lock
acquisition that refuses one row and tears nothing down. Each test is red
against the mutation of the production line it guards.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…ing check

Window 2's "must not wait" assertion raced the post-recycle reconnect,
which runs asynchronously on the I/O thread: maybeRecycleForDictReset()
refuses to touch the age guard while the loop is not yet linked, so the
timing check was passing for an unrelated reason and never exercised
armedSinceNanos or the age guard it was meant to pin. Publish a row and
await its ack first, proving the fresh loop is linked, before gating
acks and taking the timing measurement.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…ry poll

The wait snapshotted the engine once and read the epoch base afterwards,
so a monitor thread could pair the outgoing engine with the rolled base
(returning true for an FSN the fresh epoch had not published) or keep
polling the closed engine and time out. Poll through getAckedFsn's
base-first clamped read instead; pinned by a monitor parked across a real
recycle.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The epoch offset a recycle applies is memory-only, so a process restarted
on the same slot numbers from the slot's on-disk state. Say so where users
read FSN contracts: flushAndGetSequence, getAckedFsn and SenderError.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…cle lock before the owned engine's close

checkConnectionError read the volatile cursorSendLoop twice, so a monitor
thread polling awaitAckedFsn could NPE when the recycle nulls the loop.
closeRemainingResources now releases the recycle's logical lock before the
owned engine's reclaiming close, so an abandoned swap followed by close()
no longer leaves the .slot-locks pair behind. Contract docs on the engine
close and SlotLock now name the recycle as a logical-lock holder.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@jovfer

jovfer commented Sep 25, 2026

Copy link
Copy Markdown
Contributor Author

Tandem review: client #91 / OSS #7525 / ENT #1171 (level 3)

Reviewed heads: client 0496f037, OSS 186cd29978, ENT f0e373da3.

Scope: both submodule pointers point at unmerged feature-branch commits, so both are reviewed as part of this change. The OSS and ENT test files are unchanged since the previous review. The new code under review is client-side: the recycle now holds the slot's logical lock for the whole swap, awaitAckedFsn re-reads the acked position on every poll, and the docs now say FSNs are scoped to one sender instance.

Critical

None.

Moderate

[client] Recycle refuses every row forever when the slot's logical lock can never be taken. Location: QwpWebSocketSender.java:5870-5881.

  • Problem: step 0 treats a failure that will always repeat as brief contention, so the sender stays armed and refuses rows.
  • Net impact: a sender built with direct connect() and a custom rebuild factory stops ingesting once the recycle arms.
  • Evidence:
    • A probe test at 0496f037: 20 of 20 rows refused, 0 frames sent.
    • The same probe at the previous head ba54a138: 0 refused, the recycle commits.
    • Control run at 0496f037 with a normal slot path that has a parent directory: the recycle commits.

Step 0 turns any non-Error throw from SlotLock.acquireLogical(cursorEngine.sfDir()) into "recycle deferred … retried at a later row start" and leaves the sender armed. Two cases make that throw repeat on every attempt:

  • The slot path has no parent directory, for example a relative "sfslot". resolveLogicalLock returns null, so acquireLogical throws IllegalArgumentException (SlotLock.java:148-150).
  • The parent directory is not writable. acquireLogical throws SfOperationalException: could not create logical slot lock dir.

The refused row never starts, so the ring stays drained and every later table() goes straight back to step 0. An empty flush() does not re-run armIfEligible(), and setEngineRebuildFactory(null) does not disarm the sender either. On the previous head, step 3's close(true) → removeOrphanLogical silently did nothing for these paths, and the custom factory rebuilt the engine and the recycle committed.

How it is reached: a caller uses the public connect(...) with its own CursorSendEngine and installs a factory with setEngineRebuildFactory, as the new javadoc on all four connect overloads tells them to. build() is not affected: it builds slotPath as sfDir + "/" + senderId and takes acquireLogical itself at startup, so a bad parent directory fails the build instead.

Why this is Moderate rather than Critical: EngineRebuildFactory is new in this PR (not in 1.3.9, not at the base), and the failure also needs a slot path that build() never produces. No caller in any of the three repos does this.

Suggested fix:

  • Skip step 0 when resolveLogicalLock would return null, the same way removeOrphanLogical silently does nothing for such a path.
  • Keep the refuse-this-row-and-retry-later path only for SlotLockContentionException. For any other failure, disarm the recycle or latch recycleFailure instead of refusing rows forever.
  • Add a direct-connect test with a slot path that has no parent directory.

Minor

None.

Coverage

The test gate passes, with no coverage gaps admitted.

A mutation run over 232 recycle-related tests caught three changes:

  • removing the step-0 lock;
  • removing the lock release when the recycle breach terminates the sender;
  • reverting awaitAckedFsn to a single engine snapshot.

These mutants survived, but none changes user-visible behavior beyond leftover lock files:

  • keeping the lock after the recycle commits;
  • reordering the release in close();
  • closing the outgoing engine with close() instead of close(false). This behaves identically, because removeOrphanLogical takes the lock before it unlinks anything.
  • dropping the in-loop error check in awaitAckedFsn.

Verdict

  • client: approve with comments. The Moderate above should be fixed in this PR.
  • OSS: approve.
  • ENT: approve.

At review time CI was green on the client and ENT PRs. The OSS macwin run was still pending.

🤖 Generated with Claude Code

@jovfer jovfer added the READY label Sep 25, 2026
The recycle checked isLinkUp() at the barrier and stopped the loop at
step 2. A drop landing in between -- typically a server that acks the
backlog and then closes, which is also what ends the starvation wait --
let the I/O thread start its reconnect while step 0 was taking the slot
lock. close() then waited out the 30 s shutdown budget on a credential
pull or hostname resolve it cannot cancel, threw, and left a CLOSE_LOOP
resume that waited it out again on every table(), at() and flush(), so
no row was accepted until the pull returned.

connectLoop now takes a gate before anything that can block, and the
recycle stops the loop through the new closeIfLinkUp(), which claims the
same gate, so exactly one side wins. If the recycle wins, the I/O thread
never starts the walk and close() only cancels a socket. If the walk has
begun, the recycle gives back the lock it took at step 0 and stays
armed, the outcome the barrier already gives a loop it sees
reconnecting. isLinkUp() reads the gate too, so the barrier's pre-check
no longer reads a paced reconnect's pre-attempt park as a live link.

The gate is a volatile int behind a static field updater: no per-loop
allocation, and a bare allocateInstance loop starts OPEN. Every other
loop owner never claims it, so their behavior is unchanged.

Tests: CursorWebSocketSendLoopConnectGateTest pins both outcomes of the
gate deterministically; SymbolDictRecycleOutageTest reproduces the race
end to end (ack plus GOING_AWAY during the wait, uninterruptible pull)
and fails on the unfixed code. TestWebSocketServer gains
sendBinaryThenClose() to flush a frame and a close in one write.
A recycle whose deferred-close wait runs out -- an SF worker stalled
past the 35 s budget -- keeps the outgoing engine as retainedEngine.
Step 3 detaches that engine's release listener, and nothing re-attached
it. A pooled sender closed during the stall retires its slot, and when
the worker later exited, the pool was never told the slot was free: a
borrower parked on the full pool waited out its acquire timeout, or the
next housekeeper tick, instead of getting the slot at once. No borrow
failed and no capacity was lost; the re-probe still recovered the slot.

closeRemainingResources now re-attaches onSlotLockReleased to a retained
engine. Step 3 still detaches it while the sender is live, where a late
release could mark a rebuilt engine's flock as released. A closed sender
never rebuilds, so the retained engine is then the only possible holder
and its release is authoritative. The engine runs a listener registered
after its release at once, so a release landing just before the attach
is relayed too. The attach is guarded so a failing pool callback cannot
skip the dispatcher closes. The heal and breach paths retain their
engines the same way and go through the same branch.

Tests: SymbolDictRecycleDeferredCloseTest pins both orders of the
release and the close; SenderPoolSfTest parks a borrower on the full
pool and requires the slot within 10 s of the release, against a 60 s
acquire timeout. All three fail on the unfixed code.

@bluestreak01 bluestreak01 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@jovfer

Reviewing PR #91 at level 3, the full review: change-surface map, six discovery reviewers, and independent falsifiers with runtime probes at head and base.

Scope. No submodule pointer moved; core/src/main/c/share/zstd is unchanged, so no provenance check was needed. The PR metadata is in order: the title follows Conventional Commits, Fixes #80 is at the top, and the labels fit. One PR-body item is listed under Minor.

Critical

1. On some slot paths the recycle permanently stops a connect()-built sender from accepting rows

  • Problem: The recycle wedges the sender when the slot's logical lock cannot exist.
  • Net impact: connect()+factory senders on such paths stop ingesting after 100k symbols.
  • Evidence: In a probe at 17f04409, 24 of 25 rows were refused. The same app at c9968f24 accepted 25 of 25.

Where (in-diff): QwpWebSocketSender.java:5915-5928, the step-0 code in recycleForDictReset.

Trigger: two slot-path variants reproduce it.

  • Relative path: new CursorSendEngine("sf-data", …) uses a single-segment path with no parent.
  • Read-only parent: an absolute slot directory whose parent the process cannot write and which has no .slot-locks directory, e.g. a container volume mounted at /data and used directly as the slot.

In either case, the engine is passed to the public QwpWebSocketSender.connect(...) and setEngineRebuildFactory(() -> new CursorSendEngine(sameDir, …)) is installed. Every connect overload's javadoc says the factory is required for recycling (:1055, :1108, :1164, :1221).

Path:

  1. After 100k distinct symbols the recycle arms.
  2. At the first table() call that finds the backlog acknowledged and the link up, step 0 calls SlotLock.acquireLogical(sfDir).
  3. That call throws either IllegalArgumentException("slotDir must contain a parent and slot name") (SlotLock.java:148-150) or SfOperationalException("could not create logical slot lock dir") (SlotLock.java:305).
  4. Step 0 wraps the exception as "symbol dictionary recycle deferred … retried at a later row start" and leaves the sender armed.
  5. The refused row publishes nothing, so the next table() finds the backlog acknowledged again and fails the same way. sfDir is final, so this never clears.

Symptom: while the link is up, every row is refused. getSymbolDictEpoch() stays at 0 and isResetArmed() stays true. The error message calls the condition transient, but it is permanent.

Runtime evidence:

  • Probe: the public 9-arg connect, default recycle settings, 100k distinct symbols.
  • At head, both path variants accepted 1 row and refused 24 over 9 s, including an idle period longer than the 2 s max wait. The server received 1 new frame.
  • Controls at head:
    • With a writable parent the recycle commits and 25 of 25 rows are accepted.
    • Without a factory the sender never arms.

Base, same trigger: N/A, this is a new surface, since setEngineRebuildFactory does not exist at c9968f24. The identical probe minus that line accepts 25 of 25 rows there.

Why build() is unaffected: build() always uses sfDir/senderId and takes the same lock at startup (Sender.java:1664), so a path like this fails build() before any data is sent.

Net determination:

  • Population: users of the public connect(..., CursorSendEngine) API who enable recycling as documented, on a slot directory whose parent is missing or cannot hold .slot-locks.
  • Delta and magnitude: ingestion stops for good at the first acknowledged-backlog row after the threshold; at base it never stops.
  • Offsets: none documented. setEngineRebuildFactory(null) explicitly does not cancel an armed recycle.
  • Net: net-negative.

Fix:

  • At step 0, defer only on SlotLockContentionException.
  • When the logical lock cannot exist or cannot be created for this path, proceed without it. removeOrphanLogical already treats such paths as a silent no-op (SlotLock.java:184-188), and a colliding build() cannot hold the lock either, because build() takes it at startup and fails there. Step 3's close(recycleSlotLock == null) already degrades correctly in that case. Alternatively, disarm with a single WARN.
  • Add connect()-based tests for both path variants.

Moderate

2. No test pins the connect-gate CAS that closeIfLinkUp() adds

  • Problem: No test pins closeIfLinkUp()'s connect-gate CAS.
  • Net impact: Removing the CAS keeps the suite green and reopens a 30 s stall.
  • Evidence: A mutant at 17f04409 without the CAS passes all 92 recycle and gate tests, five runs out of five. A seam-free probe fails on it every time.

Location: CursorWebSocketSendLoop.java:1474-1479 (from 83bcaa4). The 83bcaa4 commit message says the new gate test "pins both outcomes of the gate deterministically", but none of the three tests checks the CAS itself:

  • testClaimedStopKeepsIoThreadOutOfConnectWalk writes STOP_CLAIMED by reflection and drives only connectLoop.
  • testCloseIfLinkUpRefusesDuringConnectWalkThenStopsLiveLoop calls closeIfLinkUp() after the walk has set lastReconnectError (:1950), so isLinkUp() refuses before the CAS runs.
  • SymbolDictRecycleOutageTest#testServerCloseBehindDrainingAckDoesNotStallProducer targets the window before the gate existed, which the isLinkUp() re-check now closes on its own.

Mutation: if (!isLinkUp()) return false;

  • The gate test plus every SymbolDictRecycle*Test (92 tests): green in 5 of 5 runs.
  • The gate and outage tests alone: green in 10 of 10 runs.
  • Six more classes that use the recycle or gate API (225 tests): green.

If this regresses: the I/O thread can enter its connect walk after isLinkUp() read the gate as OPEN but before close() clears running.

  • Step 2 then calls close() on a walk blocked in a credential pull or DNS resolve that ignores interrupts.
  • close() waits out the 30 s shutdown budget and throws, the recycle parks at CLOSE_LOOP, and the producer's table() absorbs the stall.
  • The natural window is tiny (0 hits in 20,000 forced races), which is why this is Moderate rather than Critical.

Fix (no production seam needed): add a gate test using a WebSocketClient subclass whose isConnected() does three things:

  1. drops the link,
  2. waits until the I/O thread is inside a reconnect factory that ignores interrupts,
  3. returns true.

isConnected() is the last check in isLinkUp(), after the gate read, so this forces the window. At head closeIfLinkUp() returns false in about 6 ms (5 of 5 runs). On the mutant, with setShutdownAwaitTimeoutMillis(1000), close() times out and throws (5 of 5). Probe source: /tmp/rev91/fal/f-l7-CloseIfLinkUpGateProbeTest.java.

3. The new engine test skips the suite's memory-leak check

  • Problem: A new engine test skips the suite's assertMemoryLeak wrapper.
  • Net impact: Native leaks on this path would go undetected; no user impact.
  • Evidence: CursorSendEngineTest.java:425-450 at 17f04409; every other @Test in the file wraps it.

testAdoptCountersSharesTheInstanceAndFoldsOnlyItsOwnCounter builds a disk-backed CursorSendEngine (mmap'd segments) without TestUtils.assertMemoryLeak.

Fix: wrap the body in TestUtils.assertMemoryLeak(() -> { … }).

Minor

4. The PR body's release notes announce removing an accessor that never shipped

  • Problem: The release note announces removing an accessor that never shipped.
  • Net impact: The squashed commit message tells users a nonexistent public API was removed.
  • Evidence: getSymbolDictResetsPerformed is absent at c9968f24 and in every tag from 1.0.0 to 1.3.9.

It was added (07ecb27) and later removed on this branch. Fix: drop the "getSymbolDictResetsPerformed() is gone" bullet from the PR body.

Coverage map

  • Test gate: passes. There is one admitted coverage gap, finding 2 (Moderate).
    • Search: rg -l "closeIfLinkUp|isLinkUp" core/src/test finds only CursorWebSocketSendLoopConnectGateTest.
    • Failure link: none of its assertions can fail when the CAS is removed (mutation runs above).
  • Tandem PRs: questdb/questdb#7525 (real-server e2e and fuzz) and questdb/questdb-enterprise#1171 (failover) both exist on the matching branch, are cross-linked and have green CI. Both pin client 0496f037, so 83bcaa4 (connect gate) and 17f0440 (retained-engine listener) are covered only by local tests until the pins move. That is merge mechanics, not a finding.

Summary

  • Verdict: request changes.
  • Correctness gate: fails, with one open Critical (finding 1).
  • Test gate: passes. There is no Critical coverage gap and one Moderate gap.
  • Severity: 1 Critical, 2 Moderate, 1 Minor. All four are in-diff; none is an out-of-diff breakage.
  • Submodules: none changed.
  • Also checked:
    • The Java 8 floor holds: no post-Java-8 APIs in the diff, and CI at 17f04409 is green for the JDK 8 build and tests and for the JDK 25 and 26 compile smoke.
    • No committed binaries.
    • The steady-state table() and column path adds only field reads, with no allocation.
    • The store-and-forward contract holds: after the first connect, reconnects stay on the I/O thread with no time budget. Finding 1 is a local slot-lock failure, not a transport error reaching the producer.

Falsification probes and logs are under /tmp/rev91/fal/. I removed the scratch worktrees afterwards, and the primary checkout is unchanged at 17f04409.

@bluestreak01 bluestreak01 removed the READY label Oct 1, 2026
jovfer and others added 10 commits October 3, 2026 16:21
Step 0 of the symbol-dictionary recycle wrapped every failure to take
the slot's logical lock as "deferred, retried at a later row start" and
threw it from table(). For a slot the lock can never be taken for -- a
bare relative slot name has no parent to anchor it in, and a parent the
process cannot write cannot hold .slot-locks -- the failure repeats at
every drained row start, so a connect()-built sender with a rebuild
factory refused every row from its first arm onward.

Nothing is torn down at step 0, so no failure there needs to cost the
producer a row. A contention now leaves the recycle armed and returns;
the next drained barrier retries. Any other failure logs one WARN and
stops the sender arming: swapping without the lock would let a colliding
build() take the slot between the outgoing close and the rebuild, so the
sender keeps ingesting with a growing dictionary instead, as one without
a rebuild factory does.

Tests: SymbolDictRecycleSlotOwnershipTest drives connect()-built senders
on a bare relative slot, under a read-only parent and with the lock
directory blocked by a regular file, and now requires rows to keep
flowing while the lock is contended.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Neither gate test exercised the CAS: one writes the gate by reflection
and drives only the walk's entry, the other calls closeIfLinkUp() after
the walk has already made isLinkUp() refuse. With the CAS removed the
suite stayed green.

The new test forces the window the CAS exists for. isConnected() is the
last read isLinkUp() makes, after the gate read, so a hook there drops
the link and holds the caller until the I/O thread is inside a connect
factory that ignores interrupts, then answers connected. closeIfLinkUp()
must lose the gate and refuse; without the CAS it stops a blocked walk
and waits out the shutdown budget.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
It built a disk-backed engine outside assertMemoryLeak, the only test in
the class to do so.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…e test's close budget

The comments and javadocs around recycle step 0 said the sender stops
arming only when the slot's logical lock can never be taken. The code
latches on every failure other than a contention, a transient open
failure included, because step 0 cannot tell the two apart. The text now
says so, says "disk-mode swap" where memory mode takes no lock, names
the condition as a parent where .slot-locks cannot be created or opened,
and lists step 0's outcomes in the starvation wait's javadoc. The WARN
adds that a new sender recycles again if the cause was transient.

The connect-gate test shrank the loop's shutdown budget to keep a
wrongful stop short and then ran its cleanup close under the same
budget; it now restores the default first.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Conflicts:
- Sender.java: main routes every engine build() creates through
  newCursorEngine() so it charges the shared sf_max_total_bytes budget;
  the branch had moved that construction into
  constructEngineOnSlotLocked(). The helper now takes the shared budget
  and passes it on, and the recycle's rebuild factory hands it the
  builder's budget, so a rebuilt engine charges the same budget as the
  engine it replaces.
- CursorWebSocketSendLoop.java: both sides added a field at the same
  place (externalFsnBase, gracefulStopMillis); kept both.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AaMH1ioFovgLgCa5bmwHbo
A recycle replaces the sender's cursor engine. Two pooled-sender tests,
one with sf_dir and one in memory mode, check that the replacement
charges the pool's shared sf_max_total_bytes budget and that closing
the pool releases every charge. Both fail if the rebuild factory builds
the replacement on a private budget.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AaMH1ioFovgLgCa5bmwHbo
The master connectWithCredentialSupplier overload, and the connect
overload that delegates to it, assigned the recycle threshold and the
maximum wait unchecked while the builder rejects out-of-range values. A
threshold of 0 on a sender with a rebuild factory armed a recycle at
every flush. The overload now applies the builder's ranges and messages
before it builds anything.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The threshold javadoc said a bounded live set recycles once and settles.
The re-arm bar is capped at 1,000,000, so only a static set below that
stops, after about log2(set size / threshold) recycles; a set of
1,000,000 or more, or one whose values keep changing, recycles for the
life of the sender. The javadoc and the floor comment now say so and
point the 1M-2M case at symbolDictReset(false).

The dictionary-full message said the reset acts only on build() and
fromConfig() senders; a connect() sender with a rebuild factory
recycles too.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…loop

recordFatal() leaves the connect gate open, so closeIfLinkUp() still
wins when the I/O thread latches a terminal error behind the stop's own
link check. Recycle step 2 then read only the ever-connected sticky and
dropped the loop, and checkConnectionError() polls only a live loop
reference: the error vanished, the swap committed and the sender kept
accepting rows after a terminal rejection. The resume that finishes a
failed step-2 stop dropped the loop the same way.

Both sites now go through retireStoppedLoop(), which throws the loop's
terminal error and keeps the dead loop, so every later call throws it
too. Step 2 first gives back a slot lock it took. The tests drive a
sender on a stub link that delivers a SECURITY_ERROR rejection from
inside closeIfLinkUp()'s link check, and from the second closeTraffic()
of a stop whose first attempt timed out.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The resume-arm test kept its 200 ms shutdown budget, which exists to
make the first stop time out, for the second stop as well, and that one
must succeed: a slow log write or a scheduling stall on the I/O thread's
exit path would have failed the test. It now restores the default budget
before the resume.

The connect-overload test also runs with the recycle switched off: the
ranges apply either way, as in the builder. The slot-lock field comment
lists the two give-backs in recycle step 2.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@mtopolnik

Copy link
Copy Markdown
Contributor

[PR Coverage check]

😍 pass : 562 / 634 (88.64%)

file detail

path covered line new line coverage
🔵 io/questdb/client/Sender.java 76 91 83.52%
🔵 io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java 56 64 87.50%
🔵 io/questdb/client/cutlass/qwp/client/QwpWebSocketSender.java 398 447 89.04%
🔵 io/questdb/client/impl/PooledSender.java 2 2 100.00%
🔵 io/questdb/client/cutlass/qwp/client/sf/cursor/CursorSendEngine.java 11 11 100.00%
🔵 io/questdb/client/impl/ConfigSchema.java 3 3 100.00%
🔵 io/questdb/client/cutlass/qwp/client/sf/cursor/CursorSendCounters.java 16 16 100.00%

@jovfer

jovfer commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor Author

Tandem review: client #91 / OSS #7525 / ENT #1171 (level 3)

Reviewed heads: client 6f58e888, OSS 2bffcfd467, ENT 10b5bc914.

Scope: both submodule pointers point at unmerged feature-branch commits, so the whole PR diff in each repo is in scope. Apart from their pointer lines, the OSS and ENT diffs are byte-identical to the previous review's. The new code under review is client-side:

  • the wide connect overloads validate the recycle settings;
  • the docs now state when a recycling sender stops recycling;
  • a terminal error latched while the recycle stops the I/O loop is thrown instead of discarded;
  • tests for all three.

Critical

None.

Moderate

[client] No test pins that close() stays silent after the recycle throws the loop's terminal. Location: QwpWebSocketSender.java:6233-6237.

  • Problem: the surfaced-error record on the new recycle throw is untested.
  • Net impact: a refactor could make close() rethrow, or raise IllegalArgumentException under try-with-resources.
  • Evidence:
    • At 6f58e888, the full client suite (3,595 tests) stays green with the record dropped.
    • Four probe tests go red with the record dropped and pass at head.

retireStoppedLoop() throws through cursorSendLoop.checkError(), which also records the error as already surfaced. close0() reads that record (:1750) to drop the instance the caller already owns (:1918).

Both SymbolDictRecycleLoopTerminalTest tests call table() a second time before close() (:99, :156). That call goes through checkConnectionError() → loop.checkError() and records the error anyway, so a missing record on the recycle path is invisible to them.

The regression is plausible. Step 2 already reads getTerminalError() one line before the call (:6056), and folding the throw into that read would drop the record with nothing failing.

Symptom if it regresses: the trigger is a terminal server rejection, such as SECURITY_ERROR from a read-only node, latching while the recycle stops the loop.

  • A caller that catches the terminal from table() gets it thrown again from close().
  • A try-with-resources caller gets IllegalArgumentException: Self-suppression not permitted instead of the typed LineSenderServerException.

At base this is a new surface. The same ownership contract already holds for every other terminal path; a non-recycle analog probe passes at base and at head.

Suggested fix: add one test using the existing NackOnDemandClient:

  1. Trigger the step-2 latch.
  2. Catch the terminal from the first table().
  3. Assert that close() returns without throwing.

A try-with-resources variant of the same test asserts that the thrown error is the loop's own terminal.

Minor

[client] The widest connect javadoc gives the max_frame_rejections definition to the recycle keys. Location: QwpWebSocketSender.java:1225-1230.

  • Problem: the summary sentence ties the rejection-count definition to the recycle keys.
  • Net impact: readers of the public javadoc see the wrong meaning for three keys.
  • Evidence: the same block at base a7e7db3d is correct; this PR's 3a30e4d9 inserted the recycle keys before the colon.

Suggested fix: keep the definition right after max_frame_rejections and give the recycle keys their own clause.

Coverage

  • Test gate: passes. One coverage gap is admitted, the Moderate above.
  • The new tests kill their target mutants:
    • removing checkError() from retireStoppedLoop;
    • removing the step-2 lock release;
    • reverting the CLOSE_LOOP resume arm;
    • all four validation bounds.
  • Every new test is red against the previous head 38f8ce04.
  • Flakiness: 10/10 consecutive runs passed, and 5/5 with the CPU oversubscribed 2x.

Checked and found sound

  • Terminal handling during the stop: close() returns normally only after the I/O thread has passed its shutdown countdown, and recordFatal (the only writer of terminalError) never runs after it.
    • The step-2 guard and the throw always read the same value.
    • A terminal latched inside the graceful-stop window is caught too.
  • Sender state after the throw: every entry point (table, the flush family, drain, awaitAckedFsn, ensureConnected) rethrows the same instance, and the CLOSE_LOOP resume is not re-entered.
  • Lock release: the step-0 lock is released exactly once, at step 2 or first in close().
  • close() after the throw: for default and custom handlers alike, close() is silent, as on a sender that never recycled.
  • Pool: it discards the lease.
  • Connect validation: it matches the builder's ranges and messages exactly, throws before the engine is attached, and affects no existing caller or released API.
  • Doc arithmetic: the "900k symbols → 4 recycles" example and the 1M band are correct.

Verdict

Correctness gate and test gate pass. Findings: 0 Critical, 1 Moderate, 1 Minor, both in-diff in the client.

🤖 Generated with Claude Code

@jovfer jovfer added the READY label Oct 5, 2026

@bluestreak01 bluestreak01 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewing PR #91 at level 3, the full review.

Scope. No submodule pointer moved (core/src/main/c/share/zstd is unchanged), so the whole diff is in scope. The title follows Conventional Commits, Fixes #80 is at the top of the body, and the labels fit. Reviewed head 6f58e888, merge base a7e7db3d.

Critical

None.

Moderate

1. No test pins the cap on the re-arm floor

  • Problem: No test pins the re-arm floor's 1,000,000 cap.
  • Net impact: If the cap is removed, producers with ever-growing symbol sets silently stop recycling.
  • Evidence: With the cap removed at 6f58e888, the full JDK 8 core suite still passes (3,595 tests, 0 failures). A probe fails on that mutant.

Where (in-diff): QwpWebSocketSender.java:5642-5643, the clause Math.min(dictSizeAtSwap * 2, MAX_SYMBOL_DICTIONARY_SIZE / 2).

Why this cap matters: take a producer on default settings whose distinct symbols keep growing.

  • After four recycles the floor reaches 1M (100k → 200k → 400k → 800k → 1M).
  • From then on every swap commits at 1M or more, and this clause is what holds the floor at 1M.
  • Without it, the floor jumps past 2M. armIfEligible (:5435) can then never fire again, because the dictionary is capped at 2M.
  • The producer then hits "global symbol dictionary is full", which is the #80 failure again, for exactly the users this PR is for. Base has no recycle, so this is a new surface.

What the tests reach:

  • Client floor assertions only use small values (0, 4, 8, 12). The largest dictionary at any swap in the whole suite is 18.
  • The tandem OSS and ENT tests use at most 400 symbols per epoch.

Runtime results (probe: /tmp/pr91/probes/PrReviewFloorCapProbeTest.java):

Run Head Cap removed
Probe (threshold 500k) Floor stays 1,000,000 and re-arms at each 1M; no error Floor 2,200,000, never re-arms, cap error at 2M
Default config, row API No error over 6M rows Cap error at about 5.1M rows

Fix: add a test shaped like the probe:

  1. Fill 1.1M symbols through getGlobalSymbolDictionaryForTest().
  2. Flush, drain, then call table().
  3. Assert getResetFloorSymbolsForTesting() == 1_000_000 and that the next refill to 1M re-arms.

It runs in about 1 s on the existing TestWebSocketServer.

2. Sender.table() javadoc doesn't describe its new behaviour

  • Problem: The Sender.table() javadoc omits the new pause, the reconnect and the new exceptions.
  • Net impact: WebSocket producers past 100k symbols hit table() stalls that the method's own docs don't mention.
  • Evidence: Sender.java:793-805 is identical to base. The new behaviour comes from QwpWebSocketSender.java:3650-3658.

Where (out-of-diff): the javadoc is unchanged, but the PR changes the contract it describes. On default settings table() can now:

  • park for up to 2 s (maybeBlockForStarvedReset, :5757), and for up to 30 s per wait in a deferred engine close (:185, :5468);
  • close the WebSocket and start a reconnect (:5971);
  • throw new LineSenderExceptions: transient ones marked "retried on the next send", plus the terminal breach.

None of this could happen at base. It is documented only on resetSymbolDictionary() and the three builder knobs. Sender.table(), QwpWebSocketSender.table() and PooledSender.table() say nothing.

Fix: add a WebSocket paragraph to Sender.table() that links to resetSymbolDictionary() and symbolDictResetMaxWaitMillis(long) and gives the worst-case blocking bound.

Minor

3. The master connect javadoc attaches a definition to the wrong keys

  • Problem: The master connect javadoc applies the rejection-count definition to the recycle keys.
  • Net impact: Readers of the low-level API get the wrong meaning for three keys.
  • Evidence: QwpWebSocketSender.java:1226-1231 at head, against :859-862 at base.

The recycle knobs were inserted before the colon, so "consecutive server-active rejections of the same head-of-line frame…" now reads as their definition. Fix: keep that clause next to max_frame_rejections and describe the knobs in their own sentence.

Coverage map

Test gate: passes. 1 coverage gap admitted.

Change Search Failure link
Floor cap, :5642-5643 rg getResetFloorSymbolsForTesting finds only SymbolDictRecycleArmingTest (floors 0/4/8) and SymbolDictRecycleTest (12). OSS #7525 and ENT #1171 use at most 400 symbols per epoch. None: with the cap removed, 3,595 tests pass and 0 fail

Summary

  • Verdict: approve with comments. Please address 1 (add the test) and 2 (doc sentence). Neither blocks the merge.
  • Gates: the correctness gate passes (no Critical findings) and the test gate passes.
  • Provenance: no submodule pointer changed.
  • Split: 2 in-diff (1 and 3) and 1 out-of-diff, the documentation contract in 2.
  • Severity: 0 Critical, 2 Moderate, 1 Minor.
  • Tandem PRs: questdb/questdb#7525 and questdb/questdb-enterprise#1171 are on the matching branch, cross-linked in both directions, and pinned to this head: OSS pins 6f58e888, ENT pins OSS 2bffcfd4. Client CI is green. The OSS macwin legs and ENT CI were still pending at review time.
  • Checked, nothing found:
    • Public API: a base→head javap diff removes no public or protected signature.
    • Build constraints: the Java 8 floor holds and there are no committed binaries.
    • Hot path: the per-row path only adds field reads and allocates nothing.
    • Store-and-forward contract: after initialization no transport error or reconnect time limit reaches the producer.
    • Durable-ack and transactional modes: in durable-ack mode ackedFsn moves only on durable acks. The server withholds acks for uncommitted deferred frames (QwpIngressUpgradeProcessor.java:1676-1695).
    • Local run: the 181 recycle-related tests pass on JDK 8 at head.

The primary checkout is unchanged at 6f58e888, my scratch worktrees have been removed, and I haven't posted anything to the PR.

@bluestreak01 bluestreak01 added the QUEUED FOR MERGE Approved PR in the merge queue. Do not merge master into this PR. label Oct 7, 2026

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

enhancement New feature or request QUEUED FOR MERGE Approved PR in the merge queue. Do not merge master into this PR. QWP READY tandem

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 million of district symbols limit

3 participants