Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 18 additions & 6 deletions api/gateway_chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -909,6 +909,7 @@ def _run_gateway_chat_streaming(
*,
model_provider=None,
goal_related=False,
regeneration=False,
):
"""Bridge a WebUI chat turn through Hermes Gateway's API server.

Expand Down Expand Up @@ -1279,18 +1280,29 @@ def put_gateway_event(event, data):
# same sort key; later transcript merges can then fall back to
# role/content ordering instead of turn order.
assistant_ts = now + 0.000001
user_msg = {"role": "user", "content": str(msg_text or ""), "timestamp": now}
pending_source = getattr(s, "pending_user_source", None) or "webui"
if pending_source != "webui":
user_msg["_source"] = pending_source
if attachments:
user_msg["attachments"] = list(attachments)
from api.streaming import _active_turn_authority, _materialize_active_turn_user

active_turn_identity = _active_turn_authority(s, stream_id, msg_text)
user_msg = _materialize_active_turn_user(
active_turn_identity,
str(msg_text or ""),
pending_source,
)
user_msg["timestamp"] = float(
active_turn_identity.get("timestamp") or now
)
assistant_msg = {"role": "assistant", "content": assistant_text, "timestamp": assistant_ts}
saved_reasoning = STREAM_REASONING_TEXT.get(stream_id, "")
if saved_reasoning:
assistant_msg["reasoning"] = saved_reasoning
previous_messages = list(getattr(s, "messages", None) or [])
previous_context = list(getattr(s, "context_messages", None) or getattr(s, "messages", None) or [])
stored_context = getattr(s, "context_messages", None)
previous_context = list(
stored_context
if isinstance(stored_context, list) and (regeneration or stored_context)
else getattr(s, "messages", None) or []
)
previous_process_wakeup_pause = dict(getattr(s, "process_wakeup_pause", {}) or {})
# Stamp stable ids on the two new rows (shared with the display merge
# below) so display and model-context copies share an id for the
Expand Down
45 changes: 28 additions & 17 deletions api/helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,9 @@
"_state_db_row_id",
"_db_row_id",
"state_db_row_id",
"_active_turn_token",
"_active_turn_user",
"_fork_child_turn",
})


Expand Down Expand Up @@ -1114,8 +1117,14 @@ def scrub_internal_replay_fields(
return result


def _public_message_projection(message, *, _enabled: bool):
def _public_message_projection(message, *, _enabled: bool, _active_turn_token=None):
"""Return one public transcript message without internal replay fields."""
is_active = (
isinstance(message, dict)
and message.get("role") == "user"
and _active_turn_token is not None
and message.get("_active_turn_token") == _active_turn_token
)
message = scrub_internal_replay_fields([message], message_records=True)[0]
if not isinstance(message, dict):
return _redact_value(message, _enabled=_enabled)
Expand All @@ -1131,13 +1140,22 @@ def _public_message_projection(message, *, _enabled: bool):
]
else:
item[key] = _redact_value(value, _enabled=_enabled)
if is_active:
item["_active_turn_user"] = True
return item


def _redact_messages(messages, *, _enabled: bool):
def _redact_messages(messages, *, _enabled: bool, _active_turn_token=None):
if not isinstance(messages, list):
return _redact_value(messages, _enabled=_enabled)
return [_public_message_projection(message, _enabled=_enabled) for message in messages]
return [
_public_message_projection(
message,
_enabled=_enabled,
_active_turn_token=_active_turn_token,
)
for message in messages
]


def _redact_tool_calls(tool_calls, *, _enabled: bool):
Expand All @@ -1152,10 +1170,12 @@ def _redact_nested_message_containers(value, *, _enabled: bool):
return _redact_value(scrubbed, _enabled=_enabled)
result = {}
for key, child in scrubbed.items():
if key == "messages" and isinstance(child, list):
if key in {"messages", "context_messages"} and isinstance(child, list):
result[key] = _redact_messages(child, _enabled=_enabled)
elif key == "tool_calls" and isinstance(child, list):
result[key] = _redact_tool_calls(child, _enabled=_enabled)
elif key == "runtime_journal_snapshot" and isinstance(child, dict):
result[key] = _redact_nested_message_containers(child, _enabled=_enabled)
else:
result[key] = _redact_value(child, _enabled=_enabled)
return result
Expand Down Expand Up @@ -1193,30 +1213,21 @@ def _copy_json_value(value):


def redact_session_data(session_dict: dict) -> dict:
"""Redact credentials from message content, tool data, and session sidecars.

Applies to: messages[], tool_calls[], todo_state, runtime_journal_snapshot,
and title.
The underlying session file is not modified; redaction is response-layer only.

Reads the ``api_redact_enabled`` setting ONCE for the entire response and
threads it through to avoid hundreds of settings.json reads per session
payload (a 50-message session has hundreds of nested strings). When the
setting is disabled this is also a fast path: the recursion still walks
but every string returns early.
"""
"""Redact credentials in the public session response without mutation."""
from api.config import load_settings
_enabled = bool(load_settings().get("api_redact_enabled", True))
if not isinstance(session_dict, dict):
return {}
result = {}
from api.process_event_utils import build_active_turn_token
_active_turn_token = build_active_turn_token(session_dict.get("active_stream_id"), session_dict.get("pending_started_at"))
for key, value in session_dict.items():
if key in _PUBLIC_MESSAGE_INTERNAL_FIELDS:
continue
if key == 'title' and isinstance(value, str):
result[key] = _redact_text(value, _enabled=_enabled)
elif key in {'messages', 'context_messages'}:
result[key] = _redact_messages(value, _enabled=_enabled)
result[key] = _redact_messages(value, _enabled=_enabled, _active_turn_token=_active_turn_token)
elif key == 'tool_calls' and isinstance(value, list):
result[key] = _redact_tool_calls(value, _enabled=_enabled)
elif key in {'todo_state', 'runtime_journal_snapshot'}:
Expand Down
Loading
Loading