From 01fbd437a7af7f0b585a389245fbd14be92ac62f Mon Sep 17 00:00:00 2001 From: Igor Melnichenko Date: Thu, 17 Sep 2026 16:33:42 +0300 Subject: [PATCH 1/2] Guard decoder ready handler and validate buffer size MessageDecoder accepted a non-positive maxBufferSize. Since DecodeNext only admits messages while totalAvailable is positive, such a decoder never submits anything and silently stalls the reader instead of failing fast. Construction now rejects it with IllegalArgumentException. ReadPartitionDecoder invoked readyHandler directly from EncodedMessage.decode() and from setError(). An exception thrown by the handler escaped into the decompression executor thread, where it is reported as an uncaught exception and may suppress decoding of the messages that follow. The handler is now called through notifyReady(), which logs the failure after the message has already been marked as ready. Co-Authored-By: Claude Opus 5 --- .../ydb/topic/read/impl/MessageDecoder.java | 4 ++ .../topic/read/impl/ReadPartitionDecoder.java | 12 +++- .../topic/read/impl/MessageDecoderTest.java | 58 +++++++++++++++++++ 3 files changed, 72 insertions(+), 2 deletions(-) diff --git a/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java b/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java index 8357af332..1a9d5aefe 100644 --- a/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java +++ b/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java @@ -26,6 +26,10 @@ public class MessageDecoder { private volatile boolean isStopped = false; public MessageDecoder(long maxBufferSize, Executor decompressionExecutor, CodecRegistry codecRegistry) { + if (maxBufferSize <= 0) { + throw new IllegalArgumentException("maxBufferSize must be positive, but got " + maxBufferSize); + } + this.totalAvailable = new AtomicLong(maxBufferSize); this.decompressionExecutor = decompressionExecutor; this.codecRegistry = codecRegistry; diff --git a/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java b/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java index 3e526e21f..e6070e2b1 100644 --- a/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java +++ b/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java @@ -91,6 +91,14 @@ private void release() { decoder.free(allocatedTotal.getAndSet(0)); } + private void notifyReady() { + try { + readyHandler.run(); + } catch (Throwable th) { + logger.error("[{}] Exception was thrown by the ready handler", traceID, th); + } + } + public class EncodedMessage extends MessageImpl { private final int codecCode; private final long uncompressedSize; @@ -125,7 +133,7 @@ public void setError(Throwable th) { problem = new IOException("Decompression for " + getPartitionSession() + " error", th); releaseRange(OffsetsRange.of(getOffset())); isReady = true; - readyHandler.run(); + notifyReady(); } public long allocate() { @@ -172,7 +180,7 @@ public void decode(CodecRegistry registry) { data = null; } isReady = true; - readyHandler.run(); + notifyReady(); } } } diff --git a/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java b/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java index 47ae01876..02d5da515 100644 --- a/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java +++ b/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java @@ -80,6 +80,64 @@ private static BatchMeta meta(int codec) { @Rule public final HideLoggersRule hideLogger = new HideLoggersRule(); + @Test + public void nonPositiveBufferSizeTest() { + // A decoder with a non-positive budget can never admit a message and silently stalls the reader + Assert.assertThrows(IllegalArgumentException.class, + () -> new MessageDecoder(0, Runnable::run, REGISTRY)); + Assert.assertThrows(IllegalArgumentException.class, + () -> new MessageDecoder(-1, Runnable::run, REGISTRY)); + } + + @Test + @HideLoggers(MessageDecoder.class) + public void readyHandlerThrowsOnDecodeTest() { + MessageDecoder decoder = new MessageDecoder(1000, Runnable::run, REGISTRY); + + AtomicInteger ready = new AtomicInteger(0); + ReadPartitionDecoder partition = new ReadPartitionDecoder("t1", decoder, PS1, null, () -> { + ready.incrementAndGet(); + throw new RuntimeException("ready handler is broken"); + }); + + BatchMeta meta = meta(Codec.GZIP); + MessageImpl m1 = partition.decode(meta, OffsetsRange.of(1), gzipMsg(1, 40)); + MessageImpl m2 = partition.decode(meta, OffsetsRange.of(2), gzipMsg(2, 50)); + + decoder.decodeNext(); + + // A broken handler must not escape into the decompression thread and must not stop the following messages + Assert.assertEquals(2, ready.get()); + Assert.assertTrue(m1.isReady()); + Assert.assertTrue(m2.isReady()); + Assert.assertEquals(40, m1.getData().length); + Assert.assertEquals(50, m2.getData().length); + } + + @Test + @HideLoggers(MessageDecoder.class) + public void readyHandlerThrowsOnErrorTest() { + Executor rejecting = task -> { + throw new RejectedExecutionException("executor is saturated"); + }; + MessageDecoder decoder = new MessageDecoder(1000, rejecting, REGISTRY); + + AtomicInteger ready = new AtomicInteger(0); + ReadPartitionDecoder partition = new ReadPartitionDecoder("t1", decoder, PS1, null, () -> { + ready.incrementAndGet(); + throw new RuntimeException("ready handler is broken"); + }); + + BatchMeta meta = meta(Codec.GZIP); + MessageImpl m1 = partition.decode(meta, OffsetsRange.of(1), gzipMsg(1, 40)); + + decoder.decodeNext(); + + Assert.assertEquals(1, ready.get()); + Assert.assertTrue(m1.isReady()); + assertDecompressionException("Decompression for " + PS1 + " error", m1::getData); + } + @Test public void rawDecodeTest() { MessageDecoder decoder = new MessageDecoder(10000, Runnable::run, REGISTRY); From aa48da94df9dd7dcd51647cdfabc6dcc84277554 Mon Sep 17 00:00:00 2001 From: Igor Melnichenko Date: Thu, 17 Sep 2026 16:39:36 +0300 Subject: [PATCH 2/2] Make the decompression buffer limit a hard limit DecodeNext admitted the next message whenever any budget remained, without looking at its size. A message was therefore submitted even when it did not fit, so totalAvailable went negative and the decoder held more uncompressed data than maxMemoryUsageBytes allows. A message is now admitted only when it fits into the remaining budget. A message larger than the whole budget would never fit, so it is still admitted, but only once nothing else retains the buffer, which keeps the overshoot bounded by that single message instead of letting it accumulate. The existing tests asserted the previous behavior through the negative values of getTotalAvailable(), so their expectations are updated accordingly. Co-Authored-By: Claude Opus 5 --- .../ydb/topic/read/impl/MessageDecoder.java | 20 +++- .../topic/read/impl/ReadPartitionDecoder.java | 4 + .../topic/read/impl/MessageDecoderTest.java | 93 ++++++++++++------- 3 files changed, 84 insertions(+), 33 deletions(-) diff --git a/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java b/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java index 1a9d5aefe..cea7c47a6 100644 --- a/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java +++ b/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java @@ -17,6 +17,7 @@ */ public class MessageDecoder { private static final Logger logger = LoggerFactory.getLogger(MessageDecoder.class); + private final long maxBufferSize; private final AtomicLong totalAvailable; private final Executor decompressionExecutor; @@ -30,6 +31,7 @@ public MessageDecoder(long maxBufferSize, Executor decompressionExecutor, CodecR throw new IllegalArgumentException("maxBufferSize must be positive, but got " + maxBufferSize); } + this.maxBufferSize = maxBufferSize; this.totalAvailable = new AtomicLong(maxBufferSize); this.decompressionExecutor = decompressionExecutor; this.codecRegistry = codecRegistry; @@ -65,12 +67,26 @@ void free(long bufferSize) { private final class DecodeNext implements Runnable { @Override public void run() { - while (!isStopped && totalAvailable.get() > 0) { - ReadPartitionDecoder.EncodedMessage next = decodingQueue.poll(); + while (!isStopped) { + long available = totalAvailable.get(); + if (available <= 0) { + return; + } + + // Only this runnable polls the queue and it is serialized, so peek() cannot be invalidated here + ReadPartitionDecoder.EncodedMessage next = decodingQueue.peek(); if (next == null) { return; } + // A message larger than the whole budget must still make progress, but only when nothing else + // retains the buffer. Otherwise it waits until the already admitted messages are released. + if (next.getUncompressedSize() > available && available != maxBufferSize) { + return; + } + + decodingQueue.poll(); + long size = next.allocate(); totalAvailable.addAndGet(-size); try { diff --git a/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java b/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java index e6070e2b1..3da03cf98 100644 --- a/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java +++ b/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java @@ -136,6 +136,10 @@ public void setError(Throwable th) { notifyReady(); } + long getUncompressedSize() { + return uncompressedSize; + } + public long allocate() { if (isStopped) { problem = new IOException("" + getPartitionSession() + " is already closed"); diff --git a/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java b/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java index 02d5da515..b997958bd 100644 --- a/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java +++ b/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java @@ -188,45 +188,56 @@ public void flowControlByBudgetTest() { decoder.decodeNext(); - Assert.assertEquals(-50, decoder.getTotalAvailable()); + // 40 + 50 fit into the budget, 60 does not and has to wait + Assert.assertEquals(10, decoder.getTotalAvailable()); Assert.assertTrue(m1.isReady()); Assert.assertTrue(m2.isReady()); - Assert.assertTrue(m3.isReady()); + Assert.assertFalse(m3.isReady()); Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); - Assert.assertEquals(3, ready.get()); + Assert.assertEquals(2, ready.get()); Assert.assertEquals(40, m1.getData().length); Assert.assertEquals(50, m2.getData().length); - Assert.assertEquals(60, m3.getData().length); Assert.assertEquals(1, m1.getData()[0]); Assert.assertEquals(2, m2.getData()[0]); - Assert.assertEquals(3, m3.getData()[0]); - p1.releaseRange(OffsetsRange.of(1)); // 40 is not enough to resume decoding + p1.releaseRange(OffsetsRange.of(1)); // 10 + 40 is not enough for the 60 bytes of m3 p1.releaseRange(OffsetsRange.of(4)); // that offset is not decoded yet - Assert.assertEquals(-10, decoder.getTotalAvailable()); + Assert.assertEquals(50, decoder.getTotalAvailable()); + Assert.assertFalse(m3.isReady()); Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); - Assert.assertEquals(3, ready.get()); + Assert.assertEquals(2, ready.get()); - p1.releaseRange(OffsetsRange.of(2)); - Assert.assertEquals(-30, decoder.getTotalAvailable()); - Assert.assertTrue(m4.isReady()); + p1.releaseRange(OffsetsRange.of(2)); // the whole budget is free again, m3 is admitted + Assert.assertEquals(40, decoder.getTotalAvailable()); + Assert.assertTrue(m3.isReady()); + Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); - Assert.assertEquals(4, ready.get()); + Assert.assertEquals(3, ready.get()); - Assert.assertEquals(70, m4.getData().length); - Assert.assertEquals(4, m4.getData()[0]); + Assert.assertEquals(60, m3.getData().length); + Assert.assertEquals(3, m3.getData()[0]); p1.releaseRange(OffsetsRange.of(0, 3)); // double release - Assert.assertEquals(-30, decoder.getTotalAvailable()); + Assert.assertEquals(40, decoder.getTotalAvailable()); + Assert.assertFalse(m4.isReady()); + Assert.assertEquals(3, ready.get()); + + p1.releaseRange(OffsetsRange.of(0, 5)); // releases m3, then m4 fits into the free budget + Assert.assertEquals(30, decoder.getTotalAvailable()); + Assert.assertTrue(m4.isReady()); Assert.assertFalse(m5.isReady()); Assert.assertEquals(4, ready.get()); - p1.releaseRange(OffsetsRange.of(0, 5)); + Assert.assertEquals(70, m4.getData().length); + Assert.assertEquals(4, m4.getData()[0]); + + // m5 is bigger than the whole budget, so it is admitted only when nothing else retains the buffer + p1.releaseRange(OffsetsRange.of(0, 6)); Assert.assertTrue(m5.isReady()); Assert.assertEquals(5, ready.get()); @@ -257,7 +268,8 @@ public void partitionFlowTest() { Assert.assertEquals(100, decoder.getTotalAvailable()); decoder.decodeNext(); - Assert.assertEquals(-50, decoder.getTotalAvailable()); + // 40 + 50 fit into the budget, the 60 bytes of m3 do not + Assert.assertEquals(10, decoder.getTotalAvailable()); MessageImpl m6 = p1.decode(meta, OffsetsRange.of(4), gzipMsg(4, 10)); MessageImpl m7 = p1.decode(meta, OffsetsRange.of(5), gzipMsg(5, 20)); @@ -265,11 +277,12 @@ public void partitionFlowTest() { MessageImpl m9 = p2.decode(meta, OffsetsRange.of(15), gzipMsg(15, 20)); decoder.decodeNext(); - Assert.assertEquals(-50, decoder.getTotalAvailable()); + // m3 still blocks the queue, the smaller messages behind it are not reordered + Assert.assertEquals(10, decoder.getTotalAvailable()); Assert.assertTrue(m1.isReady()); Assert.assertTrue(m2.isReady()); - Assert.assertTrue(m3.isReady()); + Assert.assertFalse(m3.isReady()); Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); Assert.assertFalse(m6.isReady()); @@ -278,8 +291,9 @@ public void partitionFlowTest() { Assert.assertFalse(m9.isReady()); Assert.assertEquals(2, r1.get()); - Assert.assertEquals(1, r2.get()); + Assert.assertEquals(0, r2.get()); + // closing p1 returns its 90 bytes, which is enough to admit m3 and then m4 p1.close(); Assert.assertEquals(0, decoder.getTotalAvailable()); @@ -328,23 +342,27 @@ public void decodeStopTest() { decoder.decodeNext(); - Assert.assertEquals(-20, decoder.getTotalAvailable()); - Assert.assertEquals(2, ready.get()); + // only m1 fits into the 70 bytes budget, the 50 bytes of m2 do not + Assert.assertEquals(30, decoder.getTotalAvailable()); + Assert.assertEquals(1, ready.get()); Assert.assertTrue(m1.isReady()); - Assert.assertTrue(m2.isReady()); + Assert.assertFalse(m2.isReady()); Assert.assertFalse(m3.isReady()); decoder.stop(); - Assert.assertEquals(-20, decoder.getTotalAvailable()); + Assert.assertEquals(30, decoder.getTotalAvailable()); - Assert.assertEquals(2, ready.get()); + Assert.assertEquals(1, ready.get()); + Assert.assertFalse(m2.isReady()); Assert.assertFalse(m3.isReady()); + // a stopped decoder neither returns the budget nor resumes the pending messages partition.releaseRange(OffsetsRange.of(0, 10)); - Assert.assertEquals(2, ready.get()); + Assert.assertEquals(1, ready.get()); + Assert.assertFalse(m2.isReady()); Assert.assertFalse(m3.isReady()); - Assert.assertEquals(-20, decoder.getTotalAvailable()); + Assert.assertEquals(30, decoder.getTotalAvailable()); } @Test @@ -377,7 +395,8 @@ public void decodesOnProvidedExecutorTest() { Assert.assertFalse(m1.isReady()); Assert.assertFalse(m2.isReady()); - Assert.assertEquals(3, decodeTasks.size()); + // 400 + 500 fit into the 1000 bytes budget, the 600 bytes of m3 do not + Assert.assertEquals(2, decodeTasks.size()); Assert.assertEquals(0, p1ready.get()); Assert.assertEquals(0, p2ready.get()); @@ -385,7 +404,7 @@ public void decodesOnProvidedExecutorTest() { Assert.assertTrue(m1.isReady()); Assert.assertFalse(m2.isReady()); - Assert.assertEquals(2, decodeTasks.size()); + Assert.assertEquals(1, decodeTasks.size()); Assert.assertEquals(1, p1ready.get()); Assert.assertEquals(0, p2ready.get()); @@ -393,18 +412,30 @@ public void decodesOnProvidedExecutorTest() { Assert.assertTrue(m1.isReady()); Assert.assertTrue(m2.isReady()); - Assert.assertEquals(1, decodeTasks.size()); + Assert.assertEquals(0, decodeTasks.size()); Assert.assertEquals(1, p1ready.get()); Assert.assertEquals(1, p2ready.get()); + // decoding a message does not return its budget, only releasing it does + p1.releaseRange(OffsetsRange.of(0, 2)); + Assert.assertEquals(0, decodeTasks.size()); + + // now the whole budget is free again and m3 and m4 are admitted + p2.releaseRange(OffsetsRange.of(0, 2)); + Assert.assertEquals(2, decodeTasks.size()); + decoder.stop(); p1.close(); p2.close(); - Assert.assertEquals(1, decodeTasks.size()); + // tasks already submitted to the executor must not decode after the partitions are closed + decodeTasks.poll().run(); decodeTasks.poll().run(); Assert.assertEquals(0, decodeTasks.size()); + Assert.assertFalse(m3.isReady()); + Assert.assertFalse(m4.isReady()); + Assert.assertFalse(m5.isReady()); Assert.assertEquals(1, p1ready.get()); Assert.assertEquals(1, p2ready.get()); }