Revert "Merge pull request #94245 from kshitijk4poor/feat/gw-event-replay"

This reverts commit df7d7f6e8d, reversing
changes made to 1a66134404.
This commit is contained in:
kshitijk4poor
2026-08-27 11:26:57 +05:30
parent df7d7f6e8d
commit 9f05b06589
16 changed files with 383 additions and 1062 deletions
+47 -104
View File
@@ -9,12 +9,11 @@ replays everything newer from the buffer, then live events resume seamlessly.
Design constraints honored:
- stdio TUI path unaffected: frames gain a ``seq`` field only on event frames;
Ink ignores unknown params keys.
- Thread safety: a single module lock guards counters + buffers, so buffer
order always matches seq order. Wire order is enforced separately by the
per-transport write path; two racing writers can briefly invert seq order
on the wire, which the client tolerates (watermarks are monotonic-max).
- Thread safety: a single module lock guards counters + buffers; write_json
already serializes per-transport writes, so stamping under the lock cannot
reorder frames relative to each other.
- Memory bound: _REPLAY_BUFFER_MAX events / _REPLAY_SESSIONS_MAX sessions,
least-recently-active session evicted first.
oldest session evicted FIFO.
"""
from __future__ import annotations
@@ -23,47 +22,34 @@ import threading
import uuid
from collections import OrderedDict, deque
# Replay ring per session. Sized for control events only — streaming token
# deltas are transient (stamped but not buffered, see _TRANSIENT_EVENT_TYPES),
# so 512 slots cover hours of durable events instead of ~one streaming burst.
# Process identity for the replay contract. Seq counters live in-process, so
# a gateway restart silently resets them to 1 while clients still hold high
# watermarks — events_since(sid, 97) then returns [] with truncated=False and
# the client believes it missed nothing (and its stale watermark makes every
# future replay empty too). The epoch lets clients detect the restart and
# reset their watermarks.
_REPLAY_EPOCH = uuid.uuid4().hex
# Replay ring per session. A long turn emits ~hundreds of token events; this
# covers several minutes of streaming plus all control events.
_REPLAY_BUFFER_MAX = 512
# Distinct sessions remembered. Desktop users rarely exceed a dozen live chats.
_REPLAY_SESSIONS_MAX = 64
# Transient event types: delivered live but never buffered for replay.
# One streaming turn emits hundreds of per-token delta frames; buffering them
# would evict every durable control event from the ring (OpenHands makes the
# same split with StreamingDeltaEvent — published, never persisted). A
# reconnecting client recovers partial streamed text from the inflight
# snapshot in ``session.resume``, not from delta replay. These frames still
# get seqs (so ordering holds live), but replay skips them and gap detection
# ignores them via the durable-seq watermark below.
_TRANSIENT_EVENT_TYPES = frozenset({
"message.delta",
"thinking.delta",
})
# Server boot epoch: lets a client detect that the seq namespace was reset
# (gateway restart) — a stale high watermark from the previous process must
# not suppress live events forever (Goose clamps; we reset via epoch).
EPOCH = uuid.uuid4().hex[:8]
_replay_lock = threading.Lock()
# sid -> deque of (seq, params_dict). params is the same dict written to the
# wire and already carries its stamped "seq" key — the replay RPC returns
# these bare event objects, matching what the client's live dispatch sees.
# Manual eviction (no maxlen): we must record the seq of every DURABLE frame
# we drop, so gap detection can distinguish "durable data lost" from "the
# missing seqs were transient deltas that were never replayable anyway".
# sid -> deque of (seq, event_object) where event_object is the frame's
# ``params`` dict (bare event: type/session_id/seq/payload) — the exact shape
# the client's dispatch path consumes.
_replay_buffers: "OrderedDict[str, deque]" = OrderedDict()
_replay_next_seq: dict[str, int] = {}
# sid -> highest DURABLE seq evicted from the ring (0 = nothing evicted).
# Precise truncation signal: durable data is lost for a client at last_seen
# iff an evicted durable frame had seq > last_seen.
_replay_evicted_seq: dict[str, int] = {}
def stamp_event(obj: dict) -> None:
def replay_epoch() -> str:
"""Opaque token identifying this server process's seq numbering."""
return _REPLAY_EPOCH
def _stamp_event(obj: dict) -> None:
"""Stamp one outgoing event frame (mutates obj in place) and record it."""
if obj.get("method") != "event":
return
@@ -75,59 +61,45 @@ def stamp_event(obj: dict) -> None:
# Session-less global events (skin.changed etc.) are re-fetchable via
# their own RPCs; no replay contract for them.
return
transient = params.get("type") in _TRANSIENT_EVENT_TYPES
with _replay_lock:
seq = _replay_next_seq.get(sid, 0) + 1
_replay_next_seq[sid] = seq
params["seq"] = seq
if transient:
# Live-only: seq stamped for wire ordering, never buffered.
return
buf = _replay_buffers.get(sid)
if buf is None:
buf = deque()
buf = deque(maxlen=_REPLAY_BUFFER_MAX)
_replay_buffers[sid] = buf
while len(_replay_buffers) > _REPLAY_SESSIONS_MAX:
_oldest_sid, _oldest_buf = _replay_buffers.popitem(last=False)
_replay_next_seq.pop(_oldest_sid, None)
_replay_evicted_seq.pop(_oldest_sid, None)
else:
# LRU, not insertion-FIFO: an actively streaming session must not
# be evicted just because it was created before idle newer ones.
_replay_buffers.move_to_end(sid)
buf.append((seq, params))
while len(buf) > _REPLAY_BUFFER_MAX:
dropped_seq, _dropped = buf.popleft()
_replay_evicted_seq[sid] = dropped_seq
def events_since(sid: str, last_seen: int) -> tuple[list[dict], int, bool]:
"""Replay contract for one session, computed atomically.
def events_since(sid: str, last_seen: int) -> list[dict]:
"""Return recorded EVENT OBJECTS with seq > last_seen for *sid*, in order.
Returns ``(events, latest_seq, truncated)``:
- ``events``: bare event params dicts (``type``/``session_id``/``seq``/
``payload``) with ``seq > last_seen``, in seq order. Transient delta
frames are never included (they were never buffered).
- ``latest_seq``: current highest stamped seq (0 when unknown).
- ``truncated``: durable data the client has not seen is unrecoverable —
either a DURABLE frame with ``seq > last_seen`` was evicted from the
ring, or ``last_seen`` is AHEAD of ``latest_seq`` (seq namespace reset
after a gateway restart / session eviction — the client should compare
``EPOCH`` and do a full state reload). Gaps consisting only of
transient delta seqs are NOT truncation: those frames were never
replayable and the resume snapshot covers their content.
Shape contract: each element is the frame's ``params`` dict — a bare event
object with top-level ``type`` / ``session_id`` / ``seq`` — because that is
exactly what the client's dispatch path consumes. Returning the full
JSON-RPC envelope here would make every replayed event fail the client's
``event.type`` gate and be silently dropped.
"""
sid = sid or ""
with _replay_lock:
latest = _replay_next_seq.get(sid, 0)
evicted = _replay_evicted_seq.get(sid, 0)
buf = _replay_buffers.get(sid)
buf = _replay_buffers.get(sid or "")
if not buf:
return [], latest, last_seen > latest or evicted > last_seen
frames = [params for seq, params in buf if seq > last_seen]
truncated = last_seen > latest or evicted > last_seen
return frames, latest, truncated
return []
return [event for seq, event in buf if seq > last_seen]
def is_truncated(sid: str, last_seen: int) -> bool:
"""True when events between *last_seen* and the ring's oldest retained
seq were evicted — the client must refetch history instead of trusting
the replay to be gap-free."""
with _replay_lock:
buf = _replay_buffers.get(sid or "")
if not buf:
return False
return last_seen + 1 < buf[0][0]
def latest_seq(sid: str) -> int:
@@ -141,42 +113,13 @@ def reset_replay_state() -> None:
with _replay_lock:
_replay_buffers.clear()
_replay_next_seq.clear()
_replay_evicted_seq.clear()
def replay_epoch() -> str:
"""This process's seq-namespace epoch (see :data:`EPOCH`)."""
return EPOCH
def replay_stats() -> dict:
"""Telemetry: buffer occupancy + per-turn timing for the ops/debug surface."""
"""Telemetry: buffer occupancy for the ops/debug surface."""
with _replay_lock:
stats = {
return {
"sessions": len(_replay_buffers),
"events": sum(len(b) for b in _replay_buffers.values()),
"max_per_session": _REPLAY_BUFFER_MAX,
"max_sessions": _REPLAY_SESSIONS_MAX,
}
# Per-turn timing from the live session table (not under the replay lock —
# the session table has its own lock).
try:
import time as _time
from tui_gateway.server import _sessions, _sessions_lock
with _sessions_lock:
active_turns = []
for sid, session in _sessions.items():
inflight = session.get("inflight_turn")
if isinstance(inflight, dict) and inflight.get("started_at"):
active_turns.append({
"session_id": sid,
"trace_id": inflight.get("trace_id"),
"elapsed_s": round(_time.time() - float(inflight["started_at"]), 2),
"streaming": inflight.get("streaming", False),
})
stats["active_turns"] = active_turns
except Exception:
stats["active_turns"] = []
return stats