From 47c06ae0412d75ad2de58f17f00fa12590fc0597 Mon Sep 17 00:00:00 2001 From: jackfanwan <61672564+jackfanwan@users.noreply.github.com> Date: Fri, 4 Nov 2022 14:40:15 +0800 Subject: [PATCH] cherry-pick when delete workflow, delete related task #12681 --- .../impl/ProcessDefinitionServiceImpl.java | 20 ++++ .../service/ProcessDefinitionServiceTest.java | 96 +++++++++++++------ .../dao/mapper/TaskDefinitionMapper.java | 8 ++ .../dao/mapper/TaskDefinitionMapper.xml | 8 ++ 4 files changed, 105 insertions(+), 27 deletions(-) diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java index 06f5fdfb95..7d276ef209 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java @@ -874,6 +874,26 @@ public class ProcessDefinitionServiceImpl extends BaseServiceImpl implements Pro return result; } } + List processTaskRelations = processTaskRelationMapper + .queryByProcessCode(project.getCode(), processDefinition.getCode()); + if (CollectionUtils.isNotEmpty(processTaskRelations)) { + Set taskCodeList = new HashSet<>(processTaskRelations.size() * 2); + for (ProcessTaskRelation processTaskRelation : processTaskRelations) { + if (processTaskRelation.getPreTaskCode() != 0) { + taskCodeList.add(processTaskRelation.getPreTaskCode()); + } + if (processTaskRelation.getPostTaskCode() != 0) { + taskCodeList.add(processTaskRelation.getPostTaskCode()); + } + } + if (CollectionUtils.isNotEmpty(taskCodeList)) { + int i = taskDefinitionMapper.deleteByBatchCodes(new ArrayList<>(taskCodeList)); + if (i != taskCodeList.size()) { + logger.error("Delete task definition error, processDefinitionCode:{}.", code); + throw new ServiceException(Status.DELETE_TASK_DEFINE_BY_CODE_ERROR); + } + } + } int delete = processDefinitionMapper.deleteById(processDefinition.getId()); if (delete == 0) { putMsg(result, Status.DELETE_PROCESS_DEFINE_BY_CODE_ERROR); diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java index a2170bd080..9eae1a6536 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java @@ -57,9 +57,11 @@ import org.apache.dolphinscheduler.dao.entity.Tenant; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.dao.mapper.DataSourceMapper; import org.apache.dolphinscheduler.dao.mapper.ProcessDefinitionMapper; +import org.apache.dolphinscheduler.dao.mapper.ProcessDefinitionLogMapper; import org.apache.dolphinscheduler.dao.mapper.ProcessTaskRelationMapper; import org.apache.dolphinscheduler.dao.mapper.ProjectMapper; import org.apache.dolphinscheduler.dao.mapper.ScheduleMapper; +import org.apache.dolphinscheduler.dao.mapper.TaskDefinitionMapper; import org.apache.dolphinscheduler.dao.mapper.TenantMapper; import org.apache.dolphinscheduler.dao.model.PageListingResult; import org.apache.dolphinscheduler.dao.repository.ProcessDefinitionDao; @@ -89,6 +91,7 @@ import org.junit.Assert; import org.junit.Ignore; import org.junit.Test; import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; import org.junit.runner.RunWith; import org.mockito.InjectMocks; import org.mockito.Mock; @@ -120,7 +123,13 @@ public class ProcessDefinitionServiceTest { private ProcessDefinitionServiceImpl processDefinitionService; @Mock - private ProcessDefinitionMapper processDefineMapper; + private ProcessDefinitionMapper processDefinitionMapper; + + @Mock + private TaskDefinitionMapper taskDefinitionMapper; + + @Mock + private ProcessDefinitionLogMapper processDefinitionLogMapper; @Mock private ProcessDefinitionDao processDefinitionDao; @@ -155,6 +164,29 @@ public class ProcessDefinitionServiceTest { @Mock private WorkFlowLineageService workFlowLineageService; + protected User user; + protected Exception exception; + protected final static long projectCode = 1L; + protected final static long projectCodeOther = 2L; + protected final static long processDefinitionCode = 11L; + protected final static String name = "testProcessDefinitionName"; + protected final static String description = "this is a description"; + protected final static String releaseState = "ONLINE"; + protected final static int warningGroupId = 1; + protected final static int timeout = 60; + protected final static String executionType = "PARALLEL"; + protected final static String tenantCode = "tenant"; + + @BeforeEach + public void before() { + User loginUser = new User(); + loginUser.setId(1); + loginUser.setTenantId(2); + loginUser.setUserType(UserType.GENERAL_USER); + loginUser.setUserName("admin"); + user = loginUser; + } + @Test public void testQueryProcessDefinitionList() { long projectCode = 1L; @@ -180,7 +212,7 @@ public class ProcessDefinitionServiceTest { .thenReturn(result); List resourceList = new ArrayList<>(); resourceList.add(getProcessDefinition()); - Mockito.when(processDefineMapper.queryAllDefinitionList(project.getCode())).thenReturn(resourceList); + Mockito.when(processDefinitionMapper.queryAllDefinitionList(project.getCode())).thenReturn(resourceList); Map checkSuccessRes = processDefinitionService.queryProcessDefinitionList(loginUser, projectCode); Assert.assertEquals(Status.SUCCESS, checkSuccessRes.get(Constants.STATUS)); @@ -265,7 +297,7 @@ public class ProcessDefinitionServiceTest { Assert.assertEquals(Status.PROCESS_DEFINE_NOT_EXIST, instanceNotexitRes.get(Constants.STATUS)); // instance exit - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(getProcessDefinition()); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(getProcessDefinition()); putMsg(result, Status.SUCCESS, projectCode); Mockito.when(projectService.checkProjectAndAuth(loginUser, project, projectCode, WORKFLOW_DEFINITION)) .thenReturn(result); @@ -300,14 +332,14 @@ public class ProcessDefinitionServiceTest { putMsg(result, Status.SUCCESS, projectCode); Mockito.when(projectService.checkProjectAndAuth(loginUser, project, projectCode, WORKFLOW_DEFINITION)) .thenReturn(result); - Mockito.when(processDefineMapper.queryByDefineName(project.getCode(), "test_def")).thenReturn(null); + Mockito.when(processDefinitionMapper.queryByDefineName(project.getCode(), "test_def")).thenReturn(null); Map instanceNotExitRes = processDefinitionService.queryProcessDefinitionByName(loginUser, projectCode, "test_def"); Assert.assertEquals(Status.PROCESS_DEFINE_NOT_EXIST, instanceNotExitRes.get(Constants.STATUS)); // instance exit - Mockito.when(processDefineMapper.queryByDefineName(project.getCode(), "test")) + Mockito.when(processDefinitionMapper.queryByDefineName(project.getCode(), "test")) .thenReturn(getProcessDefinition()); putMsg(result, Status.SUCCESS, projectCode); Mockito.when(projectService.checkProjectAndAuth(loginUser, project, projectCode, WORKFLOW_DEFINITION)) @@ -356,7 +388,7 @@ public class ProcessDefinitionServiceTest { processDefinitionList.add(definition); Set definitionCodes = Arrays.stream("46".split(Constants.COMMA)).map(Long::parseLong).collect(Collectors.toSet()); - Mockito.when(processDefineMapper.queryByCodes(definitionCodes)).thenReturn(processDefinitionList); + Mockito.when(processDefinitionMapper.queryByCodes(definitionCodes)).thenReturn(processDefinitionList); Mockito.when(processService.saveProcessDefine(loginUser, definition, Boolean.TRUE, Boolean.TRUE)).thenReturn(2); Map map3 = processDefinitionService.batchCopyProcessDefinition( loginUser, projectCode, "46", 1L); @@ -391,7 +423,7 @@ public class ProcessDefinitionServiceTest { processDefinitionList.add(definition); Set definitionCodes = Arrays.stream("46".split(Constants.COMMA)).map(Long::parseLong).collect(Collectors.toSet()); - Mockito.when(processDefineMapper.queryByCodes(definitionCodes)).thenReturn(processDefinitionList); + Mockito.when(processDefinitionMapper.queryByCodes(definitionCodes)).thenReturn(processDefinitionList); Mockito.when(processService.saveProcessDefine(loginUser, definition, Boolean.TRUE, Boolean.TRUE)).thenReturn(2); Mockito.when(processTaskRelationMapper.queryByProcessCode(projectCode, 46L)) .thenReturn(getProcessTaskRelation(projectCode)); @@ -424,7 +456,7 @@ public class ProcessDefinitionServiceTest { putMsg(result, Status.SUCCESS, projectCode); Mockito.when(projectService.checkProjectAndAuth(loginUser, project, projectCode, WORKFLOW_DEFINITION_DELETE)) .thenReturn(result); - Mockito.when(processDefineMapper.queryByCode(1L)).thenReturn(null); + Mockito.when(processDefinitionMapper.queryByCode(1L)).thenReturn(null); Map instanceNotExitRes = processDefinitionService.deleteProcessDefinitionByCode(loginUser, projectCode, 1L); Assert.assertEquals(Status.PROCESS_DEFINE_NOT_EXIST, instanceNotExitRes.get(Constants.STATUS)); @@ -435,7 +467,7 @@ public class ProcessDefinitionServiceTest { .thenReturn(result); // user no auth loginUser.setUserType(UserType.GENERAL_USER); - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(processDefinition); Map userNoAuthRes = processDefinitionService.deleteProcessDefinitionByCode(loginUser, projectCode, 46L); Assert.assertEquals(Status.USER_NO_OPERATION_PERM, userNoAuthRes.get(Constants.STATUS)); @@ -444,7 +476,7 @@ public class ProcessDefinitionServiceTest { loginUser.setUserType(UserType.ADMIN_USER); putMsg(result, Status.SUCCESS, projectCode); processDefinition.setReleaseState(ReleaseState.ONLINE); - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(processDefinition); Throwable exception = Assertions.assertThrows(ServiceException.class, () -> processDefinitionService.deleteProcessDefinitionByCode(loginUser, projectCode, 46L)); String formatter = @@ -453,11 +485,11 @@ public class ProcessDefinitionServiceTest { // scheduler list elements > 1 processDefinition.setReleaseState(ReleaseState.OFFLINE); - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(processDefinition); putMsg(result, Status.SUCCESS, projectCode); Mockito.when(scheduleMapper.queryByProcessDefinitionCode(46L)).thenReturn(getSchedule()); Mockito.when(scheduleMapper.deleteById(46)).thenReturn(1); - Mockito.when(processDefineMapper.deleteById(processDefinition.getId())).thenReturn(1); + Mockito.when(processDefinitionMapper.deleteById(processDefinition.getId())).thenReturn(1); Mockito.when(processTaskRelationMapper.deleteByCode(project.getCode(), processDefinition.getCode())) .thenReturn(1); Mockito.when(workFlowLineageService.queryTaskDepOnProcess(project.getCode(), processDefinition.getCode())) @@ -491,17 +523,25 @@ public class ProcessDefinitionServiceTest { // delete success schedule.setReleaseState(ReleaseState.OFFLINE); - Mockito.when(processDefineMapper.deleteById(46)).thenReturn(1); + Mockito.when(processTaskRelationMapper.queryByProcessCode(1, 11)) + .thenReturn(getProcessTaskRelation(projectCode)); + Mockito.when(taskDefinitionMapper.deleteByBatchCodes(Arrays.asList(100L, 200L))).thenReturn(2); + Mockito.when(processDefinitionMapper.deleteById(46)).thenReturn(1); Mockito.when(scheduleMapper.deleteById(schedule.getId())).thenReturn(1); Mockito.when(processTaskRelationMapper.deleteByCode(project.getCode(), processDefinition.getCode())) .thenReturn(1); Mockito.when(scheduleMapper.queryByProcessDefinitionCode(46L)).thenReturn(getSchedule()); Mockito.when(workFlowLineageService.queryTaskDepOnProcess(project.getCode(), processDefinition.getCode())) .thenReturn(Collections.emptySet()); - putMsg(result, Status.SUCCESS, projectCode); - Map deleteSuccess = - processDefinitionService.deleteProcessDefinitionByCode(loginUser, projectCode, 46L); - Assert.assertEquals(Status.SUCCESS, deleteSuccess.get(Constants.STATUS)); + + Assertions.assertDoesNotThrow(() -> processDefinitionService.deleteProcessDefinitionByCode(user, 46L, processDefinitionCode)); + + // delete fail + Mockito.when(taskDefinitionMapper.deleteByBatchCodes(Arrays.asList(100L, 200L))).thenReturn(1); + exception = Assertions.assertThrows(ServiceException.class, + () -> processDefinitionService.deleteProcessDefinitionByCode(user, 46L, processDefinitionCode)); + Assertions.assertEquals(Status.DELETE_TASK_DEFINE_BY_CODE_ERROR.getCode(), + ((ServiceException) exception).getCode()); } @Test @@ -525,7 +565,7 @@ public class ProcessDefinitionServiceTest { // project check auth success, processs definition online putMsg(result, Status.SUCCESS, projectCode); - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(getProcessDefinition()); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(getProcessDefinition()); List processTaskRelationList = new ArrayList<>(); ProcessTaskRelation processTaskRelation = new ProcessTaskRelation(); processTaskRelation.setProjectCode(projectCode); @@ -570,13 +610,13 @@ public class ProcessDefinitionServiceTest { // project check auth success, process not exist putMsg(result, Status.SUCCESS, projectCode); - Mockito.when(processDefineMapper.verifyByDefineName(project.getCode(), "test_pdf")).thenReturn(null); + Mockito.when(processDefinitionMapper.verifyByDefineName(project.getCode(), "test_pdf")).thenReturn(null); Map processNotExistRes = processDefinitionService.verifyProcessDefinitionName(loginUser, projectCode, "test_pdf", 0); Assert.assertEquals(Status.SUCCESS, processNotExistRes.get(Constants.STATUS)); // process exist - Mockito.when(processDefineMapper.verifyByDefineName(project.getCode(), "test_pdf")) + Mockito.when(processDefinitionMapper.verifyByDefineName(project.getCode(), "test_pdf")) .thenReturn(getProcessDefinition()); Map processExistRes = processDefinitionService.verifyProcessDefinitionName(loginUser, projectCode, "test_pdf", 0); @@ -610,7 +650,7 @@ public class ProcessDefinitionServiceTest { putMsg(result, Status.SUCCESS, projectCode); Mockito.when(projectService.checkProjectAndAuth(loginUser, project, projectCode, null)).thenReturn(result); // process definition not exist - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(null); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(null); Map processDefinitionNullRes = processDefinitionService.getTaskNodeListByDefinitionCode(loginUser, projectCode, 46L); Assert.assertEquals(Status.PROCESS_DEFINE_NOT_EXIST, processDefinitionNullRes.get(Constants.STATUS)); @@ -619,7 +659,7 @@ public class ProcessDefinitionServiceTest { ProcessDefinition processDefinition = getProcessDefinition(); putMsg(result, Status.SUCCESS, projectCode); Mockito.when(processService.genDagData(Mockito.any())).thenReturn(new DagData(processDefinition, null, null)); - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(processDefinition); Map dataNotValidRes = processDefinitionService.getTaskNodeListByDefinitionCode(loginUser, projectCode, 46L); Assert.assertEquals(Status.SUCCESS, dataNotValidRes.get(Constants.STATUS)); @@ -643,7 +683,7 @@ public class ProcessDefinitionServiceTest { String defineCodes = "46"; Set defineCodeSet = Lists.newArrayList(defineCodes.split(Constants.COMMA)).stream().map(Long::parseLong) .collect(Collectors.toSet()); - Mockito.when(processDefineMapper.queryByCodes(defineCodeSet)).thenReturn(null); + Mockito.when(processDefinitionMapper.queryByCodes(defineCodeSet)).thenReturn(null); Map processNotExistRes = processDefinitionService.getNodeListMapByDefinitionCodes(loginUser, projectCode, defineCodes); Assert.assertEquals(Status.PROCESS_DEFINE_NOT_EXIST, processNotExistRes.get(Constants.STATUS)); @@ -653,7 +693,7 @@ public class ProcessDefinitionServiceTest { List processDefinitionList = new ArrayList<>(); processDefinitionList.add(processDefinition); - Mockito.when(processDefineMapper.queryByCodes(defineCodeSet)).thenReturn(processDefinitionList); + Mockito.when(processDefinitionMapper.queryByCodes(defineCodeSet)).thenReturn(processDefinitionList); Mockito.when(processService.genDagData(Mockito.any())).thenReturn(new DagData(processDefinition, null, null)); Project project1 = getProject(projectCode); List projects = new ArrayList<>(); @@ -680,7 +720,7 @@ public class ProcessDefinitionServiceTest { ProcessDefinition processDefinition = getProcessDefinition(); List processDefinitionList = new ArrayList<>(); processDefinitionList.add(processDefinition); - Mockito.when(processDefineMapper.queryAllDefinitionList(projectCode)).thenReturn(processDefinitionList); + Mockito.when(processDefinitionMapper.queryAllDefinitionList(projectCode)).thenReturn(processDefinitionList); Map successRes = processDefinitionService.queryAllProcessDefinitionByProjectCode(loginUser, projectCode); Assert.assertEquals(Status.SUCCESS, successRes.get(Constants.STATUS)); @@ -709,7 +749,7 @@ public class ProcessDefinitionServiceTest { putMsg(result, Status.SUCCESS, projectCode); Mockito.when(projectMapper.queryByCode(1)).thenReturn(project1); Mockito.when(projectService.checkProjectAndAuth(loginUser, project1, 1, WORKFLOW_TREE_VIEW)).thenReturn(result); - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(processDefinition); Mockito.when(processService.genDagGraph(processDefinition)).thenReturn(new DAG<>()); Map taskNullRes = processDefinitionService.viewTree(loginUser, processDefinition.getProjectCode(), 46, 10); @@ -727,7 +767,7 @@ public class ProcessDefinitionServiceTest { loginUser.setId(1); loginUser.setUserType(UserType.ADMIN_USER); ProcessDefinition processDefinition = getProcessDefinition(); - Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition); + Mockito.when(processDefinitionMapper.queryByCode(46L)).thenReturn(processDefinition); Project project1 = getProject(1); Map result = new HashMap<>(); @@ -897,6 +937,8 @@ public class ProcessDefinitionServiceTest { processTaskRelation.setProjectCode(projectCode); processTaskRelation.setProcessDefinitionCode(46L); processTaskRelation.setProcessDefinitionVersion(1); + processTaskRelation.setPreTaskCode(100); + processTaskRelation.setPostTaskCode(200); processTaskRelations.add(processTaskRelation); return processTaskRelations; } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java index 4940b4e1d2..247f34aafb 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java @@ -131,4 +131,12 @@ public interface TaskDefinitionMapper extends BaseMapper { * @return task definition list */ List queryByCodeList(@Param("codes") Collection codes); + + /** + * batch delete task by task code + * + * @param taskCodeList task code list + * @return deleted row count + */ + int deleteByBatchCodes(@Param("taskCodeList") List taskCodeList); } diff --git a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.xml b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.xml index a8a09f7271..0288a30b6a 100644 --- a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.xml +++ b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.xml @@ -86,6 +86,14 @@ delete from t_ds_task_definition where code = #{code} + + + delete from t_ds_task_definition where code in + + #{taskCode} + + + insert into t_ds_task_definition (code, name, version, description, project_code, user_id, task_type, task_params, flag, task_priority, worker_group, environment_code, fail_retry_times, fail_retry_interval,