From 2e75d31d3d50dfc41173e074551ab26056abbf18 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Fri, 31 Jul 2026 00:38:22 +0800 Subject: [PATCH 1/2] Fix CodeQL concurrency warnings --- ...ffHeapNamespaceExtractionCacheManager.java | 22 +++-- .../ReadableInputStreamFrameChannel.java | 2 +- .../io/AppendableByteArrayInputStream.java | 6 +- .../AppendableByteArrayInputStreamTest.java | 85 +++++++++++++++++++ 4 files changed, 104 insertions(+), 11 deletions(-) diff --git a/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java b/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java index 0b31868a4b86..aba8dfb7bce6 100644 --- a/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java +++ b/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java @@ -115,8 +115,10 @@ void disposeManually() private void doDispose() { - if (!mmapDB.isClosed()) { - mmapDB.delete(mapDbKey); + synchronized (mmapDB) { + if (!mmapDB.isClosed()) { + mmapDB.delete(mapDbKey); + } } cacheCount.decrementAndGet(); } @@ -150,6 +152,10 @@ protected ConcurrentMap delegate() } } + /** + * MapDB synchronizes its state-changing methods on the {@link DB} instance. Check-and-use sequences must hold the + * same monitor to prevent the database from closing between the two calls. + */ private final DB mmapDB; private final File tmpFile; private AtomicLong mapDbKeyCounter = new AtomicLong(0); @@ -192,12 +198,14 @@ public void start() } @Override - public synchronized void stop() + public void stop() { - if (!mmapDB.isClosed()) { - mmapDB.close(); - if (!tmpFile.delete()) { - log.warn("Unable to delete file at [%s]", tmpFile.getAbsolutePath()); + synchronized (mmapDB) { + if (!mmapDB.isClosed()) { + mmapDB.close(); + if (!tmpFile.delete()) { + log.warn("Unable to delete file at [%s]", tmpFile.getAbsolutePath()); + } } } } diff --git a/processing/src/main/java/org/apache/druid/frame/channel/ReadableInputStreamFrameChannel.java b/processing/src/main/java/org/apache/druid/frame/channel/ReadableInputStreamFrameChannel.java index cccce0147dd3..d24c2445fd95 100644 --- a/processing/src/main/java/org/apache/druid/frame/channel/ReadableInputStreamFrameChannel.java +++ b/processing/src/main/java/org/apache/druid/frame/channel/ReadableInputStreamFrameChannel.java @@ -217,7 +217,7 @@ private void startReading() () -> { synchronized (readMonitor) { keepReading = true; - readMonitor.notify(); + readMonitor.notifyAll(); } }, Execs.directExecutor() diff --git a/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java b/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java index 6789faf56a76..deafda8632d3 100644 --- a/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java +++ b/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java @@ -51,7 +51,7 @@ public void add(byte[] bytesToAdd) synchronized (singleByteReaderDoer) { bytes.addLast(bytesToAdd); available += bytesToAdd.length; - singleByteReaderDoer.notify(); + singleByteReaderDoer.notifyAll(); } } @@ -59,7 +59,7 @@ public void done() { synchronized (singleByteReaderDoer) { done = true; - singleByteReaderDoer.notify(); + singleByteReaderDoer.notifyAll(); } } @@ -68,7 +68,7 @@ public void exceptionCaught(Throwable t) synchronized (singleByteReaderDoer) { done = true; throwable = t; - singleByteReaderDoer.notify(); + singleByteReaderDoer.notifyAll(); } } diff --git a/processing/src/test/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStreamTest.java b/processing/src/test/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStreamTest.java index b9492732c154..c337f7d9f70d 100644 --- a/processing/src/test/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStreamTest.java +++ b/processing/src/test/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStreamTest.java @@ -257,4 +257,89 @@ public byte[] call() } } + + @Test + public void testDoneUnblocksAllReaders() throws Exception + { + final AppendableByteArrayInputStream in = new AppendableByteArrayInputStream(); + final AtomicReference firstResult = new AtomicReference<>(); + final AtomicReference secondResult = new AtomicReference<>(); + final AtomicReference firstError = new AtomicReference<>(); + final AtomicReference secondError = new AtomicReference<>(); + final Thread firstReader = readerThread(in, firstResult, firstError); + final Thread secondReader = readerThread(in, secondResult, secondError); + + firstReader.start(); + secondReader.start(); + waitUntilWaiting(firstReader, secondReader); + + in.done(); + + firstReader.join(1_000); + secondReader.join(1_000); + Assertions.assertFalse(firstReader.isAlive()); + Assertions.assertFalse(secondReader.isAlive()); + Assertions.assertNull(firstError.get()); + Assertions.assertNull(secondError.get()); + Assertions.assertEquals(-1, firstResult.get()); + Assertions.assertEquals(-1, secondResult.get()); + } + + @Test + public void testExceptionUnblocksAllReaders() throws Exception + { + final AppendableByteArrayInputStream in = new AppendableByteArrayInputStream(); + final AtomicReference firstResult = new AtomicReference<>(); + final AtomicReference secondResult = new AtomicReference<>(); + final AtomicReference firstError = new AtomicReference<>(); + final AtomicReference secondError = new AtomicReference<>(); + final Thread firstReader = readerThread(in, firstResult, firstError); + final Thread secondReader = readerThread(in, secondResult, secondError); + + firstReader.start(); + secondReader.start(); + waitUntilWaiting(firstReader, secondReader); + + final Exception expected = new Exception(); + in.exceptionCaught(expected); + + firstReader.join(1_000); + secondReader.join(1_000); + Assertions.assertFalse(firstReader.isAlive()); + Assertions.assertFalse(secondReader.isAlive()); + Assertions.assertNull(firstResult.get()); + Assertions.assertNull(secondResult.get()); + Assertions.assertSame(expected, firstError.get().getCause()); + Assertions.assertSame(expected, secondError.get().getCause()); + } + + private static Thread readerThread( + final AppendableByteArrayInputStream in, + final AtomicReference result, + final AtomicReference error + ) + { + final Thread reader = new Thread(() -> { + try { + result.set(in.read()); + } + catch (Throwable t) { + error.set(t); + } + }); + reader.setDaemon(true); + return reader; + } + + private static void waitUntilWaiting(final Thread... readers) throws InterruptedException + { + for (int i = 0; i < 100; i++) { + if (Arrays.stream(readers).allMatch(reader -> reader.getState() == Thread.State.WAITING)) { + return; + } + Thread.sleep(10); + } + + Assertions.fail("Readers did not all block on the input stream"); + } } From 6c387d9faae2fb38c38844cb9040bf07010ed813 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Fri, 31 Jul 2026 06:51:38 +0800 Subject: [PATCH 2/2] fix: narrow concurrency lock and notification scope --- .../OffHeapNamespaceExtractionCacheManager.java | 14 ++++++++------ .../client/io/AppendableByteArrayInputStream.java | 2 +- 2 files changed, 9 insertions(+), 7 deletions(-) diff --git a/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java b/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java index aba8dfb7bce6..48ccc63d3f38 100644 --- a/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java +++ b/extensions-core/lookups-cached-global/src/main/java/org/apache/druid/server/lookup/namespace/cache/OffHeapNamespaceExtractionCacheManager.java @@ -158,8 +158,8 @@ protected ConcurrentMap delegate() */ private final DB mmapDB; private final File tmpFile; - private AtomicLong mapDbKeyCounter = new AtomicLong(0); - private AtomicInteger cacheCount = new AtomicInteger(0); + private final AtomicLong mapDbKeyCounter = new AtomicLong(0); + private final AtomicInteger cacheCount = new AtomicInteger(0); @Inject public OffHeapNamespaceExtractionCacheManager( @@ -200,14 +200,16 @@ public void start() @Override public void stop() { + final boolean shouldDelete; synchronized (mmapDB) { - if (!mmapDB.isClosed()) { + shouldDelete = !mmapDB.isClosed(); + if (shouldDelete) { mmapDB.close(); - if (!tmpFile.delete()) { - log.warn("Unable to delete file at [%s]", tmpFile.getAbsolutePath()); - } } } + if (shouldDelete && !tmpFile.delete()) { + log.warn("Unable to delete file at [%s]", tmpFile.getAbsolutePath()); + } } } ); diff --git a/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java b/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java index deafda8632d3..1660455ac60d 100644 --- a/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java +++ b/processing/src/main/java/org/apache/druid/java/util/http/client/io/AppendableByteArrayInputStream.java @@ -51,7 +51,7 @@ public void add(byte[] bytesToAdd) synchronized (singleByteReaderDoer) { bytes.addLast(bytesToAdd); available += bytesToAdd.length; - singleByteReaderDoer.notifyAll(); + singleByteReaderDoer.notify(); } }