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 c63ecc9bf6 [Improvement-18662][API]Aligns insertSchedule with
updateSchedule (which already allows offline workflows to save) (#18664)
c63ecc9bf6 is described below
commit c63ecc9bf68919547fe69b5266e4040b113028e8
Author: suyc <[email protected]>
AuthorDate: Wed Sep 30 16:10:27 2026 +0800
[Improvement-18662][API]Aligns insertSchedule with updateSchedule (which
already allows offline workflows to save) (#18664)
* Allow creating schedule for offline workflow definition
* Align insertSchedule validation with updateSchedule to allow scheduling
offline workflows
* Tighten tests to lock the saved-but-not-scheduled contract for offline
workflow schedules
* Preserve subworkflow validation by checking subworkflow readiness when
activating a schedule
---------
Co-authored-by: 苏义超 <[email protected]>
Co-authored-by: xiangzihao <[email protected]>
---
.../api/service/impl/SchedulerServiceImpl.java | 13 ++-
.../api/service/SchedulerServiceTest.java | 113 +++++++++++++++++++++
2 files changed, 122 insertions(+), 4 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 ad74cf4b91..ffbd850647 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
@@ -129,10 +129,12 @@ public class SchedulerServiceImpl extends BaseServiceImpl
implements SchedulerSe
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);
@@ -476,6 +478,9 @@ public class SchedulerServiceImpl extends BaseServiceImpl
implements SchedulerSe
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 67572905ac..fc0af440f4 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
@@ -324,6 +324,119 @@ public class SchedulerServiceTest extends
BaseServiceTestTool {
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<Schedule> 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
+ 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();