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

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


The following commit(s) were added to refs/heads/dev by this push:
     new db8ec27394 [Fix-18540][Master] Reset the runtime state when recreating 
a failed task instance (#18541)
db8ec27394 is described below

commit db8ec273940812069f086f0eadcadb23d2b88b81
Author: Sepuri Sai Krishna <[email protected]>
AuthorDate: Mon Aug 10 14:47:23 2026 +0530

    [Fix-18540][Master] Reset the runtime state when recreating a failed task 
instance (#18541)
---
 .../FailedRecoverTaskInstanceFactory.java          |   7 ++
 .../FailedRecoverTaskInstanceFactoryTest.java      | 125 +++++++++++++++++++++
 2 files changed, 132 insertions(+)

diff --git 
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
 
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
index 1e0b4107a4..041c338180 100644
--- 
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
+++ 
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactory.java
@@ -51,8 +51,15 @@ public class FailedRecoverTaskInstanceFactory
         taskInstance.setHost(null);
         taskInstance.setVarPool(null);
         taskInstance.setSubmitTime(new Date());
+        taskInstance.setStartTime(null);
+        taskInstance.setEndTime(null);
         taskInstance.setLogPath(null);
         taskInstance.setExecutePath(null);
+        taskInstance.setPid(0);
+        // The recreated task instance is a new attempt, it should get the 
whole retry budget back, otherwise a task
+        // which already exhausted its retry times will never be retried again 
in the recovered workflow.
+        taskInstance.setRetryTimes(0);
+        taskInstance.setAlertFlag(Flag.NO);
         taskInstanceDao.insert(taskInstance);
 
         needRecoverTaskInstance.setFlag(Flag.NO);
diff --git 
a/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactoryTest.java
 
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactoryTest.java
new file mode 100644
index 0000000000..27684fa258
--- /dev/null
+++ 
b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/engine/task/execution/FailedRecoverTaskInstanceFactoryTest.java
@@ -0,0 +1,125 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.dolphinscheduler.server.master.engine.task.execution;
+
+import static com.google.common.truth.Truth.assertThat;
+import static org.mockito.Mockito.verify;
+
+import org.apache.dolphinscheduler.common.enums.Flag;
+import org.apache.dolphinscheduler.dao.entity.TaskInstance;
+import org.apache.dolphinscheduler.dao.repository.TaskInstanceDao;
+import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus;
+
+import java.util.Date;
+
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+@ExtendWith(MockitoExtension.class)
+class FailedRecoverTaskInstanceFactoryTest {
+
+    private static final int MAX_RETRY_TIMES = 3;
+
+    @InjectMocks
+    private FailedRecoverTaskInstanceFactory failedRecoverTaskInstanceFactory;
+
+    @Mock
+    private TaskInstanceDao taskInstanceDao;
+
+    @Test
+    @DisplayName("Test the recreated task instance gets the whole retry budget 
back")
+    void testCreateTaskInstance_resetRetryTimes() {
+        final TaskInstance needRecoverTaskInstance = 
createExhaustedFailedTaskInstance();
+
+        final TaskInstance taskInstance = 
failedRecoverTaskInstanceFactory.builder()
+                .withTaskInstance(needRecoverTaskInstance)
+                .build();
+
+        assertThat(taskInstance.getRetryTimes()).isEqualTo(0);
+        assertThat(taskInstance.getMaxRetryTimes()).isEqualTo(MAX_RETRY_TIMES);
+    }
+
+    @Test
+    @DisplayName("Test the recreated task instance doesn't inherit the runtime 
state of the failed attempt")
+    void testCreateTaskInstance_clearPreviousRuntimeState() {
+        final TaskInstance needRecoverTaskInstance = 
createExhaustedFailedTaskInstance();
+
+        final TaskInstance taskInstance = 
failedRecoverTaskInstanceFactory.builder()
+                .withTaskInstance(needRecoverTaskInstance)
+                .build();
+
+        assertThat(taskInstance.getId()).isNull();
+        
assertThat(taskInstance.getState()).isEqualTo(TaskExecutionStatus.SUBMITTED_SUCCESS);
+        assertThat(taskInstance.getStartTime()).isNull();
+        assertThat(taskInstance.getEndTime()).isNull();
+        assertThat(taskInstance.getHost()).isNull();
+        assertThat(taskInstance.getLogPath()).isNull();
+        assertThat(taskInstance.getExecutePath()).isNull();
+        assertThat(taskInstance.getVarPool()).isNull();
+        assertThat(taskInstance.getPid()).isEqualTo(0);
+        // The appLink is deliberately kept: SubWorkflowLogicTask stores its 
runtime context there and needs it
+        // to recover the existing sub workflow instance instead of triggering 
a new one.
+        
assertThat(taskInstance.getAppLink()).isEqualTo("application_1717430400000_0001");
+        assertThat(taskInstance.getAlertFlag()).isEqualTo(Flag.NO);
+        assertThat(taskInstance.getSubmitTime()).isNotNull();
+    }
+
+    @Test
+    @DisplayName("Test the recovered task instance is inserted and the origin 
one is marked as invalid")
+    void testCreateTaskInstance_markOriginTaskInstanceInvalid() {
+        final TaskInstance needRecoverTaskInstance = 
createExhaustedFailedTaskInstance();
+
+        final TaskInstance taskInstance = 
failedRecoverTaskInstanceFactory.builder()
+                .withTaskInstance(needRecoverTaskInstance)
+                .build();
+
+        assertThat(taskInstance.getFlag()).isEqualTo(Flag.YES);
+        assertThat(needRecoverTaskInstance.getFlag()).isEqualTo(Flag.NO);
+        verify(taskInstanceDao).insert(taskInstance);
+        verify(taskInstanceDao).updateById(needRecoverTaskInstance);
+    }
+
+    /**
+     * Create a task instance which is failed and has already used up all its 
retry times.
+     */
+    private TaskInstance createExhaustedFailedTaskInstance() {
+        final TaskInstance taskInstance = new TaskInstance();
+        taskInstance.setId(1);
+        taskInstance.setName("A");
+        taskInstance.setState(TaskExecutionStatus.FAILURE);
+        taskInstance.setFlag(Flag.YES);
+        taskInstance.setRetryTimes(MAX_RETRY_TIMES);
+        taskInstance.setMaxRetryTimes(MAX_RETRY_TIMES);
+        taskInstance.setRetryInterval(1);
+        taskInstance.setStartTime(new Date());
+        taskInstance.setEndTime(new Date());
+        taskInstance.setSubmitTime(new Date());
+        taskInstance.setHost("127.0.0.1:1234");
+        taskInstance.setLogPath("/tmp/log/A.log");
+        taskInstance.setExecutePath("/tmp/exec/A");
+        taskInstance.setVarPool("[]");
+        taskInstance.setPid(1234);
+        taskInstance.setAppLink("application_1717430400000_0001");
+        taskInstance.setAlertFlag(Flag.YES);
+        return taskInstance;
+    }
+}

Reply via email to