diff --git a/README.md b/README.md index e945c7770..a013a15a8 100644 --- a/README.md +++ b/README.md @@ -556,6 +556,9 @@ Applied by the query pool to select and fail over between the nodes in the `addr The ingest side also accepts store-and-forward and reconnection tuning keys (`auto_flush_*`, `initial_connect_retry`, `reconnect_*`, `request_durable_ack`, `sf_*`, `max_frame_rejections`, `poison_min_escalation_window_millis`, …). +`sf_max_total_bytes` caps the unacknowledged data all pooled senders buffer together — 128 MiB of memory by default, +10 GiB of disk with `sf_dir` — so a larger `sender_pool_max` adds connections, not buffer memory. The only overshoot +is the minimum working set every live sender keeps: two segments (`2 × sf_max_segment_bytes`, 8 MiB by default). `sf_durability=periodic` checkpoints mmap-published data in the background; `sf_sync_interval_millis` defaults to `5000` in that mode. The interval is a target cadence: JVM scheduling and storage-sync latency add to the actual loss window. Use `request_durable_ack=on` when end-to-end server durability is also required. See the diff --git a/core/src/main/java/io/questdb/client/QuestDBBuilder.java b/core/src/main/java/io/questdb/client/QuestDBBuilder.java index e846ad129..11364f9ba 100644 --- a/core/src/main/java/io/questdb/client/QuestDBBuilder.java +++ b/core/src/main/java/io/questdb/client/QuestDBBuilder.java @@ -400,6 +400,11 @@ public QuestDBBuilder queryPoolSize(int size) { /** * Maximum sender-pool size. Defaults to 4. + *

+ * The pooled senders share one {@code sf_max_total_bytes} budget for data + * the server has not acknowledged yet, so a larger pool adds connections, + * not buffer memory -- apart from the minimum working set of two segments + * every live sender keeps. */ public QuestDBBuilder senderPoolMax(int max) { if (max < 1) { diff --git a/core/src/main/java/io/questdb/client/Sender.java b/core/src/main/java/io/questdb/client/Sender.java index 645d7b254..3b630da55 100644 --- a/core/src/main/java/io/questdb/client/Sender.java +++ b/core/src/main/java/io/questdb/client/Sender.java @@ -40,6 +40,7 @@ import io.questdb.client.cutlass.qwp.client.sf.cursor.CursorWebSocketSendLoop; import io.questdb.client.cutlass.qwp.client.sf.cursor.OrphanScanner; import io.questdb.client.cutlass.qwp.client.sf.cursor.PersistedSymbolDict; +import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentBudget; import io.questdb.client.cutlass.qwp.client.sf.cursor.SlotLock; import io.questdb.client.cutlass.qwp.client.sf.cursor.MmapSegmentCorruptionException; import io.questdb.client.cutlass.qwp.client.sf.cursor.SfRecoveryException; @@ -1162,6 +1163,9 @@ public int getConnectTimeout() { private SfDurability sfDurability = SfDurability.MEMORY; private long sfMaxSegmentBytes = PARAMETER_NOT_SET_EXPLICITLY; private long sfMaxTotalBytes = PARAMETER_NOT_SET_EXPLICITLY; + // Budget shared with other senders (a sender pool's), or null for a + // private budget of sfMaxTotalBytes. See storeAndForwardSharedBudget. + private SegmentBudget sfSharedBudget; private long sfSyncIntervalMillis = PARAMETER_NOT_SET_EXPLICITLY; private boolean shouldDestroyPrivKey; private boolean tlsEnabled; @@ -1484,18 +1488,8 @@ public Sender build() { // (same lock-free architecture, no disk involvement). // Durability-combination validation lives in validateParameters // so build() and no-connect validation apply the same rules. - long actualSfMaxSegmentBytes = sfMaxSegmentBytes == PARAMETER_NOT_SET_EXPLICITLY - ? DEFAULT_SEGMENT_BYTES - : sfMaxSegmentBytes; - // Default cap depends on backing: RAM (memory mode) is tight - // by default; disk (SF mode) is cheap so the default is - // generous enough that normal traffic never hits it. - long defaultMaxTotal = sfDir == null - ? DEFAULT_MAX_BYTES_MEMORY - : DEFAULT_MAX_BYTES_SF; - long actualSfMaxTotalBytes = sfMaxTotalBytes == PARAMETER_NOT_SET_EXPLICITLY - ? Math.max(defaultMaxTotal, actualSfMaxSegmentBytes * 2) - : sfMaxTotalBytes; + long actualSfMaxSegmentBytes = resolveSfMaxSegmentBytes(); + long actualSfMaxTotalBytes = resolveSfMaxTotalBytes(); long actualCloseFlushTimeoutMillis = closeFlushTimeoutMillis == CLOSE_FLUSH_TIMEOUT_NOT_SET ? DEFAULT_CLOSE_FLUSH_TIMEOUT_MILLIS : closeFlushTimeoutMillis; @@ -1624,9 +1618,9 @@ public Sender build() { CursorSendEngine cursorEngine; try { try { - cursorEngine = new CursorSendEngine( + cursorEngine = newCursorEngine( slotPath, actualSfMaxSegmentBytes, - actualSfMaxTotalBytes, actualSfAppendDeadlineNanos, + actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos, actualSfSyncIntervalNanos); } catch (SfSanitizedResidueException first) { // NOT terminal, and it must be intercepted ahead of its @@ -1639,9 +1633,9 @@ public Sender build() { LOG.info("sf slot {}: sealed residue sanitized during recovery ({}); " + "retrying over the healed chain", slotPath, first.getMessage()); - cursorEngine = new CursorSendEngine( + cursorEngine = newCursorEngine( slotPath, actualSfMaxSegmentBytes, - actualSfMaxTotalBytes, actualSfAppendDeadlineNanos, + actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos, actualSfSyncIntervalNanos); } } catch (UnreplayableSlotException | SfRecoveryException @@ -1669,7 +1663,7 @@ public Sender build() { quarantined = true; cursorEngine = quarantineTornSlot( null, e, sfDir, senderId, slotPath, actualSfMaxSegmentBytes, - actualSfMaxTotalBytes, actualSfAppendDeadlineNanos, + actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos, actualSfSyncIntervalNanos, errorHandler); } int actualErrorInboxCapacity = errorInboxCapacity != PARAMETER_NOT_SET_EXPLICITLY @@ -1742,7 +1736,7 @@ public Sender build() { quarantined = true; cursorEngine = quarantineTornSlot( cursorEngine, e, sfDir, senderId, slotPath, actualSfMaxSegmentBytes, - actualSfMaxTotalBytes, actualSfAppendDeadlineNanos, + actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos, actualSfSyncIntervalNanos, errorHandler); } catch (Throwable t) { // connect() failed before ownership of cursorEngine @@ -2899,6 +2893,19 @@ public String getConfiguredSfDir() { return sfDir; } + /** + * The {@code sf_max_total_bytes} cap {@link #build()} applies: the + * configured value, or the default for the configured backing (128 MiB + * in memory, 10 GiB with {@code sf_dir}, and never less than two + * segments). Introspection hook for the connection pool, which sizes + * the one budget all of its senders share from this value. Returns + * {@code -1} for transports other than WebSocket, which have no + * segment ring to cap. + */ + public long getResolvedSfMaxTotalBytes() { + return protocol == PROTOCOL_WEBSOCKET ? resolveSfMaxTotalBytes() : -1L; + } + /** * Excludes the connection pool's live slot set from * {@link #drainOrphans(boolean)} scanning: a sibling slot under @@ -3035,13 +3042,19 @@ public LineSenderBuilder storeAndForwardMaxSegmentBytes(long maxSegmentBytes) { } /** - * Hard cap on cursor-allocated bytes (active + spare + sealed - * segments). When the cap is reached, the producer's - * {@code Sender.flush()} blocks until ACK-driven trim frees space; - * if the cap is exhausted past the configured deadline (default 30 s), - * {@code flush()} throws. Default: {@code 128 MiB}, which applies to - * both memory-mode and SF-mode rings — for SF deployments with - * cheap disk, raise this knob explicitly. WebSocket transport only. + * Hard cap on cursor-allocated bytes: active + spare + sealed + * segments, plus the symbol-dictionary side-file with {@code sf_dir}. + * When the cap is reached, the producer's {@code Sender.flush()} + * blocks until ACK-driven trim frees space; if the cap is exhausted + * past the configured deadline (default 30 s), {@code flush()} throws. + * Default: {@code 128 MiB} of native memory without {@code sf_dir}, + * {@code 10 GiB} of disk with it. WebSocket transport only. + *

+ * The cap belongs to this sender alone unless it is built with + * {@link #storeAndForwardSharedBudget(SegmentBudget)}. The connection + * pool behind {@link QuestDB} does that for every sender it builds, + * so in a pool the cap bounds all pooled senders together, not each + * of them. */ public LineSenderBuilder storeAndForwardMaxTotalBytes(long maxTotalBytes) { if (protocol != PARAMETER_NOT_SET_EXPLICITLY && protocol != PROTOCOL_WEBSOCKET) { @@ -3054,6 +3067,29 @@ public LineSenderBuilder storeAndForwardMaxTotalBytes(long maxTotalBytes) { return this; } + /** + * Charges this sender's cursor segments to {@code budget}, shared with + * every other sender built on it, instead of to a private budget of + * {@code sf_max_total_bytes}. Together those senders never hold more + * than the budget's capacity, except that each is always granted its + * minimum working set: the active segment plus one spare. This + * sender's own {@code sf_max_total_bytes} then caps only the + * background orphan drainers it starts. The connection pool behind + * {@link QuestDB} builds every sender this way. WebSocket transport + * only. + * + * @param budget the shared budget, or {@code null} for a private one + * (the default) + * @return this instance for method chaining + */ + public LineSenderBuilder storeAndForwardSharedBudget(SegmentBudget budget) { + if (protocol != PARAMETER_NOT_SET_EXPLICITLY && protocol != PROTOCOL_WEBSOCKET) { + throw new LineSenderException("store_and_forward is only supported for WebSocket transport"); + } + this.sfSharedBudget = budget; + return this; + } + /** * Sets the target background checkpoint cadence for * {@link SfDurability#PERIODIC}. Scheduler and storage latency add to @@ -3090,6 +3126,25 @@ private static int getValue(CharSequence configurationString, int pos, StringSin return pos; } + // Every cursor engine build() creates -- including the fresh slot that + // replaces a quarantined one -- charges the shared budget when the sender + // was given one, and a private budget of sfMaxTotalBytes otherwise. + private static CursorSendEngine newCursorEngine( + String slotPath, + long sfMaxSegmentBytes, + long sfMaxTotalBytes, + SegmentBudget sfSharedBudget, + long sfAppendDeadlineNanos, + long sfSyncIntervalNanos + ) { + if (sfSharedBudget != null) { + return new CursorSendEngine(slotPath, sfMaxSegmentBytes, sfSharedBudget, + sfAppendDeadlineNanos, sfSyncIntervalNanos); + } + return new CursorSendEngine(slotPath, sfMaxSegmentBytes, sfMaxTotalBytes, + sfAppendDeadlineNanos, sfSyncIntervalNanos); + } + private static SfDurability parseDurabilityValue(@NotNull StringSink value) { if (Chars.equalsIgnoreCase("memory", value)) return SfDurability.MEMORY; if (Chars.equalsIgnoreCase("periodic", value)) return SfDurability.PERIODIC; @@ -3246,8 +3301,8 @@ private static long parseSizeValue(@NotNull StringSink value, @NotNull String na private static CursorSendEngine quarantineTornSlot( CursorSendEngine torn, RuntimeException cause, String sfDir, String senderId, String slotPath, - long sfMaxSegmentBytes, long sfMaxTotalBytes, long sfAppendDeadlineNanos, - long sfSyncIntervalNanos, + long sfMaxSegmentBytes, long sfMaxTotalBytes, SegmentBudget sfSharedBudget, + long sfAppendDeadlineNanos, long sfSyncIntervalNanos, io.questdb.client.SenderErrorHandler errorHandler ) { // The verdict, and the reason, come from the recovery seed -- the only code that @@ -3322,7 +3377,7 @@ private static CursorSendEngine quarantineTornSlot( String.valueOf(handlerFailure)); } } - return new CursorSendEngine(slotPath, sfMaxSegmentBytes, sfMaxTotalBytes, + return newCursorEngine(slotPath, sfMaxSegmentBytes, sfMaxTotalBytes, sfSharedBudget, sfAppendDeadlineNanos, sfSyncIntervalNanos); } @@ -4351,6 +4406,21 @@ private void http() { protocol = PROTOCOL_HTTP; } + private long resolveSfMaxSegmentBytes() { + return sfMaxSegmentBytes == PARAMETER_NOT_SET_EXPLICITLY ? DEFAULT_SEGMENT_BYTES : sfMaxSegmentBytes; + } + + // The default cap depends on the backing: RAM (memory mode) is tight by + // default; disk (SF mode) is cheap, so its default is generous enough + // that normal traffic never hits it. + private long resolveSfMaxTotalBytes() { + if (sfMaxTotalBytes != PARAMETER_NOT_SET_EXPLICITLY) { + return sfMaxTotalBytes; + } + long defaultMaxTotal = sfDir == null ? DEFAULT_MAX_BYTES_MEMORY : DEFAULT_MAX_BYTES_SF; + return Math.max(defaultMaxTotal, resolveSfMaxSegmentBytes() * 2); + } + private void tcp() { if (protocol != PARAMETER_NOT_SET_EXPLICITLY) { throw new LineSenderException("protocol was already configured ") diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorSendEngine.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorSendEngine.java index 66a0635ea..3ac2f262c 100644 --- a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorSendEngine.java +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorSendEngine.java @@ -281,7 +281,7 @@ public CursorSendEngine(String sfDir, long segmentSizeBytes, @TestOnly public CursorSendEngine(String sfDir, long segmentSizeBytes, long maxTotalBytes, long appendDeadlineNanos, FilesFacade dictFf) { - this(sfDir, segmentSizeBytes, null, true, maxTotalBytes, + this(sfDir, segmentSizeBytes, null, true, new SegmentBudget(maxTotalBytes), appendDeadlineNanos, 0L, dictFf); } @@ -292,8 +292,23 @@ public CursorSendEngine(String sfDir, long segmentSizeBytes, public CursorSendEngine(String sfDir, long segmentSizeBytes, long maxTotalBytes, long appendDeadlineNanos, long syncIntervalNanos) { - this(sfDir, segmentSizeBytes, null, true, maxTotalBytes, - appendDeadlineNanos, syncIntervalNanos); + this(sfDir, segmentSizeBytes, null, true, new SegmentBudget(maxTotalBytes), + appendDeadlineNanos, syncIntervalNanos, FilesFacade.INSTANCE); + } + + /** + * As {@link #CursorSendEngine(String, long, long, long, long)}, but the + * engine's private {@link SegmentManager} charges its segments to + * {@code budget} instead of to a budget of its own. Engines built on one + * budget share its capacity: together they never hold more than it, apart + * from the minimum working set (the active segment plus one spare) each + * ring is always granted. A sender pool builds every sender this way, so + * {@code sf_max_total_bytes} caps the pool rather than each pooled sender. + */ + public CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentBudget budget, + long appendDeadlineNanos, long syncIntervalNanos) { + this(sfDir, segmentSizeBytes, null, true, budget, + appendDeadlineNanos, syncIntervalNanos, FilesFacade.INSTANCE); } /** @@ -319,19 +334,16 @@ private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager man private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager manager, boolean ownsManager, long appendDeadlineNanos, long syncIntervalNanos) { - this(sfDir, segmentSizeBytes, manager, ownsManager, - SegmentManager.UNLIMITED_TOTAL_BYTES, appendDeadlineNanos, syncIntervalNanos); - } - - private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager manager, - boolean ownsManager, long maxTotalBytes, long appendDeadlineNanos, - long syncIntervalNanos) { - this(sfDir, segmentSizeBytes, manager, ownsManager, maxTotalBytes, + // A caller-supplied manager already charges its own budget. + this(sfDir, segmentSizeBytes, manager, ownsManager, null, appendDeadlineNanos, syncIntervalNanos, FilesFacade.INSTANCE); } + // budget is what an owned manager (ownsManager && manager == null) charges + // its segments to; unused, and may be null, when the caller supplies the + // manager. private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager manager, - boolean ownsManager, long maxTotalBytes, long appendDeadlineNanos, + boolean ownsManager, SegmentBudget budget, long appendDeadlineNanos, long syncIntervalNanos, FilesFacade dictFf) { this.dictFf = dictFf; // Allocate the bound callback before constructing an owned manager. @@ -355,7 +367,7 @@ private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager man } if (ownsManager && manager == null) { manager = new SegmentManager( - segmentSizeBytes, SegmentManager.DEFAULT_POLL_NANOS, maxTotalBytes); + segmentSizeBytes, SegmentManager.DEFAULT_POLL_NANOS, budget); } SlotLock acquiredLock = null; if (!memoryMode) { diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java index 6643bf036..228d3980e 100644 --- a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java @@ -87,6 +87,21 @@ */ public final class CursorWebSocketSendLoop implements QuietCloseable { + /** + * How long {@link #close()} lets a live I/O thread stop on its own before + * it breaks the connection's traffic path. After {@code running} goes false + * an idle or briefly busy I/O thread exits within microseconds, and its + * exit path closes the client in order -- a WebSocket CLOSE frame and, over + * TLS, a close_notify -- so the server sees an orderly close. Breaking + * traffic first (shutdown(SHUT_RDWR)) sent both into a dead socket, and the + * server logged every sender close as a dropped connection. Only an I/O + * thread stuck in a native send or receive outlives this window; close() + * then breaks its traffic exactly as before, so the window caps the extra + * close() latency of that stuck case. A connect walk blocked on a published + * in-flight client or credential pull skips the window: nothing is + * established to close, and only cancellation unblocks it. + */ + public static final long DEFAULT_CLOSE_GRACEFUL_STOP_MILLIS = 100L; /** * Bounded-await backstop for {@link #close()}: the maximum time close() * waits for the I/O thread to stop (count down {@code shutdownLatch}) @@ -416,6 +431,11 @@ public final class CursorWebSocketSendLoop implements QuietCloseable { // it is engine.ackedFsn() + 1, so the first replayed frame on the new // connection is wireSeq=0 and server-side cumulative ACKs still line up. private long fsnAtZero; + // Graceful-stop window for close() (see DEFAULT_CLOSE_GRACEFUL_STOP_MILLIS). + // Overridable via setGracefulStopMillis so tests can pin either outcome + // deterministically; production always uses the default. Read only on the + // owner thread inside close(). + private long gracefulStopMillis = DEFAULT_CLOSE_GRACEFUL_STOP_MILLIS; // Bounded-await backstop budget for close() (see // DEFAULT_CLOSE_SHUTDOWN_AWAIT_MILLIS). Overridable via // setShutdownAwaitTimeoutMillis so tests can exercise the timeout branch @@ -1274,7 +1294,10 @@ public synchronized void close() { // finally{shutdownLatch.countDown()} never fired — awaiting here // would block forever. isAlive()==false also covers the normal // post-exit case where the latch is already counted down. - if (t.isAlive()) { + // awaitGracefulStop() lets a thread that is not stuck stop on its + // own first, so its exit path ends the connection cleanly and + // there is no traffic left to break. + if (t.isAlive() && !awaitGracefulStop()) { // Break a native send/receive before joining. Full client close // must remain after the worker exit because it frees buffers the // worker may still access. @@ -1559,6 +1582,17 @@ public void setErrorDispatcher(SenderErrorDispatcher dispatcher) { this.errorDispatcher = dispatcher; } + /** + * Test seam: resize the {@link #close()} graceful-stop window (default + * {@link #DEFAULT_CLOSE_GRACEFUL_STOP_MILLIS}). {@code 0} disables it, so + * close() breaks traffic up front. Production never calls this. Set + * before {@link #close()}. + */ + @TestOnly + public void setGracefulStopMillis(long millis) { + this.gracefulStopMillis = millis; + } + /** * Plug an async-delivery sink for ack-watermark advances. Same lifecycle * contract as {@link #setErrorDispatcher} — set once before @@ -1696,6 +1730,30 @@ private void attemptInitialConnect() { "initial connect", 0L); } + /** + * Gives a live I/O thread {@link #gracefulStopMillis} to stop on its own + * after close() cleared {@code running}. ioLoop's exit path closes the + * client -- WebSocket CLOSE frame, then the TLS close_notify -- before it + * counts the shutdown latch down, so {@code true} means the connection + * already ended cleanly and close() has no traffic to break. Returns + * {@code false} at once while the connect walk is blocked on cancellable + * work, which only {@link ConnectCancellation#cancel()} unblocks. An + * interrupt ends the wait early and stays set, so the backstop await + * takes its usual failed-stop branch after traffic has been broken. + * Owner thread only. + */ + private boolean awaitGracefulStop() { + if (gracefulStopMillis <= 0L || connectCancellation.isConnectInFlight()) { + return false; + } + try { + return shutdownLatch.await(gracefulStopMillis, TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return shutdownLatch.getCount() == 0L; + } + } + private void clearDurableAckTracking() { if (!durableAckMode) { return; @@ -3786,6 +3844,16 @@ void cancel() { t.interrupt(); } } + + /** + * Owner-thread probe from {@link #close()}: true while the connect walk + * is blocked on work only {@link #cancel()} can break -- a published + * in-flight client or a credential pull. There is no established + * connection to close gracefully then, so close() cancels at once. + */ + boolean isConnectInFlight() { + return inFlight != null || credentialPullThread != null; + } } /** diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentBudget.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentBudget.java new file mode 100644 index 000000000..579f19001 --- /dev/null +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentBudget.java @@ -0,0 +1,141 @@ +/******************************************************************************* + * ___ _ ____ ____ + * / _ \ _ _ ___ ___| |_| _ \| __ ) + * | | | | | | |/ _ \/ __| __| | | | _ \ + * | |_| | |_| | __/\__ \ |_| |_| | |_) | + * \__\_\\__,_|\___||___/\__|____/|____/ + * + * Copyright (c) 2014-2019 Appsicle + * Copyright (c) 2019-2026 QuestDB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + ******************************************************************************/ + +package io.questdb.client.cutlass.qwp.client.sf.cursor; + +import io.questdb.client.std.ObjList; + +import java.util.function.LongSupplier; + +/** + * The {@code sf_max_total_bytes} budget: an upper bound on the bytes the + * cursor rings charged to it may hold -- every segment a ring owns (active, + * sealed and hot spare), in memory or on disk, plus each slot's live + * side-file bytes. + *

+ * A {@code Sender} built on its own charges a private budget. A sender pool + * hands one budget to every sender it builds, so {@code sf_max_total_bytes} + * caps the pool as a whole: growing the pool adds connections, not buffer + * memory. Each {@link SegmentManager} keeps its own worker thread and charges + * the budget as it provisions and trims segments. + *

+ * The cap is enforced when a manager provisions a hot spare; bytes a ring + * already owns when it registers (a recovered slot) are charged even when + * they exceed it. A ring below its minimum working set -- the active segment + * plus one spare -- is provisioned past the cap, so a budget shared by + * {@code n} rings can be exceeded by up to {@code n} such working sets. + *

+ * Thread-safe. All state is guarded by this object's monitor, a leaf lock: + * managers call in while holding their own lock, and nothing here calls back + * into a manager. Side-file gauges are read under the monitor, so they must + * be wait-free and must not throw. + */ +public final class SegmentBudget { + + private final long capacityBytes; + // Live gauges of each charged slot's .symbol-dict side-file bytes. Read on + // every cap check rather than folded into segmentBytes: a dictionary + // grows out-of-band on producer threads, so an incremental mirror would + // drift, while a live read cannot. Each gauge is a pair of volatile reads + // (PersistedSymbolDict.occupiedDiskBytes()) and must stay wait-free: a + // producer can hold its dictionary's monitor across ff.allocate and mmap, + // and a gauge that took that monitor would stall every manager charging + // this budget behind one producer's append I/O. + private final ObjList sideFileGauges = new ObjList<>(); + private long segmentBytes; + + public SegmentBudget(long capacityBytes) { + if (capacityBytes <= 0) { + throw new IllegalArgumentException("capacityBytes must be positive: " + capacityBytes); + } + this.capacityBytes = capacityBytes; + } + + /** + * Segment bytes plus live side-file bytes currently charged. + */ + public synchronized long getAccountedBytes() { + return segmentBytes + sideFileBytes(); + } + + public long getCapacityBytes() { + return capacityBytes; + } + + public synchronized long getSegmentBytes() { + return segmentBytes; + } + + public synchronized long getSideFileBytes() { + return sideFileBytes(); + } + + synchronized void addSideFileGauge(LongSupplier gauge) { + sideFileGauges.add(gauge); + } + + /** + * Charges {@code bytes} unconditionally: segments a ring already owns when + * it registers, and spares provisioned under the liveness floor. + */ + synchronized void charge(long bytes) { + segmentBytes += bytes; + } + + synchronized void release(long bytes) { + segmentBytes -= bytes; + } + + // Identity, not equals(): each registration hands in its own gauge instance. + synchronized void removeSideFileGauge(LongSupplier gauge) { + for (int i = 0, n = sideFileGauges.size(); i < n; i++) { + if (sideFileGauges.getQuick(i) == gauge) { + sideFileGauges.remove(i); + return; + } + } + } + + /** + * Charges {@code bytes} if they fit under the capacity together with + * everything already charged and the live side-file bytes. Check and charge + * are one atomic step, so managers sharing the budget can never both take + * its last free segment. + */ + synchronized boolean tryCharge(long bytes) { + if (segmentBytes + sideFileBytes() + bytes > capacityBytes) { + return false; + } + segmentBytes += bytes; + return true; + } + + private long sideFileBytes() { + long total = 0L; + for (int i = 0, n = sideFileGauges.size(); i < n; i++) { + total += sideFileGauges.getQuick(i).getAsLong(); + } + return total; + } +} diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentManager.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentManager.java index 697b3b3d2..87cda92c1 100644 --- a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentManager.java +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentManager.java @@ -82,6 +82,11 @@ public final class SegmentManager implements QuietCloseable { private static final int TRIM_RETRY_UNLINK = 2; private static final long WORKER_JOIN_TIMEOUT_MILLIS = 5_000L; + // Every segment the registered rings own, plus their live side-file bytes, + // is charged here. Private to this manager unless the caller hands the same + // budget to several managers -- a sender pool does, so sf_max_total_bytes + // caps all of its senders together. + private final SegmentBudget budget; private final AtomicLong fileGeneration = new AtomicLong(); private final FilesFacade filesFacade; // Per-ring segment bytes below which the cap check never refuses to @@ -96,9 +101,11 @@ public final class SegmentManager implements QuietCloseable { // ordinary backpressure: the ring cycles between one and two segments as // acks arrive, ingestion continues, and the cap degrades to best-effort by // exactly the dictionary's overshoot rather than stopping the pipeline. + // A budget shared with other managers needs the floor for the same reason: + // a sibling ring whose connection is down can hold the whole budget, and + // only the floor keeps this ring cycling through its own two segments. private final long livenessFloorBytes; private final Object lock = new Object(); - private final long maxTotalBytes; // Reused by the manager worker thread to build spare-segment paths // directly into native memory. Each rotation writes the path bytes plus // a trailing NUL terminator into the same buffer, and passes the @@ -133,8 +140,8 @@ public final class SegmentManager implements QuietCloseable { // never as the exact false return meaning the worker loop has exited. private volatile Runnable beforeExitCleanupRegistrationHook; // Test seam: runs on the worker thread just before the install path's - // synchronized(lock) entry (the one that performs installHotSpare + the - // totalBytes += segmentSize commit). Null in production; tests use it to + // synchronized(lock) entry (the one that performs installHotSpare and + // keeps the spare's budget charge). Null in production; tests use it to // pause after the worker has snapshotted a RingEntry and created a spare, // but before ownership/accounting commit. Callers may inject a deregister // or hold this stale worker snapshot while caller-side cleanup runs. @@ -148,7 +155,7 @@ public final class SegmentManager implements QuietCloseable { // SegmentManagerTrimDeregisterRaceTest installs it, to deterministically // inject a deregister(ring) call into the exact race window that the // entry-state check inside the trim block closes for watermark writes and - // totalBytes accounting. + // budget accounting. private volatile Runnable beforeTrimSyncHook; // Test seam invoked exactly when a retry transition/recovery is logged. // Null in production; persistent-failure tests use it to prove log bounds @@ -188,12 +195,6 @@ public final class SegmentManager implements QuietCloseable { private boolean scratchHandedToWorker; private volatile long shortestSyncIntervalNanos = Long.MAX_VALUE; private boolean workerLoopExited; - // Total bytes currently allocated across every segment owned by every - // registered ring (active + sealed + hot-spare). Mutated by the manager - // thread on provision/trim and by register/deregister callers under - // {@link #lock}; the lock covers both paths so the counter stays - // consistent across registration boundaries. - private long totalBytes; // volatile: read by awaitRingQuiescence() from arbitrary caller threads // while the @TestOnly setter may run on another. private volatile long workerJoinTimeoutMillis = WORKER_JOIN_TIMEOUT_MILLIS; @@ -203,11 +204,11 @@ public final class SegmentManager implements QuietCloseable { private volatile Thread workerThread; public SegmentManager(long segmentSizeBytes) { - this(segmentSizeBytes, DEFAULT_POLL_NANOS, UNLIMITED_TOTAL_BYTES, FilesFacade.INSTANCE, System::nanoTime); + this(segmentSizeBytes, DEFAULT_POLL_NANOS, UNLIMITED_TOTAL_BYTES); } public SegmentManager(long segmentSizeBytes, long pollNanos) { - this(segmentSizeBytes, pollNanos, UNLIMITED_TOTAL_BYTES, FilesFacade.INSTANCE, System::nanoTime); + this(segmentSizeBytes, pollNanos, UNLIMITED_TOTAL_BYTES); } /** @@ -233,12 +234,22 @@ public SegmentManager(long segmentSizeBytes, long pollNanos) { * hold an initial active plus one hot spare. */ public SegmentManager(long segmentSizeBytes, long pollNanos, long maxTotalBytes) { - this(segmentSizeBytes, pollNanos, maxTotalBytes, FilesFacade.INSTANCE, System::nanoTime); + this(segmentSizeBytes, pollNanos, new SegmentBudget(maxTotalBytes), FilesFacade.INSTANCE, System::nanoTime); + } + + /** + * As {@link #SegmentManager(long, long, long)}, but charging + * {@code budget}, which other managers may share: its capacity then caps + * the rings of all of them together. The capacity must allow at least one + * {@code segmentSizeBytes}. + */ + public SegmentManager(long segmentSizeBytes, long pollNanos, SegmentBudget budget) { + this(segmentSizeBytes, pollNanos, budget, FilesFacade.INSTANCE, System::nanoTime); } @TestOnly public SegmentManager(long segmentSizeBytes, long pollNanos, long maxTotalBytes, FilesFacade filesFacade) { - this(segmentSizeBytes, pollNanos, maxTotalBytes, filesFacade, System::nanoTime); + this(segmentSizeBytes, pollNanos, new SegmentBudget(maxTotalBytes), filesFacade, System::nanoTime); } @TestOnly @@ -248,6 +259,17 @@ public SegmentManager( long maxTotalBytes, FilesFacade filesFacade, LongSupplier ticks + ) { + this(segmentSizeBytes, pollNanos, new SegmentBudget(maxTotalBytes), filesFacade, ticks); + } + + @TestOnly + public SegmentManager( + long segmentSizeBytes, + long pollNanos, + SegmentBudget budget, + FilesFacade filesFacade, + LongSupplier ticks ) { // The pathScratch field initializer has already allocated its native // buffer by the time this body runs, so a validation throw must free @@ -257,16 +279,20 @@ public SegmentManager( pathScratch.close(); throw new IllegalArgumentException("segmentSizeBytes too small: " + segmentSizeBytes); } - if (maxTotalBytes < segmentSizeBytes) { + if (budget == null) { + pathScratch.close(); + throw new IllegalArgumentException("budget must not be null"); + } + if (budget.getCapacityBytes() < segmentSizeBytes) { pathScratch.close(); throw new IllegalArgumentException( - "maxTotalBytes (" + maxTotalBytes + ") must allow at least one segment of " + "maxTotalBytes (" + budget.getCapacityBytes() + ") must allow at least one segment of " + segmentSizeBytes + " bytes"); } + this.budget = budget; this.filesFacade = filesFacade; this.segmentSizeBytes = segmentSizeBytes; this.pollNanos = pollNanos; - this.maxTotalBytes = maxTotalBytes; // Clamp rather than multiply blind: segmentSizeBytes is user-supplied // and only bounded below, so a pathological value would wrap the // product negative and make the floor test trivially false -- silently @@ -576,7 +602,7 @@ public void deregister(SegmentRing ring) { for (int i = 0, n = rings.size(); i < n; i++) { RingEntry e = rings.get(i); if (e.ring == ring) { - // Reverse the ring's contribution to totalBytes — + // Reverse the ring's charge to the budget — // mirrors the seed in register(). Any spares the // manager provisioned during the ring's lifetime // are also part of totalSegmentBytes() now, so a @@ -584,7 +610,10 @@ public void deregister(SegmentRing ring) { // and the net manager activity (provisions minus // trims) for this ring. e.deregister(); - totalBytes -= ring.totalSegmentBytes(); + budget.release(ring.totalSegmentBytes()); + if (e.sideFileBytes != null) { + budget.removeSideFileGauge(e.sideFileBytes); + } rings.remove(i); return; } @@ -636,10 +665,10 @@ public void register(SegmentRing ring, String dir, AckWatermark watermark, long * {@code .sfa} segments; a {@code null} gauge contributes zero (memory * mode, degraded full-dict sessions). *

- * The gauge is invoked with the manager's internal lock held, so it must - * be wait-free and must not throw: a throwing gauge terminates the - * manager worker for every registered ring, and a blocking gauge stalls - * register/deregister for all slots. + * The gauge is invoked under the {@link SegmentBudget}'s monitor, so it + * must be wait-free and must not throw: a throwing gauge terminates the + * worker of every manager charging that budget, and a blocking gauge + * stalls register/deregister and provisioning for all of their slots. */ public void register(SegmentRing ring, String dir, AckWatermark watermark, long syncIntervalNanos, LongSupplier sideFileBytes) { if (syncIntervalNanos < 0L) { @@ -650,7 +679,7 @@ public void register(SegmentRing ring, String dir, AckWatermark watermark, long } // Account for bytes the ring already owns when it joins. A recovered // ring (post-restart, orphan adoption) can come up at-or-above the cap; - // without this seed, totalBytes stays at 0 and the per-tick cap check + // without this seed, the budget stays at 0 and the per-tick cap check // at serviceRing would let the manager keep provisioning new spares on // top of the recovered set, effectively doubling the documented cap. long ringBytes = ring.totalSegmentBytes(); @@ -663,13 +692,25 @@ public void register(SegmentRing ring, String dir, AckWatermark watermark, long Runnable managerWakeup = this::wakeWorker; RingEntry e = new RingEntry(ring, dir, watermark, sideFileBytes, syncIntervalNanos, ticks.getAsLong()); // ObjList.add either throws before storing e or makes the entry visible. - // Once visible, only non-throwing state commits may remain. + // Once visible, only non-throwing state commits may remain, so the gauge + // (whose list can also throw on growth) is added first and withdrawn if + // the entry cannot be published. synchronized (lock) { if (dir != null) { advanceFileGeneration(minNextGeneration); } - rings.add(e); - totalBytes += ringBytes; + if (sideFileBytes != null) { + budget.addSideFileGauge(sideFileBytes); + } + try { + rings.add(e); + } catch (Throwable t) { + if (sideFileBytes != null) { + budget.removeSideFileGauge(sideFileBytes); + } + throw t; + } + budget.charge(ringBytes); if (syncIntervalNanos > 0L) { ring.enablePeriodicSync(); if (syncIntervalNanos < shortestSyncIntervalNanos) { @@ -697,40 +738,14 @@ public SegmentRing getInServiceRingForTesting() { return entry == null ? null : entry.ring; } - // Callers must hold `lock` (the rings list is mutated under it). The - // side-file bytes are read live from each slot's gauge instead of being - // folded into the incremental totalBytes counter: the dictionary grows - // out-of-band on producer threads, so an incremental mirror would - // drift, while a live read cannot. Each gauge is WAIT-FREE and takes no - // lock at all -- PersistedSymbolDict.occupiedDiskBytes() is a pair of - // volatile reads -- and it must stay that way. This runs with `lock` held, - // on the worker that drives provisioning and trim for every registered ring, - // while a producer can hold that dictionary's monitor across ff.allocate - // and mmap. A gauge that took the monitor would park the whole manager - // behind one producer's append I/O. - private long sideFileBytesLocked() { - long total = 0L; - for (int i = 0, n = rings.size(); i < n; i++) { - LongSupplier gauge = rings.get(i).sideFileBytes; - if (gauge != null) { - total += gauge.getAsLong(); - } - } - return total; - } - @TestOnly public long getCapAccountedBytesForTesting() { - synchronized (lock) { - return totalBytes + sideFileBytesLocked(); - } + return budget.getAccountedBytes(); } @TestOnly public long getTotalBytesForTesting() { - synchronized (lock) { - return totalBytes; - } + return budget.getSegmentBytes(); } @TestOnly @@ -958,55 +973,59 @@ private boolean serviceRing0(RingEntry e) { // DISK_FULL_LOG_THROTTLE_NANOS so a sustained-disk-full state // doesn't drown the log. if (e.ring.needsHotSpare()) { - // Snapshot totalBytes under lock -- register/deregister can mutate - // it from caller threads -- and add the live side-file bytes of - // every registered slot, so .symbol-dict growth counts against the - // cap. Heavy provisioning I/O happens outside the lock; the - // post-install commit re-acquires it. - long observedTotal; - long observedSideFileBytes; - synchronized (lock) { - observedSideFileBytes = sideFileBytesLocked(); - observedTotal = totalBytes + observedSideFileBytes; - } - boolean withinCap = observedTotal + segmentSizeBytes <= maxTotalBytes; + // Charge the spare BEFORE provisioning it, as one atomic check-and- + // charge against everything else on the budget, live side-file bytes + // included. A snapshot-then-commit check would let two managers that + // share the budget both see room for its last segment and both take + // it. Installing the spare keeps the charge; every other outcome + // below releases it. + boolean withinCap = budget.tryCharge(segmentSizeBytes); // Liveness floor. Refusing on segment bytes is productive -- an ack // trims a sealed segment and the shortfall clears. Refusing on // side-file bytes is not: the dictionary never shrinks, so a ring // held below its minimum working set by them would never rotate - // again, on this run or any later one. Provision anyway while this - // ring is under the floor, and account it honestly below. + // again, on this run or any later one. Nor is refusing on the + // segments of OTHER rings that share the budget: their acks may + // never come (a sibling whose connection is down). Provision anyway + // while this ring is under the floor, and charge it honestly. long ringSegmentBytes = withinCap ? 0L : e.ring.totalSegmentBytes(); boolean belowLivenessFloor = !withinCap && ringSegmentBytes < livenessFloorBytes; + if (belowLivenessFloor) { + budget.charge(segmentSizeBytes); + } if (!withinCap && !belowLivenessFloor) { long now = System.nanoTime(); if (now - lastDiskFullLogNs >= DISK_FULL_LOG_THROTTLE_NANOS) { + long sideFileBytes = budget.getSideFileBytes(); LOG.warn("SF {}: cannot provision spare in {} " + "(totalBytes={}, sideFileBytes={}, cap={}, segmentSize={}). " + "Producer is backpressured until ACK-driven trim frees segment " + "space; side-file bytes are not reclaimed by trim.", memoryMode ? "memory cap reached" : "disk-full", - memoryMode ? "" : e.dir, observedTotal, observedSideFileBytes, - maxTotalBytes, segmentSizeBytes); + memoryMode ? "" : e.dir, budget.getSegmentBytes() + sideFileBytes, + sideFileBytes, budget.getCapacityBytes(), segmentSizeBytes); lastDiskFullLogNs = now; } } else { if (belowLivenessFloor) { // Exceeding a configured cap is worth saying out loud, and // saying WHY: the operator's remedy is to raise - // sf_max_total_bytes or shrink the symbol dictionary, never - // to wait for a trim. + // sf_max_total_bytes, or shrink the symbol dictionary or the + // number of senders sharing the budget, never to wait for a + // trim. long now = System.nanoTime(); if (now - lastDiskFullLogNs >= DISK_FULL_LOG_THROTTLE_NANOS) { + long sideFileBytes = budget.getSideFileBytes(); LOG.warn("SF {}: provisioning past sf_max_total_bytes to keep the slot " + "usable (totalBytes={}, sideFileBytes={}, cap={}, " + "segmentSize={}, ringSegmentBytes={}, minWorkingSet={}). " - + "The symbol dictionary alone leaves no room for the " - + "active segment plus one spare, and trim cannot reclaim " - + "it; raise sf_max_total_bytes or reduce symbol " - + "cardinality.", - memoryMode ? "" : e.dir, observedTotal, - observedSideFileBytes, maxTotalBytes, segmentSizeBytes, + + "The rest of the budget -- symbol dictionaries, which trim " + + "cannot reclaim, and the segments of other senders sharing " + + "it -- leaves no room for this slot's active segment plus " + + "one spare; raise sf_max_total_bytes, or reduce symbol " + + "cardinality or the number of pooled senders.", + memoryMode ? "" : e.dir, budget.getSegmentBytes() + sideFileBytes, + sideFileBytes, budget.getCapacityBytes(), segmentSizeBytes, ringSegmentBytes, livenessFloorBytes); lastDiskFullLogNs = now; } @@ -1051,22 +1070,21 @@ private boolean serviceRing0(RingEntry e) { "could not sync hot-spare directory " + e.dir); } } - // Install + commit atomically under the manager lock. - // If `e.ring` was deregistered between the snapshot - // above and now, abandoning the spare here is the only - // way to keep totalBytes consistent: deregister already - // subtracted ring.totalSegmentBytes() (without the - // spare, since it wasn't installed yet) so a commit at - // this point would inflate totalBytes by one segment - // with no future subtractor. By holding `lock` across - // installHotSpare AND the += commit AND the registration - // check, deregister is forced to either - // observe the spare in the ring (and subtract it) or - // run before installation (so no install happens). + // Install under the manager lock, gated on the entry + // still being registered. If `e.ring` was deregistered + // since the charge above, abandoning the spare here is + // the only way to keep the budget consistent: + // deregister already released ring.totalSegmentBytes() + // (without the spare, since it wasn't installed yet), so + // installing now would keep the spare's charge with no + // future releaser. By holding `lock` across the + // registration check AND installHotSpare, deregister is + // forced to either observe the spare in the ring (and + // release it) or run before installation (so no install + // happens and the !installed path releases the charge). synchronized (lock) { if (e.isRegistered()) { e.ring.installHotSpare(spare); - totalBytes += segmentSizeBytes; installed = true; } } @@ -1076,6 +1094,7 @@ private boolean serviceRing0(RingEntry e) { memoryMode ? "" : e.dir, t); } if (!installed) { + budget.release(segmentSizeBytes); if (spare != null) { try { spare.close(); @@ -1144,7 +1163,7 @@ private boolean serviceRing0(RingEntry e) { synchronized (lock) { long removedBytes = e.ring.commitPendingTrims(trimBatch, closed); if (e.isRegistered()) { - totalBytes -= removedBytes; + budget.release(removedBytes); } } } @@ -1259,7 +1278,7 @@ private boolean serviceRing0(RingEntry e) { synchronized (lock) { long removedBytes = e.ring.commitPendingTrims(trimBatch, unlinked); if (e.isRegistered()) { - totalBytes -= removedBytes; + budget.release(removedBytes); } } } catch (Throwable t) { @@ -1445,8 +1464,8 @@ private static final class RingEntry { final SegmentRing ring; // Live gauge of the slot's .symbol-dict side-file bytes, or null when // the slot has no dictionary (memory mode, degraded full-dict - // sessions). Read at the provisioning cap check only -- never folded - // into totalBytes, which stays a segments-only counter. + // sessions). Lent to the budget for its cap checks between register + // and deregister -- never folded into its segments-only counter. final LongSupplier sideFileBytes; final long syncIntervalNanos; // Engine-owned ack watermark for this slot, or null in memory diff --git a/core/src/main/java/io/questdb/client/impl/SenderPool.java b/core/src/main/java/io/questdb/client/impl/SenderPool.java index 1a4390269..bc93309d3 100644 --- a/core/src/main/java/io/questdb/client/impl/SenderPool.java +++ b/core/src/main/java/io/questdb/client/impl/SenderPool.java @@ -33,6 +33,7 @@ import io.questdb.client.cutlass.qwp.client.QwpWebSocketSender; import io.questdb.client.cutlass.qwp.client.sf.cursor.BackgroundDrainerListener; import io.questdb.client.cutlass.qwp.client.sf.cursor.OrphanScanner; +import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentBudget; import io.questdb.client.cutlass.qwp.client.sf.cursor.SenderErrorDispatcher; import io.questdb.client.cutlass.qwp.client.sf.cursor.SlotLock; import io.questdb.client.cutlass.qwp.client.sf.cursor.SlotLockContentionException; @@ -71,6 +72,13 @@ * check ({@code allSize + inFlightCreations + closingSlots + leakedSlots < * maxSize}) stays correct under concurrent borrows. *

+ * Shared buffer budget. Every WebSocket sender the pool builds charges + * its cursor segments -- unacknowledged data, in memory or under + * {@code sf_dir} -- to one {@link SegmentBudget} sized from the configured + * {@code sf_max_total_bytes}. That key therefore caps the pool as a whole: + * growing the pool adds connections, not buffer memory. The only overshoot is + * the minimum working set (two segments) each live sender is always granted. + *

* Store-and-forward slots. When the configuration enables SF * ({@code sf_dir} set), every sender owns an exclusive on-disk slot * {@code /} guarded by a {@code flock}. A pool reuses one @@ -193,6 +201,10 @@ public final class SenderPool implements AutoCloseable { // enabled; null otherwise. Each pooled sender's slot id is // {@code slotBaseId + "-" + slotIndex}. private final String slotBaseId; + // The sf_max_total_bytes budget shared by every sender this pool builds, + // live and recovery delegates alike; null for non-WebSocket transports, + // which have no cursor ring. + private final SegmentBudget segmentBudget; // SF group root (sf_dir) when SF is enabled; null otherwise. Used to // locate this pool's own managed slot dirs /- // for startup recovery of unacked data left by a previous run. @@ -543,6 +555,8 @@ private SenderPool( probe.httpTokenProvider(tokenProvider); } this.storeAndForward = probe.isStoreAndForwardEnabled(); + long budgetBytes = probe.getResolvedSfMaxTotalBytes(); + this.segmentBudget = budgetBytes > 0 ? new SegmentBudget(budgetBytes) : null; this.slotBaseId = this.storeAndForward ? probe.getConfiguredSenderId() : null; this.sfDir = this.storeAndForward ? probe.getConfiguredSfDir() : null; this.slotInUse = this.storeAndForward ? new boolean[maxSize] : null; @@ -1388,6 +1402,11 @@ public long getRetiredSlotProbeCountForTesting() { } } + @TestOnly + public SegmentBudget getSegmentBudgetForTesting() { + return segmentBudget; + } + @TestOnly public Thread getStartupRecoveryThreadForTesting() { return startupRecoveryThread; @@ -2043,6 +2062,10 @@ public void onError(SenderError error) { return builder; } + private Sender.LineSenderBuilder applySegmentBudget(Sender.LineSenderBuilder builder) { + return segmentBudget == null ? builder : builder.storeAndForwardSharedBudget(segmentBudget); + } + private Sender.LineSenderBuilder applyTokenProvider(Sender.LineSenderBuilder builder) { if (tokenProvider != null) { builder.httpTokenProvider(tokenProvider); @@ -2078,7 +2101,7 @@ private static boolean isRecoveryEventUserRelevant(SenderError e) { private Sender buildManagedSlotSender(int slotIndex, boolean forRecovery) { if (!storeAndForward) { - return applyUserCallbacks(applyTokenProvider(Sender.builder(configurationString))).build(); + return applyUserCallbacks(applyTokenProvider(applySegmentBudget(Sender.builder(configurationString)))).build(); } // Give this pooled sender its own slot dir /- // so concurrent SF senders sharing one sf_dir never collide on @@ -2101,7 +2124,7 @@ private Sender buildManagedSlotSender(int slotIndex, boolean forRecovery) { // per-sender drainer is an additional path that only runs when // drain_orphans=on; foreign leftovers under other names are drained // only by that path. - Sender.LineSenderBuilder builder = Sender.builder(configurationString) + Sender.LineSenderBuilder builder = applySegmentBudget(Sender.builder(configurationString)) .senderId(slotBaseId + "-" + slotIndex) .orphanDrainExcludeManagedSlots(slotBaseId, maxSize); if (forRecovery) { diff --git a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/CursorSendEngineTest.java b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/CursorSendEngineTest.java index 998b6e585..c93de6902 100644 --- a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/CursorSendEngineTest.java +++ b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/CursorSendEngineTest.java @@ -33,6 +33,7 @@ import io.questdb.client.cutlass.qwp.client.GlobalSymbolDictionary; import io.questdb.client.cutlass.qwp.client.sf.cursor.MmapSegment; import io.questdb.client.cutlass.qwp.client.sf.cursor.MmapSegmentException; +import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentBudget; import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentManager; import io.questdb.client.cutlass.qwp.client.sf.cursor.SlotLock; import io.questdb.client.std.Files; @@ -994,12 +995,13 @@ private static void assertSlotCanBeReacquired(String sfDir) { private static Throwable invokeOwnedPrivateConstructorExpectingFailure( String sfDir, long segmentSizeBytes, SegmentManager manager) throws Exception { Constructor ctor = CursorSendEngine.class.getDeclaredConstructor( - String.class, long.class, SegmentManager.class, boolean.class, long.class, + String.class, long.class, SegmentManager.class, boolean.class, SegmentBudget.class, long.class, long.class, FilesFacade.class); ctor.setAccessible(true); try { - ctor.newInstance(sfDir, segmentSizeBytes, manager, true, - SegmentManager.UNLIMITED_TOTAL_BYTES, + // The budget is only consulted when the engine builds its own manager; + // here the (owned) manager is supplied, so none is needed. + ctor.newInstance(sfDir, segmentSizeBytes, manager, true, null, CursorSendEngine.DEFAULT_APPEND_DEADLINE_NANOS, 0L, FilesFacade.INSTANCE); fail("expected constructor failure"); return null; diff --git a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoopGracefulCloseTest.java b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoopGracefulCloseTest.java new file mode 100644 index 000000000..ba947f9d7 --- /dev/null +++ b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoopGracefulCloseTest.java @@ -0,0 +1,238 @@ +/*+***************************************************************************** + * ___ _ ____ ____ + * / _ \ _ _ ___ ___| |_| _ \| __ ) + * | | | | | | |/ _ \/ __| __| | | | _ \ + * | |_| | |_| | __/\__ \ |_| |_| | |_) | + * \__\_\\__,_|\___||___/\__|____/|____/ + * + * Copyright (c) 2014-2019 Appsicle + * Copyright (c) 2019-2026 QuestDB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + ******************************************************************************/ + +package io.questdb.client.test.cutlass.qwp.client.sf.cursor; + +import io.questdb.client.DefaultHttpClientConfiguration; +import io.questdb.client.cutlass.http.client.WebSocketClient; +import io.questdb.client.cutlass.http.client.WebSocketFrameHandler; +import io.questdb.client.cutlass.line.LineSenderException; +import io.questdb.client.cutlass.qwp.client.sf.cursor.CursorSendEngine; +import io.questdb.client.cutlass.qwp.client.sf.cursor.CursorWebSocketSendLoop; +import io.questdb.client.network.PlainSocketFactory; +import io.questdb.client.test.tools.TestUtils; +import org.junit.Assert; +import org.junit.Test; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +/** + * Pins the order of {@link CursorWebSocketSendLoop#close()}: an I/O thread that + * is not stuck must stop on its own and close its client -- the WebSocket CLOSE + * frame and TLS close_notify go out on that path -- before close() would break + * the connection's traffic. Breaking traffic first made the server see every + * sender close as a dropped connection. A connect blocked on cancellable work + * must still be cancelled at once, not after the graceful-stop window. + */ +public class CursorWebSocketSendLoopGracefulCloseTest { + + @Test(timeout = 30_000L) + public void testCloseCancelsInFlightConnectWithoutWaitingOutGracefulStop() throws Exception { + TestUtils.assertMemoryLeak(() -> { + final AtomicReference published = new AtomicReference<>(); + final CountDownLatch connectEntered = new CountDownLatch(1); + final CursorWebSocketSendLoop.ReconnectFactory factory = new CursorWebSocketSendLoop.ReconnectFactory() { + @Override + public WebSocketClient reconnect() { + return reconnect(null); + } + + @Override + public WebSocketClient reconnect(CursorWebSocketSendLoop.ConnectCancellation cancellation) { + InFlightConnectClient c = new InFlightConnectClient(); + published.set(c); + if (cancellation != null) { + cancellation.publish(c); + } + connectEntered.countDown(); + c.awaitTrafficBreak(); + c.close(); + throw new LineSenderException("connect cancelled by closeTraffic()"); + } + }; + final CursorSendEngine engine = new CursorSendEngine(null, 64 * 1024); + final CursorWebSocketSendLoop loop = new CursorWebSocketSendLoop( + null, + engine, + 0L, + CursorWebSocketSendLoop.DEFAULT_PARK_NANOS, + factory, + 1_000L, + 5_000L, + false + ); + // Far beyond the test timeout: waiting it out here would fail the test. + loop.setGracefulStopMillis(TimeUnit.MINUTES.toMillis(5)); + try { + loop.start(); + Assert.assertTrue("I/O worker never entered the blocking connect", + connectEntered.await(5, TimeUnit.SECONDS)); + + long startNanos = System.nanoTime(); + loop.close(); + long elapsedMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos); + + Assert.assertEquals("close() must cancel the in-flight connect exactly once", + 1, published.get().trafficCloseCount.get()); + Assert.assertTrue("close() waited out the graceful-stop window for a connect it had to cancel, took " + + elapsedMillis + "ms", elapsedMillis < TimeUnit.SECONDS.toMillis(10)); + Assert.assertNull("ordinary connect-phase close must not manufacture a terminal error", + loop.getTerminalError()); + } finally { + InFlightConnectClient inFlight = published.get(); + if (inFlight != null) { + inFlight.trafficBroken.countDown(); + } + loop.close(); + engine.close(); + } + }); + } + + @Test(timeout = 30_000L) + public void testCloseLetsIdleWorkerCloseClientBeforeBreakingTraffic() throws Exception { + TestUtils.assertMemoryLeak(() -> { + final IdleClient client = new IdleClient(); + final CursorSendEngine engine = new CursorSendEngine(null, 64 * 1024); + final CursorWebSocketSendLoop loop = new CursorWebSocketSendLoop( + client, + engine, + 0L, + CursorWebSocketSendLoop.DEFAULT_PARK_NANOS, + null, + 1_000L, + 5_000L, + false + ); + // Far above any scheduling delay, so a loaded machine cannot flip the outcome. + loop.setGracefulStopMillis(TimeUnit.SECONDS.toMillis(20)); + try { + loop.start(); + Assert.assertTrue("I/O worker never polled for ACKs", client.polled.await(5, TimeUnit.SECONDS)); + + loop.close(); + + Assert.assertEquals("close() broke the traffic of an idle worker instead of letting it close cleanly", + 0, client.trafficCloseCount.get()); + Assert.assertNotNull("the client was never closed", client.firstCloseThread.get()); + Assert.assertEquals("the I/O worker must close its own client on the way out", + client.pollThread.get(), client.firstCloseThread.get()); + Assert.assertNull("ordinary close must not manufacture a terminal error", loop.getTerminalError()); + } finally { + loop.close(); + engine.close(); + client.close(); + } + }); + } + + /** + * Connected client with nothing to receive. Records which thread closes it + * first and whether anyone broke its traffic. + */ + private static final class IdleClient extends WebSocketClient { + final AtomicReference firstCloseThread = new AtomicReference<>(); + final AtomicReference pollThread = new AtomicReference<>(); + final CountDownLatch polled = new CountDownLatch(1); + final AtomicInteger trafficCloseCount = new AtomicInteger(); + + private IdleClient() { + super(DefaultHttpClientConfiguration.INSTANCE, PlainSocketFactory.INSTANCE); + } + + @Override + public void close() { + firstCloseThread.compareAndSet(null, Thread.currentThread()); + super.close(); + } + + @Override + public void closeTraffic() { + trafficCloseCount.incrementAndGet(); + } + + @Override + public boolean tryReceiveFrame(WebSocketFrameHandler handler) { + pollThread.compareAndSet(null, Thread.currentThread()); + polled.countDown(); + return false; + } + + @Override + protected void ioWait(int timeout, int op) { + throw new UnsupportedOperationException("stub: no socket"); + } + + @Override + protected void setupIoWait() { + // no-op + } + } + + /** + * Fake in-flight connect: parks until {@link #closeTraffic()} releases it, + * like a native connect that neither unpark nor interrupt can cancel. + */ + private static final class InFlightConnectClient extends WebSocketClient { + final CountDownLatch trafficBroken = new CountDownLatch(1); + final AtomicInteger trafficCloseCount = new AtomicInteger(); + + private InFlightConnectClient() { + super(DefaultHttpClientConfiguration.INSTANCE, PlainSocketFactory.INSTANCE); + } + + @Override + public void closeTraffic() { + trafficCloseCount.incrementAndGet(); + trafficBroken.countDown(); + } + + @Override + protected void ioWait(int timeout, int op) { + throw new UnsupportedOperationException("stub: no socket"); + } + + @Override + protected void setupIoWait() { + // no-op + } + + void awaitTrafficBreak() { + boolean interrupted = false; + while (trafficBroken.getCount() != 0L) { + try { + trafficBroken.await(); + } catch (InterruptedException e) { + interrupted = true; + } + } + if (interrupted) { + Thread.currentThread().interrupt(); + } + } + } +} diff --git a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentManagerSharedBudgetTest.java b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentManagerSharedBudgetTest.java new file mode 100644 index 000000000..7efcc48ba --- /dev/null +++ b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentManagerSharedBudgetTest.java @@ -0,0 +1,231 @@ +/******************************************************************************* + * ___ _ ____ ____ + * / _ \ _ _ ___ ___| |_| _ \| __ ) + * | | | | | | |/ _ \/ __| __| | | | _ \ + * | |_| | |_| | __/\__ \ |_| |_| | |_) | + * \__\_\\__,_|\___||___/\__|____/|____/ + * + * Copyright (c) 2014-2019 Appsicle + * Copyright (c) 2019-2026 QuestDB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + ******************************************************************************/ + +package io.questdb.client.test.cutlass.qwp.client.sf.cursor; + +import io.questdb.client.cutlass.qwp.client.sf.cursor.MmapSegment; +import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentBudget; +import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentManager; +import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentRing; +import io.questdb.client.std.MemoryTag; +import io.questdb.client.std.Unsafe; +import io.questdb.client.test.tools.TestUtils; +import org.junit.Assert; +import org.junit.Test; + +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BooleanSupplier; + +/** + * Several {@link SegmentManager}s charging one {@link SegmentBudget}: what a + * sender pool does so that {@code sf_max_total_bytes} caps the pool as a whole + * instead of each pooled sender. + */ +public class SegmentManagerSharedBudgetTest { + + private static final int FRAME_PAYLOAD = 64; + private static final long POLL_NANOS = 200_000L; + // Exactly two frames per segment, so the rotation arithmetic below is exact. + private static final long SEGMENT_SIZE = MmapSegment.HEADER_SIZE + + 2 * (MmapSegment.FRAME_HEADER_SIZE + FRAME_PAYLOAD); + + @Test + public void testConcurrentManagersStayWithinBudgetAndReleaseEveryCharge() throws Exception { + TestUtils.assertMemoryLeak(() -> { + // Producers on four managers race for one budget while acks trim behind + // them. tryCharge must never let the budget exceed its capacity by more + // than the minimum working sets the managers grant unconditionally, and + // once every ring is deregistered nothing may remain charged. + final int ringCount = 4; + final SegmentBudget budget = new SegmentBudget(6 * SEGMENT_SIZE); + final long bound = budget.getCapacityBytes() + ringCount * 2 * SEGMENT_SIZE; + final SegmentRing[] rings = new SegmentRing[ringCount]; + final SegmentManager[] managers = new SegmentManager[ringCount]; + final Thread[] producers = new Thread[ringCount]; + final AtomicBoolean stop = new AtomicBoolean(); + final AtomicReference failure = new AtomicReference<>(); + try { + for (int i = 0; i < ringCount; i++) { + rings[i] = newRing(); + managers[i] = new SegmentManager(SEGMENT_SIZE, POLL_NANOS, budget); + managers[i].start(); + managers[i].register(rings[i], null); + } + for (int i = 0; i < ringCount; i++) { + final SegmentRing ring = rings[i]; + final int lag = i; + producers[i] = new Thread(() -> { + long buf = Unsafe.malloc(FRAME_PAYLOAD, MemoryTag.NATIVE_DEFAULT); + try { + while (!stop.get()) { + long fsn = ring.appendOrFsn(buf, FRAME_PAYLOAD); + if (fsn >= 0) { + // Ack with a per-ring lag so the rings trim, and + // therefore release, at different rates. + if (fsn > lag) { + ring.acknowledge(fsn - lag); + } + } else { + Assert.assertEquals(SegmentRing.BACKPRESSURE_NO_SPARE, fsn); + Thread.yield(); + } + } + } catch (Throwable t) { + failure.compareAndSet(null, t); + } finally { + Unsafe.free(buf, FRAME_PAYLOAD, MemoryTag.NATIVE_DEFAULT); + } + }); + producers[i].start(); + } + long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(500); + while (System.nanoTime() < deadline && failure.get() == null) { + long charged = budget.getSegmentBytes(); + if (charged > bound) { + Assert.fail("budget overshot its capacity by more than the minimum working sets: " + + charged + " > " + bound); + } + } + } finally { + stop.set(true); + for (Thread producer : producers) { + if (producer != null) { + producer.join(); + } + } + for (int i = 0; i < ringCount; i++) { + if (managers[i] != null && rings[i] != null) { + managers[i].deregister(rings[i]); + if (!managers[i].awaitRingQuiescence(rings[i])) { + failure.compareAndSet(null, new AssertionError("ring " + i + " did not quiesce")); + } + } + } + for (int i = 0; i < ringCount; i++) { + if (managers[i] != null) { + managers[i].close(); + } + if (rings[i] != null) { + rings[i].close(); + } + } + } + if (failure.get() != null) { + throw new AssertionError("producer or teardown failed", failure.get()); + } + Assert.assertEquals("every charge must be released once all rings are deregistered", + 0, budget.getSegmentBytes()); + }); + } + + @Test + public void testRingsOfDifferentManagersShareOneBudget() throws Exception { + TestUtils.assertMemoryLeak(() -> { + SegmentBudget budget = new SegmentBudget(4 * SEGMENT_SIZE); + long buf = Unsafe.malloc(FRAME_PAYLOAD, MemoryTag.NATIVE_DEFAULT); + try (SegmentRing busy = newRing(); + SegmentRing sibling = newRing(); + SegmentManager busyManager = new SegmentManager(SEGMENT_SIZE, POLL_NANOS, budget); + SegmentManager siblingManager = new SegmentManager(SEGMENT_SIZE, POLL_NANOS, budget)) { + busyManager.start(); + siblingManager.start(); + + // The busy ring grows to the whole budget: three segments of frames + // plus the active one -- eight frames -- and no further spare. + busyManager.register(busy, null); + appendFrames(busy, buf, 8); + assertCapped(busy, buf); + Assert.assertEquals(4 * SEGMENT_SIZE, busy.totalSegmentBytes()); + Assert.assertEquals(4 * SEGMENT_SIZE, budget.getSegmentBytes()); + + // A ring on ANOTHER manager finds the budget spent. It still gets its + // minimum working set -- the active segment plus one spare, four + // frames -- but nothing beyond: a private budget of the same size + // would have let it grow to four segments of its own. + siblingManager.register(sibling, null); + appendFrames(sibling, buf, 4); + assertCapped(sibling, buf); + Assert.assertEquals(2 * SEGMENT_SIZE, sibling.totalSegmentBytes()); + Assert.assertEquals("the budget may be exceeded only by the sibling's minimum working set", + 6 * SEGMENT_SIZE, budget.getSegmentBytes()); + + // Once the busy ring gives its segments back, the sibling grows into + // the room they leave -- up to the whole budget, and no further. + busyManager.deregister(busy); + Assert.assertTrue(busyManager.awaitRingQuiescence(busy)); + Assert.assertTrue("the sibling must get a spare once the budget has room", + waitFor(() -> !sibling.needsHotSpare())); + Assert.assertEquals(3 * SEGMENT_SIZE, budget.getSegmentBytes()); + appendFrames(sibling, buf, 4); + assertCapped(sibling, buf); + Assert.assertEquals(4 * SEGMENT_SIZE, sibling.totalSegmentBytes()); + Assert.assertEquals(4 * SEGMENT_SIZE, budget.getSegmentBytes()); + + siblingManager.deregister(sibling); + Assert.assertTrue(siblingManager.awaitRingQuiescence(sibling)); + Assert.assertEquals(0, budget.getSegmentBytes()); + } finally { + Unsafe.free(buf, FRAME_PAYLOAD, MemoryTag.NATIVE_DEFAULT); + } + }); + } + + // Appends `frames` frames, waiting out the moments a rotation finds the + // spare not yet provisioned. Every frame is expected to fit eventually. + private static void appendFrames(SegmentRing ring, long buf, int frames) throws InterruptedException { + for (int i = 0; i < frames; i++) { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + long fsn; + while ((fsn = ring.appendOrFsn(buf, FRAME_PAYLOAD)) == SegmentRing.BACKPRESSURE_NO_SPARE) { + Assert.assertTrue("no spare arrived for frame " + i, System.nanoTime() < deadline); + Thread.sleep(1); + } + Assert.assertTrue("append failed: " + fsn, fsn >= 0); + } + } + + private static void assertCapped(SegmentRing ring, long buf) throws InterruptedException { + // Ample time for a manager worker to (wrongly) provision past the budget. + Thread.sleep(100); + Assert.assertTrue("the budget must refuse this ring another spare", ring.needsHotSpare()); + Assert.assertEquals(SegmentRing.BACKPRESSURE_NO_SPARE, ring.appendOrFsn(buf, FRAME_PAYLOAD)); + } + + private static SegmentRing newRing() { + return new SegmentRing(MmapSegment.createInMemory(0, SEGMENT_SIZE), SEGMENT_SIZE); + } + + private static boolean waitFor(BooleanSupplier condition) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + while (!condition.getAsBoolean()) { + if (System.nanoTime() > deadline) { + return false; + } + Thread.sleep(1); + } + return true; + } +} diff --git a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentManagerUnlinkFailureTest.java b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentManagerUnlinkFailureTest.java index b8c23901e..b4b60d6fb 100644 --- a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentManagerUnlinkFailureTest.java +++ b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentManagerUnlinkFailureTest.java @@ -39,7 +39,6 @@ import org.junit.Before; import org.junit.Test; -import java.lang.reflect.Field; import java.nio.file.Paths; import java.util.Arrays; import java.util.concurrent.CountDownLatch; @@ -136,7 +135,7 @@ public void testFailedUnlinkRetainsBookkeepingAndUsesSuccessorPath() throws Exce Assert.assertTrue("manager never retried the injected unlink", facade.removeRetried.await(5, TimeUnit.SECONDS)); Assert.assertEquals("failed unlink must retain conservative registered bytes", - ring.totalSegmentBytes(), readTotalBytes(manager)); + ring.totalSegmentBytes(), manager.getTotalBytesForTesting()); } Assert.assertTrue("failed unlink path must remain observable", Files.exists(failedPath)); @@ -211,17 +210,6 @@ private static void fill(long address, int len, byte value) { } } - private static long readTotalBytes(SegmentManager manager) throws Exception { - Field field = SegmentManager.class.getDeclaredField("totalBytes"); - field.setAccessible(true); - Field lockField = SegmentManager.class.getDeclaredField("lock"); - lockField.setAccessible(true); - Object lock = lockField.get(manager); - synchronized (lock) { - return field.getLong(manager); - } - } - private static void removeRecursive(String dir) { long find = Files.findFirst(dir); if (find > 0) { diff --git a/core/src/test/java/io/questdb/client/test/impl/SenderPoolSharedBudgetTest.java b/core/src/test/java/io/questdb/client/test/impl/SenderPoolSharedBudgetTest.java new file mode 100644 index 000000000..184ca8696 --- /dev/null +++ b/core/src/test/java/io/questdb/client/test/impl/SenderPoolSharedBudgetTest.java @@ -0,0 +1,182 @@ +/******************************************************************************* + * ___ _ ____ ____ + * / _ \ _ _ ___ ___| |_| _ \| __ ) + * | | | | | | |/ _ \/ __| __| | | | _ \ + * | |_| | |_| | __/\__ \ |_| |_| | |_) | + * \__\_\\__,_|\___||___/\__|____/|____/ + * + * Copyright (c) 2014-2019 Appsicle + * Copyright (c) 2019-2026 QuestDB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + ******************************************************************************/ + +package io.questdb.client.test.impl; + +import io.questdb.client.Sender; +import io.questdb.client.cutlass.line.LineSenderException; +import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentBudget; +import io.questdb.client.impl.SenderPool; +import io.questdb.client.std.MemoryTag; +import io.questdb.client.std.Unsafe; +import io.questdb.client.test.cutlass.qwp.websocket.TestWebSocketServer; +import io.questdb.client.test.tools.TestUtils; +import org.junit.Assert; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.TemporaryFolder; + +import java.util.concurrent.TimeUnit; + +/** + * Every sender a pool builds charges one {@code sf_max_total_bytes} budget, so + * the key caps the pool rather than each pooled sender: growing the pool adds + * connections, not buffer memory. Before, each pooled sender buffered up to the + * whole cap on its own, and a server that stopped acknowledging grew the + * client's native memory by the cap times the pool size. + */ +public class SenderPoolSharedBudgetTest { + + private static final long SEGMENT_BYTES = 64 * 1024; + private static final long BUDGET_BYTES = 16 * SEGMENT_BYTES; + // Every live sender is always granted its active segment plus one spare, + // budget or not. + private static final long MIN_WORKING_SET_BYTES = 2 * SEGMENT_BYTES; + private static final int POOL_SIZE = 4; + + @Rule + public final TemporaryFolder temp = TemporaryFolder.builder().assureDeletion().build(); + + @Test + public void testBudgetIsSizedFromTheResolvedSfMaxTotalBytes() throws Exception { + TestUtils.assertMemoryLeak(() -> { + // min=0 builds no sender, so nothing connects to these addresses. + try (SenderPool pool = newIdlePool("ws::addr=127.0.0.1:1;")) { + Assert.assertEquals("memory-mode default", + 128L * 1024 * 1024, pool.getSegmentBudgetForTesting().getCapacityBytes()); + } + try (SenderPool pool = newIdlePool("ws::addr=127.0.0.1:1;sf_max_total_bytes=5m;")) { + Assert.assertEquals(5L * 1024 * 1024, pool.getSegmentBudgetForTesting().getCapacityBytes()); + } + String sfDir = temp.getRoot().toPath().resolve("sf").toString(); + try (SenderPool pool = newIdlePool("ws::addr=127.0.0.1:1;sf_dir=" + sfDir + ";")) { + Assert.assertEquals("store-and-forward default", + 10L * 1024 * 1024 * 1024, pool.getSegmentBudgetForTesting().getCapacityBytes()); + } + try (SenderPool pool = newIdlePool("http::addr=127.0.0.1:1;protocol_version=2;")) { + Assert.assertNull("an HTTP sender has no segment ring to budget", + pool.getSegmentBudgetForTesting()); + } + }); + } + + @Test + public void testPooledSendersShareOneBudget() throws Exception { + TestUtils.assertMemoryLeak(() -> { + // The server never acknowledges, so everything the senders write stays + // buffered client-side -- the outage that used to cost the cap times + // the pool size. + try (TestWebSocketServer server = new TestWebSocketServer(new TestWebSocketServer.WebSocketServerHandler() { + })) { + server.start(); + Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS)); + String cfg = "ws::addr=localhost:" + server.getPort() + ";" + + "sf_max_segment_bytes=" + SEGMENT_BYTES + ";" + + "sf_max_total_bytes=" + BUDGET_BYTES + ";" + + "sf_append_deadline_millis=200;" + + "close_flush_timeout_millis=100;"; + SegmentBudget budget; + try (SenderPool pool = new SenderPool(cfg, POOL_SIZE, POOL_SIZE, 5_000, Long.MAX_VALUE, Long.MAX_VALUE)) { + budget = pool.getSegmentBudgetForTesting(); + Assert.assertEquals(BUDGET_BYTES, budget.getCapacityBytes()); + Sender[] senders = new Sender[POOL_SIZE]; + try { + for (int i = 0; i < POOL_SIZE; i++) { + senders[i] = pool.borrow(); + } + // Every sender is built and holds its starting segments, so + // from here on native memory grows only by buffered data. + long nativeBefore = Unsafe.getMemUsedByTag(MemoryTag.NATIVE_DEFAULT); + long bufferedBefore = budget.getSegmentBytes(); + for (int i = 0; i < POOL_SIZE; i++) { + writeUntilBackpressured(senders[i]); + } + // Measured independently of the budget's own accounting: with a + // cap per sender this grew by about POOL_SIZE budgets. + long nativeGrowth = Unsafe.getMemUsedByTag(MemoryTag.NATIVE_DEFAULT) - nativeBefore; + Assert.assertTrue("native memory must grow by no more than the rest of the one budget " + + "[growth=" + nativeGrowth + ", bufferedBefore=" + bufferedBefore + ']', + nativeGrowth <= BUDGET_BYTES - bufferedBefore + 4 * SEGMENT_BYTES); + long buffered = budget.getSegmentBytes(); + Assert.assertTrue("the senders together must fill the budget [buffered=" + buffered + ']', + buffered >= BUDGET_BYTES - SEGMENT_BYTES); + // With a cap per sender, the first sender alone would have + // buffered the whole budget and each of the others up to + // another one: POOL_SIZE budgets in all. + Assert.assertTrue("the pool must buffer at most one budget plus each sender's minimum " + + "working set [buffered=" + buffered + ']', + buffered <= BUDGET_BYTES + POOL_SIZE * MIN_WORKING_SET_BYTES); + } finally { + for (Sender sender : senders) { + closeQuietly(sender); + } + } + } + Assert.assertEquals("closing the pool must release every charge", 0, budget.getSegmentBytes()); + } + }); + } + + private static void closeQuietly(Sender sender) { + if (sender == null) { + return; + } + try { + // The ring is still full, so the flush in close() is backpressured too + // and the pool discards the sender instead of recycling it. + sender.close(); + } catch (LineSenderException ignored) { + } + } + + private static SenderPool newIdlePool(String cfg) { + return new SenderPool(cfg, 0, POOL_SIZE, 1_000, Long.MAX_VALUE, Long.MAX_VALUE); + } + + // Writes ~16 KiB batches until the sender's flush gives up waiting for room + // in its segment ring. + private static void writeUntilBackpressured(Sender sender) { + StringBuilder payload = new StringBuilder(1024); + for (int i = 0; i < 1024; i++) { + payload.append((char) ('a' + i % 26)); + } + try { + for (int batch = 0; batch < 10_000; batch++) { + for (int row = 0; row < 16; row++) { + sender.table("budget").stringColumn("payload", payload).longColumn("row", row).atNow(); + } + sender.flush(); + } + } catch (LineSenderException e) { + // The append failure arrives wrapped; its cause names the deadline. + for (Throwable t = e; t != null; t = t.getCause()) { + if (t.getMessage() != null && t.getMessage().contains("backpressured")) { + return; + } + } + throw new AssertionError("expected a backpressure failure", e); + } + Assert.fail("the sender never ran out of room to buffer"); + } +}