Skip to content
Open
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
160 changes: 160 additions & 0 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -9262,6 +9262,155 @@ def _clear_goal_pending_continuations(self, session_key: str, adapter: Any) -> i
_q_state.conversation.queued_events = kept
return removed

def _replace_queued_message(
self,
session_key: str,
adapter: Any,
message_id: str,
new_text: str,
) -> bool:
"""Replace the text of the first queued event matching ``message_id``.

Searches the primary pending slot first, then the overflow FIFO.
Mutates the matching event's ``.text`` in place and preserves all
other fields and FIFO order. Returns True if a match was replaced,
False if no match was found in either queue level.

Modeled on ``_clear_goal_pending_continuations`` (same two-level
search pattern).
"""
# Primary slot
pending_slot = getattr(adapter, "_pending_messages", None) if adapter is not None else None
if isinstance(pending_slot, dict):
pending_event = pending_slot.get(session_key)
if pending_event is not None and getattr(pending_event, "message_id", None) == message_id:
pending_event.text = new_text
return True

# Overflow FIFO
_q_state = self._peek_session_state(session_key)
overflow = _q_state.conversation.queued_events if _q_state else []
for ev in overflow:
if getattr(ev, "message_id", None) == message_id:
ev.text = new_text
return True

return False

def _handle_edit_supersede(
self,
event: "MessageEvent",
session_key: str,
adapter: Any,
) -> bool:
"""Correlate an inbound edit with queued or in-flight messages.

Called from the synchronous edit-supersede block in
``_handle_message`` (#35535). Returns True if the event was handled
(queue-replaced, in-flight redirected/steered, or silently dropped);
the caller should return ``None`` to skip normal dispatch. Returns
False for non-edit events that should continue through normal
processing.

Design intent (per issue #35535 design review):
* Queue match -> replace text in-place, silent (no ack).
* In-flight match (active_message_id) -> redirect/steer with
framing text. Neither primitive cancels in-flight tools — this
is the best available approach without transcript mutation.
* No match (uncorrelated edit) -> silently dropped. An isolated
out-of-context correction would confuse the model.

SYNCHRONOUS: this method and ``_replace_queued_message`` are only
called from the gateway's single asyncio event loop with **no await
points** between the queue check, the in-flight check, and the
mutation — cooperative scheduling therefore cannot interleave a queue
promotion or turn completion inside the correlation block, and no lock
is required. Cross-vendor review (Gemini/GPT-OSS) verified this
invariant.
"""
_is_edit = bool(event.metadata.get("is_edit", False) if event.metadata else False)
_msg_id = event.message_id
if not _is_edit or not _msg_id:
return False # not an edit — continue normal dispatch

# SECURITY: correlation is scoped by session_key (derived from
# platform + chat_id, distinct for group vs DM sessions). Platform
# message ids are unique per chat, so an edit from user A cannot
# supersede a queued message from user B in a shared group — each
# group session is separate.
# 1. Try queued supersede
if self._replace_queued_message(session_key, adapter, _msg_id, event.text or ""):
logger.info(
"Edit supersede — replaced queued message_id=%s for session %s",
_msg_id, session_key,
)
return True

# 2. Try in-flight redirect / steer
_q_state = self._peek_session_state(session_key)
if _q_state and _q_state.turn.active_message_id == _msg_id:
running_agent = _q_state.turn.agent
if running_agent is not None and running_agent is not _AGENT_PENDING_SENTINEL:
_framing = f'[User edited their earlier message. Corrected message: "{event.text}"]'
_handled = False
# Prefer redirect() over steer() — redirect delivers as
# user-side input next model iteration, steer appends to
# last tool result. Neither cancels in-flight tools.
if (
getattr(running_agent, "_supports_active_turn_redirect", False) is True
and hasattr(running_agent, "redirect")
):
try:
_handled = bool(running_agent.redirect(_framing))
except Exception as exc:
logger.warning(
"Edit supersede — redirect failed for session %s: %s",
session_key, exc,
)
if not _handled and hasattr(running_agent, "steer"):
try:
_handled = bool(running_agent.steer(_framing))
except Exception as exc:
logger.warning(
"Edit supersede — steer failed for session %s: %s",
session_key, exc,
)
if _handled:
logger.info(
"Edit supersede — redirected/steered in-flight turn "
"message_id=%s for session %s",
_msg_id, session_key,
)
return True
# Fallback: queue the edit as a normal pending event
logger.warning(
"Edit supersede — redirect/steer unavailable for session %s, "
"queueing edit as normal pending event",
session_key,
)
self._queue_or_replace_pending_event(session_key, event)
return True

elif running_agent is _AGENT_PENDING_SENTINEL:
logger.info(
"Edit supersede — edit matched in-flight turn but agent "
"still pending (sentinel); edit dropped for session %s",
session_key,
)
return True

