This is an automated email from the ASF dual-hosted git repository.

FrankChen021 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/master by this push:
     new 9068d2abdf6 fix: address CodeQL concurrency warnings (#19827)
9068d2abdf6 is described below

commit 9068d2abdf6a997e11c8cb0faf4591a4960577d8
Author: Frank Chen <[email protected]>
AuthorDate: Mon Aug 31 14:17:27 2026 +0800

    fix: address CodeQL concurrency warnings (#19827)
    
    * Fix CodeQL concurrency warnings
    
    * fix: narrow concurrency lock and notification scope
---
 .../OffHeapNamespaceExtractionCacheManager.java    | 28 ++++---
 .../channel/ReadableInputStreamFrameChannel.java   |  2 +-
 .../client/io/AppendableByteArrayInputStream.java  |  4 +-
 .../io/AppendableByteArrayInputStreamTest.java     | 85 ++++++++++++++++++++++
 4 files changed, 107 insertions(+), 12 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 0b31868a4b8..48ccc63d3f3 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 @@ public class OffHeapNamespaceExtractionCacheManager 
extends NamespaceExtractionC
 
     private void doDispose()
     {
-      if (!mmapDB.isClosed()) {
-        mmapDB.delete(mapDbKey);
+      synchronized (mmapDB) {
+        if (!mmapDB.isClosed()) {
+          mmapDB.delete(mapDbKey);
+        }
       }
       cacheCount.decrementAndGet();
     }
@@ -150,10 +152,14 @@ public class OffHeapNamespaceExtractionCacheManager 
extends NamespaceExtractionC
     }
   }
 
+  /**
+   * 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);
-  private AtomicInteger cacheCount = new AtomicInteger(0);
+  private final AtomicLong mapDbKeyCounter = new AtomicLong(0);
+  private final AtomicInteger cacheCount = new AtomicInteger(0);
 
   @Inject
   public OffHeapNamespaceExtractionCacheManager(
@@ -192,14 +198,18 @@ public class OffHeapNamespaceExtractionCacheManager 
extends NamespaceExtractionC
             }
 
             @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());
+              final boolean shouldDelete;
+              synchronized (mmapDB) {
+                shouldDelete = !mmapDB.isClosed();
+                if (shouldDelete) {
+                  mmapDB.close();
                 }
               }
+              if (shouldDelete && !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 cccce0147dd..d24c2445fd9 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 @@ public class ReadableInputStreamFrameChannel implements 
ReadableFrameChannel
                       () -> {
                         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 6789faf56a7..1660455ac60 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
@@ -59,7 +59,7 @@ public class AppendableByteArrayInputStream extends 
InputStream
   {
     synchronized (singleByteReaderDoer) {
       done = true;
-      singleByteReaderDoer.notify();
+      singleByteReaderDoer.notifyAll();
     }
   }
 
@@ -68,7 +68,7 @@ public class AppendableByteArrayInputStream extends 
InputStream
     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 b9492732c15..c337f7d9f70 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 class AppendableByteArrayInputStreamTest
     }
 
   }
+
+  @Test
+  public void testDoneUnblocksAllReaders() throws Exception
+  {
+    final AppendableByteArrayInputStream in = new 
AppendableByteArrayInputStream();
+    final AtomicReference<Integer> firstResult = new AtomicReference<>();
+    final AtomicReference<Integer> secondResult = new AtomicReference<>();
+    final AtomicReference<Throwable> firstError = new AtomicReference<>();
+    final AtomicReference<Throwable> 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<Integer> firstResult = new AtomicReference<>();
+    final AtomicReference<Integer> secondResult = new AtomicReference<>();
+    final AtomicReference<Throwable> firstError = new AtomicReference<>();
+    final AtomicReference<Throwable> 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<Integer> result,
+      final AtomicReference<Throwable> 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");
+  }
 }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to