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");
+ }
+}