feat(stream): redact secret values from SSE workflow stream responses

Adds `secret_values` param to `WorkflowResponseConverter.__init__` and a
`_redact` helper that delegates to `redact_secret_values` from
`core.workflow.secret_scrub`. Applies redaction to event-sourced
inputs/process_data/outputs in all four stream builders:
`workflow_node_finish_to_stream_response`,
`workflow_node_retry_to_stream_response`,
`workflow_finish_to_stream_response`, and
`workflow_pause_to_stream_response`. The `workflow_run.outputs_dict` path
(reconnect/resume) is intentionally untouched (already redacted by Plan 2).
Empty registry is a strict no-op, so existing tests pass unchanged.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Charles Yao
2026-07-15 06:07:11 -07:00
parent 5ca7d9cee0
commit 7b55944a8c
2 changed files with 97 additions and 8 deletions
@@ -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)
@@ -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)