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

Jackie-Jiang pushed a commit to branch fix-group-commit-empty-queue-exception
in repository https://gitbox.apache.org/repos/asf/pinot.git


The following commit(s) were added to 
refs/heads/fix-group-commit-empty-queue-exception by this push:
     new 0e22b1f3a29 Test group commit failure after queue drains
0e22b1f3a29 is described below

commit 0e22b1f3a299aecd91bcdf6b8bdb22a2686b0f6c
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Fri Sep 25 18:06:09 2026 -0700

    Test group commit failure after queue drains
---
 .../helix/IdealStateGroupCommitTest.java           | 51 ++++++++++++++++++++++
 1 file changed, 51 insertions(+)

diff --git 
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
 
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
index 257a03c880a..a4c036ce627 100644
--- 
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
+++ 
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
@@ -331,6 +331,57 @@ public class IdealStateGroupCommitTest {
     }
   }
 
+  /// A follower can acquire the queue after a failed leader removes its entry 
but before it signals the owner.
+  @Test
+  @SuppressWarnings("unchecked")
+  public void testProcessedFailurePropagatesWhenQueueIsEmpty()
+      throws Exception {
+    String tableName = TABLE_NAME_PREFIX + "empty_OFFLINE";
+    IdealStateGroupCommit commit = new IdealStateGroupCommit();
+    Field queuesField = 
IdealStateGroupCommit.class.getDeclaredField("_queues");
+    queuesField.setAccessible(true);
+    Object[] queues = (Object[]) queuesField.get(commit);
+    Object queue = queues[(tableName.hashCode() & Integer.MAX_VALUE) % 
queues.length];
+
+    Field runningField = queue.getClass().getDeclaredField("_running");
+    runningField.setAccessible(true);
+    AtomicReference<Thread> running = (AtomicReference<Thread>) 
runningField.get(queue);
+    Field pendingField = queue.getClass().getDeclaredField("_pending");
+    pendingField.setAccessible(true);
+    ConcurrentLinkedQueue<Object> pending = (ConcurrentLinkedQueue<Object>) 
pendingField.get(queue);
+
+    AtomicReference<Throwable> observedFailure = new AtomicReference<>();
+    RuntimeException batchFailure = new RuntimeException("simulated failed 
batch");
+    running.set(Thread.currentThread());
+    Thread owner = new Thread(() -> {
+      try {
+        commit.commit(TEST_INSTANCE.getHelixManager(), tableName, is -> is,
+            RetryPolicies.noDelayRetryPolicy(1), false);
+      } catch (Throwable t) {
+        observedFailure.set(t);
+      }
+    });
+    owner.start();
+    try {
+      long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
+      while (pending.peek() == null && System.nanoTime() < deadline) {
+        Thread.sleep(1);
+      }
+      Object entry = pending.peek();
+      Assert.assertNotNull(entry, "Owner did not enqueue its entry");
+      Field exceptionField = entry.getClass().getDeclaredField("_exception");
+      exceptionField.setAccessible(true);
+      exceptionField.set(entry, batchFailure);
+      Assert.assertTrue(pending.remove(entry));
+    } finally {
+      running.set(null);
+      owner.join(TimeUnit.SECONDS.toMillis(10));
+    }
+    Assert.assertFalse(owner.isAlive(), "Owner did not finish");
+    Assert.assertNotNull(observedFailure.get(), "Processed entry's failure was 
lost");
+    Assert.assertSame(observedFailure.get().getCause(), batchFailure);
+  }
+
   /// Integration-style test that walks the same chain the orphan-creating 
production bug walked:
   ///
   ///   - Step 2 of commitSegmentMetadataInternal: write the new consuming 
segment's ZK metadata


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

Reply via email to