# 3. Uncoupled edit (#35535 deliberate policy): no queued match,
# no in-flight match. An isolated out-of-context correction
# would confuse the model — drop it silently.
logger.info(
"Edit supersede — uncorrelated edit dropped "
"(platform=%s, chat=%s, message_id=%s)",
getattr(getattr(event.source, "platform", None), "value", "?"),
getattr(event.source, "chat_id", "?"),
_msg_id,
)
return True

def _goal_still_active_for_session(self, session_id: str) -> bool:
"""Best-effort fresh DB check before running a queued continuation."""
if not session_id:
Expand Down Expand Up @@ -17703,6 +17852,16 @@ async def _handle_message(self, event: MessageEvent) -> Optional[str]:
)
self._release_running_agent_state(_quick_key)

# ── Edit supersede (#35535) ────────────────────────────────────
# SYNCHRONOUS block: no await between queued-replace, in-flight
# check, and decision. An inbound EDIT of a previously-sent message
# supersedes the queued original or the in-flight turn; an isolated
# out-of-context correction is silently dropped.
_edit_adapter = self._adapter_for_source(source)
if _edit_adapter is not None:
if self._handle_edit_supersede(event, _quick_key, _edit_adapter):
return None

if self._is_session_running(_quick_key):
# Resolve the command once; every command's mid-run behavior is
# declared on its CommandDef (busy_policy / busy_handler in
Expand Down Expand Up @@ -18677,6 +18836,7 @@ async def _do_undo():
_claim_state.turn.lease = _active_session_lease
_claim_state.turn.agent = _AGENT_PENDING_SENTINEL
_claim_state.turn.started_ts = time.time()
_claim_state.turn.active_message_id = event.message_id or None
self._persist_active_agents()
_run_generation = self._begin_session_run_generation(_quick_key)

Expand Down
4 changes: 4 additions & 0 deletions gateway/session_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,9 @@ class TurnState:
# preserves that: release/rebind only match when generation is current.
lease_token: Any = None
lease_generation: Optional[int] = None
# MessageEvent.message_id of the event that started this turn; used by
# edit-supersede (#35535) to correlate inbound edits with in-flight turns.
active_message_id: Optional[str] = None

def clear(self) -> None:
"""Reset the per-turn slot (agent / start ts / lease / busy-ack).
Expand All @@ -85,6 +88,7 @@ def clear(self) -> None:
self.started_ts = 0.0
self.lease = None
self.busy_ack_ts = 0.0
self.active_message_id = None


@dataclass
Expand Down
35 changes: 32 additions & 3 deletions plugins/platforms/simplex/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -440,6 +440,25 @@ async def _handle_event(self, event: dict) -> None:
logger.exception("SimpleX: error processing chat item")
return

# Edited messages — simplex-chat sends "chatItemUpdated" with the
# full updated chatItem wrapper. The itemId is stable across edits,
# so we use it as message_id and mark metadata["is_edit"] so
# consumers can react differently (e.g. re-summarize instead of
# append). Outgoing-direction guard inside _handle_chat_item
# naturally drops edits of the bot's own messages.
# NOTE: edited media items carry only their caption text on the edit
# event — file/image payloads are not reconstructed on edits
# (intentional; matches newChatItems behavior for caption-only
# edits).
if resp_type == "chatItemUpdated":
try:
await self._handle_chat_item(
resp.get("chatItem", {}), is_edit=True
)
except Exception:
logger.exception("SimpleX: error processing chat item update")
return

# File transfer completion — deliver any deferred chat item
if resp_type == "rcvFileComplete":
chat_item = resp.get("chatItem", {}) or {}
Expand Down Expand Up @@ -471,8 +490,8 @@ async def _handle_event(self, event: dict) -> None:
if resp_type:
logger.debug("SimpleX: unhandled event type: %s", resp_type)

async def _handle_chat_item(self, chat_item: dict) -> None:
"""Process a single chat item from a newChatItems event."""
async def _handle_chat_item(self, chat_item: dict, is_edit: bool = False) -> None:
"""Process a single chat item from a newChatItems or chatItemUpdated event."""
chat_info = chat_item.get("chatInfo", {}) or {}
chat_item_data = chat_item.get("chatItem", {}) or {}

Expand Down Expand Up @@ -643,7 +662,11 @@ async def _handle_chat_item(self, chat_item: dict) -> None:
except (ValueError, AttributeError):
timestamp = datetime.now(tz=timezone.utc)

msg_event = MessageEvent(
# Extract the stable item ID for edit tracking
item_id = meta.get("itemId")
msg_id = str(item_id) if item_id is not None else None

msg_kwargs: Dict[str, Any] = dict(
source=source,
text=text or "",
message_type=msg_type,
Expand All @@ -652,6 +675,12 @@ async def _handle_chat_item(self, chat_item: dict) -> None:
timestamp=timestamp,
raw_message=chat_item,
)
if msg_id is not None:
msg_kwargs["message_id"] = msg_id
if is_edit:
msg_kwargs["metadata"] = {"is_edit": True}

msg_event = MessageEvent(**msg_kwargs)

logger.debug(
"SimpleX: message from %s in %s: %s",
Expand Down
Loading
Loading