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 771e47ccdbb2..b95a5cba8908 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 @@ -264,7 +264,13 @@ public WorkflowDefinition createWorkflowDefinition(User loginUser, List taskDefinitionLogs = generateTaskDefinitionList(taskDefinitionJson); List taskRelationList = generateTaskRelationList(taskRelationJson, taskDefinitionLogs); - SensitivePropertyUtils.requireNoPlaceholder(GlobalParameterUtils.deserializeGlobalParameter(globalParams)); + List submittedGlobalParams = GlobalParameterUtils.deserializeGlobalParameter(globalParams); + if (CollectionUtils.isNotEmpty(submittedGlobalParams)) { + globalParams = GlobalParameterUtils.serializeGlobalParameter( + SensitivePropertyUtils.encodeForCreate(submittedGlobalParams)); + } else { + SensitivePropertyUtils.requireNoPlaceholder(submittedGlobalParams); + } for (TaskDefinitionLog taskDefinitionLog : CollectionUtils.emptyIfNull(taskDefinitionLogs)) { taskDefinitionLog.setTaskParams( SensitivePropertyUtils.mergeLocalParams(taskDefinitionLog.getTaskParams(), null)); @@ -669,7 +675,7 @@ public WorkflowDefinition updateWorkflowDefinition(User loginUser, GlobalParameterUtils.deserializeGlobalParameter(globalParams); if (CollectionUtils.isNotEmpty(submittedGlobalParams)) { globalParams = GlobalParameterUtils.serializeGlobalParameter( - SensitivePropertyUtils.merge(submittedGlobalParams, + SensitivePropertyUtils.mergeAndEncode(submittedGlobalParams, GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams()))); } List versionKeys = new ArrayList<>(); 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 bf4356d1f46e..9ceacfa83fb6 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 @@ -413,12 +413,13 @@ public WorkflowDefinition updateWorkflowInstance(User loginUser, long projectCod } List submittedGlobalParams = GlobalParameterUtils.deserializeGlobalParameter(globalParams); + String instanceGlobalParams = globalParams; if (CollectionUtils.isNotEmpty(submittedGlobalParams)) { - globalParams = GlobalParameterUtils.serializeGlobalParameter( + instanceGlobalParams = GlobalParameterUtils.serializeGlobalParameter( SensitivePropertyUtils.merge(submittedGlobalParams, GlobalParameterUtils.deserializeGlobalParameter(workflowInstance.getGlobalParams()))); } - setWorkflowInstance(workflowInstance, scheduleTime, globalParams, timeout, timezoneId); + setWorkflowInstance(workflowInstance, scheduleTime, instanceGlobalParams, timeout, timezoneId); List taskDefinitionLogs = JSONUtils.toList(taskDefinitionJson, TaskDefinitionLog.class); if (taskDefinitionLogs.isEmpty()) { log.warn("Parameter taskDefinitionJson is empty"); @@ -466,8 +467,14 @@ public WorkflowDefinition updateWorkflowInstance(User loginUser, long projectCod // check workflow json is valid (throws ServiceException on validation failures) workflowDefinitionService.checkWorkflowNodeList(taskRelationJson, taskDefinitionLogs); + String definitionGlobalParams = globalParams; + if (CollectionUtils.isNotEmpty(submittedGlobalParams)) { + definitionGlobalParams = GlobalParameterUtils.serializeGlobalParameter( + SensitivePropertyUtils.mergeAndEncode(submittedGlobalParams, + GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams()))); + } workflowDefinition.set(projectCode, workflowDefinition.getName(), workflowDefinition.getDescription(), - globalParams, locations, timeout); + definitionGlobalParams, locations, timeout); workflowDefinition.setUpdateTime(new Date()); int insertVersion = processService.saveWorkflowDefine(loginUser, workflowDefinition, syncDefine, Boolean.FALSE); if (insertVersion == 0) { 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 index a94766f1747b..8b96e7eab501 100644 --- 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 @@ -36,6 +36,7 @@ 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.datasource.api.utils.PasswordUtils; 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; @@ -62,9 +63,8 @@ /** * 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. + * Definition write path: merge keep-original then encode with {@link PasswordUtils} when enabled. + * Query masks values as {@code ******}. Start / execution use decrypt copies. */ @UtilityClass public class SensitivePropertyUtils { @@ -80,9 +80,28 @@ public void requireNoPlaceholder(List properties) { } /** - * Update: restore DB plaintext for keep-original {@code ******}. + * Create definition params: reject {@code ******}, then encode sensitive plaintext. + */ + public List encodeForCreate(List properties) { + requireNoPlaceholder(properties); + return encodeNewSensitivePlaintext(properties, Collections.emptyList()); + } + + /** + * Update definition params: restore DB value for keep-original {@code ******}, then encode. * Empty / null is a real empty value. {@code false→true} + {@code ******} is allowed; * {@code true→false} + {@code ******} is rejected. + * Instance {@code global_params} stay plaintext via {@link #merge}. Copies written to a workflow + * definition must use this method against the definition's stored values, not the instance plaintext. + */ + public List mergeAndEncode(List submittedProperties, List existingProperties) { + List merged = merge(submittedProperties, existingProperties); + return encodeNewSensitivePlaintext(merged, existingProperties); + } + + /** + * Update: restore DB value for keep-original {@code ******} without encoding. + * Used for runtime instance global_params (plaintext materialization). */ public List merge(List submittedProperties, List existingProperties) { if (CollectionUtils.isEmpty(submittedProperties)) { @@ -100,6 +119,7 @@ public List merge(List submittedProperties, List e /** * Build start/command params before persist. + * Decrypts definition globals first, then: *
    *
  • Same name as a workflow global: inherit all global attributes, only override {@code value} * ({@code ******} keeps the global plaintext).
  • @@ -110,7 +130,8 @@ public List restoreStartParams(List startParams, List globals = CollectionUtils.emptyIfNull(globalParams).stream() + List decryptedGlobals = decodeSensitiveValues(globalParams); + Map globals = CollectionUtils.emptyIfNull(decryptedGlobals).stream() .filter(Objects::nonNull) .filter(property -> property.getProp() != null) .collect(Collectors.toMap(Property::getProp, Function.identity(), (left, right) -> right)); @@ -137,16 +158,68 @@ public List restoreStartParams(List startParams, List decodeSensitiveValues(List properties) { + return PropertySensitiveUtils.transformSensitiveValues(properties, PasswordUtils::decodePassword); + } + + /** + * Merge task {@code localParams} and encode sensitive values for definition persist. + * Instance edits also persist task-definition snapshots, so this always encodes. + */ public String mergeLocalParams(String submittedTaskParams, String existingTaskParams) { return rewriteLocalParams(submittedTaskParams, submitted -> { if (StringUtils.isEmpty(existingTaskParams)) { - requireNoPlaceholder(submitted); - return submitted; + return encodeForCreate(submitted); } - return merge(submitted, getLocalParams(existingTaskParams)); + return mergeAndEncode(submitted, getLocalParams(existingTaskParams)); }); } + /** + * Encode sensitive values that are new plaintext for definition persist. + * Keep-original from an already-sensitive existing value is written as-is (no re-encode). + * {@code false→true} keep-original merges non-sensitive plaintext and then encodes. + */ + private List encodeNewSensitivePlaintext(List mergedProperties, + List existingProperties) { + if (CollectionUtils.isEmpty(mergedProperties)) { + return mergedProperties; + } + Map existingMap = CollectionUtils.emptyIfNull(existingProperties).stream() + .filter(Objects::nonNull) + .filter(property -> property.getProp() != null) + .collect(Collectors.toMap(Property::getProp, Function.identity(), (left, right) -> right)); + List encoded = new ArrayList<>(mergedProperties.size()); + for (Property property : mergedProperties) { + Property copy = PropertySensitiveUtils.copy(property); + if (!PropertySensitiveUtils.isSensitive(copy)) { + Property existing = existingMap.get(copy.getProp()); + // true→false: if a non-UI client re-sent ciphertext, decode back to plaintext. + if (existing != null && existing.isSensitive() + && StringUtils.isNotEmpty(copy.getValue()) + && Objects.equals(copy.getValue(), existing.getValue())) { + copy.setValue(PasswordUtils.decodePassword(copy.getValue())); + } + encoded.add(copy); + continue; + } + if (StringUtils.isEmpty(copy.getValue())) { + encoded.add(copy); + continue; + } + Property existing = existingMap.get(copy.getProp()); + if (existing != null && existing.isSensitive() + && Objects.equals(copy.getValue(), existing.getValue())) { + // true→true keep-original (or identical resubmit): already persisted form. + encoded.add(copy); + continue; + } + copy.setValue(PasswordUtils.encodePassword(copy.getValue())); + encoded.add(copy); + } + return encoded; + } + /** * 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). 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 97bf50e73100..7187fcee7859 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 @@ -49,6 +49,7 @@ import org.apache.dolphinscheduler.common.model.TaskNodeRelation; import org.apache.dolphinscheduler.common.utils.DateUtils; import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.common.utils.PropertyUtils; import org.apache.dolphinscheduler.dao.AlertDao; import org.apache.dolphinscheduler.dao.entity.DependentResultTaskInstanceContext; import org.apache.dolphinscheduler.dao.entity.Project; @@ -74,6 +75,8 @@ 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.datasource.api.constants.DataSourceConstants; +import org.apache.dolphinscheduler.plugin.datasource.api.utils.PasswordUtils; 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; @@ -81,6 +84,7 @@ import org.apache.dolphinscheduler.plugin.task.api.enums.Direct; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import org.apache.dolphinscheduler.plugin.task.api.model.Property; +import org.apache.dolphinscheduler.plugin.task.api.utils.GlobalParameterUtils; import org.apache.dolphinscheduler.service.expand.CuringParamsService; import org.apache.dolphinscheduler.service.model.TaskNode; import org.apache.dolphinscheduler.service.process.ProcessService; @@ -806,6 +810,118 @@ public void testUpdateWorkflowInstancePersistsPlaintextSensitiveGlobalParams() { Assertions.assertSame(updated, persisted.getValue()); } + @Test + public void testUpdateWorkflowInstanceEncryptsDefinitionSnapshotsForBothSyncDefine() { + try ( + MockedStatic propertyUtils = Mockito.mockStatic(PropertyUtils.class); + MockedStatic taskPluginManager = Mockito.mockStatic(TaskPluginManager.class)) { + propertyUtils.when(() -> PropertyUtils.getBoolean(DataSourceConstants.DATASOURCE_ENCRYPTION_ENABLE, false)) + .thenReturn(true); + propertyUtils.when(() -> PropertyUtils.getString(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT, + DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT)) + .thenReturn(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT); + taskPluginManager.when(() -> TaskPluginManager.checkTaskParameters(any(), any())).thenReturn(true); + when(curingGlobalParamsService.curingGlobalParams(any(), any(), any(), any(), any(), any())) + .thenAnswer(invocation -> GlobalParameterUtils.serializeGlobalParameter(invocation.getArgument(2))); + + String keptGlobalCipher = PasswordUtils.encodePassword("kept-secret"); + String keptLocalCipher = PasswordUtils.encodePassword("kept-local"); + assertInstanceEditEncryptsDefinitionSnapshot(Boolean.TRUE, keptGlobalCipher, keptLocalCipher); + assertInstanceEditEncryptsDefinitionSnapshot(Boolean.FALSE, keptGlobalCipher, keptLocalCipher); + } + } + + private void assertInstanceEditEncryptsDefinitionSnapshot(Boolean syncDefine, String keptGlobalCipher, + String keptLocalCipher) { + Mockito.clearInvocations(processService); + 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\":\"kept\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"kept-secret\",\"sensitive\":true}]"); + WorkflowDefinition workflowDefinition = getProcessDefinition(); + workflowDefinition.setProjectCode(projectCode); + workflowDefinition.setGlobalParams( + "[{\"prop\":\"kept\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"" + keptGlobalCipher + + "\",\"sensitive\":true}]"); + + 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); + 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); + + TaskDefinitionLog existingTaskLog = new TaskDefinitionLog(); + existingTaskLog.setCode(4254862762304L); + existingTaskLog.setVersion(1); + existingTaskLog.setTaskParams( + "{\"localParams\":[{\"prop\":\"keptLocal\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"" + + keptLocalCipher + "\",\"sensitive\":true}],\"rawScript\":\"echo 1\"}"); + when(taskDefinitionLogMapper.queryByTaskDefinitions(any())) + .thenReturn(Collections.singletonList(existingTaskLog)); + + String submittedGlobalParams = + "[{\"prop\":\"kept\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"******\",\"sensitive\":true}," + + "{\"prop\":\"pwd\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"new-global\",\"sensitive\":true}]"; + String submittedTaskDefinitionJson = + "[{\"code\":4254862762304,\"name\":\"test1\",\"version\":1,\"description\":\"\",\"delayTime\":0," + + "\"taskType\":\"SHELL\",\"taskParams\":{\"resourceList\":[],\"localParams\":[" + + "{\"prop\":\"keptLocal\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"******\",\"sensitive\":true}," + + "{\"prop\":\"apiKey\",\"direct\":\"IN\",\"type\":\"VARCHAR\",\"value\":\"new-local\",\"sensitive\":true}" + + "],\"rawScript\":\"echo 1\"},\"flag\":\"YES\",\"taskPriority\":\"MEDIUM\",\"workerGroup\":\"default\"," + + "\"failRetryTimes\":0,\"failRetryInterval\":1,\"timeoutFlag\":\"CLOSE\",\"timeoutNotifyStrategy\":null," + + "\"timeout\":0,\"environmentCode\":-1}]"; + + ArgumentCaptor savedTaskLogs = ArgumentCaptor.forClass(List.class); + when(processService.saveTaskDefine(any(), Mockito.anyLong(), savedTaskLogs.capture(), eq(syncDefine))) + .thenReturn(1); + when(processService.saveWorkflowDefine(any(), eq(workflowDefinition), eq(syncDefine), eq(Boolean.FALSE))) + .thenReturn(1); + + workflowInstanceService.updateWorkflowInstance(loginUser, projectCode, 1, + taskRelationJson, submittedTaskDefinitionJson, "2020-02-21 00:00:00", syncDefine, + submittedGlobalParams, "", 0); + + List instanceGlobals = + GlobalParameterUtils.deserializeGlobalParameter(workflowInstance.getGlobalParams()); + Assertions.assertEquals("kept-secret", findProperty(instanceGlobals, "kept").getValue()); + Assertions.assertEquals("new-global", findProperty(instanceGlobals, "pwd").getValue()); + + List definitionGlobals = + GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams()); + Assertions.assertEquals(keptGlobalCipher, findProperty(definitionGlobals, "kept").getValue()); + Property encodedGlobal = findProperty(definitionGlobals, "pwd"); + Assertions.assertNotEquals("new-global", encodedGlobal.getValue()); + Assertions.assertEquals("new-global", PasswordUtils.decodePassword(encodedGlobal.getValue())); + Assertions.assertFalse(workflowDefinition.getGlobalParams().contains("kept-secret")); + + TaskDefinitionLog savedTask = ((List) savedTaskLogs.getValue()).get(0); + List localParams = JSONUtils.toList( + JSONUtils.getNodeString(savedTask.getTaskParams(), "localParams"), Property.class); + Assertions.assertEquals(keptLocalCipher, findProperty(localParams, "keptLocal").getValue()); + Property encodedLocal = findProperty(localParams, "apiKey"); + Assertions.assertNotEquals("new-local", encodedLocal.getValue()); + Assertions.assertEquals("new-local", PasswordUtils.decodePassword(encodedLocal.getValue())); + verify(processService).saveTaskDefine(eq(loginUser), eq(projectCode), any(), eq(syncDefine)); + verify(processService).saveWorkflowDefine(any(), eq(workflowDefinition), eq(syncDefine), eq(Boolean.FALSE)); + } + + private static Property findProperty(List properties, String prop) { + return properties.stream() + .filter(property -> prop.equals(property.getProp())) + .findFirst() + .orElseThrow(() -> new AssertionError(prop)); + } + @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 index 368622976718..c7ad4176ca3b 100644 --- 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 @@ -19,6 +19,7 @@ import org.apache.dolphinscheduler.api.exceptions.ServiceException; import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.common.utils.PropertyUtils; import org.apache.dolphinscheduler.dao.entity.Command; import org.apache.dolphinscheduler.dao.entity.ErrorCommand; import org.apache.dolphinscheduler.dao.entity.TaskInstance; @@ -27,6 +28,8 @@ 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.datasource.api.constants.DataSourceConstants; +import org.apache.dolphinscheduler.plugin.datasource.api.utils.PasswordUtils; 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; @@ -39,6 +42,8 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.mockito.MockedStatic; +import org.mockito.Mockito; class SensitivePropertyUtilsTest { @@ -280,6 +285,91 @@ void maskErrorCommandMasksSensitiveCommandParam() { Assertions.assertNotSame(original, masked); } + @Test + void encodeForCreateEncryptsWhenEnabled() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mockEncryptionEnabled(mocked); + List encoded = SensitivePropertyUtils.encodeForCreate( + Collections.singletonList(sensitive("pwd", "Secret123"))); + Assertions.assertNotEquals("Secret123", encoded.get(0).getValue()); + Assertions.assertEquals("Secret123", PasswordUtils.decodePassword(encoded.get(0).getValue())); + } + } + + @Test + void encodeForCreateSkipsWhenDisabled() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mockEncryptionDisabled(mocked); + List encoded = SensitivePropertyUtils.encodeForCreate( + Collections.singletonList(sensitive("pwd", "Secret123"))); + Assertions.assertEquals("Secret123", encoded.get(0).getValue()); + } + } + + @Test + void mergeAndEncodeDoesNotDoubleEncryptKeepOriginal() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mockEncryptionEnabled(mocked); + String ciphertext = PasswordUtils.encodePassword("Secret123"); + List encoded = SensitivePropertyUtils.mergeAndEncode( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(sensitive("pwd", ciphertext))); + Assertions.assertEquals(ciphertext, encoded.get(0).getValue()); + } + } + + @Test + void mergeAndEncodeFalseToTrueWithPlaceholderThenEncrypts() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mockEncryptionEnabled(mocked); + List encoded = SensitivePropertyUtils.mergeAndEncode( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(nonSensitive("pwd", "plain"))); + Assertions.assertNotEquals("plain", encoded.get(0).getValue()); + Assertions.assertEquals("plain", PasswordUtils.decodePassword(encoded.get(0).getValue())); + Assertions.assertTrue(encoded.get(0).isSensitive()); + } + } + + @Test + void mergeAndEncodeTrueToFalseDecodesWhenCiphertextResubmitted() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mockEncryptionEnabled(mocked); + String ciphertext = PasswordUtils.encodePassword("Secret123"); + List encoded = SensitivePropertyUtils.mergeAndEncode( + Collections.singletonList(nonSensitive("pwd", ciphertext)), + Collections.singletonList(sensitive("pwd", ciphertext))); + Assertions.assertEquals("Secret123", encoded.get(0).getValue()); + Assertions.assertFalse(encoded.get(0).isSensitive()); + } + } + + @Test + void restoreStartParamsDecryptsDefinitionGlobals() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mockEncryptionEnabled(mocked); + String ciphertext = PasswordUtils.encodePassword("Secret123"); + List restored = SensitivePropertyUtils.restoreStartParams( + Collections.singletonList(sensitive("pwd", TaskConstants.SENSITIVE_DATA_MASK)), + Collections.singletonList(sensitive("pwd", ciphertext))); + Assertions.assertEquals("Secret123", restored.get(0).getValue()); + Assertions.assertTrue(restored.get(0).isSensitive()); + } + } + + private static void mockEncryptionEnabled(MockedStatic mocked) { + mocked.when(() -> PropertyUtils.getBoolean(DataSourceConstants.DATASOURCE_ENCRYPTION_ENABLE, false)) + .thenReturn(true); + mocked.when(() -> PropertyUtils.getString(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT, + DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT)) + .thenReturn(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT); + } + + private static void mockEncryptionDisabled(MockedStatic mocked) { + mocked.when(() -> PropertyUtils.getBoolean(DataSourceConstants.DATASOURCE_ENCRYPTION_ENABLE, false)) + .thenReturn(false); + } + private static Property sensitive(String prop, String value) { return Property.builder() .prop(prop) diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/command/handler/RunWorkflowCommandHandler.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/command/handler/RunWorkflowCommandHandler.java index 414921a6ba98..e149d4387888 100644 --- a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/command/handler/RunWorkflowCommandHandler.java +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/command/handler/RunWorkflowCommandHandler.java @@ -35,6 +35,7 @@ import org.apache.dolphinscheduler.server.master.engine.task.execution.TaskExecution; import org.apache.dolphinscheduler.server.master.engine.task.execution.TaskExecutionBuilder; import org.apache.dolphinscheduler.server.master.runner.WorkflowExecuteContext.WorkflowExecuteContextBuilder; +import org.apache.dolphinscheduler.server.master.utils.SensitivePropertyCryptoUtils; import org.apache.commons.collections4.CollectionUtils; @@ -125,7 +126,8 @@ private String mergeCommandParamsWithWorkflowParams(final Command command, .map(ICommandParam::getCommandParams) .orElse(null); final List globalParamsList = - GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams()); + SensitivePropertyCryptoUtils.decodeSensitiveValues( + GlobalParameterUtils.deserializeGlobalParameter(workflowDefinition.getGlobalParams())); Map finalParams = new HashMap<>(); if (CollectionUtils.isNotEmpty(globalParamsList)) { globalParamsList.forEach(globalParam -> finalParams.put(globalParam.getProp(), globalParam)); diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/TaskExecutionContextBuilder.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/TaskExecutionContextBuilder.java index 69765d690e18..6338ba481e78 100644 --- a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/TaskExecutionContextBuilder.java +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/engine/task/execution/TaskExecutionContextBuilder.java @@ -29,6 +29,7 @@ import org.apache.dolphinscheduler.plugin.task.api.enums.TaskTimeoutStrategy; import org.apache.dolphinscheduler.plugin.task.api.model.Property; import org.apache.dolphinscheduler.plugin.task.api.parameters.resource.ResourceParametersHelper; +import org.apache.dolphinscheduler.server.master.utils.SensitivePropertyCryptoUtils; import java.util.Map; import java.util.concurrent.TimeUnit; @@ -81,7 +82,8 @@ public TaskExecutionContextBuilder buildTaskDefinitionRelatedInfo(final TaskDefi (int) Math.min(TimeUnit.MINUTES.toSeconds(taskDefinition.getTimeout()), Integer.MAX_VALUE)); } } - taskExecutionContext.setTaskParams(taskDefinition.getTaskParams()); + taskExecutionContext.setTaskParams( + SensitivePropertyCryptoUtils.decodeLocalParamsInTaskParams(taskDefinition.getTaskParams())); return this; } diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/TaskExecutionContextFactory.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/TaskExecutionContextFactory.java index 653088911514..d6627150a499 100644 --- a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/TaskExecutionContextFactory.java +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/TaskExecutionContextFactory.java @@ -46,6 +46,7 @@ import org.apache.dolphinscheduler.server.master.engine.task.execution.ITaskExecution; import org.apache.dolphinscheduler.server.master.engine.task.execution.TaskExecutionContextBuilder; import org.apache.dolphinscheduler.server.master.engine.task.execution.TaskExecutionContextCreateRequest; +import org.apache.dolphinscheduler.server.master.utils.SensitivePropertyCryptoUtils; import org.apache.dolphinscheduler.service.expand.CuringParamsService; import org.apache.dolphinscheduler.service.process.ProcessService; @@ -170,7 +171,7 @@ private Map getPrepareParams(final TaskInstance taskInstance, final Project project) { final AbstractParameters baseParam = TaskPluginManager.parseTaskParameters( taskInstance.getTaskType(), - taskInstance.getTaskParams()); + SensitivePropertyCryptoUtils.decodeLocalParamsInTaskParams(taskInstance.getTaskParams())); return curingParamsService.paramParsingPreparation( taskInstance, diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/utils/SensitivePropertyCryptoUtils.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/utils/SensitivePropertyCryptoUtils.java new file mode 100644 index 000000000000..c4f23d375a32 --- /dev/null +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/utils/SensitivePropertyCryptoUtils.java @@ -0,0 +1,43 @@ +/* + * 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.server.master.utils; + +import org.apache.dolphinscheduler.plugin.datasource.api.utils.PasswordUtils; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; +import org.apache.dolphinscheduler.plugin.task.api.utils.PropertySensitiveUtils; + +import java.util.List; + +import lombok.experimental.UtilityClass; + +/** + * Decrypt definition-time sensitive values for execution. + * API/UI still mask; Worker receives plaintext copies only. + */ +@UtilityClass +public class SensitivePropertyCryptoUtils { + + public List decodeSensitiveValues(List properties) { + return PropertySensitiveUtils.transformSensitiveValues(properties, PasswordUtils::decodePassword); + } + + public String decodeLocalParamsInTaskParams(String taskParams) { + return PropertySensitiveUtils.transformLocalParamsInTaskParams(taskParams, + SensitivePropertyCryptoUtils::decodeSensitiveValues); + } +} diff --git a/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/utils/SensitivePropertyCryptoUtilsTest.java b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/utils/SensitivePropertyCryptoUtilsTest.java new file mode 100644 index 000000000000..87921d98083a --- /dev/null +++ b/dolphinscheduler-master/src/test/java/org/apache/dolphinscheduler/server/master/utils/SensitivePropertyCryptoUtilsTest.java @@ -0,0 +1,76 @@ +/* + * 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.server.master.utils; + +import org.apache.dolphinscheduler.common.utils.PropertyUtils; +import org.apache.dolphinscheduler.plugin.datasource.api.constants.DataSourceConstants; +import org.apache.dolphinscheduler.plugin.datasource.api.utils.PasswordUtils; +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.List; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.MockedStatic; +import org.mockito.Mockito; + +class SensitivePropertyCryptoUtilsTest { + + @Test + void decodeSensitiveValuesDecryptsWhenEnabled() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mocked.when(() -> PropertyUtils.getBoolean(DataSourceConstants.DATASOURCE_ENCRYPTION_ENABLE, false)) + .thenReturn(true); + mocked.when(() -> PropertyUtils.getString(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT, + DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT)) + .thenReturn(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT); + + String ciphertext = PasswordUtils.encodePassword("Secret123"); + List decoded = SensitivePropertyCryptoUtils.decodeSensitiveValues( + Collections.singletonList(Property.builder() + .prop("pwd") + .direct(Direct.IN) + .type(DataType.VARCHAR) + .value(ciphertext) + .sensitive(true) + .build())); + Assertions.assertEquals("Secret123", decoded.get(0).getValue()); + } + } + + @Test + void decodeLocalParamsInTaskParams() { + try (MockedStatic mocked = Mockito.mockStatic(PropertyUtils.class)) { + mocked.when(() -> PropertyUtils.getBoolean(DataSourceConstants.DATASOURCE_ENCRYPTION_ENABLE, false)) + .thenReturn(true); + mocked.when(() -> PropertyUtils.getString(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT, + DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT)) + .thenReturn(DataSourceConstants.DATASOURCE_ENCRYPTION_SALT_DEFAULT); + + String ciphertext = PasswordUtils.encodePassword("token-value"); + String taskParams = "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"" + ciphertext + "\",\"sensitive\":true}]}"; + String decoded = SensitivePropertyCryptoUtils.decodeLocalParamsInTaskParams(taskParams); + Assertions.assertTrue(decoded.contains("token-value")); + Assertions.assertFalse(decoded.contains(ciphertext)); + } + } +} 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 index 2302206d5edb..3b15d81c60ab 100644 --- 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 @@ -17,10 +17,12 @@ package org.apache.dolphinscheduler.plugin.task.api.utils; +import org.apache.dolphinscheduler.common.utils.JSONUtils; import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; import org.apache.dolphinscheduler.plugin.task.api.model.Property; import org.apache.commons.collections4.CollectionUtils; +import org.apache.commons.lang3.StringUtils; import java.util.Collections; import java.util.List; @@ -31,6 +33,9 @@ import lombok.experimental.UtilityClass; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; + /** * Mask / merge helpers for {@link Property#isSensitive()}. * Keep-original marker is only {@code ******}; empty / null is a real empty value. @@ -85,6 +90,28 @@ public List maskSensitiveValues(List properties) { .collect(Collectors.toList()); } + /** + * Deep-copy then apply {@code transformer} to each non-empty {@code sensitive=true} value. + * Used for definition-time encode / execution-time decode without pulling crypto deps into task-api. + */ + public List transformSensitiveValues(List properties, + Function transformer) { + if (CollectionUtils.isEmpty(properties)) { + return Collections.emptyList(); + } + return properties.stream() + .map(property -> transformSensitiveValue(property, transformer)) + .collect(Collectors.toList()); + } + + public Property transformSensitiveValue(Property property, Function transformer) { + Property copied = copy(property); + if (isSensitive(copied) && StringUtils.isNotEmpty(copied.getValue())) { + copied.setValue(transformer.apply(copied.getValue())); + } + return copied; + } + public List mergeSensitiveValuePlaceholders(List submittedProperties, List existingProperties) { if (CollectionUtils.isEmpty(submittedProperties)) { @@ -157,4 +184,26 @@ private Map toPropMap(List properties) { .filter(property -> property.getProp() != null) .collect(Collectors.toMap(Property::getProp, Function.identity(), (left, right) -> right)); } + + /** + * Rewrite {@code localParams} inside a taskParams JSON string. + * Returns the original string when missing / unparsable. + */ + public String transformLocalParamsInTaskParams(String taskParams, + Function, List> transform) { + if (StringUtils.isEmpty(taskParams) || transform == null) { + return taskParams; + } + ObjectNode taskParamsNode = JSONUtils.parseObject(taskParams); + if (taskParamsNode == null) { + return taskParams; + } + JsonNode localParamsNode = taskParamsNode.findValue("localParams"); + if (localParamsNode == null || localParamsNode.isNull()) { + return taskParams; + } + List localParams = JSONUtils.toList(localParamsNode.toString(), Property.class); + taskParamsNode.set("localParams", JSONUtils.toJsonNode(transform.apply(localParams))); + return JSONUtils.toJsonString(taskParamsNode); + } } 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 index 7227d9d31cb6..f2fa90a5c807 100644 --- 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 @@ -145,6 +145,27 @@ void serializationRoundTripKeepsSensitiveFlag() { Assertions.assertEquals("abc", parsed.get(0).getValue()); } + @Test + void transformSensitiveValuesAppliesTransformer() { + Property original = sensitive("pwd", "Secret123"); + List transformed = PropertySensitiveUtils.transformSensitiveValues( + Collections.singletonList(original), value -> "enc:" + value); + + Assertions.assertEquals("enc:Secret123", transformed.get(0).getValue()); + Assertions.assertEquals("Secret123", original.getValue()); + } + + @Test + void transformLocalParamsInTaskParamsDecryptsSensitive() { + String taskParams = "{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," + + "\"value\":\"cipher\",\"sensitive\":true}],\"rawScript\":\"echo 1\"}"; + String decoded = PropertySensitiveUtils.transformLocalParamsInTaskParams(taskParams, + props -> PropertySensitiveUtils.transformSensitiveValues(props, value -> "plain")); + Assertions.assertTrue(decoded.contains("\"plain\"")); + Assertions.assertTrue(decoded.contains("rawScript")); + Assertions.assertTrue(taskParams.contains("cipher")); + } + private static Property sensitive(String prop, String value) { return Property.builder() .prop(prop)