From 4fddb9233c66d75f8590099a4dc2af6279d2bcde Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=8B=8F=E4=B9=89=E8=B6=85?= Date: Wed, 23 Sep 2026 09:38:59 +0800 Subject: [PATCH 1/4] Allow creating schedule for offline workflow definition --- .../service/impl/SchedulerServiceImpl.java | 14 ++--- .../api/service/SchedulerServiceTest.java | 56 ++++++++++++++++++- 2 files changed, 59 insertions(+), 11 deletions(-) diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java index ad74cf4b91b2..b10e3e5c0990 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java @@ -22,7 +22,6 @@ import org.apache.dolphinscheduler.api.dto.ScheduleParam; import org.apache.dolphinscheduler.api.enums.Status; import org.apache.dolphinscheduler.api.exceptions.ServiceException; -import org.apache.dolphinscheduler.api.service.ExecutorService; import org.apache.dolphinscheduler.api.service.ProjectService; import org.apache.dolphinscheduler.api.service.SchedulerService; import org.apache.dolphinscheduler.api.utils.PageInfo; @@ -78,9 +77,6 @@ public class SchedulerServiceImpl extends BaseServiceImpl implements SchedulerSe @Autowired private ProjectService projectService; - @Autowired - private ExecutorService executorService; - @Autowired private ScheduleDao scheduleDao; @@ -129,10 +125,12 @@ public Schedule insertSchedule(User loginUser, projectService.checkHasProjectWritePermissionThrowException(loginUser, project); - // check workflow define release state - WorkflowDefinition workflowDefinition = workflowDefinitionDao.queryByCode(workflowDefinitionCode).orElse(null); - executorService.checkWorkflowDefinitionValid(projectCode, workflowDefinition, workflowDefinitionCode, - workflowDefinition.getVersion()); + // check workflow definition exists + WorkflowDefinition workflowDefinition = workflowDefinitionDao.queryByCode(workflowDefinitionCode) + .orElseThrow(() -> new ServiceException(Status.WORKFLOW_DEFINITION_NOT_EXIST, workflowDefinitionCode)); + if (projectCode != workflowDefinition.getProjectCode()) { + throw new ServiceException(Status.WORKFLOW_DEFINITION_NOT_EXIST, workflowDefinitionCode); + } Schedule scheduleExists = scheduleDao.queryByWorkflowDefinitionCode(workflowDefinitionCode); diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java index 67572905accb..4d5edabf66dc 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java @@ -72,9 +72,6 @@ public class SchedulerServiceTest extends BaseServiceTestTool { @Mock private ProjectService projectService; - @Mock - private ExecutorService executorService; - @Mock private TenantExistValidator tenantExistValidator; @@ -324,6 +321,59 @@ public void testDeleteSchedules() { Assertions.assertDoesNotThrow(() -> schedulerService.deleteSchedulesById(user, scheduleId)); } + @Test + public void testInsertScheduleWorkflowNotExists() { + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject()); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.empty()); + + exception = Assertions.assertThrows(ServiceException.class, + () -> schedulerService.insertSchedule( + user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType.NONE, 0, + FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode)); + Assertions.assertEquals(Status.WORKFLOW_DEFINITION_NOT_EXIST.getCode(), + ((ServiceException) exception).getCode()); + } + + @Test + public void testInsertScheduleWorkflowFromAnotherProject() { + Project project = this.getProject(); + WorkflowDefinition workflowDefinition = this.getProcessDefinition(); + workflowDefinition.setProjectCode(999L); + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(project); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.of(workflowDefinition)); + + exception = Assertions.assertThrows(ServiceException.class, + () -> schedulerService.insertSchedule( + user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType.NONE, 0, + FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode)); + Assertions.assertEquals(Status.WORKFLOW_DEFINITION_NOT_EXIST.getCode(), + ((ServiceException) exception).getCode()); + } + + @Test + public void testInsertScheduleOfflineWorkflow() { + Project project = this.getProject(); + WorkflowDefinition workflowDefinition = this.getProcessDefinition(); + workflowDefinition.setReleaseState(ReleaseState.OFFLINE); + Schedule insertedSchedule = new Schedule(); + insertedSchedule.setId(scheduleId); + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(project); + Mockito.when(scheduleDao.queryByWorkflowDefinitionCode(processDefinitionCode)).thenReturn(null); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.of(workflowDefinition)); + Mockito.when(scheduleDao.queryById(Mockito.any())).thenReturn(insertedSchedule); + + Schedule result = schedulerService.insertSchedule( + user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType.NONE, 0, + FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode); + + ArgumentCaptor scheduleCaptor = ArgumentCaptor.forClass(Schedule.class); + Mockito.verify(scheduleDao).insert(scheduleCaptor.capture()); + Assertions.assertSame(insertedSchedule, result); + } + @Test public void testReadOnlyUserCannotChangeSchedule() { Project project = this.getProject(); From d7773ca07f6a8185efe12e686d9f7d85f50b5c7a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=8B=8F=E4=B9=89=E8=B6=85?= Date: Wed, 23 Sep 2026 09:38:59 +0800 Subject: [PATCH 2/4] Align insertSchedule validation with updateSchedule to allow scheduling offline workflows --- .../service/impl/SchedulerServiceImpl.java | 14 ++--- .../api/service/SchedulerServiceTest.java | 56 ++++++++++++++++++- 2 files changed, 59 insertions(+), 11 deletions(-) diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java index ad74cf4b91b2..b10e3e5c0990 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java @@ -22,7 +22,6 @@ import org.apache.dolphinscheduler.api.dto.ScheduleParam; import org.apache.dolphinscheduler.api.enums.Status; import org.apache.dolphinscheduler.api.exceptions.ServiceException; -import org.apache.dolphinscheduler.api.service.ExecutorService; import org.apache.dolphinscheduler.api.service.ProjectService; import org.apache.dolphinscheduler.api.service.SchedulerService; import org.apache.dolphinscheduler.api.utils.PageInfo; @@ -78,9 +77,6 @@ public class SchedulerServiceImpl extends BaseServiceImpl implements SchedulerSe @Autowired private ProjectService projectService; - @Autowired - private ExecutorService executorService; - @Autowired private ScheduleDao scheduleDao; @@ -129,10 +125,12 @@ public Schedule insertSchedule(User loginUser, projectService.checkHasProjectWritePermissionThrowException(loginUser, project); - // check workflow define release state - WorkflowDefinition workflowDefinition = workflowDefinitionDao.queryByCode(workflowDefinitionCode).orElse(null); - executorService.checkWorkflowDefinitionValid(projectCode, workflowDefinition, workflowDefinitionCode, - workflowDefinition.getVersion()); + // check workflow definition exists + WorkflowDefinition workflowDefinition = workflowDefinitionDao.queryByCode(workflowDefinitionCode) + .orElseThrow(() -> new ServiceException(Status.WORKFLOW_DEFINITION_NOT_EXIST, workflowDefinitionCode)); + if (projectCode != workflowDefinition.getProjectCode()) { + throw new ServiceException(Status.WORKFLOW_DEFINITION_NOT_EXIST, workflowDefinitionCode); + } Schedule scheduleExists = scheduleDao.queryByWorkflowDefinitionCode(workflowDefinitionCode); diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java index 67572905accb..4d5edabf66dc 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java @@ -72,9 +72,6 @@ public class SchedulerServiceTest extends BaseServiceTestTool { @Mock private ProjectService projectService; - @Mock - private ExecutorService executorService; - @Mock private TenantExistValidator tenantExistValidator; @@ -324,6 +321,59 @@ public void testDeleteSchedules() { Assertions.assertDoesNotThrow(() -> schedulerService.deleteSchedulesById(user, scheduleId)); } + @Test + public void testInsertScheduleWorkflowNotExists() { + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject()); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.empty()); + + exception = Assertions.assertThrows(ServiceException.class, + () -> schedulerService.insertSchedule( + user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType.NONE, 0, + FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode)); + Assertions.assertEquals(Status.WORKFLOW_DEFINITION_NOT_EXIST.getCode(), + ((ServiceException) exception).getCode()); + } + + @Test + public void testInsertScheduleWorkflowFromAnotherProject() { + Project project = this.getProject(); + WorkflowDefinition workflowDefinition = this.getProcessDefinition(); + workflowDefinition.setProjectCode(999L); + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(project); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.of(workflowDefinition)); + + exception = Assertions.assertThrows(ServiceException.class, + () -> schedulerService.insertSchedule( + user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType.NONE, 0, + FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode)); + Assertions.assertEquals(Status.WORKFLOW_DEFINITION_NOT_EXIST.getCode(), + ((ServiceException) exception).getCode()); + } + + @Test + public void testInsertScheduleOfflineWorkflow() { + Project project = this.getProject(); + WorkflowDefinition workflowDefinition = this.getProcessDefinition(); + workflowDefinition.setReleaseState(ReleaseState.OFFLINE); + Schedule insertedSchedule = new Schedule(); + insertedSchedule.setId(scheduleId); + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(project); + Mockito.when(scheduleDao.queryByWorkflowDefinitionCode(processDefinitionCode)).thenReturn(null); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.of(workflowDefinition)); + Mockito.when(scheduleDao.queryById(Mockito.any())).thenReturn(insertedSchedule); + + Schedule result = schedulerService.insertSchedule( + user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType.NONE, 0, + FailureStrategy.CONTINUE, Priority.MEDIUM, "default", "tenantCode", environmentCode); + + ArgumentCaptor scheduleCaptor = ArgumentCaptor.forClass(Schedule.class); + Mockito.verify(scheduleDao).insert(scheduleCaptor.capture()); + Assertions.assertSame(insertedSchedule, result); + } + @Test public void testReadOnlyUserCannotChangeSchedule() { Project project = this.getProject(); From 00c1da074023531d32371b56d149b383707cc546 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=8B=8F=E4=B9=89=E8=B6=85?= Date: Tue, 29 Sep 2026 16:16:25 +0800 Subject: [PATCH 3/4] Tighten tests to lock the saved-but-not-scheduled contract for offline workflow schedules --- .../api/service/SchedulerServiceTest.java | 20 +++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java index 4d5edabf66dc..c46b976e88ee 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java @@ -372,6 +372,26 @@ user, projectCode, processDefinitionCode, scheduleExpression(null), WarningType. ArgumentCaptor scheduleCaptor = ArgumentCaptor.forClass(Schedule.class); Mockito.verify(scheduleDao).insert(scheduleCaptor.capture()); Assertions.assertSame(insertedSchedule, result); + Assertions.assertEquals(ReleaseState.OFFLINE, scheduleCaptor.getValue().getReleaseState()); + Mockito.verifyNoInteractions(schedulerApi); + } + + @Test + public void testOnlineSchedulerRejectsOfflineWorkflow() { + Schedule schedule = this.getSchedule(); + schedule.setReleaseState(ReleaseState.OFFLINE); + WorkflowDefinition workflowDefinition = this.getProcessDefinition(); + workflowDefinition.setReleaseState(ReleaseState.OFFLINE); + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject()); + Mockito.when(scheduleDao.queryById(scheduleId)).thenReturn(schedule); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.of(workflowDefinition)); + + ServiceException ex = Assertions.assertThrows(ServiceException.class, + () -> schedulerService.onlineScheduler(user, projectCode, scheduleId)); + Assertions.assertEquals(Status.WORKFLOW_DEFINITION_NOT_RELEASE.getCode(), ex.getCode()); + Mockito.verify(scheduleDao, Mockito.never()).updateById(Mockito.any()); + Mockito.verifyNoInteractions(schedulerApi); } @Test From 484cb791c59a402ddae7a0376138f3cc54d80607 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=8B=8F=E4=B9=89=E8=B6=85?= Date: Tue, 29 Sep 2026 16:49:11 +0800 Subject: [PATCH 4/4] Preserve subworkflow validation by checking subworkflow readiness when activating a schedule --- .../service/impl/SchedulerServiceImpl.java | 7 +++ .../api/service/SchedulerServiceTest.java | 43 +++++++++++++++++++ 2 files changed, 50 insertions(+) diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java index b10e3e5c0990..ffbd85064725 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/SchedulerServiceImpl.java @@ -22,6 +22,7 @@ import org.apache.dolphinscheduler.api.dto.ScheduleParam; import org.apache.dolphinscheduler.api.enums.Status; import org.apache.dolphinscheduler.api.exceptions.ServiceException; +import org.apache.dolphinscheduler.api.service.ExecutorService; import org.apache.dolphinscheduler.api.service.ProjectService; import org.apache.dolphinscheduler.api.service.SchedulerService; import org.apache.dolphinscheduler.api.utils.PageInfo; @@ -77,6 +78,9 @@ public class SchedulerServiceImpl extends BaseServiceImpl implements SchedulerSe @Autowired private ProjectService projectService; + @Autowired + private ExecutorService executorService; + @Autowired private ScheduleDao scheduleDao; @@ -474,6 +478,9 @@ private void doOnlineScheduler(Schedule schedule) { if (!ReleaseState.ONLINE.equals(workflowDefinition.getReleaseState())) { throw new ServiceException(Status.WORKFLOW_DEFINITION_NOT_RELEASE, workflowDefinition.getName()); } + if (!executorService.checkSubWorkflowDefinitionValid(workflowDefinition)) { + throw new ServiceException(Status.SUB_WORKFLOW_DEFINITION_NOT_RELEASE); + } schedule.setReleaseState(ReleaseState.ONLINE); schedule.setUpdateTime(new Date()); diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java index c46b976e88ee..fc0af440f43e 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/SchedulerServiceTest.java @@ -72,6 +72,9 @@ public class SchedulerServiceTest extends BaseServiceTestTool { @Mock private ProjectService projectService; + @Mock + private ExecutorService executorService; + @Mock private TenantExistValidator tenantExistValidator; @@ -394,6 +397,46 @@ public void testOnlineSchedulerRejectsOfflineWorkflow() { Mockito.verifyNoInteractions(schedulerApi); } + @Test + public void testOnlineSchedulerRejectsOfflineSubWorkflow() { + Schedule schedule = this.getSchedule(); + schedule.setReleaseState(ReleaseState.OFFLINE); + WorkflowDefinition workflowDefinition = this.getProcessDefinition(); + workflowDefinition.setReleaseState(ReleaseState.ONLINE); + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(this.getProject()); + Mockito.when(scheduleDao.queryById(scheduleId)).thenReturn(schedule); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.of(workflowDefinition)); + Mockito.when(executorService.checkSubWorkflowDefinitionValid(workflowDefinition)).thenReturn(false); + + ServiceException ex = Assertions.assertThrows(ServiceException.class, + () -> schedulerService.onlineScheduler(user, projectCode, scheduleId)); + Assertions.assertEquals(Status.SUB_WORKFLOW_DEFINITION_NOT_RELEASE.getCode(), ex.getCode()); + Assertions.assertEquals(ReleaseState.OFFLINE, schedule.getReleaseState()); + Mockito.verify(scheduleDao, Mockito.never()).updateById(Mockito.any()); + Mockito.verifyNoInteractions(schedulerApi); + } + + @Test + public void testOnlineSchedulerSucceedsWhenSubWorkflowOnline() { + Schedule schedule = this.getSchedule(); + schedule.setReleaseState(ReleaseState.OFFLINE); + WorkflowDefinition workflowDefinition = this.getProcessDefinition(); + workflowDefinition.setReleaseState(ReleaseState.ONLINE); + Project project = this.getProject(); + Mockito.when(projectDao.queryByCode(projectCode)).thenReturn(project); + Mockito.when(scheduleDao.queryById(scheduleId)).thenReturn(schedule); + Mockito.when(workflowDefinitionDao.queryByCode(processDefinitionCode)) + .thenReturn(Optional.of(workflowDefinition)); + Mockito.when(executorService.checkSubWorkflowDefinitionValid(workflowDefinition)).thenReturn(true); + + schedulerService.onlineScheduler(user, projectCode, scheduleId); + + Assertions.assertEquals(ReleaseState.ONLINE, schedule.getReleaseState()); + Mockito.verify(scheduleDao).updateById(schedule); + Mockito.verify(schedulerApi).insertOrUpdateScheduleTask(project.getId(), schedule); + } + @Test public void testReadOnlyUserCannotChangeSchedule() { Project project = this.getProject();