fix(gateway): internal events must not re-key the session-context / channel prompt pin

Internal events (kanban wakes, delegation completions, watch
notifications) carry a source rebuilt from the persisted origin: no
chat_name/user_name/message_id/parent_chat_id, and channel_prompt=None.
_pinned_session_context_prompt re-rendered from that source, so every
internal turn re-keyed the pin and the next human turn re-keyed it back
(A->B->A). Each flip rewrote already-sent system bytes and collapsed the
prompt cache to the static prefix. The same toggle happened through
channel_prompt and parent-keyed channel_overrides in the ephemeral
system prompt.

- _pinned_session_context_prompt(internal=True) reuses the existing pin
  verbatim and never re-pins.
- Human turns record (channel_prompt, parent_chat_id) in
  ConversationState.channel_pin; internal turns reuse them (main run and
  the queued follow-up path).

tests/gateway/test_internal_event_pin_wiring.py drives the real
_handle_message_with_agent human->internal->human with _run_agent
stubbed. Both tests fail on main and pass here; dropping internal= at
the call site makes both fail again.

(cherry picked from commit 3b38527ae7f0f498ca265fe9f64fcbc1c8c7765c)
This commit is contained in:
Kyzcreig
2026-09-24 17:03:17 -07:00
committed by kshitij
parent e30e61f6ff
commit 9c18b383de
4 changed files with 264 additions and 7 deletions

View File

@@ -4,6 +4,7 @@ for GatewayRunner (MRO mixin). ``gateway.run`` internals are imported lazily ins
from __future__ import annotations
import dataclasses
import importlib
import logging
import threading
@@ -618,12 +619,22 @@ class GatewayAgentCacheMixin:
return None
return f"[Voice channel now: {vc_now or 'not connected to a voice channel'}]"
def _pinned_session_context_prompt(self, context, redact_pii: bool, session_key: Optional[str]) -> str:
def _pinned_session_context_prompt(
self, context, redact_pii: bool, session_key: Optional[str], *, internal: bool = False,
) -> str:
"""Session-context prompt pinned per session: key hit → pinned bytes reused VERBATIM (immune
to renderer nondeterminism); key miss → re-render and re-pin (rename, topic edit, /sethome)."""
_eph_key = self._ephemeral_change_key(context, redact_pii)
to renderer nondeterminism); key miss → re-render and re-pin (rename, topic edit, /sethome).
``internal`` events (kanban wakes, delegation completions, watch notifications) carry a
source rebuilt from the persisted origin, without chat_name/user_name/message_id. Rendering
from it re-keyed the pin, and the next human turn re-keyed it back (A→B→A), rewriting
already-sent system bytes each time. An internal event is never a real metadata change, so
it reuses the existing pin verbatim and never re-pins."""
_pin_state = self._peek_session_state(session_key) if session_key else None
_eph_pin = _pin_state.conversation.ephemeral_pin if _pin_state else None
if internal:
return _eph_pin[1] if _eph_pin is not None else build_session_context_prompt(context, redact_pii=redact_pii)
_eph_key = self._ephemeral_change_key(context, redact_pii)
if _eph_pin is not None and _eph_pin[0] == _eph_key:
return _eph_pin[1]
text = build_session_context_prompt(context, redact_pii=redact_pii)
@@ -631,6 +642,31 @@ class GatewayAgentCacheMixin:
self._session_state(session_key).conversation.ephemeral_pin = (_eph_key, text)
return text
def _pinned_channel_inputs(self, session_key, event, source):
"""``(channel_prompt, source)`` for this turn's agent run.
The ephemeral system prompt also appends ``channel_prompt`` and the ``channel_overrides``
prompt (looked up by chat/thread/``parent_chat_id``). Internal events carry
``channel_prompt=None`` and a source without ``parent_chat_id``, so they dropped both and
toggled the system prompt like the context pin did. Human turns record their inputs;
internal turns reuse them."""
channel_prompt = getattr(event, "channel_prompt", None)
if not session_key:
return channel_prompt, source
if not getattr(event, "internal", False):
self._session_state(session_key).conversation.channel_pin = (
channel_prompt, getattr(source, "parent_chat_id", None),
)
return channel_prompt, source
state = self._peek_session_state(session_key)
pin = state.conversation.channel_pin if state else None
if pin is None:
return channel_prompt, source
pinned_prompt, pinned_parent = pin
if pinned_parent and not getattr(source, "parent_chat_id", None):
source = dataclasses.replace(source, parent_chat_id=pinned_parent)
return pinned_prompt, source
@staticmethod
def _ephemeral_change_key(context, redact_pii: bool) -> str:
"""Hash the exact inputs ``build_session_context_prompt`` renders. Invariant

