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

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 06dae2b7ab [Test][Zeta] Stabilize coordinator retryable operation test 
(#11425)
06dae2b7ab is described below

commit 06dae2b7abe0bad8ce3b75f931158df004f05f70
Author: Daniel <[email protected]>
AuthorDate: Wed Jul 15 22:32:43 2026 +0800

    [Test][Zeta] Stabilize coordinator retryable operation test (#11425)
    
    Co-authored-by: DanielLeens <[email protected]>
---
 .../engine/server/CoordinatorServiceTest.java      | 40 ++++++++++++++++------
 .../operation/ReturnRetryTimesOperation.java       | 18 ++++++++++
 2 files changed, 48 insertions(+), 10 deletions(-)

diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
index 8c934a9b5b..286ba67fa6 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
@@ -71,6 +71,7 @@ import com.hazelcast.logging.ILogger;
 import com.hazelcast.map.IMap;
 import com.hazelcast.spi.exception.RetryableHazelcastException;
 import com.hazelcast.spi.impl.NodeEngineImpl;
+import com.hazelcast.spi.properties.ClusterProperty;
 import lombok.extern.slf4j.Slf4j;
 
 import java.lang.reflect.Field;
@@ -596,26 +597,45 @@ public class CoordinatorServiceTest {
     @Test
     public void 
testSeaTunnelEngineRetryableExceptionOperationCanBeRetryByHazelcast() {
 
-        HazelcastInstanceImpl instance =
-                SeaTunnelServerStarter.createHazelcastInstance(
+        int maxRetryCount = 3;
+        ReturnRetryTimesOperation.resetRetryTimes();
+        SeaTunnelConfig seaTunnelConfig = 
ConfigProvider.locateAndGetSeaTunnelConfig();
+        seaTunnelConfig
+                .getHazelcastConfig()
+                .setClusterName(
                         TestUtils.getClusterName(
                                 
"CoordinatorServiceTest_testSeaTunnelEngineRetryableExceptionOperationCanBeRetryByHazelcast"));
+        // Keep the production retry contract intact while using a small 
test-only budget so CI
+        // can still verify terminal exception propagation without waiting for 
250 retries.
+        seaTunnelConfig
+                .getHazelcastConfig()
+                .setProperty(
+                        ClusterProperty.INVOCATION_MAX_RETRY_COUNT.getName(),
+                        String.valueOf(maxRetryCount));
+        seaTunnelConfig
+                .getHazelcastConfig()
+                .setProperty(ClusterProperty.INVOCATION_RETRY_PAUSE.getName(), 
"1");
+        HazelcastInstanceImpl instance =
+                
SeaTunnelServerStarter.createHazelcastInstance(seaTunnelConfig);
         try {
             CompletionException exception =
                     Assertions.assertThrows(
                             CompletionException.class,
-                            () -> {
-                                NodeEngineUtil.sendOperationToMemberNode(
-                                                instance.node.getNodeEngine(),
-                                                new 
ReturnRetryTimesOperation(),
-                                                
instance.getCluster().getLocalMember().getAddress())
-                                        .join();
-                            });
+                            () ->
+                                    NodeEngineUtil.sendOperationToMemberNode(
+                                                    
instance.node.getNodeEngine(),
+                                                    new 
ReturnRetryTimesOperation(),
+                                                    instance.getCluster()
+                                                            .getLocalMember()
+                                                            .getAddress())
+                                            .join());
             Assertions.assertTrue(
                     exception
                             .getCause()
                             .getMessage()
-                            .contains("Retryable exception occurred, retry 
times: 250"));
+                            .contains(
+                                    "Retryable exception occurred, retry 
times: " + maxRetryCount));
+            Assertions.assertEquals(maxRetryCount, 
ReturnRetryTimesOperation.getRetryTimes());
         } finally {
             instance.shutdown();
         }
diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
index 9f3a1f46fe..d90a4c4165 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
@@ -30,6 +30,24 @@ public class ReturnRetryTimesOperation extends Operation
 
     private static final AtomicInteger retryTimes = new AtomicInteger(0);
 
+    /**
+     * Reset the static retry counter before starting a new retry-focused test.
+     *
+     * <p>The operation is reused across Hazelcast retries in the same test 
JVM.
+     */
+    public static void resetRetryTimes() {
+        retryTimes.set(0);
+    }
+
+    /**
+     * Return the number of operation attempts observed by the current test.
+     *
+     * <p>The terminal assertion uses this value to verify retry exhaustion.
+     */
+    public static int getRetryTimes() {
+        return retryTimes.get();
+    }
+
     @Override
     public void run() {
         retryTimes.getAndIncrement();

Reply via email to