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