View File

@@ -2087,7 +2087,9 @@ class GatewayTurnMixin:
# The context prompt render is pinned per session, keyed by a hash of the renderer inputs, so
# the system prompt cannot drift turn-over-turn; a miss (thread rename, /sethome) re-renders.
context_prompt = self._pinned_session_context_prompt(context, _redact_pii, session_key)
context_prompt = self._pinned_session_context_prompt(
context, _redact_pii, session_key, internal=bool(getattr(event, "internal", False)),
)
# Per-turn notes ride the user message via the api_content sidecar, NOT context_prompt
# (appending to the ephemeral system prompt forced a full agent rebuild).
@@ -2205,12 +2207,14 @@ class GatewayTurnMixin:
# Admission/typing is not execution. All routing, authorization and
# turn preparation gates have passed when the agent runner is entered.
event._heartbeat_execution_started = True
# Internal events reuse the last human turn's channel inputs (see _pinned_channel_inputs).
_turn_channel_prompt, _turn_source = self._pinned_channel_inputs(session_key, event, source)
agent_result = await self._run_agent(
message=message_text, context_prompt=prepared.context_prompt, history=history, source=source,
message=message_text, context_prompt=prepared.context_prompt, history=history, source=_turn_source,
session_id=_run_start_session_id, session_key=session_key,
run_generation=run_generation, event_message_id=self._reply_anchor_for_event(event),
inbound_message_id=str(event.message_id) if event.message_id else None,
channel_prompt=event.channel_prompt, moa_config=getattr(event, "_moa_config", None),
channel_prompt=_turn_channel_prompt, moa_config=getattr(event, "_moa_config", None),
persist_user_message=prepared.persist_user_message,
persist_user_timestamp=prepared.persist_user_timestamp,
persist_user_display_kind=prepared.persist_user_display_kind,
@@ -3872,7 +3876,9 @@ class GatewayTurnMixin:
next_persist_message = strip_discord_triggering_note(pending_event, next_message)
next_message_id = self._reply_anchor_for_event(pending_event)
next_inbound_id = str(pending_event.message_id) if getattr(pending_event, "message_id", None) else None
next_channel_prompt = getattr(pending_event, "channel_prompt", None)
next_channel_prompt, next_source = self._pinned_channel_inputs(
next_session_key, pending_event, next_source,
)
next_message_type = getattr(pending_event, "message_type", None)
# Clear the prior turn's streaming-TTS completion marker so the recursive turn isn't suppressed.

View File

@@ -50,6 +50,8 @@ class ConversationState:
queued_events: List[Any] = field(default_factory=list) # /queue overflow FIFO (head in adapter)
sidecar_notes: List[str] = field(default_factory=list) # one-shot must-deliver notes
ephemeral_pin: Optional[Tuple[Any, ...]] = None # pinned session-context (change_key, text)
# (channel_prompt, parent_chat_id) of the last non-internal turn; internal events reuse it
channel_pin: Optional[Tuple[Optional[str], Optional[str]]] = None
vc_last: Optional[str] = None # last voice-channel context delivered
def clear(self) -> None: