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..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 @@ -129,10 +129,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); @@ -476,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 67572905accb..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 @@ -324,6 +324,119 @@ 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); + 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();