diff --git a/api/core/app/apps/common/workflow_response_converter.py b/api/core/app/apps/common/workflow_response_converter.py index 236b8c58d3a..97f9d882a7b 100644 --- a/api/core/app/apps/common/workflow_response_converter.py +++ b/api/core/app/apps/common/workflow_response_converter.py @@ -61,6 +61,7 @@ from core.workflow.human_input_policy import ( resolve_human_input_pause_reason_inputs, ) from core.workflow.nodes.human_input.pause_reason import HumanInputRequired +from core.workflow.secret_scrub import redact_secret_values from core.workflow.system_variables import SystemVariableKey, system_variables_to_mapping from core.workflow.workflow_entry import WorkflowEntry from extensions.ext_database import db @@ -130,10 +131,12 @@ class WorkflowResponseConverter: application_generate_entity: Union[AdvancedChatAppGenerateEntity, WorkflowAppGenerateEntity], user: Union[Account, EndUser], system_variables: Sequence[Variable], + secret_values: tuple[str, ...] = (), ): self._application_generate_entity = application_generate_entity self._user = user self._system_variables = system_variables_to_mapping(system_variables) + self._secret_values = secret_values self._workflow_inputs = self._prepare_workflow_inputs() # Disable truncation for SERVICE_API calls to keep backward compatibility. @@ -234,6 +237,9 @@ class WorkflowResponseConverter: converter = WorkflowRuntimeTypeConverter() return converter.to_json_encodable(outputs) + def _redact(self, value: object) -> object: + return redact_secret_values(value, self._secret_values) + def workflow_start_to_stream_response( self, *, @@ -279,7 +285,7 @@ class WorkflowResponseConverter: elapsed_time = (finished_at - started_at).total_seconds() outputs_mapping = graph_runtime_state.outputs or {} - encoded_outputs = WorkflowRuntimeTypeConverter().to_json_encodable(outputs_mapping) + encoded_outputs = self._redact(WorkflowRuntimeTypeConverter().to_json_encodable(outputs_mapping)) created_by: CreatedByDict | dict[str, object] = {} user = self._user @@ -331,7 +337,7 @@ class WorkflowResponseConverter: ) paused_at = naive_utc_now() elapsed_time = (paused_at - started_at).total_seconds() - encoded_outputs = self._encode_outputs(event.outputs) or {} + encoded_outputs = self._redact(self._encode_outputs(event.outputs) or {}) if self._application_generate_entity.invoke_from == InvokeFrom.SERVICE_API: encoded_outputs = {} variable_pool = graph_runtime_state.variable_pool @@ -575,9 +581,9 @@ class WorkflowResponseConverter: finished_at = event.finished_at or naive_utc_now() elapsed_time = (finished_at - start_at).total_seconds() - inputs, inputs_truncated = self._truncate_mapping(event.inputs) - process_data, process_data_truncated = self._truncate_mapping(event.process_data) - encoded_outputs = self._encode_outputs(event.outputs) + inputs, inputs_truncated = self._truncate_mapping(self._redact(event.inputs)) + process_data, process_data_truncated = self._truncate_mapping(self._redact(event.process_data)) + encoded_outputs = self._redact(self._encode_outputs(event.outputs)) outputs, outputs_truncated = self._truncate_mapping(encoded_outputs) metadata = self._merge_metadata(event.execution_metadata, snapshot) @@ -635,9 +641,9 @@ class WorkflowResponseConverter: finished_at = naive_utc_now() elapsed_time = (finished_at - event.start_at).total_seconds() - inputs, inputs_truncated = self._truncate_mapping(event.inputs) - process_data, process_data_truncated = self._truncate_mapping(event.process_data) - encoded_outputs = self._encode_outputs(event.outputs) + inputs, inputs_truncated = self._truncate_mapping(self._redact(event.inputs)) + process_data, process_data_truncated = self._truncate_mapping(self._redact(event.process_data)) + encoded_outputs = self._redact(self._encode_outputs(event.outputs)) outputs, outputs_truncated = self._truncate_mapping(encoded_outputs) metadata = self._merge_metadata(event.execution_metadata, snapshot) diff --git a/api/tests/unit_tests/core/app/apps/common/test_workflow_response_converter.py b/api/tests/unit_tests/core/app/apps/common/test_workflow_response_converter.py index dd6cd0e919d..31649a19f07 100644 --- a/api/tests/unit_tests/core/app/apps/common/test_workflow_response_converter.py +++ b/api/tests/unit_tests/core/app/apps/common/test_workflow_response_converter.py @@ -1,8 +1,17 @@ from collections.abc import Mapping, Sequence +from unittest.mock import Mock from core.app.apps.common.workflow_response_converter import WorkflowResponseConverter +from core.app.entities.app_invoke_entities import InvokeFrom, WorkflowAppGenerateEntity +from core.app.entities.queue_entities import QueueNodeStartedEvent, QueueNodeSucceededEvent +from core.workflow.secret_scrub import SECRET_PLACEHOLDER +from core.workflow.system_variables import build_system_variables +from graphon.entities import WorkflowStartReason +from graphon.enums import BuiltinNodeTypes from graphon.file import FILE_MODEL_IDENTITY, File, FileTransferMethod, FileType from graphon.variables.segments import ArrayFileSegment, FileSegment +from libs.datetime_utils import naive_utc_now +from models import Account class TestWorkflowResponseConverterFetchFilesFromVariableValue: @@ -256,3 +265,77 @@ class TestWorkflowResponseConverterFetchFilesFromVariableValue: assert len(result) == 2 assert result[0]["id"] == "complex_file" assert result[1]["id"] == "complex_obj" + + +class TestWorkflowResponseConverterSecretRedaction: + """Test that WorkflowResponseConverter redacts secret values from SSE stream responses.""" + + def _make_converter(self, secret_values: tuple[str, ...] = ()) -> WorkflowResponseConverter: + mock_entity = Mock(spec=WorkflowAppGenerateEntity) + mock_app_config = Mock() + mock_app_config.tenant_id = "test-tenant-id" + mock_entity.invoke_from = InvokeFrom.WEB_APP + mock_entity.app_config = mock_app_config + mock_entity.inputs = {} + + mock_user = Mock(spec=Account) + mock_user.id = "test-user-id" + mock_user.name = "Test User" + mock_user.email = "test@example.com" + + system_variables = build_system_variables(workflow_id="wf-id", workflow_execution_id="run-id") + return WorkflowResponseConverter( + application_generate_entity=mock_entity, + user=mock_user, + system_variables=system_variables, + secret_values=secret_values, + ) + + def _prime_converter(self, converter: WorkflowResponseConverter) -> QueueNodeStartedEvent: + """Seed workflow run id and node snapshot so finish events are processable.""" + converter.workflow_start_to_stream_response( + task_id="t1", + workflow_run_id="run-id", + workflow_id="wf-id", + reason=WorkflowStartReason.INITIAL, + ) + start_event = QueueNodeStartedEvent( + node_execution_id="exec-1", + node_id="node-1", + node_title="Test Node", + node_type=BuiltinNodeTypes.CODE, + start_at=naive_utc_now(), + in_iteration_id=None, + in_loop_id=None, + provider_type="built-in", + provider_id="code", + ) + converter.workflow_node_start_to_stream_response(event=start_event, task_id="t1") + return start_event + + def test_node_finish_redacts_inputs_process_data_outputs(self): + """node-finish stream response must replace secret values with SECRET_PLACEHOLDER.""" + secret = "supersecretvalue123" + converter = self._make_converter(secret_values=(secret,)) + start_event = self._prime_converter(converter) + + succeeded_event = QueueNodeSucceededEvent( + node_execution_id=start_event.node_execution_id, + node_id="node-1", + node_type=BuiltinNodeTypes.CODE, + start_at=naive_utc_now(), + inputs={"api_key": secret}, + process_data={"raw_response": f"token={secret}"}, + outputs={"result": f"hello {secret} world"}, + ) + + response = converter.workflow_node_finish_to_stream_response(event=succeeded_event, task_id="t1") + + assert response is not None + # Secret must be replaced in all three fields + assert SECRET_PLACEHOLDER in str(response.data.inputs) + assert secret not in str(response.data.inputs) + assert SECRET_PLACEHOLDER in str(response.data.process_data) + assert secret not in str(response.data.process_data) + assert SECRET_PLACEHOLDER in str(response.data.outputs) + assert secret not in str(response.data.outputs)