diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java index 8c96bff9c8de..32bd866ec532 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java @@ -28,13 +28,17 @@ import org.apache.dolphinscheduler.api.audit.enums.AuditType; import org.apache.dolphinscheduler.api.exceptions.ApiException; import org.apache.dolphinscheduler.api.service.TaskDefinitionService; +import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.api.vo.TaskDefinitionVO; import org.apache.dolphinscheduler.common.constants.Constants; import org.apache.dolphinscheduler.common.enums.ReleaseState; +import org.apache.dolphinscheduler.dao.entity.TaskDefinitionLog; import org.apache.dolphinscheduler.dao.entity.User; import java.util.List; +import java.util.stream.Collectors; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.HttpStatus; @@ -88,7 +92,17 @@ public Result queryTaskDefinitionVersions(@Parameter(hidden = true) @RequestAttr @RequestParam(value = "pageNo") int pageNo, @RequestParam(value = "pageSize") int pageSize) { checkPageParams(pageNo, pageSize); - return taskDefinitionService.queryTaskDefinitionVersions(loginUser, projectCode, code, pageNo, pageSize); + Result result = taskDefinitionService.queryTaskDefinitionVersions(loginUser, projectCode, code, pageNo, + pageSize); + @SuppressWarnings("unchecked") + PageInfo pageInfo = (PageInfo) result.getData(); + if (pageInfo != null && pageInfo.getTotalList() != null) { + pageInfo.setTotalList(pageInfo.getTotalList().stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); + } + result.setData(pageInfo); + return result; } /** @@ -163,7 +177,7 @@ public Result queryTaskDefinitionDetail(@Parameter(hidden = tr @PathVariable(value = "code") long code) { TaskDefinitionVO taskDefinitionVO = taskDefinitionService.queryTaskDefinitionDetail(loginUser, projectCode, code); - return Result.success(taskDefinitionVO); + return Result.success(SensitivePropertyUtils.mask(taskDefinitionVO)); } /** diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java index 6ca1d7b34950..c26037172ed7 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java @@ -26,13 +26,18 @@ import org.apache.dolphinscheduler.api.audit.enums.AuditType; import org.apache.dolphinscheduler.api.exceptions.ApiException; import org.apache.dolphinscheduler.api.service.TaskInstanceService; +import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.common.constants.Constants; import org.apache.dolphinscheduler.common.enums.TaskExecuteType; +import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import org.apache.dolphinscheduler.plugin.task.api.utils.ParameterUtils; +import java.util.stream.Collectors; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.HttpStatus; import org.springframework.web.bind.annotation.GetMapping; @@ -112,7 +117,7 @@ public Result queryTaskListPaging(@Parameter(hidden = true) @RequestAttribute(va @RequestParam("pageSize") Integer pageSize) { checkPageParams(pageNo, pageSize); searchVal = ParameterUtils.handleEscapes(searchVal); - return taskInstanceService.queryTaskListPaging( + Result result = taskInstanceService.queryTaskListPaging( loginUser, projectCode, workflowInstanceId, @@ -129,6 +134,15 @@ public Result queryTaskListPaging(@Parameter(hidden = true) @RequestAttribute(va taskExecuteType, pageNo, pageSize); + @SuppressWarnings("unchecked") + PageInfo pageInfo = (PageInfo) result.getData(); + if (pageInfo != null && pageInfo.getTotalList() != null) { + pageInfo.setTotalList(pageInfo.getTotalList().stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); + } + result.setData(pageInfo); + return result; } /** diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionController.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionController.java index 29885ed4ac4e..941048535f39 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionController.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionController.java @@ -43,6 +43,7 @@ import org.apache.dolphinscheduler.api.service.WorkflowDefinitionService; import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.common.constants.Constants; import org.apache.dolphinscheduler.common.enums.ReleaseState; import org.apache.dolphinscheduler.common.enums.WorkflowExecutionTypeEnum; @@ -51,10 +52,14 @@ import org.apache.dolphinscheduler.dao.entity.TaskDefinition; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; +import org.apache.dolphinscheduler.dao.entity.WorkflowDefinitionLog; import org.apache.dolphinscheduler.plugin.task.api.utils.ParameterUtils; +import org.apache.dolphinscheduler.plugin.task.api.utils.PropertySensitiveUtils; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; @@ -128,7 +133,7 @@ public Result createWorkflowDefinition(@Parameter(hidden = t WorkflowDefinition workflowDefinition = workflowDefinitionService.createWorkflowDefinition(loginUser, projectCode, name, description, globalParams, locations, timeout, taskRelationJson, taskDefinitionJson, otherParamsJson, executionType); - return Result.success(workflowDefinition); + return Result.success(SensitivePropertyUtils.mask(workflowDefinition)); } /** @@ -254,7 +259,7 @@ public Result updateWorkflowDefinition(@Parameter(hidden = t if (releaseState == ReleaseState.ONLINE) { workflowDefinitionService.onlineWorkflowDefinition(loginUser, projectCode, code); } - return Result.success(workflowDefinition); + return Result.success(SensitivePropertyUtils.mask(workflowDefinition)); } /** @@ -283,8 +288,17 @@ public Result queryWorkflowDefinitionVersions(@Parameter(hidden = true) @Request @PathVariable(value = "code") long code) { checkPageParams(pageNo, pageSize); - return workflowDefinitionService.queryWorkflowDefinitionVersions(loginUser, projectCode, pageNo, pageSize, - code); + Result result = workflowDefinitionService.queryWorkflowDefinitionVersions(loginUser, projectCode, pageNo, + pageSize, code); + @SuppressWarnings("unchecked") + PageInfo pageInfo = (PageInfo) result.getData(); + if (pageInfo != null && pageInfo.getTotalList() != null) { + pageInfo.setTotalList(pageInfo.getTotalList().stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); + } + result.setData(pageInfo); + return result; } /** @@ -386,7 +400,7 @@ public Result queryWorkflowDefinitionByCode(@Parameter(hidden = true) @ @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, @PathVariable(value = "code", required = true) long code) { DagData dagData = workflowDefinitionService.queryWorkflowDefinitionByCode(loginUser, projectCode, code); - return Result.success(dagData); + return Result.success(SensitivePropertyUtils.mask(dagData)); } /** @@ -408,7 +422,7 @@ public Result queryWorkflowDefinitionByName(@Parameter(hidden = true) @ @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, @RequestParam("name") String name) { DagData dagData = workflowDefinitionService.queryWorkflowDefinitionByName(loginUser, projectCode, name); - return Result.success(dagData); + return Result.success(SensitivePropertyUtils.mask(dagData)); } /** @@ -424,7 +438,13 @@ public Result queryWorkflowDefinitionByName(@Parameter(hidden = true) @ @ApiException(QUERY_WORKFLOW_DEFINITION_LIST) public Result> queryWorkflowDefinitionList(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser, @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode) { - return Result.success(workflowDefinitionService.queryWorkflowDefinitionList(loginUser, projectCode)); + List dagDataList = workflowDefinitionService.queryWorkflowDefinitionList(loginUser, projectCode); + if (dagDataList == null) { + return Result.success(null); + } + return Result.success(dagDataList.stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); } /** @@ -480,6 +500,11 @@ public Result> queryWorkflowDefinitionListPaging( PageInfo pageInfo = workflowDefinitionService.queryWorkflowDefinitionListPaging( loginUser, projectCode, searchVal, otherParamsJson, userId, pageNo, pageSize); + if (pageInfo != null && pageInfo.getTotalList() != null) { + pageInfo.setTotalList(pageInfo.getTotalList().stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); + } return Result.success(pageInfo); } @@ -526,8 +551,14 @@ public Result viewTree(@Parameter(hidden = true) @RequestAttribute( public Result> getNodeListByDefinitionCode(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser, @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, @PathVariable("code") long code) { - return Result.success( - workflowDefinitionService.getTaskNodeListByDefinitionCode(loginUser, projectCode, code)); + List taskDefinitions = + workflowDefinitionService.getTaskNodeListByDefinitionCode(loginUser, projectCode, code); + if (taskDefinitions == null) { + return Result.success(null); + } + return Result.success(taskDefinitions.stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); } /** @@ -548,8 +579,18 @@ public Result> getNodeListByDefinitionCode(@Parameter(hidde public Result>> getNodeListMapByDefinitionCodes(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser, @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, @RequestParam("codes") String codes) { - return Result.success( - workflowDefinitionService.getNodeListMapByDefinitionCodes(loginUser, projectCode, codes)); + Map> taskMap = + workflowDefinitionService.getNodeListMapByDefinitionCodes(loginUser, projectCode, codes); + if (taskMap == null) { + return Result.success(null); + } + Map> masked = new LinkedHashMap<>(); + for (Map.Entry> entry : taskMap.entrySet()) { + List tasks = entry.getValue(); + masked.put(entry.getKey(), tasks == null ? null + : tasks.stream().map(SensitivePropertyUtils::mask).collect(Collectors.toList())); + } + return Result.success(masked); } /** @@ -645,8 +686,14 @@ public Result batchDeleteWorkflowDefinitionByCodes(@Parameter(hidden = tru @ApiException(QUERY_WORKFLOW_DEFINITION_LIST) public Result> queryAllWorkflowDefinitionByProjectCode(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser, @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode) { - return Result.success( - workflowDefinitionService.queryAllWorkflowDefinitionByProjectCode(loginUser, projectCode)); + List dagDataList = + workflowDefinitionService.queryAllWorkflowDefinitionByProjectCode(loginUser, projectCode); + if (dagDataList == null) { + return Result.success(null); + } + return Result.success(dagDataList.stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); } /** @@ -666,7 +713,14 @@ public Result> queryAllWorkflowDefinitionByProjectCode(@Parameter( public Result viewVariables(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser, @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, @PathVariable("code") Long code) { - return Result.success(workflowDefinitionService.viewVariables(loginUser, projectCode, code)); + WorkflowDefinitionVariablesDTO variables = + workflowDefinitionService.viewVariables(loginUser, projectCode, code); + if (variables == null) { + return Result.success(null); + } + return Result.success(new WorkflowDefinitionVariablesDTO( + PropertySensitiveUtils.maskSensitiveValues(variables.getGlobalParams()), + SensitivePropertyUtils.mask(variables.getLocalParams()))); } } diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceController.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceController.java index 20b450dd3dc4..f87a23a578a8 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceController.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceController.java @@ -29,6 +29,7 @@ import org.apache.dolphinscheduler.api.service.WorkflowInstanceService; import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.api.vo.WorkflowInstanceSummaryVO; import org.apache.dolphinscheduler.common.constants.Constants; import org.apache.dolphinscheduler.common.enums.WorkflowExecutionStatus; @@ -36,6 +37,7 @@ import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; import org.apache.dolphinscheduler.plugin.task.api.utils.ParameterUtils; +import org.apache.dolphinscheduler.plugin.task.api.utils.PropertySensitiveUtils; import org.apache.commons.lang3.StringUtils; @@ -43,6 +45,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; @@ -146,7 +149,15 @@ public Result queryTaskListByWorkflowInstanceId(@Pa @PathVariable("id") Integer id) { WorkflowInstanceTaskListDTO taskList = workflowInstanceService.queryTaskListByWorkflowInstanceId(loginUser, projectCode, id); - return Result.success(taskList); + if (taskList == null) { + return Result.success(null); + } + return Result.success(new WorkflowInstanceTaskListDTO( + taskList.getWorkflowInstanceState(), + taskList.getTaskList() == null ? null + : taskList.getTaskList().stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList()))); } /** @@ -189,7 +200,7 @@ public Result updateWorkflowInstance(@Parameter(hidden = tru @RequestParam(value = "timeout", required = false, defaultValue = "0") int timeout) { WorkflowDefinition workflowDefinition = workflowInstanceService.updateWorkflowInstance(loginUser, projectCode, id, taskRelationJson, taskDefinitionJson, scheduleTime, syncDefine, globalParams, locations, timeout); - return Result.success(workflowDefinition); + return Result.success(SensitivePropertyUtils.mask(workflowDefinition)); } /** @@ -212,7 +223,7 @@ public Result queryWorkflowInstanceById(@Parameter(hidden = tr @PathVariable("id") Integer id) { WorkflowInstance workflowInstance = workflowInstanceService.queryWorkflowInstanceById(loginUser, projectCode, id); - return Result.success(workflowInstance); + return Result.success(SensitivePropertyUtils.mask(workflowInstance)); } /** @@ -333,7 +344,12 @@ public Result viewVariables(@Parameter(hidden = tr @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, @PathVariable("id") Integer id) { WorkflowInstanceVariablesDTO variables = workflowInstanceService.viewVariables(loginUser, projectCode, id); - return Result.success(variables); + if (variables == null) { + return Result.success(null); + } + return Result.success(new WorkflowInstanceVariablesDTO( + PropertySensitiveUtils.maskSensitiveValues(variables.getGlobalParams()), + SensitivePropertyUtils.mask(variables.getLocalParams()))); } /** @@ -420,4 +436,5 @@ public Result> queryWorkflowInstancesByTriggerCo workflowInstanceService.queryByTriggerCode(loginUser, projectCode, triggerCode); return Result.success(workflowInstances); } + } 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 4470b8f1d2a3..f841c24c3015 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 @@ -49,6 +49,7 @@ import org.apache.dolphinscheduler.api.utils.CheckUtils; import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.api.validator.GlobalParamsValidator; import org.apache.dolphinscheduler.common.constants.Constants; import org.apache.dolphinscheduler.common.enums.ReleaseState; @@ -98,6 +99,7 @@ import org.apache.dolphinscheduler.plugin.task.api.parameters.ConditionsParameters; import org.apache.dolphinscheduler.plugin.task.api.parameters.DependentParameters; import org.apache.dolphinscheduler.plugin.task.api.parameters.SwitchParameters; +import org.apache.dolphinscheduler.plugin.task.api.utils.GlobalParameterUtils; import org.apache.dolphinscheduler.plugin.task.api.utils.TaskTypeUtils; import org.apache.dolphinscheduler.service.model.TaskNode; import org.apache.dolphinscheduler.service.process.ProcessService; @@ -262,6 +264,11 @@ public WorkflowDefinition createWorkflowDefinition(User loginUser, List taskDefinitionLogs = generateTaskDefinitionList(taskDefinitionJson); List taskRelationList = generateTaskRelationList(taskRelationJson, taskDefinitionLogs); + SensitivePropertyUtils.requireNoPlaceholder(GlobalParameterUtils.deserializeGlobalParameter(globalParams)); + for (TaskDefinitionLog taskDefinitionLog : CollectionUtils.emptyIfNull(taskDefinitionLogs)) { + taskDefinitionLog.setTaskParams( + SensitivePropertyUtils.mergeLocalParams(taskDefinitionLog.getTaskParams(), null)); + } long workflowDefinitionCode = CodeGenerateUtils.genCode(); WorkflowDefinition workflowDefinition = @@ -656,6 +663,34 @@ public WorkflowDefinition updateWorkflowDefinition(User loginUser, throw new ServiceException(Status.WORKFLOW_DEFINITION_NAME_EXIST, name); } } + List submittedGlobalParams = + GlobalParameterUtils.deserializeGlobalParameter(globalParams); + if (CollectionUtils.isNotEmpty(submittedGlobalParams)) { + globalParams = GlobalParameterUtils.serializeGlobalParameter( + SensitivePropertyUtils.merge(submittedGlobalParams, + GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams()))); + } + List versionKeys = new ArrayList<>(); + for (TaskDefinitionLog taskDefinitionLog : CollectionUtils.emptyIfNull(taskDefinitionLogs)) { + if (taskDefinitionLog.getCode() > 0 && taskDefinitionLog.getVersion() > 0) { + versionKeys.add(new TaskDefinition(taskDefinitionLog.getCode(), taskDefinitionLog.getVersion())); + } + } + List existingTaskLogs = CollectionUtils.isEmpty(versionKeys) + ? Collections.emptyList() + : taskDefinitionLogMapper.queryByTaskDefinitions(versionKeys); + Map existingTaskParamsMap = existingTaskLogs.stream() + .collect(Collectors.toMap( + log -> log.getCode() + "_" + log.getVersion(), + TaskDefinitionLog::getTaskParams, + (left, right) -> left)); + for (TaskDefinitionLog submitted : CollectionUtils.emptyIfNull(taskDefinitionLogs)) { + String existingTaskParams = submitted.getCode() <= 0 || submitted.getVersion() <= 0 + ? null + : existingTaskParamsMap.get(submitted.getCode() + "_" + submitted.getVersion()); + submitted.setTaskParams(SensitivePropertyUtils.mergeLocalParams( + submitted.getTaskParams(), existingTaskParams)); + } WorkflowDefinition workflowDefinitionDeepCopy = JSONUtils.parseObject(JSONUtils.toJsonString(workflowDefinition), WorkflowDefinition.class); workflowDefinition.set(projectCode, name, description, globalParams, locations, timeout); 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 f3dd1179ba12..6d95bf4ba5e6 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 @@ -41,6 +41,7 @@ import org.apache.dolphinscheduler.api.service.WorkflowInstanceService; import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.api.vo.WorkflowInstanceSummaryVO; import org.apache.dolphinscheduler.common.constants.Constants; import org.apache.dolphinscheduler.common.enums.ContextType; @@ -411,6 +412,12 @@ public WorkflowDefinition updateWorkflowInstance(User loginUser, long projectCod timezoneId = commandParam.getTimeZone(); } + List submittedGlobalParams = GlobalParameterUtils.deserializeGlobalParameter(globalParams); + if (CollectionUtils.isNotEmpty(submittedGlobalParams)) { + globalParams = GlobalParameterUtils.serializeGlobalParameter( + SensitivePropertyUtils.merge(submittedGlobalParams, + GlobalParameterUtils.deserializeGlobalParameter(workflowInstance.getGlobalParams()))); + } setWorkflowInstance(workflowInstance, scheduleTime, globalParams, timeout, timezoneId); List taskDefinitionLogs = JSONUtils.toList(taskDefinitionJson, TaskDefinitionLog.class); if (taskDefinitionLogs.isEmpty()) { @@ -423,6 +430,27 @@ public WorkflowDefinition updateWorkflowInstance(User loginUser, long projectCod throw new ServiceException(Status.WORKFLOW_NODE_S_PARAMETER_INVALID, taskDefinitionLog.getName()); } } + List versionKeys = new ArrayList<>(); + for (TaskDefinitionLog taskDefinitionLog : taskDefinitionLogs) { + if (taskDefinitionLog.getCode() > 0 && taskDefinitionLog.getVersion() > 0) { + versionKeys.add(new TaskDefinition(taskDefinitionLog.getCode(), taskDefinitionLog.getVersion())); + } + } + List existingTaskLogs = CollectionUtils.isEmpty(versionKeys) + ? Collections.emptyList() + : taskDefinitionLogMapper.queryByTaskDefinitions(versionKeys); + Map existingTaskParamsMap = existingTaskLogs.stream() + .collect(Collectors.toMap( + log -> log.getCode() + "_" + log.getVersion(), + TaskDefinitionLog::getTaskParams, + (left, right) -> left)); + for (TaskDefinitionLog submitted : taskDefinitionLogs) { + String existingTaskParams = submitted.getCode() <= 0 || submitted.getVersion() <= 0 + ? null + : existingTaskParamsMap.get(submitted.getCode() + "_" + submitted.getVersion()); + submitted.setTaskParams(SensitivePropertyUtils.mergeLocalParams( + submitted.getTaskParams(), existingTaskParams)); + } taskDatasourcePermissionChecker.checkPermission(loginUser, taskDefinitionLogs); taskSubWorkflowPermissionChecker.checkPermission(loginUser, taskDefinitionLogs); int saveTaskResult = processService.saveTaskDefine(loginUser, projectCode, taskDefinitionLogs, syncDefine); diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/utils/SensitivePropertyUtils.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/utils/SensitivePropertyUtils.java new file mode 100644 index 000000000000..40b732b77314 --- /dev/null +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/utils/SensitivePropertyUtils.java @@ -0,0 +1,299 @@ +/* + * 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.utils; + +import static org.apache.dolphinscheduler.common.constants.Constants.LOCAL_PARAMS; +import static org.apache.dolphinscheduler.plugin.task.api.TaskConstants.LOCAL_PARAMS_LIST; + +import org.apache.dolphinscheduler.api.enums.Status; +import org.apache.dolphinscheduler.api.exceptions.ServiceException; +import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.dao.entity.DagData; +import org.apache.dolphinscheduler.dao.entity.TaskDefinition; +import org.apache.dolphinscheduler.dao.entity.TaskInstance; +import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; +import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; +import org.apache.dolphinscheduler.extract.master.command.AbstractCommandParam; +import org.apache.dolphinscheduler.extract.master.command.ICommandParam; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; +import org.apache.dolphinscheduler.plugin.task.api.utils.GlobalParameterUtils; +import org.apache.dolphinscheduler.plugin.task.api.utils.PropertySensitiveUtils; +import org.apache.dolphinscheduler.plugin.task.api.utils.VarPoolUtils; + +import org.apache.commons.collections4.CollectionUtils; +import org.apache.commons.lang3.StringUtils; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.function.Function; +import java.util.stream.Collectors; + +import lombok.experimental.UtilityClass; + +import org.springframework.beans.BeanUtils; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; + +/** + * HTTP-layer helpers for {@link Property#isSensitive()}. + * Values stay plaintext in DB. Query masks them as {@code ******}. + * Create rejects {@code ******}; Update {@link #merge}s it back to the old value. + * Start: {@code ******} is replaced with the definition global plaintext. + */ +@UtilityClass +public class SensitivePropertyUtils { + + /** Create: {@code ******} is never a real value. */ + public void requireNoPlaceholder(List properties) { + String invalidProp = PropertySensitiveUtils.findPlaceholderProp(properties); + if (invalidProp != null) { + throw new ServiceException(Status.REQUEST_PARAMS_NOT_VALID_ERROR, + "parameter '" + invalidProp + + "' cannot use ****** when creating; please re-enter the value"); + } + } + + /** + * Update: restore DB plaintext for keep-original {@code ******}. + * Empty / null is a real empty value. {@code false→true} + {@code ******} is allowed; + * {@code true→false} + {@code ******} is rejected. + */ + public List merge(List submittedProperties, List existingProperties) { + if (CollectionUtils.isEmpty(submittedProperties)) { + return submittedProperties; + } + String invalidProp = PropertySensitiveUtils.findInvalidSensitivePlaceholderProp( + submittedProperties, existingProperties); + if (invalidProp != null) { + throw new ServiceException(Status.REQUEST_PARAMS_NOT_VALID_ERROR, + "parameter '" + invalidProp + + "' cannot use ****** when creating, enabling, or disabling sensitive; please re-enter the value"); + } + return PropertySensitiveUtils.mergeSensitiveValuePlaceholders(submittedProperties, existingProperties); + } + + /** Start: {@code ******} is replaced with the matching definition global plaintext. */ + public List restoreStartParams(List startParams, List globalParams) { + if (CollectionUtils.isEmpty(startParams)) { + return startParams; + } + Map globals = CollectionUtils.emptyIfNull(globalParams).stream() + .filter(Objects::nonNull) + .filter(property -> property.getProp() != null) + .collect(Collectors.toMap(Property::getProp, Function.identity(), (left, right) -> right)); + List restored = new ArrayList<>(); + for (Property startParam : startParams) { + if (startParam == null) { + continue; + } + if (!PropertySensitiveUtils.isSensitiveValuePlaceholder(startParam.getValue())) { + restored.add(startParam); + continue; + } + Property global = globals.get(startParam.getProp()); + if (global != null) { + restored.add(PropertySensitiveUtils.copy(global)); + } + } + return restored; + } + + public String mergeLocalParams(String submittedTaskParams, String existingTaskParams) { + return rewriteLocalParams(submittedTaskParams, submitted -> { + if (StringUtils.isEmpty(existingTaskParams)) { + requireNoPlaceholder(submitted); + return submitted; + } + return merge(submitted, getLocalParams(existingTaskParams)); + }); + } + + /** + * Copy then mask for HTTP responses. Do not mask Service return values in place. + * Nested {@link DagData} is remasked from the original (BeanUtils is shallow). + */ + public Map> mask(Map> source) { + if (source == null) { + return null; + } + return maskLocalParamsMap(source); + } + + public DagData mask(DagData source) { + if (source == null) { + return null; + } + DagData copy = copyBean(source); + copy.setWorkflowDefinition(mask(source.getWorkflowDefinition())); + if (source.getTaskDefinitionList() != null) { + copy.setTaskDefinitionList(source.getTaskDefinitionList().stream() + .map(SensitivePropertyUtils::mask) + .collect(Collectors.toList())); + } + return copy; + } + + public WorkflowInstance mask(WorkflowInstance source) { + if (source == null) { + return null; + } + WorkflowInstance copy = copyBean(source); + copy.setGlobalParams(maskGlobalParams(copy.getGlobalParams())); + copy.setVarPool(maskVarPool(copy.getVarPool())); + copy.setCommandParam(maskCommandParam(copy.getCommandParam())); + copy.setDagData(mask(source.getDagData())); + return copy; + } + + public T mask(T source) { + if (source == null) { + return null; + } + T copy = copyBean(source); + copy.setGlobalParams(maskGlobalParams(copy.getGlobalParams())); + copy.setGlobalParamMap(null); + return copy; + } + + public T mask(T source) { + if (source == null) { + return null; + } + T copy = copyBean(source); + copy.setTaskParams(rewriteLocalParams(copy.getTaskParams(), + PropertySensitiveUtils::maskSensitiveValues)); + copy.setTaskParamMap(null); + return copy; + } + + public T mask(T source) { + if (source == null) { + return null; + } + T copy = copyBean(source); + copy.setTaskParams(rewriteLocalParams(copy.getTaskParams(), + PropertySensitiveUtils::maskSensitiveValues)); + copy.setVarPool(maskVarPool(copy.getVarPool())); + return copy; + } + + private String maskGlobalParams(String globalParams) { + List properties = GlobalParameterUtils.deserializeGlobalParameter(globalParams); + if (CollectionUtils.isEmpty(properties)) { + return globalParams; + } + return GlobalParameterUtils.serializeGlobalParameter(PropertySensitiveUtils.maskSensitiveValues(properties)); + } + + private String maskVarPool(String varPool) { + if (StringUtils.isEmpty(varPool)) { + return varPool; + } + List properties = VarPoolUtils.deserializeVarPool(varPool); + if (CollectionUtils.isEmpty(properties)) { + return varPool; + } + return VarPoolUtils.serializeVarPool(PropertySensitiveUtils.maskSensitiveValues(properties)); + } + + /** + * Mask {@link ICommandParam#getCommandParams()} in the response copy. + * Start restores {@code ******} to plaintext before Master persists {@code commandParam}; + * query must hide those values again. + */ + private String maskCommandParam(String commandParam) { + if (StringUtils.isEmpty(commandParam)) { + return commandParam; + } + ICommandParam parsed = JSONUtils.parseObject(commandParam, ICommandParam.class); + if (!(parsed instanceof AbstractCommandParam)) { + return commandParam; + } + AbstractCommandParam abstractCommandParam = (AbstractCommandParam) parsed; + if (CollectionUtils.isEmpty(abstractCommandParam.getCommandParams())) { + return commandParam; + } + abstractCommandParam.setCommandParams( + PropertySensitiveUtils.maskSensitiveValues(abstractCommandParam.getCommandParams())); + return JSONUtils.toJsonString(abstractCommandParam); + } + + private Map> maskLocalParamsMap(Map> localParams) { + Map> masked = new LinkedHashMap<>(); + for (Map.Entry> entry : localParams.entrySet()) { + Map inner = entry.getValue(); + if (inner == null) { + masked.put(entry.getKey(), null); + continue; + } + Map copied = new LinkedHashMap<>(inner); + Object localParamsList = copied.get(LOCAL_PARAMS_LIST); + if (localParamsList instanceof List) { + @SuppressWarnings("unchecked") + List properties = (List) localParamsList; + copied.put(LOCAL_PARAMS_LIST, PropertySensitiveUtils.maskSensitiveValues(properties)); + } + masked.put(entry.getKey(), copied); + } + return masked; + } + + private List getLocalParams(String taskParams) { + if (StringUtils.isEmpty(taskParams)) { + return Collections.emptyList(); + } + String localParams = JSONUtils.getNodeString(taskParams, LOCAL_PARAMS); + if (StringUtils.isEmpty(localParams)) { + return Collections.emptyList(); + } + return JSONUtils.toList(localParams, Property.class); + } + + @SuppressWarnings("unchecked") + private T copyBean(T source) { + try { + T copy = (T) source.getClass().getDeclaredConstructor().newInstance(); + BeanUtils.copyProperties(source, copy); + return copy; + } catch (ReflectiveOperationException e) { + throw new IllegalStateException("Failed to copy " + source.getClass().getSimpleName() + " for masking", e); + } + } + + private String rewriteLocalParams(String taskParams, Function, List> transform) { + if (StringUtils.isEmpty(taskParams)) { + return taskParams; + } + ObjectNode taskParamsNode = JSONUtils.parseObject(taskParams); + if (taskParamsNode == null) { + return taskParams; + } + JsonNode localParamsNode = taskParamsNode.findValue(LOCAL_PARAMS); + if (localParamsNode == null || localParamsNode.isNull()) { + return taskParams; + } + List localParams = JSONUtils.toList(localParamsNode.toString(), Property.class); + taskParamsNode.set(LOCAL_PARAMS, JSONUtils.toJsonNode(transform.apply(localParams))); + return JSONUtils.toJsonString(taskParamsNode); + } +} diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/BackfillWorkflowRequestTransformer.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/BackfillWorkflowRequestTransformer.java index 8fea31a95385..b66f6b2a00ef 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/BackfillWorkflowRequestTransformer.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/BackfillWorkflowRequestTransformer.java @@ -19,6 +19,7 @@ import org.apache.dolphinscheduler.api.dto.workflow.WorkflowBackFillRequest; import org.apache.dolphinscheduler.api.exceptions.ServiceException; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.api.utils.WorkflowUtils; import org.apache.dolphinscheduler.api.validator.ITransformer; import org.apache.dolphinscheduler.common.utils.CodeGenerateUtils; @@ -26,6 +27,7 @@ import org.apache.dolphinscheduler.dao.entity.Schedule; import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; import org.apache.dolphinscheduler.dao.repository.WorkflowDefinitionDao; +import org.apache.dolphinscheduler.plugin.task.api.utils.GlobalParameterUtils; import org.apache.dolphinscheduler.plugin.task.api.utils.PropertyUtils; import org.apache.dolphinscheduler.service.cron.CronUtils; import org.apache.dolphinscheduler.service.process.ProcessService; @@ -54,6 +56,10 @@ public class BackfillWorkflowRequestTransformer implements ITransformer new ServiceException( + "Cannot find the workflow: " + workflowBackFillRequest.getWorkflowDefinitionCode())); final BackfillWorkflowDTO.BackfillParamsDTO backfillParams = transformBackfillParamsDTO(workflowBackFillRequest); final BackfillWorkflowDTO backfillWorkflowDTO = BackfillWorkflowDTO.builder() @@ -69,18 +75,13 @@ public BackfillWorkflowDTO transform(WorkflowBackFillRequest workflowBackFillReq .workerGroup(workflowBackFillRequest.getWorkerGroup()) .tenantCode(workflowBackFillRequest.getTenantCode()) .environmentCode(workflowBackFillRequest.getEnvironmentCode()) - .startParamList( - PropertyUtils.startParamsTransformPropertyList(workflowBackFillRequest.getStartParamList())) + .startParamList(SensitivePropertyUtils.restoreStartParams( + PropertyUtils.startParamsTransformPropertyList(workflowBackFillRequest.getStartParamList()), + GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams()))) .dryRun(workflowBackFillRequest.getDryRun()) .triggerCode(CodeGenerateUtils.genCode()) .backfillParams(backfillParams) .build(); - - WorkflowDefinition workflowDefinition = workflowDefinitionDao - .queryByCode(workflowBackFillRequest.getWorkflowDefinitionCode()) - .orElseThrow(() -> new ServiceException( - "Cannot find the workflow: " + workflowBackFillRequest.getWorkflowDefinitionCode())); - backfillWorkflowDTO.setWorkflowDefinition(workflowDefinition); return backfillWorkflowDTO; } diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/TriggerWorkflowRequestTransformer.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/TriggerWorkflowRequestTransformer.java index c1c85acbcdaa..748756c73396 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/TriggerWorkflowRequestTransformer.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/validator/workflow/TriggerWorkflowRequestTransformer.java @@ -19,10 +19,12 @@ import org.apache.dolphinscheduler.api.dto.workflow.WorkflowTriggerRequest; import org.apache.dolphinscheduler.api.exceptions.ServiceException; +import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; import org.apache.dolphinscheduler.api.utils.WorkflowUtils; import org.apache.dolphinscheduler.api.validator.ITransformer; import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; import org.apache.dolphinscheduler.dao.repository.WorkflowDefinitionDao; +import org.apache.dolphinscheduler.plugin.task.api.utils.GlobalParameterUtils; import org.apache.dolphinscheduler.plugin.task.api.utils.PropertyUtils; import lombok.extern.slf4j.Slf4j; @@ -39,6 +41,11 @@ public class TriggerWorkflowRequestTransformer implements ITransformer new ServiceException( + "Cannot find the workflow: " + workflowTriggerRequest.getWorkflowDefinitionCode())); + TriggerWorkflowDTO triggerWorkflowDTO = TriggerWorkflowDTO.builder() .loginUser(workflowTriggerRequest.getLoginUser()) .startNodes(WorkflowUtils.parseStartNodeList(workflowTriggerRequest.getStartNodes())) @@ -51,16 +58,11 @@ public TriggerWorkflowDTO transform(WorkflowTriggerRequest workflowTriggerReques .workerGroup(workflowTriggerRequest.getWorkerGroup()) .tenantCode(workflowTriggerRequest.getTenantCode()) .environmentCode(workflowTriggerRequest.getEnvironmentCode()) - .startParamList( - PropertyUtils.startParamsTransformPropertyList(workflowTriggerRequest.getStartParamList())) + .startParamList(SensitivePropertyUtils.restoreStartParams( + PropertyUtils.startParamsTransformPropertyList(workflowTriggerRequest.getStartParamList()), + GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams()))) .dryRun(workflowTriggerRequest.getDryRun()) .build(); - - WorkflowDefinition workflowDefinition = workflowDefinitionDao - .queryByCode(workflowTriggerRequest.getWorkflowDefinitionCode()) - .orElseThrow(() -> new ServiceException( - "Cannot find the workflow: " + workflowTriggerRequest.getWorkflowDefinitionCode())); - triggerWorkflowDTO.setWorkflowDefinition(workflowDefinition); return triggerWorkflowDTO; } diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionControllerTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionControllerTest.java new file mode 100644 index 000000000000..3f75f9193c2a --- /dev/null +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionControllerTest.java @@ -0,0 +1,104 @@ +/* + * 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.controller; + +import org.apache.dolphinscheduler.api.enums.Status; +import org.apache.dolphinscheduler.api.service.TaskDefinitionService; +import org.apache.dolphinscheduler.api.utils.PageInfo; +import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.vo.TaskDefinitionVO; +import org.apache.dolphinscheduler.common.enums.UserType; +import org.apache.dolphinscheduler.dao.entity.TaskDefinitionLog; +import org.apache.dolphinscheduler.dao.entity.User; +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; + +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.InjectMocks; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.jupiter.MockitoExtension; + +@ExtendWith(MockitoExtension.class) +public class TaskDefinitionControllerTest { + + @InjectMocks + private TaskDefinitionController taskDefinitionController; + + @Mock + private TaskDefinitionService taskDefinitionService; + + private User user; + + @BeforeEach + public void before() { + User loginUser = new User(); + loginUser.setId(1); + loginUser.setUserType(UserType.GENERAL_USER); + loginUser.setUserName("admin"); + user = loginUser; + } + + @Test + public void testQueryTaskDefinitionDetailMasksSensitiveLocalParams() { + TaskDefinitionVO taskDefinitionVO = new TaskDefinitionVO(); + taskDefinitionVO.setTaskParams(sensitiveTaskParams()); + Mockito.when(taskDefinitionService.queryTaskDefinitionDetail(user, 1L, 2L)) + .thenReturn(taskDefinitionVO); + + Result response = taskDefinitionController.queryTaskDefinitionDetail(user, 1L, 2L); + + Assertions.assertEquals(Status.SUCCESS.getCode(), response.getCode().intValue()); + Assertions.assertTrue(response.getData().getTaskParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertFalse(response.getData().getTaskParams().contains("abc")); + Assertions.assertTrue(taskDefinitionVO.getTaskParams().contains("abc")); + } + + @Test + public void testQueryTaskDefinitionVersionsMasksSensitiveLocalParams() { + TaskDefinitionLog taskDefinitionLog = new TaskDefinitionLog(); + taskDefinitionLog.setTaskParams(sensitiveTaskParams()); + PageInfo pageInfo = new PageInfo<>(1, 10); + pageInfo.setTotalList(Collections.singletonList(taskDefinitionLog)); + Result result = new Result(); + result.setCode(Status.SUCCESS.getCode()); + result.setMsg(Status.SUCCESS.getMsg()); + result.setData(pageInfo); + Mockito.when(taskDefinitionService.queryTaskDefinitionVersions(user, 1L, 2L, 1, 10)) + .thenReturn(result); + + Result response = taskDefinitionController.queryTaskDefinitionVersions(user, 1L, 2L, 1, 10); + + Assertions.assertEquals(Status.SUCCESS.getCode(), response.getCode().intValue()); + @SuppressWarnings("unchecked") + PageInfo maskedPage = (PageInfo) response.getData(); + Assertions.assertTrue( + maskedPage.getTotalList().get(0).getTaskParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertFalse(maskedPage.getTotalList().get(0).getTaskParams().contains("abc")); + Assertions.assertTrue(taskDefinitionLog.getTaskParams().contains("abc")); + } + + private static String sensitiveTaskParams() { + return "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"abc\",\"sensitive\":true}]}"; + } +} diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java index b58944537b87..03512c9b4f93 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java @@ -34,8 +34,11 @@ import org.apache.dolphinscheduler.common.utils.JSONUtils; import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.User; +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; +import java.util.Collections; + import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; @@ -61,7 +64,11 @@ public void testQueryTaskListPaging() { Result result = new Result(); Integer pageNo = 1; Integer pageSize = 20; - PageInfo pageInfo = new PageInfo(pageNo, pageSize); + TaskInstance taskInstance = new TaskInstance(); + taskInstance.setTaskParams("{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"abc\",\"sensitive\":true}]}"); + PageInfo pageInfo = new PageInfo<>(pageNo, pageSize); + pageInfo.setTotalList(Collections.singletonList(taskInstance)); result.setData(pageInfo); result.setCode(Status.SUCCESS.getCode()); result.setMsg(Status.SUCCESS.getMsg()); @@ -74,6 +81,12 @@ public void testQueryTaskListPaging() { "", 1L, "", TaskExecutionStatus.SUCCESS, "192.168.xx.xx", "2020-01-01 00:00:00", "2020-01-02 00:00:00", TaskExecuteType.BATCH, pageNo, pageSize); Assertions.assertEquals(Integer.valueOf(Status.SUCCESS.getCode()), taskResult.getCode()); + @SuppressWarnings("unchecked") + PageInfo maskedPage = (PageInfo) taskResult.getData(); + Assertions.assertTrue( + maskedPage.getTotalList().get(0).getTaskParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertFalse(maskedPage.getTotalList().get(0).getTaskParams().contains("abc")); + Assertions.assertTrue(taskInstance.getTaskParams().contains("abc")); } @Disabled diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionControllerTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionControllerTest.java index a65213038a97..66242295fdf8 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionControllerTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowDefinitionControllerTest.java @@ -30,15 +30,23 @@ import org.apache.dolphinscheduler.common.enums.UserType; import org.apache.dolphinscheduler.common.enums.WorkflowExecutionTypeEnum; import org.apache.dolphinscheduler.dao.entity.DagData; +import org.apache.dolphinscheduler.dao.entity.TaskDefinition; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; import org.apache.dolphinscheduler.dao.entity.WorkflowDefinitionLog; +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; +import org.apache.dolphinscheduler.plugin.task.api.enums.DataType; +import org.apache.dolphinscheduler.plugin.task.api.enums.Direct; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; +import org.apache.dolphinscheduler.plugin.task.api.utils.GlobalParameterUtils; import java.text.MessageFormat; import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; @@ -82,6 +90,7 @@ public void testCreateWorkflowDefinition() { WorkflowDefinition workflowDefinition = new WorkflowDefinition(); workflowDefinition.setName(name); + workflowDefinition.setGlobalParams(sensitiveGlobalParams()); Mockito.when(processDefinitionService.createWorkflowDefinition(user, projectCode, name, description, globalParams, locations, timeout, relationJson, taskDefinitionJson, "", @@ -92,6 +101,8 @@ public void testCreateWorkflowDefinition() { name, description, globalParams, locations, timeout, relationJson, taskDefinitionJson, "", WorkflowExecutionTypeEnum.PARALLEL); Assertions.assertEquals(Status.SUCCESS.getCode(), response.getCode().intValue()); + assertMaskedAndOriginalUnchanged(workflowDefinition.getGlobalParams(), + response.getData().getGlobalParams()); } public void putMsg(Result result, Status status, Object... statusParams) { @@ -129,6 +140,7 @@ public void updateWorkflowDefinition() { WorkflowDefinition workflowDefinition = new WorkflowDefinition(); workflowDefinition.setCode(code); + workflowDefinition.setGlobalParams(sensitiveGlobalParams()); Mockito.when(processDefinitionService.updateWorkflowDefinition(user, projectCode, name, code, description, globalParams, locations, timeout, relationJson, taskDefinitionJson, @@ -138,6 +150,8 @@ public void updateWorkflowDefinition() { name, code, description, globalParams, locations, timeout, relationJson, taskDefinitionJson, WorkflowExecutionTypeEnum.PARALLEL, ReleaseState.OFFLINE); Assertions.assertEquals(Status.SUCCESS.getCode(), response.getCode().intValue()); + assertMaskedAndOriginalUnchanged(workflowDefinition.getGlobalParams(), + response.getData().getGlobalParams()); } @Test @@ -159,13 +173,21 @@ public void testQueryWorkflowDefinitionByCode() { WorkflowDefinition workflowDefinition = new WorkflowDefinition(); workflowDefinition.setCode(code); - DagData dagData = new DagData(workflowDefinition, Collections.emptyList(), Collections.emptyList()); + workflowDefinition.setGlobalParams(sensitiveGlobalParams()); + TaskDefinition taskDefinition = new TaskDefinition(); + taskDefinition.setTaskParams(sensitiveTaskParams()); + DagData dagData = new DagData(workflowDefinition, Collections.emptyList(), + Collections.singletonList(taskDefinition)); Mockito.when(processDefinitionService.queryWorkflowDefinitionByCode(user, projectCode, code)) .thenReturn(dagData); Result response = workflowDefinitionController.queryWorkflowDefinitionByCode(user, projectCode, code); Assertions.assertEquals(Status.SUCCESS.getCode(), response.getCode().intValue()); + assertMaskedAndOriginalUnchanged(workflowDefinition.getGlobalParams(), + response.getData().getWorkflowDefinition().getGlobalParams()); + assertMaskedAndOriginalUnchanged(taskDefinition.getTaskParams(), + response.getData().getTaskDefinitionList().get(0).getTaskParams()); } @Test @@ -281,7 +303,10 @@ public void testQueryWorkflowDefinitionListPaging() { String searchVal = ""; int userId = 1; + WorkflowDefinition workflowDefinition = new WorkflowDefinition(); + workflowDefinition.setGlobalParams(sensitiveGlobalParams()); PageInfo pageInfo = new PageInfo<>(1, 10); + pageInfo.setTotalList(Collections.singletonList(workflowDefinition)); Mockito.when( processDefinitionService.queryWorkflowDefinitionListPaging(user, projectCode, searchVal, "", userId, @@ -291,6 +316,8 @@ public void testQueryWorkflowDefinitionListPaging() { .queryWorkflowDefinitionListPaging(user, projectCode, searchVal, "", userId, pageNo, pageSize); Assertions.assertTrue(response != null && response.isSuccess()); + assertMaskedAndOriginalUnchanged(workflowDefinition.getGlobalParams(), + response.getData().getTotalList().get(0).getGlobalParams()); } @Test @@ -299,7 +326,11 @@ public void testQueryWorkflowDefinitionVersions() { long projectCode = 1L; Result resultMap = new Result(); putMsg(resultMap, Status.SUCCESS); - resultMap.setData(new PageInfo(1, 10)); + WorkflowDefinitionLog workflowDefinitionLog = new WorkflowDefinitionLog(); + workflowDefinitionLog.setGlobalParams(sensitiveGlobalParams()); + PageInfo pageInfo = new PageInfo<>(1, 10); + pageInfo.setTotalList(Collections.singletonList(workflowDefinitionLog)); + resultMap.setData(pageInfo); Mockito.when(processDefinitionService.queryWorkflowDefinitionVersions( user, projectCode, 1, 10, 1)) .thenReturn(resultMap); @@ -307,6 +338,10 @@ public void testQueryWorkflowDefinitionVersions() { user, projectCode, 1, 10, 1); Assertions.assertEquals(Status.SUCCESS.getCode(), (int) result.getCode()); + @SuppressWarnings("unchecked") + PageInfo maskedPage = (PageInfo) result.getData(); + assertMaskedAndOriginalUnchanged(workflowDefinitionLog.getGlobalParams(), + maskedPage.getTotalList().get(0).getGlobalParams()); } @Test @@ -332,12 +367,55 @@ public void testDeleteWorkflowDefinitionVersion() { public void testViewVariables() { long projectCode = 1L; + Map localParam = new LinkedHashMap<>(); + localParam.put(TaskConstants.LOCAL_PARAMS_LIST, Collections.singletonList(sensitive("token", "abc"))); + WorkflowDefinitionVariablesDTO variables = new WorkflowDefinitionVariablesDTO( + Collections.singletonList(sensitive("pwd", "Secret123")), + Collections.singletonMap("shell-1", localParam)); Mockito.when(processDefinitionService.viewVariables(user, projectCode, 1L)) - .thenReturn(new WorkflowDefinitionVariablesDTO()); + .thenReturn(variables); Result result = workflowDefinitionController.viewVariables(user, projectCode, 1L); Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode().intValue()); + WorkflowDefinitionVariablesDTO masked = (WorkflowDefinitionVariablesDTO) result.getData(); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, masked.getGlobalParams().get(0).getValue()); + Assertions.assertEquals("Secret123", variables.getGlobalParams().get(0).getValue()); + @SuppressWarnings("unchecked") + List maskedLocalParams = + (List) masked.getLocalParams().get("shell-1").get(TaskConstants.LOCAL_PARAMS_LIST); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, maskedLocalParams.get(0).getValue()); + @SuppressWarnings("unchecked") + List originalLocalParams = + (List) variables.getLocalParams().get("shell-1").get(TaskConstants.LOCAL_PARAMS_LIST); + Assertions.assertEquals("abc", originalLocalParams.get(0).getValue()); + } + + private static String sensitiveGlobalParams() { + return GlobalParameterUtils.serializeGlobalParameter(Collections.singletonList(sensitive("pwd", "Secret123"))); + } + + private static String sensitiveTaskParams() { + return "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"abc\",\"sensitive\":true}]}"; + } + + private static Property sensitive(String prop, String value) { + return Property.builder() + .prop(prop) + .direct(Direct.IN) + .type(DataType.VARCHAR) + .value(value) + .sensitive(true) + .build(); + } + + private static void assertMaskedAndOriginalUnchanged(String original, String masked) { + Assertions.assertTrue(original.contains("Secret123") || original.contains("\"abc\"")); + Assertions.assertFalse(original.contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertTrue(masked.contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertFalse(masked.contains("Secret123")); + Assertions.assertFalse(masked.contains("\"abc\"")); } } diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceControllerTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceControllerTest.java index 4182f4c43066..c5fea82ad5e7 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceControllerTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/WorkflowInstanceControllerTest.java @@ -24,6 +24,7 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.content; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import org.apache.dolphinscheduler.api.dto.workflowInstance.WorkflowInstanceTaskListDTO; import org.apache.dolphinscheduler.api.dto.workflowInstance.WorkflowInstanceVariablesDTO; import org.apache.dolphinscheduler.api.enums.Status; import org.apache.dolphinscheduler.api.exceptions.ServiceException; @@ -33,13 +34,23 @@ import org.apache.dolphinscheduler.api.vo.WorkflowInstanceSummaryVO; import org.apache.dolphinscheduler.common.enums.WorkflowExecutionStatus; import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.dao.entity.DagData; +import org.apache.dolphinscheduler.dao.entity.TaskDefinition; +import org.apache.dolphinscheduler.dao.entity.TaskInstanceDependentDetails; import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; import org.apache.dolphinscheduler.dao.model.WorkflowInstanceSummaryDto; +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; +import org.apache.dolphinscheduler.plugin.task.api.enums.DataType; +import org.apache.dolphinscheduler.plugin.task.api.enums.Direct; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; +import org.apache.dolphinscheduler.plugin.task.api.utils.GlobalParameterUtils; import java.util.ArrayList; import java.util.Collections; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -163,9 +174,37 @@ public void testQueryTaskListByWorkflowInstanceId() throws Exception { Assertions.assertEquals(Status.PROJECT_NOT_FOUND.getCode(), result.getCode().intValue()); } + @Test + public void testQueryTaskListByWorkflowInstanceIdMasksSensitiveTaskParams() throws Exception { + TaskInstanceDependentDetails taskInstance = new TaskInstanceDependentDetails(); + taskInstance.setTaskParams(sensitiveTaskParams()); + WorkflowInstanceTaskListDTO dto = new WorkflowInstanceTaskListDTO("SUCCESS", + Collections.singletonList(taskInstance)); + Mockito.when(workflowInstanceService.queryTaskListByWorkflowInstanceId(Mockito.any(), Mockito.anyLong(), + Mockito.any())) + .thenReturn(dto); + + MvcResult mvcResult = mockMvc + .perform(get("/projects/{projectCode}/workflow-instances/{id}/tasks", "1113", "123") + .header(SESSION_ID, sessionId)) + .andExpect(status().isOk()) + .andExpect(content().contentType(MediaType.APPLICATION_JSON)) + .andReturn(); + + String responseBody = mvcResult.getResponse().getContentAsString(); + Result result = JSONUtils.parseObject(responseBody, Result.class); + Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode().intValue()); + com.fasterxml.jackson.databind.node.ObjectNode root = JSONUtils.parseObject(responseBody); + String taskParams = root.path("data").path("taskList").path(0).path("taskParams").toString(); + Assertions.assertTrue(taskParams.contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertFalse(taskParams.contains("abc")); + Assertions.assertTrue(taskInstance.getTaskParams().contains("abc")); + } + @Test public void testUpdateWorkflowInstance() throws Exception { WorkflowDefinition mockResult = new WorkflowDefinition(); + mockResult.setGlobalParams(sensitiveGlobalParams()); Mockito.when(workflowInstanceService .updateWorkflowInstance(Mockito.any(), Mockito.anyLong(), Mockito.anyInt(), Mockito.anyString(), Mockito.anyString(), Mockito.anyString(), Mockito.anyBoolean(), Mockito.anyString(), @@ -197,13 +236,25 @@ public void testUpdateWorkflowInstance() throws Exception { Result result = JSONUtils.parseObject(mvcResult.getResponse().getContentAsString(), Result.class); Assertions.assertNotNull(result); Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode().intValue()); + com.fasterxml.jackson.databind.node.ObjectNode updateRoot = + JSONUtils.parseObject(mvcResult.getResponse().getContentAsString()); + String maskedGlobalParams = updateRoot.path("data").path("globalParams").asText(); + Assertions.assertTrue(maskedGlobalParams.contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertFalse(maskedGlobalParams.contains("Secret123")); + Assertions.assertTrue(mockResult.getGlobalParams().contains("Secret123")); } @Test public void testQueryWorkflowInstanceById() throws Exception { + WorkflowInstance workflowInstance = new WorkflowInstance(); + workflowInstance.setGlobalParams(sensitiveGlobalParams()); + TaskDefinition taskDefinition = new TaskDefinition(); + taskDefinition.setTaskParams(sensitiveTaskParams()); + workflowInstance.setDagData(new DagData(new WorkflowDefinition(), Collections.emptyList(), + Collections.singletonList(taskDefinition))); Mockito.when( workflowInstanceService.queryWorkflowInstanceById(Mockito.any(), Mockito.anyLong(), Mockito.anyInt())) - .thenReturn(new WorkflowInstance()); + .thenReturn(workflowInstance); MvcResult mvcResult = mockMvc.perform(get("/projects/{projectCode}/workflow-instances/{id}", "1113", "123") .header(SESSION_ID, sessionId)) .andExpect(status().isOk()) @@ -213,6 +264,14 @@ public void testQueryWorkflowInstanceById() throws Exception { Result result = JSONUtils.parseObject(mvcResult.getResponse().getContentAsString(), Result.class); Assertions.assertNotNull(result); Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode().intValue()); + com.fasterxml.jackson.databind.node.ObjectNode instanceRoot = + JSONUtils.parseObject(mvcResult.getResponse().getContentAsString()); + String maskedPayload = instanceRoot.path("data").toString(); + Assertions.assertTrue(maskedPayload.contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertFalse(maskedPayload.contains("Secret123")); + Assertions.assertFalse(maskedPayload.contains("\"abc\"")); + Assertions.assertTrue(workflowInstance.getGlobalParams().contains("Secret123")); + Assertions.assertTrue(taskDefinition.getTaskParams().contains("abc")); } @Test @@ -255,8 +314,11 @@ public void testQueryParentInstanceBySubId() throws Exception { @Test public void testViewVariables() throws Exception { + Map localParam = new LinkedHashMap<>(); + localParam.put(TaskConstants.LOCAL_PARAMS_LIST, Collections.singletonList(sensitive("token", "abc"))); WorkflowInstanceVariablesDTO mockResult = - new WorkflowInstanceVariablesDTO(Collections.emptyList(), Collections.emptyMap()); + new WorkflowInstanceVariablesDTO(Collections.singletonList(sensitive("pwd", "Secret123")), + Collections.singletonMap("shell-1", localParam)); Mockito.when(workflowInstanceService.viewVariables(Mockito.any(), Mockito.eq(1113L), Mockito.eq(123))) .thenReturn(mockResult); MvcResult mvcResult = mockMvc @@ -268,6 +330,14 @@ public void testViewVariables() throws Exception { Result result = JSONUtils.parseObject(mvcResult.getResponse().getContentAsString(), Result.class); Assertions.assertNotNull(result); Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode().intValue()); + com.fasterxml.jackson.databind.node.ObjectNode variablesRoot = + JSONUtils.parseObject(mvcResult.getResponse().getContentAsString()); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, + variablesRoot.path("data").path("globalParams").path(0).path("value").asText()); + Assertions.assertEquals("Secret123", mockResult.getGlobalParams().get(0).getValue()); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, + variablesRoot.path("data").path("localParams").path("shell-1") + .path(TaskConstants.LOCAL_PARAMS_LIST).path(0).path("value").asText()); } @Test @@ -372,4 +442,23 @@ public void testFromSummaryDto_MapsAllFields() { Assertions.assertEquals("admin", result.getExecutorName()); Assertions.assertEquals("1h 2m", result.getDuration()); } + + private static String sensitiveGlobalParams() { + return GlobalParameterUtils.serializeGlobalParameter(Collections.singletonList(sensitive("pwd", "Secret123"))); + } + + private static String sensitiveTaskParams() { + return "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"abc\",\"sensitive\":true}]}"; + } + + private static Property sensitive(String prop, String value) { + return Property.builder() + .prop(prop) + .direct(Direct.IN) + .type(DataType.VARCHAR) + .value(value) + .sensitive(true) + .build(); + } } 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 8e055c74df61..c5af336eff84 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 @@ -75,6 +75,7 @@ import org.apache.dolphinscheduler.dao.repository.WorkflowDefinitionLogDao; import org.apache.dolphinscheduler.dao.repository.WorkflowTaskRelationDao; import org.apache.dolphinscheduler.dao.utils.WorkerGroupUtils; +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; import org.apache.dolphinscheduler.plugin.task.api.model.ConditionDependentItem; import org.apache.dolphinscheduler.plugin.task.api.model.ConditionDependentTaskModel; import org.apache.dolphinscheduler.plugin.task.api.model.SwitchResultVo; @@ -112,6 +113,7 @@ 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.InjectMocks; import org.mockito.Mock; import org.mockito.Mockito; @@ -934,6 +936,49 @@ public void testCreateWorkflowDefinitionShouldSyncVersionToResponse() { Assertions.assertEquals(1, workflowDefinition.getVersion()); } + @Test + public void testCreateWorkflowDefinitionPersistsPlaintextSensitiveGlobalParams() { + 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()); + when(processService.saveTaskDefine(eq(user), eq(projectCode), anyList(), eq(Boolean.TRUE))).thenReturn(1); + when(processService.saveTaskRelation(eq(user), eq(projectCode), anyLong(), eq(1), anyList(), anyList(), + eq(Boolean.TRUE))).thenReturn(Constants.EXIT_CODE_SUCCESS); + + ArgumentCaptor persisted = ArgumentCaptor.forClass(WorkflowDefinition.class); + when(processService.saveWorkflowDefine(any(User.class), persisted.capture(), eq(Boolean.TRUE), + eq(Boolean.TRUE))).thenReturn(1); + + String sensitiveGlobalParams = + "[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"Secret123\",\"sensitive\":true}]"; + WorkflowDefinition created = workflowDefinitionService.createWorkflowDefinition( + user, projectCode, name, description, sensitiveGlobalParams, "[]", timeout, + taskRelationJson, taskDefinitionJson, null, WorkflowExecutionTypeEnum.PARALLEL); + + Assertions.assertTrue(created.getGlobalParams().contains("Secret123")); + Assertions.assertTrue(persisted.getValue().getGlobalParams().contains("Secret123")); + Assertions.assertSame(created, persisted.getValue()); + } + + @Test + public void testCreateWorkflowDefinitionRejectsSensitivePlaceholder() { + 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); + + String placeholderGlobalParams = + "[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"******\",\"sensitive\":true}]"; + Assertions.assertThrows(ServiceException.class, + () -> workflowDefinitionService.createWorkflowDefinition( + user, projectCode, name, description, placeholderGlobalParams, "[]", timeout, + taskRelationJson, taskDefinitionJson, null, WorkflowExecutionTypeEnum.PARALLEL)); + Mockito.verify(processService, Mockito.never()) + .saveWorkflowDefine(any(User.class), any(WorkflowDefinition.class), eq(Boolean.TRUE), eq(Boolean.TRUE)); + } + @Test public void testCreateWorkflowDefinitionShouldRejectUnauthorizedDatasource() { Project project = getProject(projectCode); @@ -1098,6 +1143,93 @@ public void testUpdateWorkflowDefinitionShouldSyncVersionToResponse() { Assertions.assertEquals(2, resultDefinition.getVersion()); } + @Test + public void testUpdateWorkflowDefinitionMergesSensitiveLocalParamsFromMatchingTaskVersion() { + Project project = getProject(projectCode); + WorkflowDefinition workflowDefinition = getWorkflowDefinition(); + workflowDefinition.setName("origin-name"); + when(projectDao.queryByCode(projectCode)).thenReturn(project); + Mockito.doNothing().when(projectService).checkHasProjectWritePermissionThrowException(eq(user), eq(project)); + when(processService.transformTask(anyList(), anyList())).thenReturn(getTaskNodeList()); + when(workflowDefinitionDao.queryByCode(processDefinitionCode)).thenReturn(Optional.of(workflowDefinition)); + when(workflowDefinitionDao.verifyByDefineName(projectCode, name)).thenReturn(null); + when(processService.saveWorkflowDefine(any(User.class), any(WorkflowDefinition.class), eq(Boolean.TRUE), + eq(Boolean.TRUE))).thenReturn(2); + when(workflowTaskRelationDao.queryByWorkflowDefinitionCode(processDefinitionCode)) + .thenReturn(Collections.emptyList()); + when(processService.saveTaskRelation(eq(user), eq(projectCode), eq(processDefinitionCode), eq(2), anyList(), + anyList(), eq(Boolean.TRUE))).thenReturn(Constants.EXIT_CODE_SUCCESS); + + TaskDefinitionLog version1Log = new TaskDefinitionLog(); + version1Log.setCode(123456789L); + version1Log.setVersion(1); + version1Log.setTaskParams( + "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"old-secret\",\"sensitive\":true}],\"rawScript\":\"echo 1\"}"); + when(taskDefinitionLogMapper.queryByTaskDefinitions(anyList())) + .thenReturn(Collections.singletonList(version1Log)); + + String submittedTaskDefinitionJson = + "[{\"code\":123456789,\"name\":\"test1\",\"version\":1,\"description\":\"\",\"delayTime\":0,\"taskType\":\"SHELL\"," + + "\"taskParams\":{\"resourceList\":[],\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"******\",\"sensitive\":true}],\"rawScript\":\"echo 1\",\"dependence\":{},\"conditionResult\":{\"successNode\":[],\"failedNode\":[]},\"waitStartTimeout\":{}," + + "\"switchResult\":{}},\"flag\":\"YES\",\"taskPriority\":\"MEDIUM\",\"workerGroup\":\"default\",\"failRetryTimes\":0,\"failRetryInterval\":1,\"timeoutFlag\":\"CLOSE\"," + + "\"timeoutNotifyStrategy\":null,\"timeout\":0,\"environmentCode\":-1},{\"code\":123451234,\"name\":\"test2\",\"version\":1,\"description\":\"\",\"delayTime\":0,\"taskType\":\"SHELL\"," + + "\"taskParams\":{\"resourceList\":[],\"localParams\":[],\"rawScript\":\"echo 2\",\"dependence\":{},\"conditionResult\":{\"successNode\":[],\"failedNode\":[]},\"waitStartTimeout\":{}," + + "\"switchResult\":{}},\"flag\":\"YES\",\"taskPriority\":\"MEDIUM\",\"workerGroup\":\"default\",\"failRetryTimes\":0,\"failRetryInterval\":1,\"timeoutFlag\":\"CLOSE\"," + + "\"timeoutNotifyStrategy\":\"WARN\",\"timeout\":0,\"environmentCode\":-1}]"; + + ArgumentCaptor savedTaskLogs = ArgumentCaptor.forClass(List.class); + when(processService.saveTaskDefine(eq(user), eq(projectCode), savedTaskLogs.capture(), eq(Boolean.TRUE))) + .thenReturn(1); + + workflowDefinitionService.updateWorkflowDefinition( + user, projectCode, name, processDefinitionCode, description, "[]", "[]", timeout, + taskRelationJson, submittedTaskDefinitionJson, WorkflowExecutionTypeEnum.PARALLEL); + + TaskDefinitionLog merged = ((List) savedTaskLogs.getValue()).stream() + .filter(log -> log.getCode() == 123456789L) + .findFirst() + .orElseThrow(AssertionError::new); + Assertions.assertTrue(merged.getTaskParams().contains("old-secret")); + Assertions.assertFalse(merged.getTaskParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + verify(taskDefinitionDao, Mockito.never()).queryByCodes(any()); + verify(taskDefinitionLogMapper).queryByTaskDefinitions(anyList()); + } + + @Test + public void testUpdateWorkflowDefinitionPersistsPlaintextNewSensitiveGlobalParams() { + Project project = getProject(projectCode); + WorkflowDefinition workflowDefinition = getWorkflowDefinition(); + workflowDefinition.setName("origin-name"); + workflowDefinition.setGlobalParams( + "[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"old-secret\",\"sensitive\":true}]"); + when(projectDao.queryByCode(projectCode)).thenReturn(project); + Mockito.doNothing().when(projectService).checkHasProjectWritePermissionThrowException(eq(user), eq(project)); + when(processService.transformTask(anyList(), anyList())).thenReturn(getTaskNodeList()); + when(workflowDefinitionDao.queryByCode(processDefinitionCode)).thenReturn(Optional.of(workflowDefinition)); + when(workflowDefinitionDao.verifyByDefineName(projectCode, name)).thenReturn(null); + when(processService.saveTaskDefine(eq(user), eq(projectCode), anyList(), eq(Boolean.TRUE))).thenReturn(1); + when(workflowTaskRelationDao.queryByWorkflowDefinitionCode(processDefinitionCode)) + .thenReturn(Collections.emptyList()); + when(processService.saveTaskRelation(eq(user), eq(projectCode), eq(processDefinitionCode), eq(2), anyList(), + anyList(), eq(Boolean.TRUE))).thenReturn(Constants.EXIT_CODE_SUCCESS); + + ArgumentCaptor persisted = ArgumentCaptor.forClass(WorkflowDefinition.class); + when(processService.saveWorkflowDefine(any(User.class), persisted.capture(), eq(Boolean.TRUE), + eq(Boolean.TRUE))).thenReturn(2); + + String submittedGlobalParams = + "[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"new-secret\",\"sensitive\":true}]"; + WorkflowDefinition updated = workflowDefinitionService.updateWorkflowDefinition( + user, projectCode, name, processDefinitionCode, description, submittedGlobalParams, "[]", timeout, + taskRelationJson, taskDefinitionJson, WorkflowExecutionTypeEnum.PARALLEL); + + Assertions.assertTrue(updated.getGlobalParams().contains("new-secret")); + Assertions.assertFalse(updated.getGlobalParams().contains("old-secret")); + Assertions.assertTrue(persisted.getValue().getGlobalParams().contains("new-secret")); + Assertions.assertSame(updated, persisted.getValue()); + } + @Test public void testUpdateWorkflowDefinitionShouldRejectUnavailableSubWorkflow() { 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 da71cd1621d8..97bf50e73100 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 @@ -20,9 +20,12 @@ import static org.apache.dolphinscheduler.api.AssertionsHelper.assertThrowsServiceException; import static org.apache.dolphinscheduler.api.constants.ApiFuncIdentificationConstant.WORKFLOW_INSTANCE; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import org.apache.dolphinscheduler.api.dto.workflowInstance.WorkflowInstanceTaskListDTO; @@ -59,6 +62,7 @@ import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; import org.apache.dolphinscheduler.dao.entity.WorkflowDefinitionLog; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; +import org.apache.dolphinscheduler.dao.mapper.TaskDefinitionLogMapper; import org.apache.dolphinscheduler.dao.mapper.WorkflowDefinitionLogMapper; import org.apache.dolphinscheduler.dao.model.WorkflowInstanceSummaryDto; import org.apache.dolphinscheduler.dao.repository.ProjectDao; @@ -70,6 +74,7 @@ import org.apache.dolphinscheduler.dao.repository.WorkflowInstanceDao; import org.apache.dolphinscheduler.dao.repository.WorkflowInstanceMapDao; import org.apache.dolphinscheduler.extract.master.command.RunWorkflowCommandParam; +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; import org.apache.dolphinscheduler.plugin.task.api.TaskPluginManager; import org.apache.dolphinscheduler.plugin.task.api.enums.DataType; import org.apache.dolphinscheduler.plugin.task.api.enums.DependResult; @@ -82,6 +87,7 @@ import java.util.ArrayList; import java.util.Calendar; +import java.util.Collections; import java.util.Date; import java.util.List; import java.util.Map; @@ -90,6 +96,7 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.MockedStatic; @@ -144,6 +151,9 @@ public class WorkflowInstanceServiceTest { @Mock TaskDefinitionDao taskDefinitionDao; + @Mock + TaskDefinitionLogMapper taskDefinitionLogMapper; + @Mock TaskDatasourcePermissionChecker taskDatasourcePermissionChecker; @@ -677,6 +687,125 @@ public void testUpdateWorkflowInstanceWithMismatchedProjectCode() { Mockito.verifyNoInteractions(workflowDefinitionDao); } + @Test + public void testUpdateWorkflowInstanceMergesSensitiveLocalParamsFromMatchingTaskVersion() { + long projectCode = 1L; + User loginUser = getAdminUser(); + WorkflowInstance workflowInstance = getProcessInstance(); + workflowInstance.setProjectCode(projectCode); + workflowInstance.setState(WorkflowExecutionStatus.SUCCESS); + workflowInstance.setTimeout(3000); + workflowInstance.setCommandType(CommandType.STOP); + workflowInstance.setWorkflowDefinitionCode(46L); + workflowInstance.setWorkflowDefinitionVersion(1); + WorkflowDefinition workflowDefinition = getProcessDefinition(); + workflowDefinition.setProjectCode(projectCode); + + doNothing().when(projectService).checkHasProjectWritePermissionThrowException(loginUser, projectCode); + when(processService.findWorkflowInstanceDetailById(1)).thenReturn(Optional.of(workflowInstance)); + when(workflowDefinitionDao.queryByCode(46L)).thenReturn(Optional.of(workflowDefinition)); + when(workflowInstanceDao.updateById(workflowInstance)).thenReturn(true); + when(processService.saveWorkflowDefine(loginUser, workflowDefinition, Boolean.TRUE, Boolean.FALSE)) + .thenReturn(1); + Mockito.doNothing().when(workflowDefinitionService).checkWorkflowNodeList(any(), any()); + Mockito.doNothing().when(taskDatasourcePermissionChecker).checkPermission(any(), any()); + + TaskDefinitionLog version1Log = new TaskDefinitionLog(); + version1Log.setCode(4254862762304L); + version1Log.setVersion(1); + version1Log.setTaskParams( + "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"old-secret\",\"sensitive\":true}],\"rawScript\":\"echo 1\"}"); + when(taskDefinitionLogMapper.queryByTaskDefinitions(any())).thenReturn(Collections.singletonList(version1Log)); + + TaskDefinition currentTask = new TaskDefinition(); + currentTask.setCode(4254862762304L); + currentTask.setTaskParams( + "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"new-secret\",\"sensitive\":true}]}"); + when(taskDefinitionDao.queryByCodes(any())).thenReturn(Collections.singletonList(currentTask)); + + String submittedTaskDefinitionJson = + "[{\"code\":4254862762304,\"name\":\"test1\",\"version\":1,\"description\":\"\",\"delayTime\":0,\"taskType\":\"SHELL\",\"taskParams\":{\"resourceList\":[]," + + "\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"******\",\"sensitive\":true}],\"rawScript\":\"echo 1\",\"dependence\":{},\"conditionResult\":{\"successNode\":[],\"failedNode\":[]},\"waitStartTimeout\":{},\"switchResult\":{}},\"flag\":\"YES\"," + + "\"taskPriority\":\"MEDIUM\",\"workerGroup\":\"default\",\"failRetryTimes\":0,\"failRetryInterval\":1,\"timeoutFlag\":\"CLOSE\",\"timeoutNotifyStrategy\":null,\"timeout\":0," + + "\"environmentCode\":-1},{\"code\":4254865123776,\"name\":\"test2\",\"version\":1,\"description\":\"\",\"delayTime\":0,\"taskType\":\"SHELL\",\"taskParams\":{\"resourceList\":[]," + + "\"localParams\":[],\"rawScript\":\"echo 2\",\"dependence\":{},\"conditionResult\":{\"successNode\":[],\"failedNode\":[]},\"waitStartTimeout\":{},\"switchResult\":{}},\"flag\":\"YES\"," + + "\"taskPriority\":\"MEDIUM\",\"workerGroup\":\"default\",\"failRetryTimes\":0,\"failRetryInterval\":1,\"timeoutFlag\":\"CLOSE\",\"timeoutNotifyStrategy\":\"WARN\",\"timeout\":0," + + "\"environmentCode\":-1}]"; + + ArgumentCaptor savedTaskLogs = ArgumentCaptor.forClass(List.class); + when(processService.saveTaskDefine(any(), Mockito.anyLong(), savedTaskLogs.capture(), any())).thenReturn(1); + + try ( + MockedStatic taskPluginManagerMockedStatic = + Mockito.mockStatic(TaskPluginManager.class)) { + taskPluginManagerMockedStatic + .when(() -> TaskPluginManager.checkTaskParameters(any(), any())) + .thenReturn(true); + workflowInstanceService.updateWorkflowInstance(loginUser, projectCode, 1, + taskRelationJson, submittedTaskDefinitionJson, "2020-02-21 00:00:00", true, "", "", 0); + } + + TaskDefinitionLog merged = ((List) savedTaskLogs.getValue()).stream() + .filter(log -> log.getCode() == 4254862762304L) + .findFirst() + .orElseThrow(AssertionError::new); + Assertions.assertTrue(merged.getTaskParams().contains("old-secret")); + Assertions.assertFalse(merged.getTaskParams().contains("new-secret")); + Assertions.assertFalse(merged.getTaskParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + verify(taskDefinitionDao, never()).queryByCodes(any()); + verify(taskDefinitionLogMapper, never()).queryByDefinitionCodeAndVersion(Mockito.anyLong(), Mockito.anyInt()); + verify(taskDefinitionLogMapper).queryByTaskDefinitions(any()); + } + + @Test + public void testUpdateWorkflowInstancePersistsPlaintextSensitiveGlobalParams() { + long projectCode = 1L; + User loginUser = getAdminUser(); + WorkflowInstance workflowInstance = getProcessInstance(); + workflowInstance.setProjectCode(projectCode); + workflowInstance.setState(WorkflowExecutionStatus.SUCCESS); + workflowInstance.setTimeout(3000); + workflowInstance.setCommandType(CommandType.STOP); + workflowInstance.setWorkflowDefinitionCode(46L); + workflowInstance.setWorkflowDefinitionVersion(1); + workflowInstance.setGlobalParams( + "[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"old-secret\",\"sensitive\":true}]"); + WorkflowDefinition workflowDefinition = getProcessDefinition(); + workflowDefinition.setProjectCode(projectCode); + + doNothing().when(projectService).checkHasProjectWritePermissionThrowException(loginUser, projectCode); + when(processService.findWorkflowInstanceDetailById(1)).thenReturn(Optional.of(workflowInstance)); + when(workflowDefinitionDao.queryByCode(46L)).thenReturn(Optional.of(workflowDefinition)); + when(workflowInstanceDao.updateById(workflowInstance)).thenReturn(true); + when(processService.saveTaskDefine(any(), Mockito.anyLong(), any(), any())).thenReturn(1); + Mockito.doNothing().when(workflowDefinitionService).checkWorkflowNodeList(any(), any()); + Mockito.doNothing().when(taskDatasourcePermissionChecker).checkPermission(any(), any()); + when(processService.saveTaskRelation(any(), Mockito.anyLong(), Mockito.anyLong(), Mockito.anyInt(), any(), + any(), any())).thenReturn(Constants.EXIT_CODE_SUCCESS); + + ArgumentCaptor persisted = ArgumentCaptor.forClass(WorkflowDefinition.class); + when(processService.saveWorkflowDefine(any(), persisted.capture(), eq(Boolean.TRUE), eq(Boolean.FALSE))) + .thenReturn(1); + + String submittedGlobalParams = + "[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"new-secret\",\"sensitive\":true}]"; + WorkflowDefinition updated; + try ( + MockedStatic taskPluginManagerMockedStatic = + Mockito.mockStatic(TaskPluginManager.class)) { + taskPluginManagerMockedStatic + .when(() -> TaskPluginManager.checkTaskParameters(any(), any())) + .thenReturn(true); + updated = workflowInstanceService.updateWorkflowInstance(loginUser, projectCode, 1, + taskRelationJson, taskDefinitionJson, "2020-02-21 00:00:00", true, submittedGlobalParams, "", 0); + } + + Assertions.assertTrue(updated.getGlobalParams().contains("new-secret")); + Assertions.assertFalse(updated.getGlobalParams().contains("old-secret")); + Assertions.assertTrue(persisted.getValue().getGlobalParams().contains("new-secret")); + Assertions.assertSame(updated, persisted.getValue()); + } + @Test public void testQueryParentInstanceBySubId() { long projectCode = 1L; diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/utils/SensitivePropertyUtilsTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/utils/SensitivePropertyUtilsTest.java new file mode 100644 index 000000000000..53cc2818b522 --- /dev/null +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/utils/SensitivePropertyUtilsTest.java @@ -0,0 +1,241 @@ +/* + * 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.utils; + +import org.apache.dolphinscheduler.api.exceptions.ServiceException; +import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.dao.entity.TaskInstance; +import org.apache.dolphinscheduler.dao.entity.TaskInstanceDependentDetails; +import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; +import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; +import org.apache.dolphinscheduler.extract.master.command.ICommandParam; +import org.apache.dolphinscheduler.extract.master.command.RunWorkflowCommandParam; +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; +import org.apache.dolphinscheduler.plugin.task.api.enums.DataType; +import org.apache.dolphinscheduler.plugin.task.api.enums.Direct; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; + +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +class SensitivePropertyUtilsTest { + + @Test + void requireNoPlaceholderRejectsMaskOnCreate() { + Assertions.assertThrows(ServiceException.class, + () -> SensitivePropertyUtils.requireNoPlaceholder( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)))); + } + + @Test + void mergeKeepsOriginalOnPlaceholder() { + List merged = SensitivePropertyUtils.merge( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(sensitive("pwd", "Secret123"))); + + Assertions.assertEquals("Secret123", merged.get(0).getValue()); + Assertions.assertTrue(merged.get(0).isSensitive()); + } + + @Test + void mergeKeepsOriginalWhenFalseToTrueWithPlaceholder() { + List merged = SensitivePropertyUtils.merge( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(nonSensitive("pwd", "plain"))); + + Assertions.assertEquals("plain", merged.get(0).getValue()); + Assertions.assertTrue(merged.get(0).isSensitive()); + } + + @Test + void mergeRejectsTrueToFalseWithPlaceholder() { + Assertions.assertThrows(ServiceException.class, + () -> SensitivePropertyUtils.merge( + Collections.singletonList(nonSensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(sensitive("pwd", "Secret123")))); + } + + @Test + void restoreStartParamsUsesGlobalPlaintext() { + List restored = SensitivePropertyUtils.restoreStartParams( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(sensitive("pwd", "Secret123"))); + + Assertions.assertEquals("Secret123", restored.get(0).getValue()); + Assertions.assertTrue(restored.get(0).isSensitive()); + } + + @Test + void restoreStartParamsKeepsUserOverride() { + List restored = SensitivePropertyUtils.restoreStartParams( + Collections.singletonList(sensitive("pwd", "new-secret")), + Collections.singletonList(sensitive("pwd", "Secret123"))); + + Assertions.assertEquals("new-secret", restored.get(0).getValue()); + } + + @Test + void restoreStartParamsUsesGlobalWhenSensitiveFlagMissing() { + List restored = SensitivePropertyUtils.restoreStartParams( + Collections.singletonList(new Property("pwd", Direct.IN, DataType.VARCHAR, + TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(sensitive("pwd", "Secret123"))); + + Assertions.assertEquals("Secret123", restored.get(0).getValue()); + Assertions.assertTrue(restored.get(0).isSensitive()); + } + + @Test + void restoreStartParamsDropsPlaceholderWithoutGlobal() { + List restored = SensitivePropertyUtils.restoreStartParams( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.emptyList()); + + Assertions.assertTrue(restored.isEmpty()); + } + + @Test + void emptyStringIsPersistedAsEmpty() { + List merged = SensitivePropertyUtils.merge( + Collections.singletonList(sensitive("pwd", "")), + Collections.singletonList(sensitive("pwd", "Secret123"))); + + Assertions.assertEquals("", merged.get(0).getValue()); + } + + @Test + void mergeLocalParamsRejectsPlaceholderWhenNoExisting() { + String submitted = + "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"" + + TaskConstants.SENSITIVE_DATA_MASK + "\",\"sensitive\":true}]}"; + + Assertions.assertThrows(ServiceException.class, + () -> SensitivePropertyUtils.mergeLocalParams(submitted, null)); + } + + @Test + void mergeLocalParamsRejectsTrueToFalseWithPlaceholder() { + String existing = + "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"abc\",\"sensitive\":true}]}"; + String submitted = + "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"" + + TaskConstants.SENSITIVE_DATA_MASK + "\",\"sensitive\":false}]}"; + + Assertions.assertThrows(ServiceException.class, + () -> SensitivePropertyUtils.mergeLocalParams(submitted, existing)); + } + + @Test + @SuppressWarnings("unchecked") + void maskDoesNotMutateLocalParamsMap() { + Map inner = new LinkedHashMap<>(); + inner.put(TaskConstants.LOCAL_PARAMS_LIST, Collections.singletonList(sensitive("token", "abc"))); + Map> localParams = new LinkedHashMap<>(); + localParams.put("shell-1", inner); + + Map> masked = SensitivePropertyUtils.mask(localParams); + + Assertions.assertEquals("abc", + ((List) localParams.get("shell-1").get(TaskConstants.LOCAL_PARAMS_LIST)).get(0).getValue()); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, + ((List) masked.get("shell-1").get(TaskConstants.LOCAL_PARAMS_LIST)).get(0).getValue()); + } + + @Test + void maskDoesNotMutateTaskInstanceDependentDetails() { + TaskInstanceDependentDetails original = new TaskInstanceDependentDetails<>(); + original.setTaskParams("{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"abc\",\"sensitive\":true}]}"); + original.setVarPool("[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"Secret123\",\"sensitive\":true}]"); + + TaskInstance masked = SensitivePropertyUtils.mask(original); + + Assertions.assertTrue(masked.getTaskParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertTrue(masked.getVarPool().contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertTrue(original.getTaskParams().contains("abc")); + Assertions.assertTrue(original.getVarPool().contains("Secret123")); + Assertions.assertFalse(masked.getTaskParams().contains("\"abc\"")); + Assertions.assertNotSame(original, masked); + } + + @Test + void maskDoesNotMutateWorkflowDefinition() { + WorkflowDefinition original = new WorkflowDefinition(); + original.setGlobalParams("[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"Secret123\",\"sensitive\":true}]"); + + WorkflowDefinition masked = SensitivePropertyUtils.mask(original); + + Assertions.assertTrue(original.getGlobalParams().contains("Secret123")); + Assertions.assertTrue(masked.getGlobalParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, masked.getGlobalParamMap().get("pwd")); + Assertions.assertEquals("Secret123", original.getGlobalParamMap().get("pwd")); + } + + @Test + void maskWorkflowInstanceMasksCommandParam() { + String plaintextCommandParam = JSONUtils.toJsonString(RunWorkflowCommandParam.builder() + .commandParams(Collections.singletonList(sensitive("pwd", "Secret123"))) + .timeZone("UTC") + .build()); + WorkflowInstance original = new WorkflowInstance(); + original.setCommandParam(plaintextCommandParam); + original.setGlobalParams("[{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"Secret123\",\"sensitive\":true}]"); + + WorkflowInstance masked = SensitivePropertyUtils.mask(original); + + Assertions.assertEquals(plaintextCommandParam, original.getCommandParam()); + Assertions.assertTrue(original.getCommandParam().contains("Secret123")); + Assertions.assertFalse(masked.getCommandParam().contains("Secret123")); + ICommandParam maskedCommandParam = + JSONUtils.parseObject(masked.getCommandParam(), ICommandParam.class); + Assertions.assertNotNull(maskedCommandParam); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, + maskedCommandParam.getCommandParams().get(0).getValue()); + Assertions.assertTrue(maskedCommandParam.getCommandParams().get(0).isSensitive()); + Assertions.assertTrue(masked.getGlobalParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); + Assertions.assertNotSame(original, masked); + } + + private static Property sensitive(String prop, String value) { + return Property.builder() + .prop(prop) + .direct(Direct.IN) + .type(DataType.VARCHAR) + .value(value) + .sensitive(true) + .build(); + } + + private static Property nonSensitive(String prop, String value) { + return Property.builder() + .prop(prop) + .direct(Direct.IN) + .type(DataType.VARCHAR) + .value(value) + .sensitive(false) + .build(); + } +} diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/model/Property.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/model/Property.java index e96f4020d3a0..4fee144846f4 100644 --- a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/model/Property.java +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/model/Property.java @@ -27,6 +27,8 @@ import lombok.Data; import lombok.NoArgsConstructor; +import com.fasterxml.jackson.annotation.JsonInclude; + @Data @Builder @NoArgsConstructor @@ -51,4 +53,18 @@ public class Property implements Serializable { private String value; + /** + * sensitive flag + */ + @Builder.Default + @JsonInclude(JsonInclude.Include.NON_DEFAULT) + private boolean sensitive = false; + + public Property(String prop, Direct direct, DataType type, String value) { + this.prop = prop; + this.direct = direct; + this.type = type; + this.value = value; + } + } diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/utils/PropertySensitiveUtils.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/utils/PropertySensitiveUtils.java new file mode 100644 index 000000000000..2302206d5edb --- /dev/null +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/utils/PropertySensitiveUtils.java @@ -0,0 +1,160 @@ +/* + * 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.plugin.task.api.utils; + +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; + +import org.apache.commons.collections4.CollectionUtils; + +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.function.Function; +import java.util.stream.Collectors; + +import lombok.experimental.UtilityClass; + +/** + * Mask / merge helpers for {@link Property#isSensitive()}. + * Keep-original marker is only {@code ******}; empty / null is a real empty value. + */ +@UtilityClass +public class PropertySensitiveUtils { + + public boolean isSensitive(Property property) { + return property != null && property.isSensitive(); + } + + public boolean isSensitiveValuePlaceholder(String value) { + return TaskConstants.SENSITIVE_DATA_MASK.equals(value); + } + + public Property copy(Property property) { + if (property == null) { + return null; + } + return Property.builder() + .prop(property.getProp()) + .direct(property.getDirect()) + .type(property.getType()) + .value(property.getValue()) + .sensitive(property.isSensitive()) + .build(); + } + + public List copy(List properties) { + if (CollectionUtils.isEmpty(properties)) { + return Collections.emptyList(); + } + return properties.stream() + .map(PropertySensitiveUtils::copy) + .collect(Collectors.toList()); + } + + public Property maskSensitiveValue(Property property) { + Property maskedProperty = copy(property); + if (isSensitive(maskedProperty) && maskedProperty.getValue() != null) { + maskedProperty.setValue(TaskConstants.SENSITIVE_DATA_MASK); + } + return maskedProperty; + } + + public List maskSensitiveValues(List properties) { + if (CollectionUtils.isEmpty(properties)) { + return Collections.emptyList(); + } + return properties.stream() + .map(PropertySensitiveUtils::maskSensitiveValue) + .collect(Collectors.toList()); + } + + public List mergeSensitiveValuePlaceholders(List submittedProperties, + List existingProperties) { + if (CollectionUtils.isEmpty(submittedProperties)) { + return Collections.emptyList(); + } + Map existingPropertyMap = toPropMap(existingProperties); + + return submittedProperties.stream() + .map(PropertySensitiveUtils::copy) + .peek(property -> mergeSensitiveValuePlaceholder(property, existingPropertyMap)) + .collect(Collectors.toList()); + } + + /** Create: {@code ******} is never a real value. */ + public String findPlaceholderProp(List properties) { + if (CollectionUtils.isEmpty(properties)) { + return null; + } + for (Property property : properties) { + if (property != null && isSensitiveValuePlaceholder(property.getValue())) { + return property.getProp(); + } + } + return null; + } + + /** + * Update: {@code ******} is keep-original only when an existing property can be merged. + * {@code false→true} + {@code ******} is allowed; {@code true→false} + {@code ******} is rejected. + */ + public String findInvalidSensitivePlaceholderProp(List submittedProperties, + List existingProperties) { + if (CollectionUtils.isEmpty(submittedProperties)) { + return null; + } + Map existingPropertyMap = toPropMap(existingProperties); + for (Property submitted : submittedProperties) { + if (!isSensitiveValuePlaceholder(submitted.getValue())) { + continue; + } + Property existing = existingPropertyMap.get(submitted.getProp()); + if (!isSensitive(submitted)) { + if (existing != null && existing.isSensitive()) { + return submitted.getProp(); + } + continue; + } + if (existing == null) { + return submitted.getProp(); + } + } + return null; + } + + private void mergeSensitiveValuePlaceholder(Property submittedProperty, + Map existingPropertyMap) { + if (!isSensitive(submittedProperty) || !isSensitiveValuePlaceholder(submittedProperty.getValue())) { + return; + } + Property existingProperty = existingPropertyMap.get(submittedProperty.getProp()); + if (existingProperty != null) { + submittedProperty.setValue(existingProperty.getValue()); + } + } + + private Map toPropMap(List properties) { + return CollectionUtils.emptyIfNull(properties) + .stream() + .filter(Objects::nonNull) + .filter(property -> property.getProp() != null) + .collect(Collectors.toMap(Property::getProp, Function.identity(), (left, right) -> right)); + } +} diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/test/java/org/apache/dolphinscheduler/plugin/task/api/utils/PropertySensitiveUtilsTest.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/test/java/org/apache/dolphinscheduler/plugin/task/api/utils/PropertySensitiveUtilsTest.java new file mode 100644 index 000000000000..7227d9d31cb6 --- /dev/null +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/test/java/org/apache/dolphinscheduler/plugin/task/api/utils/PropertySensitiveUtilsTest.java @@ -0,0 +1,167 @@ +/* + * 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.plugin.task.api.utils; + +import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; +import org.apache.dolphinscheduler.plugin.task.api.enums.DataType; +import org.apache.dolphinscheduler.plugin.task.api.enums.Direct; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +class PropertySensitiveUtilsTest { + + @Test + void maskSensitiveValuesShouldDeepCopyAndMask() { + Property original = sensitive("pwd", "Secret123"); + List masked = PropertySensitiveUtils.maskSensitiveValues(Collections.singletonList(original)); + + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, masked.get(0).getValue()); + Assertions.assertEquals("Secret123", original.getValue()); + Assertions.assertTrue(masked.get(0).isSensitive()); + } + + @Test + void maskShouldNotChangeNonSensitive() { + Property original = nonSensitive("name", "alice"); + List masked = PropertySensitiveUtils.maskSensitiveValues(Collections.singletonList(original)); + Assertions.assertEquals("alice", masked.get(0).getValue()); + Assertions.assertFalse(masked.get(0).isSensitive()); + } + + @Test + void mergeShouldReplacePlaceholderWithExistingValue() { + Property submitted = sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK); + Property existing = sensitive("pwd", "Secret123"); + + List merged = PropertySensitiveUtils.mergeSensitiveValuePlaceholders( + Collections.singletonList(submitted), Collections.singletonList(existing)); + + Assertions.assertEquals("Secret123", merged.get(0).getValue()); + Assertions.assertEquals(TaskConstants.SENSITIVE_DATA_MASK, submitted.getValue()); + } + + @Test + void emptyStringIsRealEmptyValueNotKeepOriginal() { + Property submitted = sensitive("pwd", ""); + Property existing = sensitive("pwd", "Secret123"); + + List merged = PropertySensitiveUtils.mergeSensitiveValuePlaceholders( + Collections.singletonList(submitted), Collections.singletonList(existing)); + + Assertions.assertEquals("", merged.get(0).getValue()); + } + + @Test + void findPlaceholderPropOnCreate() { + Property submitted = sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK); + + Assertions.assertEquals("pwd", + PropertySensitiveUtils.findPlaceholderProp(Collections.singletonList(submitted))); + } + + @Test + void falseToTrueWithPlaceholderIsKeepOriginal() { + Property submitted = sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK); + Property existing = nonSensitive("pwd", "plain"); + + Assertions.assertNull(PropertySensitiveUtils.findInvalidSensitivePlaceholderProp( + Collections.singletonList(submitted), Collections.singletonList(existing))); + List merged = PropertySensitiveUtils.mergeSensitiveValuePlaceholders( + Collections.singletonList(submitted), Collections.singletonList(existing)); + Assertions.assertEquals("plain", merged.get(0).getValue()); + Assertions.assertTrue(merged.get(0).isSensitive()); + } + + @Test + void findInvalidPlaceholderWhenNewSensitiveWithPlaceholder() { + Property submitted = sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK); + + Assertions.assertEquals("pwd", + PropertySensitiveUtils.findInvalidSensitivePlaceholderProp( + Collections.singletonList(submitted), Collections.emptyList())); + } + + @Test + void validPlaceholderWhenTrueToTrue() { + Property submitted = sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK); + Property existing = sensitive("pwd", "Secret123"); + + Assertions.assertNull(PropertySensitiveUtils.findInvalidSensitivePlaceholderProp( + Collections.singletonList(submitted), Collections.singletonList(existing))); + } + + @Test + void findInvalidPlaceholderWhenTrueToFalse() { + Property submitted = nonSensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK); + Property existing = sensitive("pwd", "Secret123"); + + Assertions.assertEquals("pwd", + PropertySensitiveUtils.findInvalidSensitivePlaceholderProp( + Collections.singletonList(submitted), Collections.singletonList(existing))); + } + + @Test + void validPlaceholderAfterToggleBackToSensitive() { + // Uncheck then check again: final submit is still sensitive + ****** (keep-original). + Property submitted = sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK); + Property existing = sensitive("pwd", "Secret123"); + + Assertions.assertNull(PropertySensitiveUtils.findInvalidSensitivePlaceholderProp( + Collections.singletonList(submitted), Collections.singletonList(existing))); + List merged = PropertySensitiveUtils.mergeSensitiveValuePlaceholders( + Collections.singletonList(submitted), Collections.singletonList(existing)); + Assertions.assertEquals("Secret123", merged.get(0).getValue()); + Assertions.assertTrue(merged.get(0).isSensitive()); + } + + @Test + void serializationRoundTripKeepsSensitiveFlag() { + Property property = sensitive("token", "abc"); + String json = org.apache.dolphinscheduler.common.utils.JSONUtils.toJsonString( + Arrays.asList(property)); + List parsed = org.apache.dolphinscheduler.common.utils.JSONUtils.toList(json, Property.class); + Assertions.assertTrue(parsed.get(0).isSensitive()); + Assertions.assertEquals("abc", parsed.get(0).getValue()); + } + + private static Property sensitive(String prop, String value) { + return Property.builder() + .prop(prop) + .direct(Direct.IN) + .type(DataType.VARCHAR) + .value(value) + .sensitive(true) + .build(); + } + + private static Property nonSensitive(String prop, String value) { + return Property.builder() + .prop(prop) + .direct(Direct.IN) + .type(DataType.VARCHAR) + .value(value) + .sensitive(false) + .build(); + } +} diff --git a/dolphinscheduler-ui/src/components/form/fields/checkbox.ts b/dolphinscheduler-ui/src/components/form/fields/checkbox.ts index a8af013e6950..23c523af731a 100644 --- a/dolphinscheduler-ui/src/components/form/fields/checkbox.ts +++ b/dolphinscheduler-ui/src/components/form/fields/checkbox.ts @@ -28,8 +28,13 @@ export function renderCheckbox( if (!options) { return h(NCheckbox, { ...props, - value: fields[field], - onUpdateChecked: (checked: boolean) => void (fields[field] = checked) + checked: !!fields[field], + onUpdateChecked: (checked: boolean) => { + fields[field] = checked + if (props && typeof props.onUpdateChecked === 'function') { + props.onUpdateChecked(checked) + } + } }) } return h( diff --git a/dolphinscheduler-ui/src/components/form/fields/custom-parameters.ts b/dolphinscheduler-ui/src/components/form/fields/custom-parameters.ts index 49bb371b5f46..ae1f99a46537 100644 --- a/dolphinscheduler-ui/src/components/form/fields/custom-parameters.ts +++ b/dolphinscheduler-ui/src/components/form/fields/custom-parameters.ts @@ -92,7 +92,7 @@ const getDefaultValue = (children: IJsonItem[]) => { } return } else { - parent[mergedChild.field] = mergedChild.value || null + parent[mergedChild.field] = mergedChild.value ?? null if (mergedChild.validate) ruleParent[mergedChild.field] = formatValidate(mergedChild.validate) } diff --git a/dolphinscheduler-ui/src/locales/en_US/project.ts b/dolphinscheduler-ui/src/locales/en_US/project.ts index 8084f192d7d2..e40e4601f2e7 100644 --- a/dolphinscheduler-ui/src/locales/en_US/project.ts +++ b/dolphinscheduler-ui/src/locales/en_US/project.ts @@ -351,6 +351,7 @@ export default { minute: 'Minute', key: 'Key', value: 'Value', + sensitive: 'Sensitive', success: 'Success', delete_cell: 'Delete selected edges and nodes', online_directly: 'Whether to go online the workflow definition', @@ -460,6 +461,7 @@ export default { prop_repeat: 'prop is repeat', value_tips: 'value(optional)', value_required_tips: 'value(required)', + sensitive: 'Sensitive', custom_labels: 'Customized labels', node_selectors: 'Node Selectors', label_repeat: 'repeated label', diff --git a/dolphinscheduler-ui/src/locales/zh_CN/project.ts b/dolphinscheduler-ui/src/locales/zh_CN/project.ts index 34901950bdad..9d48dcfeb039 100644 --- a/dolphinscheduler-ui/src/locales/zh_CN/project.ts +++ b/dolphinscheduler-ui/src/locales/zh_CN/project.ts @@ -345,6 +345,7 @@ export default { minute: '分', key: '键', value: '值', + sensitive: '敏感', success: '成功', delete_cell: '删除选中的线或节点', online_directly: '是否上线工作流定义', @@ -446,6 +447,7 @@ export default { prop_repeat: 'prop中有重复', value_tips: 'value(选填)', value_required_tips: 'value(必填)', + sensitive: '敏感', custom_labels: '自定义标签', node_selectors: '节点选择器', label_repeat: 'label中有重复', diff --git a/dolphinscheduler-ui/src/utils/sensitive-param.ts b/dolphinscheduler-ui/src/utils/sensitive-param.ts new file mode 100644 index 000000000000..c1e4ebc328d8 --- /dev/null +++ b/dolphinscheduler-ui/src/utils/sensitive-param.ts @@ -0,0 +1,29 @@ +/* + * 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. + */ + +export const SENSITIVE_VALUE_MASK = '******' + +/** Unchecking does not clear value; empty + re-check restores ****** (keep-original). */ +export function applySensitiveToggle( + param: { value?: string; sensitive?: boolean }, + checked: boolean +) { + param.sensitive = checked + if (checked && (param.value === '' || param.value == null)) { + param.value = SENSITIVE_VALUE_MASK + } +} diff --git a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-custom-params.ts b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-custom-params.ts index 2052f524fe16..926cb9f564b4 100644 --- a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-custom-params.ts +++ b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-custom-params.ts @@ -16,6 +16,7 @@ */ import { Ref } from 'vue' import { useI18n } from 'vue-i18n' +import { applySensitiveToggle } from '@/utils/sensitive-param' import type { IJsonItem } from '../types' export function useCustomParams({ @@ -44,7 +45,7 @@ export function useCustomParams({ { type: 'input', field: 'prop', - span: 6, + span: 5, class: 'input-param-key', props: { placeholder: t('project.node.prop_tips'), @@ -81,7 +82,7 @@ export function useCustomParams({ { type: 'select', field: 'type', - span: 6, + span: 5, options: TYPE_LIST, value: 'VARCHAR', props: { @@ -91,13 +92,29 @@ export function useCustomParams({ { type: 'input', field: 'value', - span: 6, + span: 5, class: 'input-param-value', props: { placeholder: t('project.node.value_tips'), maxLength: 256 } - } + }, + (i?: number) => ({ + type: 'checkbox' as const, + field: 'sensitive', + name: t('project.node.sensitive'), + span: 3, + value: false, + props: { + onUpdateChecked: (checked: boolean) => { + const rows = model[field] + if (i === undefined || !rows || !rows[i]) { + return + } + applySensitiveToggle(rows[i], checked) + } + } + }) ] } ] diff --git a/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts b/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts index 4ed68ce942d1..49d88de66920 100644 --- a/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts +++ b/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts @@ -64,6 +64,7 @@ interface ILocalParam { direct?: string type?: string value?: string + sensitive?: boolean } interface ILabel { diff --git a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx index 1e39a6d37e4e..204dc5ab0a0b 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx +++ b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx @@ -41,6 +41,7 @@ import { useRoute } from 'vue-router' import { verifyName } from '@/service/modules/workflow-definition' import './x6-style.scss' import { positiveIntegerRegex } from '@/utils/regex' +import { applySensitiveToggle } from '@/utils/sensitive-param' import type { SaveForm, WorkflowDefinition, WorkflowInstance } from './types' const props = { @@ -156,7 +157,8 @@ export default defineComponent({ key: param.prop, value: param.value, direct: param.direct, - type: param.type + type: param.type, + sensitive: param.sensitive || false }) ) } @@ -243,7 +245,8 @@ export default defineComponent({ key: '', direct: 'IN', type: 'VARCHAR', - value: '' + value: '', + sensitive: false } }} class='input-global-params' @@ -255,10 +258,11 @@ export default defineComponent({ direct: string type: string value: string + sensitive?: boolean } }) => ( - + - + - + + + + applySensitiveToggle(param.value, checked) + } + > + {t('project.dag.sensitive')} + + ) }} diff --git a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts index f136fa5be08f..5e7497503eb7 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts +++ b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts @@ -140,6 +140,7 @@ export interface GlobalParam { direct: string type: string value: string + sensitive?: boolean } export interface SaveForm { diff --git a/dolphinscheduler-ui/src/views/projects/workflow/definition/create/index.tsx b/dolphinscheduler-ui/src/views/projects/workflow/definition/create/index.tsx index 511f830c110d..39f703aef2ea 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/definition/create/index.tsx +++ b/dolphinscheduler-ui/src/views/projects/workflow/definition/create/index.tsx @@ -59,7 +59,8 @@ export default defineComponent({ prop: p.key, value: p.value, direct: p.direct, - type: p.type + type: p.type, + sensitive: p.sensitive || false } }) diff --git a/dolphinscheduler-ui/src/views/projects/workflow/definition/detail/index.tsx b/dolphinscheduler-ui/src/views/projects/workflow/definition/detail/index.tsx index 2831014c55a6..30ccc0df5b40 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/definition/detail/index.tsx +++ b/dolphinscheduler-ui/src/views/projects/workflow/definition/detail/index.tsx @@ -85,7 +85,8 @@ export default defineComponent({ prop: p.key, value: p.value, direct: p.direct, - type: p.type + type: p.type, + sensitive: p.sensitive || false } }) diff --git a/dolphinscheduler-ui/src/views/projects/workflow/instance/components/variables-view.tsx b/dolphinscheduler-ui/src/views/projects/workflow/instance/components/variables-view.tsx index dcd37ad03a1f..0bcdf1404f8c 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/instance/components/variables-view.tsx +++ b/dolphinscheduler-ui/src/views/projects/workflow/instance/components/variables-view.tsx @@ -66,17 +66,23 @@ export default defineComponent({ ctx.emit('copy', text) } + const shouldDisplayParamField = (key: string, taskType: string) => { + if (key === 'sensitive') { + return false + } + return ( + !(taskType === 'SQL' || taskType === 'PROCEDURE') || + (key !== 'direct' && key !== 'type') + ) + } + /** * Copyed text processing */ const rtClipboard = (el: any, taskType: string) => { const arr: Array = [] Object.keys(el).forEach((key) => { - if (taskType === 'SQL' || taskType === 'PROCEDURE') { - if (key !== 'direct' && key !== 'type') { - arr.push(`${key}=${el[key]}`) - } - } else { + if (shouldDisplayParamField(key, taskType)) { arr.push(`${key}=${el[key]}`) } }) @@ -92,21 +98,14 @@ export default defineComponent({ onClick={() => handleCopy(rtClipboard(el, taskType))} > {Object.keys(el).map((key: string) => { - if (taskType === 'SQL' || taskType === 'PROCEDURE') { - return key !== 'direct' && key !== 'type' ? ( - - {key} = {el[key]} - - ) : ( - '' - ) - } else { + if (shouldDisplayParamField(key, taskType)) { return ( {key} = {el[key]} ) } + return '' })} ) @@ -120,7 +119,8 @@ export default defineComponent({ globalParams, localParams, localButton, - handleCopy + handleCopy, + shouldDisplayParamField } }, render() { diff --git a/dolphinscheduler-ui/src/views/projects/workflow/instance/detail/index.tsx b/dolphinscheduler-ui/src/views/projects/workflow/instance/detail/index.tsx index 1957a6fe9120..7e9654067172 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/instance/detail/index.tsx +++ b/dolphinscheduler-ui/src/views/projects/workflow/instance/detail/index.tsx @@ -82,7 +82,8 @@ export default defineComponent({ prop: p.key, value: p.value, direct: p.direct, - type: p.type + type: p.type, + sensitive: p.sensitive || false } })