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]