xxubai commented on code in PR #4349:
URL: https://github.com/apache/amoro/pull/4349#discussion_r3911207767


##########
amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultProcessService.java:
##########
@@ -392,6 +395,66 @@ public void testRecoverProcessFailSafe() {
     }
   }
 
+  /**
+   * Verify that a process which keeps failing is retried a bounded number of 
times and then
+   * dropped, instead of being resubmitted in an infinite loop. The process 
store must be created
+   * with the retry limit rather than the current retry count, otherwise a 
FAILED process is treated
+   * as terminal immediately (retryNumber >= maxRetryTime), every 
RETRY_REQUESTED transition is
+   * rejected, retryNumber never increases and the process is resubmitted 
forever.
+   */
+  @Test(timeout = 60_000)
+  public void testFailedProcessRetryIsBounded() throws InterruptedException {
+    AlwaysFailingExecuteEngine failingEngine = new 
AlwaysFailingExecuteEngine();
+    processServiceService().unInstallAllExecuteEngines();
+    processServiceService().installExecuteEngine(failingEngine);
+    processServiceService().unInstallAllActionCoordinators();
+    processServiceService().installActionCoordinator(new 
MockActionCoordinator(failingEngine));
+    try {
+      createTable();
+
+      // Wait until the failing process has been submitted at least once.
+      awaitCondition(() -> failingEngine.getSubmitAttempts() >= 1, 
WAIT_TIMEOUT_MS, 100L);
+
+      // The process must exhaust its bounded retries and be untracked. With 
the infinite
+      // retry loop, the active process map never empties and this wait times 
out.
+      awaitCondition(
+          () -> processServiceService().getActiveTableProcess().isEmpty(), 
15_000L, 100L);
+
+      Assert.assertEquals(
+          1 + ActionCoordinatorScheduler.PROCESS_MAX_RETRY_NUMBER,
+          failingEngine.getSubmitAttempts());
+
+      dropTable();
+    } catch (Throwable t) {
+      throw new RuntimeException(t);
+    }
+  }
+
+  /** Execute engine whose processes always fail immediately after submission. 
*/
+  private static class AlwaysFailingExecuteEngine extends MockExecuteEngine {
+    private final AtomicInteger submitAttempts = new AtomicInteger();
+    private final Set<String> failedIdentifiers = 
ConcurrentHashMap.newKeySet();
+
+    @Override
+    public String submitTableProcess(org.apache.amoro.process.TableProcess 
tableProcess) {
+      String identifier = "failing-" + submitAttempts.incrementAndGet();
+      failedIdentifiers.add(identifier);
+      return identifier;
+    }
+
+    @Override
+    public ProcessStatus getStatus(String processIdentifier) {
+      if (failedIdentifiers.contains(processIdentifier)) {

Review Comment:
   nit: Please make the fake engine fail directly from `submitTableProcess()`. 
   
   Returning `FAILED` immediately takes the fast-terminal path and produces 
redundant `COMPLETE_FAILED` transition errors.
   
   
   ```suggestion
         throw new IllegalStateException("Submission failure");
   ```



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to