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 51b29690c4 [Fix-18565][API] Validate datasource access for task 
definitions (#18566)
51b29690c4 is described below

commit 51b29690c429d220c5b5a96fd9355250be66412f
Author: Wenjun Ruan <[email protected]>
AuthorDate: Fri Aug 21 15:26:33 2026 +0800

    [Fix-18565][API] Validate datasource access for task definitions (#18566)
---
 .../TaskDatasourcePermissionChecker.java           |  87 +++++++++++++++++
 .../service/impl/TaskDefinitionServiceImpl.java    |  22 +++--
 .../impl/WorkflowDefinitionServiceImpl.java        |  31 ++++--
 .../service/impl/WorkflowInstanceServiceImpl.java  |   5 +
 .../TaskDatasourcePermissionCheckerTest.java       | 108 +++++++++++++++++++++
 .../api/service/TaskDefinitionServiceImplTest.java |  26 +++++
 .../api/service/WorkflowDefinitionServiceTest.java |  74 ++++++++++++++
 .../api/service/WorkflowInstanceServiceTest.java   |   6 ++
 8 files changed, 342 insertions(+), 17 deletions(-)

diff --git 
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/permission/TaskDatasourcePermissionChecker.java
 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/permission/TaskDatasourcePermissionChecker.java
new file mode 100644
index 0000000000..de40ee2594
--- /dev/null
+++ 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/permission/TaskDatasourcePermissionChecker.java
@@ -0,0 +1,87 @@
+/*
+ * 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.api.permission;
+
+import org.apache.dolphinscheduler.api.enums.Status;
+import org.apache.dolphinscheduler.api.exceptions.ServiceException;
+import org.apache.dolphinscheduler.common.enums.AuthorizationType;
+import org.apache.dolphinscheduler.common.enums.UserType;
+import org.apache.dolphinscheduler.dao.entity.TaskDefinition;
+import org.apache.dolphinscheduler.dao.entity.User;
+import org.apache.dolphinscheduler.plugin.task.api.TaskPluginManager;
+import org.apache.dolphinscheduler.plugin.task.api.enums.ResourceType;
+import 
org.apache.dolphinscheduler.plugin.task.api.parameters.AbstractParameters;
+import 
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.AbstractResourceParameters;
+import 
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.ResourceParametersHelper;
+
+import java.util.Collection;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
+
+import lombok.extern.slf4j.Slf4j;
+
+import org.springframework.stereotype.Component;
+
+@Slf4j
+@Component
+public class TaskDatasourcePermissionChecker {
+
+    private final ResourcePermissionCheckService 
resourcePermissionCheckService;
+
+    public TaskDatasourcePermissionChecker(ResourcePermissionCheckService 
resourcePermissionCheckService) {
+        this.resourcePermissionCheckService = resourcePermissionCheckService;
+    }
+
+    public void checkPermission(User loginUser, Collection<? extends 
TaskDefinition> taskDefinitions) {
+        if (taskDefinitions == null || taskDefinitions.isEmpty()) {
+            return;
+        }
+
+        Set<Integer> datasourceIds = new HashSet<>();
+        for (TaskDefinition taskDefinition : taskDefinitions) {
+            AbstractParameters taskParameters = 
TaskPluginManager.parseTaskParameters(
+                    taskDefinition.getTaskType(), 
taskDefinition.getTaskParams());
+            ResourceParametersHelper resources = taskParameters.getResources();
+            if (resources == null) {
+                continue;
+            }
+            Map<Integer, AbstractResourceParameters> datasourceResources =
+                    resources.getResourceMap(ResourceType.DATASOURCE);
+            if (datasourceResources == null) {
+                continue;
+            }
+            datasourceResources.keySet().stream()
+                    .filter(datasourceId -> datasourceId != null && 
datasourceId > 0)
+                    .forEach(datasourceIds::add);
+        }
+
+        if (datasourceIds.isEmpty()) {
+            return;
+        }
+
+        int userId = loginUser.getUserType() == UserType.ADMIN_USER ? 0 : 
loginUser.getId();
+        Integer[] datasourceIdArray = datasourceIds.toArray(new Integer[0]);
+        if (!resourcePermissionCheckService.resourcePermissionCheck(
+                AuthorizationType.DATASOURCE, datasourceIdArray, userId, log)) 
{
+            log.warn("User does not have permission to use datasource 
referenced by task, userId:{}.",
+                    loginUser.getId());
+            throw new 
ServiceException(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION);
+        }
+    }
+}
diff --git 
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
index 5c6be5f86c..b39b0de17c 100644
--- 
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
+++ 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
@@ -24,6 +24,7 @@ import static 
org.apache.dolphinscheduler.plugin.task.api.TaskPluginManager.chec
 import org.apache.dolphinscheduler.api.enums.Status;
 import org.apache.dolphinscheduler.api.exceptions.ServiceException;
 import org.apache.dolphinscheduler.api.permission.PermissionCheck;
+import 
org.apache.dolphinscheduler.api.permission.TaskDatasourcePermissionChecker;
 import org.apache.dolphinscheduler.api.service.ProjectService;
 import org.apache.dolphinscheduler.api.service.TaskDefinitionService;
 import org.apache.dolphinscheduler.api.service.WorkflowTaskRelationService;
@@ -116,6 +117,9 @@ public class TaskDefinitionServiceImpl extends 
BaseServiceImpl implements TaskDe
     @Autowired
     private WorkflowDefinitionLogMapper workflowDefinitionLogMapper;
 
+    @Autowired
+    private TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
+
     /**
      * query task definition
      *
@@ -222,13 +226,6 @@ public class TaskDefinitionServiceImpl extends 
BaseServiceImpl implements TaskDe
         }
         TaskDefinitionLog taskDefinitionToUpdate =
                 JSONUtils.parseObject(taskDefinitionJsonObj, 
TaskDefinitionLog.class);
-        if (TimeoutFlag.CLOSE == taskDefinition.getTimeoutFlag()) {
-            taskDefinition.setTimeoutNotifyStrategy(null);
-        }
-        if (taskDefinition.equals(taskDefinitionToUpdate)) {
-            log.warn("Task definition does not need update because no change, 
taskDefinitionCode:{}.", taskCode);
-            return null;
-        }
         if (taskDefinitionToUpdate == null) {
             log.warn("Parameter taskDefinitionJson is invalid.");
             throw new ServiceException(Status.DATA_IS_NOT_VALID, 
taskDefinitionJsonObj);
@@ -236,6 +233,14 @@ public class TaskDefinitionServiceImpl extends 
BaseServiceImpl implements TaskDe
         if (!checkTaskParameters(taskDefinitionToUpdate.getTaskType(), 
taskDefinitionToUpdate.getTaskParams())) {
             throw new 
ServiceException(Status.WORKFLOW_NODE_S_PARAMETER_INVALID, 
taskDefinitionToUpdate.getName());
         }
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
Collections.singletonList(taskDefinitionToUpdate));
+        if (TimeoutFlag.CLOSE == taskDefinition.getTimeoutFlag()) {
+            taskDefinition.setTimeoutNotifyStrategy(null);
+        }
+        if (taskDefinition.equals(taskDefinitionToUpdate)) {
+            log.warn("Task definition does not need update because no change, 
taskDefinitionCode:{}.", taskCode);
+            return null;
+        }
         Integer version = 
taskDefinitionLogMapper.queryMaxVersionForDefinition(taskCode);
         if (version == null || version == 0) {
             log.error("Max version task definitionLog can not be found in 
database, taskDefinitionCode:{}.",
@@ -516,6 +521,7 @@ public class TaskDefinitionServiceImpl extends 
BaseServiceImpl implements TaskDe
         }
         TaskDefinitionLog taskDefinitionUpdate =
                 
taskDefinitionLogMapper.queryByDefinitionCodeAndVersion(taskCode, version);
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
Collections.singletonList(taskDefinitionUpdate));
         taskDefinitionUpdate.setUserId(loginUser.getId());
         taskDefinitionUpdate.setUpdateTime(new Date());
         taskDefinitionUpdate.setId(taskDefinition.getId());
@@ -662,6 +668,8 @@ public class TaskDefinitionServiceImpl extends 
BaseServiceImpl implements TaskDe
                 taskDefinitionLog.setFlag(Flag.NO);
                 break;
             case ONLINE:
+                taskDatasourcePermissionChecker.checkPermission(loginUser,
+                        Collections.singletonList(taskDefinitionLog));
                 String resourceIds = taskDefinition.getResourceIds();
                 if (StringUtils.isNotBlank(resourceIds)) {
                     Integer[] resourceIdArray =
diff --git 
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowDefinitionServiceImpl.java
 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowDefinitionServiceImpl.java
index d777635ac1..538a35160f 100644
--- 
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowDefinitionServiceImpl.java
+++ 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowDefinitionServiceImpl.java
@@ -37,6 +37,7 @@ import 
org.apache.dolphinscheduler.api.dto.treeview.TreeViewDto;
 import 
org.apache.dolphinscheduler.api.dto.workflow.WorkflowDefinitionVariablesDTO;
 import org.apache.dolphinscheduler.api.enums.Status;
 import org.apache.dolphinscheduler.api.exceptions.ServiceException;
+import 
org.apache.dolphinscheduler.api.permission.TaskDatasourcePermissionChecker;
 import org.apache.dolphinscheduler.api.service.ProjectService;
 import org.apache.dolphinscheduler.api.service.SchedulerService;
 import org.apache.dolphinscheduler.api.service.TaskDefinitionLogService;
@@ -206,6 +207,9 @@ public class WorkflowDefinitionServiceImpl extends 
BaseServiceImpl implements Wo
     @Autowired
     private GlobalParamsValidator globalParamsValidator;
 
+    @Autowired
+    private TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
+
     /**
      * create workflow definition
      *
@@ -268,6 +272,7 @@ public class WorkflowDefinitionServiceImpl extends 
BaseServiceImpl implements Wo
                                                  List<WorkflowTaskRelationLog> 
taskRelationList,
                                                  WorkflowDefinition 
workflowDefinition,
                                                  List<TaskDefinitionLog> 
taskDefinitionLogs) {
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
taskDefinitionLogs);
         int saveTaskResult = processService.saveTaskDefine(loginUser, 
workflowDefinition.getProjectCode(),
                 taskDefinitionLogs, Boolean.TRUE);
         if (saveTaskResult == Constants.EXIT_CODE_SUCCESS) {
@@ -691,6 +696,7 @@ public class WorkflowDefinitionServiceImpl extends 
BaseServiceImpl implements Wo
                                                  WorkflowDefinition 
workflowDefinition,
                                                  WorkflowDefinition 
workflowDefinitionDeepCopy,
                                                  List<TaskDefinitionLog> 
taskDefinitionLogs) {
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
taskDefinitionLogs);
         int saveTaskResult = processService.saveTaskDefine(loginUser, 
workflowDefinition.getProjectCode(),
                 taskDefinitionLogs, Boolean.TRUE);
         if (saveTaskResult == Constants.EXIT_CODE_SUCCESS) {
@@ -1590,14 +1596,6 @@ public class WorkflowDefinitionServiceImpl extends 
BaseServiceImpl implements Wo
                     
Status.SWITCH_WORKFLOW_DEFINITION_VERSION_NOT_EXIST_WORKFLOW_DEFINITION_VERSION_ERROR,
                     workflowDefinition.getCode(), version);
         }
-        int switchVersion = processService.switchVersion(workflowDefinition, 
workflowDefinitionLog);
-        if (switchVersion <= 0) {
-            log.error(
-                    "Switch workflow definition version error, projectCode:{}, 
workflowDefinitionCode:{}, version:{}.",
-                    projectCode, code, version);
-            throw new 
ServiceException(Status.SWITCH_WORKFLOW_DEFINITION_VERSION_ERROR);
-        }
-
         List<WorkflowTaskRelation> workflowTaskRelationList = 
workflowTaskRelationDao
                 
.queryWorkflowTaskRelationsByWorkflowDefinitionCode(workflowDefinitionLog.getCode(),
                         workflowDefinitionLog.getVersion());
@@ -1610,6 +1608,16 @@ public class WorkflowDefinitionServiceImpl extends 
BaseServiceImpl implements Wo
                             
taskDefinitionLog.setVersion(taskCodeVersionDto.getVersion());
                             return Stream.of(taskDefinitionLog);
                         }).collect(Collectors.toList()));
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
taskDefinitionLogList);
+
+        int switchVersion = processService.switchVersion(workflowDefinition, 
workflowDefinitionLog);
+        if (switchVersion <= 0) {
+            log.error(
+                    "Switch workflow definition version error, projectCode:{}, 
workflowDefinitionCode:{}, version:{}.",
+                    projectCode, code, version);
+            throw new 
ServiceException(Status.SWITCH_WORKFLOW_DEFINITION_VERSION_ERROR);
+        }
+
         saveWorkflowLineage(workflowDefinitionLog.getProjectCode(), 
workflowDefinitionLog.getCode(),
                 workflowDefinitionLog.getVersion(), taskDefinitionLogList);
 
@@ -1755,7 +1763,7 @@ public class WorkflowDefinitionServiceImpl extends 
BaseServiceImpl implements Wo
             return;
         }
 
-        checkWorkflowDefinitionIsValidated(workflowDefinition.getCode());
+        checkWorkflowDefinitionIsValidated(loginUser, 
workflowDefinition.getCode());
         checkAllSubWorkflowDefinitionIsOnline(workflowDefinition.getCode());
 
         workflowDefinition.setReleaseState(ReleaseState.ONLINE);
@@ -1851,13 +1859,16 @@ public class WorkflowDefinitionServiceImpl extends 
BaseServiceImpl implements Wo
         return localUserDefParams;
     }
 
-    private void checkWorkflowDefinitionIsValidated(Long 
workflowDefinitionCode) {
+    private void checkWorkflowDefinitionIsValidated(User loginUser, Long 
workflowDefinitionCode) {
         // todo: build dag check if the dag is validated
         List<WorkflowTaskRelation> workflowTaskRelations =
                 
workflowTaskRelationDao.queryByWorkflowDefinitionCode(workflowDefinitionCode);
         if (CollectionUtils.isEmpty(workflowTaskRelations)) {
             throw new ServiceException(Status.WORKFLOW_DAG_IS_EMPTY);
         }
+        List<TaskDefinitionLog> taskDefinitionLogs =
+                
taskDefinitionLogDao.queryTaskDefineLogList(workflowTaskRelations);
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
taskDefinitionLogs);
         // todo : check Workflow is validate
     }
 
diff --git 
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java
 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java
index 96cfde7061..7c79613da0 100644
--- 
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java
+++ 
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java
@@ -32,6 +32,7 @@ import 
org.apache.dolphinscheduler.api.dto.workflowInstance.WorkflowInstanceTask
 import 
org.apache.dolphinscheduler.api.dto.workflowInstance.WorkflowInstanceVariablesDTO;
 import org.apache.dolphinscheduler.api.enums.Status;
 import org.apache.dolphinscheduler.api.exceptions.ServiceException;
+import 
org.apache.dolphinscheduler.api.permission.TaskDatasourcePermissionChecker;
 import org.apache.dolphinscheduler.api.service.ProjectService;
 import org.apache.dolphinscheduler.api.service.TaskInstanceService;
 import org.apache.dolphinscheduler.api.service.UsersService;
@@ -160,6 +161,9 @@ public class WorkflowInstanceServiceImpl extends 
BaseServiceImpl implements Work
     @Autowired
     private TaskInstanceContextDao taskInstanceContextDao;
 
+    @Autowired
+    private TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
+
     @Override
     public List<WorkflowInstance> queryTopNLongestRunningWorkflowInstance(User 
loginUser, long projectCode, int size,
                                                                           
String startTime, String endTime) {
@@ -406,6 +410,7 @@ public class WorkflowInstanceServiceImpl extends 
BaseServiceImpl implements Work
                 throw new 
ServiceException(Status.WORKFLOW_NODE_S_PARAMETER_INVALID, 
taskDefinitionLog.getName());
             }
         }
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
taskDefinitionLogs);
         int saveTaskResult = processService.saveTaskDefine(loginUser, 
projectCode, taskDefinitionLogs, syncDefine);
         if (saveTaskResult == Constants.DEFINITION_FAILURE) {
             log.error("Update task definition error, projectCode:{}, 
workflowInstanceId:{}", projectCode,
diff --git 
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/permission/TaskDatasourcePermissionCheckerTest.java
 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/permission/TaskDatasourcePermissionCheckerTest.java
new file mode 100644
index 0000000000..3eef903517
--- /dev/null
+++ 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/permission/TaskDatasourcePermissionCheckerTest.java
@@ -0,0 +1,108 @@
+/*
+ * 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.api.permission;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+
+import org.apache.dolphinscheduler.api.enums.Status;
+import org.apache.dolphinscheduler.api.exceptions.ServiceException;
+import org.apache.dolphinscheduler.common.enums.AuthorizationType;
+import org.apache.dolphinscheduler.common.enums.UserType;
+import org.apache.dolphinscheduler.dao.entity.TaskDefinition;
+import org.apache.dolphinscheduler.dao.entity.User;
+
+import java.util.Collections;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mock;
+import org.mockito.Mockito;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.slf4j.Logger;
+
+@ExtendWith(MockitoExtension.class)
+public class TaskDatasourcePermissionCheckerTest {
+
+    private static final int USER_ID = 1;
+    private static final int DATASOURCE_ID = 42;
+
+    @Mock
+    private ResourcePermissionCheckService resourcePermissionCheckService;
+
+    private TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
+
+    private User loginUser;
+
+    @BeforeEach
+    public void setUp() {
+        taskDatasourcePermissionChecker = new 
TaskDatasourcePermissionChecker(resourcePermissionCheckService);
+        loginUser = new User();
+        loginUser.setId(USER_ID);
+        loginUser.setUserType(UserType.GENERAL_USER);
+    }
+
+    @Test
+    public void shouldRejectUnauthorizedDatasourceReferencedByTask() {
+        Mockito.when(resourcePermissionCheckService.resourcePermissionCheck(
+                eq(AuthorizationType.DATASOURCE), any(Object[].class), 
eq(USER_ID), any(Logger.class)))
+                .thenReturn(false);
+
+        ServiceException exception = 
Assertions.assertThrows(ServiceException.class,
+                () -> taskDatasourcePermissionChecker.checkPermission(
+                        loginUser, 
Collections.singletonList(getRemoteShellTaskDefinition())));
+
+        
Assertions.assertEquals(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION.getCode(), 
exception.getCode());
+        ArgumentCaptor<Object[]> datasourceIdsCaptor = 
ArgumentCaptor.forClass(Object[].class);
+        Mockito.verify(resourcePermissionCheckService).resourcePermissionCheck(
+                eq(AuthorizationType.DATASOURCE), 
datasourceIdsCaptor.capture(), eq(USER_ID), any(Logger.class));
+        Assertions.assertArrayEquals(new Integer[]{DATASOURCE_ID}, 
datasourceIdsCaptor.getValue());
+    }
+
+    @Test
+    public void shouldAllowAuthorizedDatasourceReferencedByTask() {
+        Mockito.when(resourcePermissionCheckService.resourcePermissionCheck(
+                eq(AuthorizationType.DATASOURCE), any(Object[].class), 
eq(USER_ID), any(Logger.class)))
+                .thenReturn(true);
+
+        Assertions.assertDoesNotThrow(() -> 
taskDatasourcePermissionChecker.checkPermission(
+                loginUser, 
Collections.singletonList(getRemoteShellTaskDefinition())));
+    }
+
+    @Test
+    public void shouldSkipPermissionCheckForTaskWithoutDatasource() {
+        TaskDefinition taskDefinition = new TaskDefinition();
+        taskDefinition.setTaskType("SHELL");
+        taskDefinition.setTaskParams("{\"rawScript\":\"echo test\"}");
+
+        taskDatasourcePermissionChecker.checkPermission(loginUser, 
Collections.singletonList(taskDefinition));
+
+        Mockito.verifyNoInteractions(resourcePermissionCheckService);
+    }
+
+    private TaskDefinition getRemoteShellTaskDefinition() {
+        TaskDefinition taskDefinition = new TaskDefinition();
+        taskDefinition.setTaskType("REMOTESHELL");
+        taskDefinition.setTaskParams(
+                "{\"rawScript\":\"echo 
test\",\"type\":\"SSH\",\"datasource\":" + DATASOURCE_ID + "}");
+        return taskDefinition;
+    }
+}
diff --git 
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
index d85ee6ef3f..f34c1c7e81 100644
--- 
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
+++ 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
@@ -21,6 +21,7 @@ import static 
org.apache.dolphinscheduler.api.AssertionsHelper.assertThrowsServi
 import static 
org.apache.dolphinscheduler.api.constants.ApiFuncIdentificationConstant.TASK_DEFINITION;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyList;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.ArgumentMatchers.isA;
 import static org.mockito.Mockito.doNothing;
@@ -29,6 +30,7 @@ import static org.mockito.Mockito.when;
 
 import org.apache.dolphinscheduler.api.enums.Status;
 import org.apache.dolphinscheduler.api.exceptions.ServiceException;
+import 
org.apache.dolphinscheduler.api.permission.TaskDatasourcePermissionChecker;
 import org.apache.dolphinscheduler.api.service.impl.ProjectServiceImpl;
 import org.apache.dolphinscheduler.api.service.impl.TaskDefinitionServiceImpl;
 import org.apache.dolphinscheduler.common.constants.Constants;
@@ -108,6 +110,9 @@ public class TaskDefinitionServiceImplTest {
     @Mock
     private WorkflowDefinitionLogMapper workflowDefinitionLogMapper;
 
+    @Mock
+    private TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
+
     @Mock
     private WorkflowTaskRelationLogMapper workflowTaskRelationLogMapper;
 
@@ -182,6 +187,26 @@ public class TaskDefinitionServiceImplTest {
                 () -> taskDefinitionService.switchVersion(user, PROJECT_CODE, 
TASK_CODE, VERSION));
     }
 
+    @Test
+    public void switchVersionShouldRejectUnauthorizedDatasource() {
+        Project project = getProject();
+        when(projectDao.queryByCode(PROJECT_CODE)).thenReturn(project);
+        
Mockito.doNothing().when(projectService).checkHasProjectWritePermissionThrowException(user,
 project);
+
+        TaskDefinition taskDefinition = new TaskDefinition();
+        taskDefinition.setProjectCode(PROJECT_CODE);
+        
when(taskDefinitionDao.queryByCode(TASK_CODE)).thenReturn(taskDefinition);
+        
when(taskDefinitionLogMapper.queryByDefinitionCodeAndVersion(TASK_CODE, 
VERSION))
+                .thenReturn(new TaskDefinitionLog());
+        doThrow(new 
ServiceException(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION))
+                
.when(taskDatasourcePermissionChecker).checkPermission(eq(user), anyList());
+
+        
assertThrowsServiceException(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION,
+                () -> taskDefinitionService.switchVersion(user, PROJECT_CODE, 
TASK_CODE, VERSION));
+
+        Mockito.verify(taskDefinitionDao, 
Mockito.never()).updateById(any(TaskDefinition.class));
+    }
+
     @Test
     public void deleteByCodeAndVersion() {
         Project project = getProject();
@@ -397,6 +422,7 @@ public class TaskDefinitionServiceImplTest {
             Long updatedTaskCode = 
taskDefinitionService.updateTaskWithUpstream(user, PROJECT_CODE, TASK_CODE,
                     taskDefinitionJson, UPSTREAM_CODE);
             assertEquals(TASK_CODE, updatedTaskCode);
+            
Mockito.verify(taskDatasourcePermissionChecker).checkPermission(eq(user), 
anyList());
             user.setUserType(UserType.GENERAL_USER);
         }
     }
diff --git 
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowDefinitionServiceTest.java
 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowDefinitionServiceTest.java
index b26d0da000..dbf1dd5787 100644
--- 
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowDefinitionServiceTest.java
+++ 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowDefinitionServiceTest.java
@@ -35,6 +35,7 @@ import static org.mockito.Mockito.when;
 
 import org.apache.dolphinscheduler.api.enums.Status;
 import org.apache.dolphinscheduler.api.exceptions.ServiceException;
+import 
org.apache.dolphinscheduler.api.permission.TaskDatasourcePermissionChecker;
 import org.apache.dolphinscheduler.api.service.impl.ProjectServiceImpl;
 import 
org.apache.dolphinscheduler.api.service.impl.WorkflowDefinitionServiceImpl;
 import org.apache.dolphinscheduler.api.utils.PageInfo;
@@ -57,6 +58,7 @@ import org.apache.dolphinscheduler.dao.entity.TaskMainInfo;
 import org.apache.dolphinscheduler.dao.entity.User;
 import org.apache.dolphinscheduler.dao.entity.UserWithWorkflowDefinitionCode;
 import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition;
+import org.apache.dolphinscheduler.dao.entity.WorkflowDefinitionLog;
 import org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelation;
 import org.apache.dolphinscheduler.dao.mapper.TaskDefinitionLogMapper;
 import org.apache.dolphinscheduler.dao.mapper.WorkflowDefinitionLogMapper;
@@ -180,6 +182,9 @@ public class WorkflowDefinitionServiceTest extends 
BaseServiceTestTool {
     @Mock
     private TaskDefinitionLogMapper taskDefinitionLogMapper;
 
+    @Mock
+    private TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
+
     @Mock
     private TaskDefinitionService taskDefinitionService;
 
@@ -925,6 +930,75 @@ public class WorkflowDefinitionServiceTest extends 
BaseServiceTestTool {
         Assertions.assertEquals(1, workflowDefinition.getVersion());
     }
 
+    @Test
+    public void 
testCreateWorkflowDefinitionShouldRejectUnauthorizedDatasource() {
+        Project project = getProject(projectCode);
+        when(projectDao.queryByCode(projectCode)).thenReturn(project);
+        
Mockito.doNothing().when(projectService).checkHasProjectWritePermissionThrowException(eq(user),
 eq(project));
+        when(workflowDefinitionDao.verifyByDefineName(projectCode, 
name)).thenReturn(null);
+        when(processService.transformTask(anyList(), 
anyList())).thenReturn(getTaskNodeList());
+        doThrow(new 
ServiceException(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION))
+                
.when(taskDatasourcePermissionChecker).checkPermission(eq(user), anyList());
+
+        ServiceException exception = 
Assertions.assertThrows(ServiceException.class,
+                () -> workflowDefinitionService.createWorkflowDefinition(
+                        user, projectCode, name, description, "[]", "[]", 
timeout,
+                        taskRelationJson, taskDefinitionJson, null, 
WorkflowExecutionTypeEnum.PARALLEL));
+
+        
Assertions.assertEquals(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION.getCode(), 
exception.getCode());
+        Mockito.verify(processService, Mockito.never())
+                .saveTaskDefine(eq(user), eq(projectCode), anyList(), 
eq(Boolean.TRUE));
+    }
+
+    @Test
+    public void 
testSwitchWorkflowDefinitionVersionShouldRejectUnauthorizedDatasource() {
+        WorkflowDefinition workflowDefinition = getWorkflowDefinition();
+        WorkflowDefinitionLog workflowDefinitionLog = new 
WorkflowDefinitionLog();
+        workflowDefinitionLog.setCode(processDefinitionCode);
+        workflowDefinitionLog.setProjectCode(projectCode);
+        workflowDefinitionLog.setVersion(1);
+        WorkflowTaskRelation workflowTaskRelation =
+                getWorkflowTaskRelation(1, 1, projectCode, 
processDefinitionCode, 0, 0, 123456789L, 1);
+
+        
when(workflowDefinitionDao.queryByCode(processDefinitionCode)).thenReturn(Optional.of(workflowDefinition));
+        
when(workflowDefinitionLogMapper.queryByDefinitionCodeAndVersion(processDefinitionCode,
 1))
+                .thenReturn(workflowDefinitionLog);
+        
when(workflowTaskRelationDao.queryWorkflowTaskRelationsByWorkflowDefinitionCode(processDefinitionCode,
 1))
+                .thenReturn(Collections.singletonList(workflowTaskRelation));
+        when(taskDefinitionLogMapper.queryByTaskDefinitions(anyList()))
+                .thenReturn(Collections.singletonList(new 
TaskDefinitionLog()));
+        doThrow(new 
ServiceException(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION))
+                
.when(taskDatasourcePermissionChecker).checkPermission(eq(user), anyList());
+
+        ServiceException exception = 
Assertions.assertThrows(ServiceException.class,
+                () -> 
workflowDefinitionService.switchWorkflowDefinitionVersion(
+                        user, projectCode, processDefinitionCode, 1));
+
+        
Assertions.assertEquals(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION.getCode(), 
exception.getCode());
+        Mockito.verify(processService, Mockito.never())
+                .switchVersion(any(WorkflowDefinition.class), 
any(WorkflowDefinitionLog.class));
+    }
+
+    @Test
+    public void 
testOnlineWorkflowDefinitionShouldRejectUnauthorizedDatasource() {
+        WorkflowDefinition workflowDefinition = getWorkflowDefinition();
+        WorkflowTaskRelation workflowTaskRelation =
+                getWorkflowTaskRelation(1, 1, projectCode, 
processDefinitionCode, 0, 0, 123456789L, 1);
+        
when(workflowDefinitionDao.queryByCode(processDefinitionCode)).thenReturn(Optional.of(workflowDefinition));
+        
when(workflowTaskRelationDao.queryByWorkflowDefinitionCode(processDefinitionCode))
+                .thenReturn(Collections.singletonList(workflowTaskRelation));
+        when(taskDefinitionLogDao.queryTaskDefineLogList(anyList()))
+                .thenReturn(Collections.singletonList(new 
TaskDefinitionLog()));
+        doThrow(new 
ServiceException(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION))
+                
.when(taskDatasourcePermissionChecker).checkPermission(eq(user), anyList());
+
+        ServiceException exception = 
Assertions.assertThrows(ServiceException.class,
+                () -> workflowDefinitionService.onlineWorkflowDefinition(user, 
projectCode, processDefinitionCode));
+
+        
Assertions.assertEquals(Status.RESOURCE_NOT_EXIST_OR_NO_PERMISSION.getCode(), 
exception.getCode());
+        Mockito.verify(workflowDefinitionDao, 
Mockito.never()).updateById(any(WorkflowDefinition.class));
+    }
+
     @Test
     public void testUpdateWorkflowDefinitionShouldSyncVersionToResponse() {
         Project project = getProject(projectCode);
diff --git 
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowInstanceServiceTest.java
 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowInstanceServiceTest.java
index b066cad2a8..aa694aa4ca 100644
--- 
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowInstanceServiceTest.java
+++ 
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/WorkflowInstanceServiceTest.java
@@ -29,6 +29,7 @@ import 
org.apache.dolphinscheduler.api.dto.workflowInstance.WorkflowInstanceTask
 import 
org.apache.dolphinscheduler.api.dto.workflowInstance.WorkflowInstanceVariablesDTO;
 import org.apache.dolphinscheduler.api.enums.Status;
 import org.apache.dolphinscheduler.api.exceptions.ServiceException;
+import 
org.apache.dolphinscheduler.api.permission.TaskDatasourcePermissionChecker;
 import org.apache.dolphinscheduler.api.service.impl.LoggerServiceImpl;
 import org.apache.dolphinscheduler.api.service.impl.ProjectServiceImpl;
 import 
org.apache.dolphinscheduler.api.service.impl.WorkflowInstanceServiceImpl;
@@ -140,6 +141,9 @@ public class WorkflowInstanceServiceTest {
     @Mock
     TaskDefinitionDao taskDefinitionDao;
 
+    @Mock
+    TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
+
     @Mock
     private TaskInstanceContextDao taskInstanceContextDao;
 
@@ -615,6 +619,8 @@ public class WorkflowInstanceServiceTest {
                     taskRelationJson, taskDefinitionJson, "2020-02-21 
00:00:00", Boolean.FALSE, "", "", 0);
             Assertions.assertNotNull(successRes);
         }
+        Mockito.verify(taskDatasourcePermissionChecker, Mockito.times(3))
+                .checkPermission(eq(loginUser), Mockito.anyList());
     }
 
     @Test

Reply via email to