Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -264,7 +264,13 @@ public WorkflowDefinition createWorkflowDefinition(User loginUser,

List<TaskDefinitionLog> taskDefinitionLogs = generateTaskDefinitionList(taskDefinitionJson);
List<WorkflowTaskRelationLog> taskRelationList = generateTaskRelationList(taskRelationJson, taskDefinitionLogs);
SensitivePropertyUtils.requireNoPlaceholder(GlobalParameterUtils.deserializeGlobalParameter(globalParams));
List<Property> 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));
Expand Down Expand Up @@ -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<TaskDefinition> versionKeys = new ArrayList<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -413,12 +413,13 @@ public WorkflowDefinition updateWorkflowInstance(User loginUser, long projectCod
}

List<Property> 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<TaskDefinitionLog> taskDefinitionLogs = JSONUtils.toList(taskDefinitionJson, TaskDefinitionLog.class);
if (taskDefinitionLogs.isEmpty()) {
log.warn("Parameter taskDefinitionJson is empty");
Expand Down Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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 {
Expand All @@ -80,9 +80,28 @@ public void requireNoPlaceholder(List<Property> properties) {
}

/**
* Update: restore DB plaintext for keep-original {@code ******}.
* Create definition params: reject {@code ******}, then encode sensitive plaintext.
*/
public List<Property> encodeForCreate(List<Property> 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<Property> mergeAndEncode(List<Property> submittedProperties, List<Property> existingProperties) {
List<Property> 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<Property> merge(List<Property> submittedProperties, List<Property> existingProperties) {
if (CollectionUtils.isEmpty(submittedProperties)) {
Expand All @@ -100,6 +119,7 @@ public List<Property> merge(List<Property> submittedProperties, List<Property> e

/**
* Build start/command params before persist.
* Decrypts definition globals first, then:
* <ul>
* <li>Same name as a workflow global: inherit all global attributes, only override {@code value}
* ({@code ******} keeps the global plaintext).</li>
Expand All @@ -110,7 +130,8 @@ public List<Property> restoreStartParams(List<Property> startParams, List<Proper
if (CollectionUtils.isEmpty(startParams)) {
return startParams;
}
Map<String, Property> globals = CollectionUtils.emptyIfNull(globalParams).stream()
List<Property> decryptedGlobals = decodeSensitiveValues(globalParams);
Map<String, Property> globals = CollectionUtils.emptyIfNull(decryptedGlobals).stream()
.filter(Objects::nonNull)
.filter(property -> property.getProp() != null)
.collect(Collectors.toMap(Property::getProp, Function.identity(), (left, right) -> right));
Expand All @@ -137,16 +158,68 @@ public List<Property> restoreStartParams(List<Property> startParams, List<Proper
return restored;
}

public List<Property> decodeSensitiveValues(List<Property> 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<Property> encodeNewSensitivePlaintext(List<Property> mergedProperties,
List<Property> existingProperties) {
if (CollectionUtils.isEmpty(mergedProperties)) {
return mergedProperties;
}
Map<String, Property> existingMap = CollectionUtils.emptyIfNull(existingProperties).stream()
.filter(Objects::nonNull)
.filter(property -> property.getProp() != null)
.collect(Collectors.toMap(Property::getProp, Function.identity(), (left, right) -> right));
List<Property> 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).
Expand Down
Loading
Loading