fix(tui): bound subscriber delivery without blocking shared turns
This commit is contained in:
@@ -5300,7 +5300,7 @@ def test_close_transport_rebinds_session_to_remaining_viewer(monkeypatch):
|
||||
The rebind #83716 added is gone; multi-client fan-out subsumes it. Both
|
||||
windows are attached to the slot at once, so the pop-out is a fan-out peer
|
||||
rather than a viewer waiting to be promoted, and closing it detaches that
|
||||
peer and collapses the slot back onto the main window. This pins the same
|
||||
peer while retaining the surviving ordered mailbox. This pins the same
|
||||
guarantee through the mechanism that replaced the rebind: the session is
|
||||
not parked, not reaped, not handed to the orphan reaper, and the surviving
|
||||
window keeps receiving frames.
|
||||
@@ -5311,9 +5311,11 @@ def test_close_transport_rebinds_session_to_remaining_viewer(monkeypatch):
|
||||
class _LiveTransport:
|
||||
def __init__(self):
|
||||
self.frames = []
|
||||
self.received = threading.Event()
|
||||
|
||||
def write(self, obj=None, *a, **k):
|
||||
self.frames.append(obj)
|
||||
self.received.set()
|
||||
return True
|
||||
|
||||
main = _LiveTransport()
|
||||
@@ -5332,14 +5334,14 @@ def test_close_transport_rebinds_session_to_remaining_viewer(monkeypatch):
|
||||
reaped, detached = server._close_sessions_for_transport(popout)
|
||||
|
||||
assert reaped == 0 and detached == 0
|
||||
# One peer left, so the fan-out collapses back to the bare transport —
|
||||
# the slot is indistinguishable from a session that never fanned out.
|
||||
assert session["transport"] is main
|
||||
assert server._session_transport_contains(session, main)
|
||||
assert not server._session_transport_contains(session, popout)
|
||||
assert "multi-sid" not in reap_calls
|
||||
assert server._ws_session_is_orphaned(session) is False
|
||||
|
||||
# And it is still a working stream, not just a surviving reference.
|
||||
server._emit("message.delta", "multi-sid", {"text": "still here"})
|
||||
assert main.received.wait(timeout=5)
|
||||
assert [(f.get("params") or {}).get("type") for f in main.frames] == [
|
||||
"message.delta"
|
||||
]
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -7,10 +7,12 @@ from tui_gateway import server
|
||||
class Peer:
|
||||
def __init__(self):
|
||||
self.frames = []
|
||||
self.received = threading.Event()
|
||||
self._closed = False
|
||||
|
||||
def write(self, frame):
|
||||
self.frames.append(frame)
|
||||
self.received.set()
|
||||
return not self._closed
|
||||
|
||||
def close(self):
|
||||
@@ -24,10 +26,14 @@ def test_reattach_preserves_terminal_delivery(monkeypatch):
|
||||
with session["history_lock"]:
|
||||
server._rebind_live_transport("shared", session, second)
|
||||
server._emit("message.complete", "shared", {"text": "finished"})
|
||||
assert first.received.wait(timeout=5)
|
||||
assert second.received.wait(timeout=5)
|
||||
assert first.frames == second.frames
|
||||
assert len(first.frames) == 1
|
||||
second.close()
|
||||
assert server._close_sessions_for_transport(second) == (0, 0)
|
||||
assert second not in session.get("viewers", {})
|
||||
first.received.clear()
|
||||
server._emit("message.complete", "shared", {"text": "still attached"})
|
||||
assert first.received.wait(timeout=5)
|
||||
assert len(first.frames) == 2
|
||||
|
||||
@@ -4,146 +4,71 @@ from __future__ import annotations
|
||||
import threading
|
||||
from tui_gateway.method_ctx import bind_module
|
||||
|
||||
# ── multi-client fan-out ─────────────────────────────────────────────────────
|
||||
#
|
||||
# A session's ``transport`` slot holds either ONE transport — the historical,
|
||||
# single-client shape, byte-identical to pre-fan-out behaviour — or a
|
||||
# ``FanoutTransport`` wrapping several. Because the fan-out satisfies the same
|
||||
# Transport protocol, ``write_json`` and every other reader of the slot are
|
||||
# untouched; only the attach/detach ladder below knows the difference.
|
||||
#
|
||||
# ``_session_transport_lock`` serializes the read-modify-write of that one slot.
|
||||
# It is a LEAF lock: nothing under it acquires ``_sessions_lock`` or a session's
|
||||
# ``history_lock`` (the fan-out's own lock is likewise a leaf), so it is safe to
|
||||
# take while holding either of those.
|
||||
_session_transport_lock = threading.Lock()
|
||||
# Leaf lock: callers may hold sessions/history locks, never acquire them here.
|
||||
_session_transport_lock = threading.RLock()
|
||||
|
||||
|
||||
def _transport_is_live_peer(transport) -> bool:
|
||||
"""True when *transport* is a real, currently attached client.
|
||||
|
||||
Excluded: the parked drop sentinel (by definition clientless) and stdio,
|
||||
which is the process-wide fallback sink — a standalone ``hermes --tui``
|
||||
writes there, but nothing "attaches" to it and a second client can never
|
||||
share it.
|
||||
"""
|
||||
if transport is None:
|
||||
return False
|
||||
if transport is _detached_ws_transport or transport is _stdio_transport:
|
||||
return False
|
||||
if isinstance(transport, (_DropTransport, StdioTransport)):
|
||||
return False
|
||||
# A socket that already latched ``_closed`` is a departed client, not a live
|
||||
# peer: without this, a stale WSTransport keeps a session out of the park /
|
||||
# reap path and keeps answering "another client is still here".
|
||||
# ``_transport_is_dead`` is the deadness predicate this module shares with
|
||||
# the reaper (defined in ``session_reaper``, published onto this namespace
|
||||
# by ``method_ctx.bind_module``; module globals resolve at call time).
|
||||
return not _transport_is_dead(transport)
|
||||
"""Exclude the process fallback sink, parked sentinel, and closed peers."""
|
||||
return (transport is not None
|
||||
and transport is not _detached_ws_transport
|
||||
and transport is not _stdio_transport
|
||||
and not isinstance(transport, (_DropTransport, StdioTransport))
|
||||
and not _transport_is_dead(transport))
|
||||
|
||||
|
||||
def _session_transport_contains(session: dict | None, transport) -> bool:
|
||||
"""True when *transport* is attached to *session*, directly or via fan-out."""
|
||||
if not session or transport is None:
|
||||
if not session or transport is None or _transport_is_dead(transport):
|
||||
return False
|
||||
existing = session.get("transport")
|
||||
if existing is transport:
|
||||
return True
|
||||
if isinstance(existing, FanoutTransport):
|
||||
return existing.contains(transport)
|
||||
return False
|
||||
return existing is transport or (
|
||||
isinstance(existing, FanoutTransport) and existing.contains(transport))
|
||||
|
||||
|
||||
def _session_live_transports(session: dict | None) -> list:
|
||||
"""Every live client attached to *session* (empty for a parked/stdio slot)."""
|
||||
existing = (session or {}).get("transport")
|
||||
if isinstance(existing, FanoutTransport):
|
||||
return [t for t in existing.transports() if _transport_is_live_peer(t)]
|
||||
return [existing] if _transport_is_live_peer(existing) else []
|
||||
peers = existing.transports() if isinstance(existing, FanoutTransport) else [existing]
|
||||
return [peer for peer in peers if _transport_is_live_peer(peer)]
|
||||
|
||||
|
||||
def _session_has_live_transport(session: dict | None, *, excluding=None) -> bool:
|
||||
"""True when *session* still has a live client attached, ignoring *excluding*.
|
||||
|
||||
``excluding`` answers whether another live client remains once one peer is
|
||||
ignored, which is what the detach and orphan paths need: whether anything
|
||||
survives the departing transport.
|
||||
"""
|
||||
return any(t is not excluding for t in _session_live_transports(session))
|
||||
return any(peer is not excluding for peer in _session_live_transports(session))
|
||||
|
||||
|
||||
def _attach_session_transport(session: dict | None, transport) -> bool:
|
||||
"""Attach *transport* to *session* ADDITIVELY — never steal the slot.
|
||||
|
||||
A ``FanoutTransport`` argument is FLATTENED before the ladder runs, whatever
|
||||
the slot holds: each of its peers is attached as a leaf, so the slot never
|
||||
nests one fan-out inside another. The queued-prompt paths make that
|
||||
reachable — a busy submit records ``session["transport"]`` as the prompt's
|
||||
transport, which is the fan-out itself when two clients are attached, and
|
||||
the drain hands it back here after a disconnect has collapsed the slot to a
|
||||
single client. Nesting would blind every single-level reader of the slot
|
||||
(``FanoutTransport.contains``, the steer-authority scan, detach) to the
|
||||
inner peers, so a client inside the inner fan-out could never be found,
|
||||
parked, reaped, or granted authority over its own turn.
|
||||
|
||||
The ladder, for the leaf newcomer that always reaches it:
|
||||
|
||||
* same object already in the slot → no-op;
|
||||
* the slot already fans out → attach into it;
|
||||
* the slot is empty / stdio / the parked drop sentinel → take the slot,
|
||||
which is exactly what the old rebind did;
|
||||
* otherwise → wrap the incumbent and the newcomer in a ``FanoutTransport``
|
||||
so both clients keep streaming.
|
||||
|
||||
A non-peer newcomer (stdio, drop sentinel) never displaces a live client:
|
||||
an activate/resume dispatched without a bound websocket must not silence the
|
||||
socket that owns the session.
|
||||
|
||||
Returns ``True`` when *transport* is attached afterwards.
|
||||
"""
|
||||
"""Add live peers; flatten captured queued fanouts without nesting authority."""
|
||||
if not session or transport is None:
|
||||
return False
|
||||
if isinstance(transport, FanoutTransport):
|
||||
# Flatten OUTSIDE the lock — _session_transport_lock is not reentrant.
|
||||
# Each peer then walks the ladder on its own, so a peer whose socket
|
||||
# died while the prompt sat in the queue is dropped by the liveness rung
|
||||
# instead of being pinned back into the slot.
|
||||
attached = False
|
||||
for peer in transport.transports():
|
||||
if _attach_session_transport(session, peer):
|
||||
attached = True
|
||||
return attached
|
||||
with _session_transport_lock:
|
||||
if isinstance(transport, FanoutTransport):
|
||||
# Snapshot and attach share detach's lock: a queued fanout cannot
|
||||
# resurrect a still-open peer removed during flattening.
|
||||
attached = [_attach_session_transport(session, peer) for peer in transport.transports()]
|
||||
return any(attached)
|
||||
existing = session.get("transport")
|
||||
if existing is transport:
|
||||
return True
|
||||
if isinstance(existing, FanoutTransport):
|
||||
if isinstance(transport, FanoutTransport):
|
||||
for member in transport.transports():
|
||||
existing.attach(member)
|
||||
else:
|
||||
existing.attach(transport)
|
||||
return True
|
||||
if _transport_is_dead(transport):
|
||||
if isinstance(existing, FanoutTransport):
|
||||
existing.detach(transport)
|
||||
return False
|
||||
if not _transport_is_live_peer(transport):
|
||||
if _transport_is_live_peer(existing):
|
||||
if _session_has_live_transport(session):
|
||||
return False
|
||||
session["transport"] = transport
|
||||
return True
|
||||
if not _transport_is_live_peer(existing):
|
||||
session["transport"] = transport
|
||||
if existing is transport:
|
||||
return True
|
||||
session["transport"] = FanoutTransport(existing, transport)
|
||||
if isinstance(existing, FanoutTransport):
|
||||
existing.attach(transport)
|
||||
return existing.contains(transport)
|
||||
elif _transport_is_live_peer(existing):
|
||||
session["transport"] = FanoutTransport(existing, transport)
|
||||
else:
|
||||
session["transport"] = transport
|
||||
return True
|
||||
|
||||
|
||||
def _detach_session_transport(session: dict | None, transport) -> bool:
|
||||
"""Detach *transport* from *session*'s slot.
|
||||
|
||||
Returns ``True`` when a live client OTHER than *transport* remains — i.e.
|
||||
the session must keep streaming and must NOT be parked or reaped. A
|
||||
single-client session leaves the departing transport in the slot, exactly as
|
||||
before fan-out existed; its caller parks the drop sentinel over it.
|
||||
"""
|
||||
"""Remove membership; return whether another live client prevents parking."""
|
||||
if not session:
|
||||
return False
|
||||
with _session_transport_lock:
|
||||
@@ -151,33 +76,27 @@ def _detach_session_transport(session: dict | None, transport) -> bool:
|
||||
existing = session.get("transport")
|
||||
if isinstance(existing, FanoutTransport):
|
||||
existing.detach(transport)
|
||||
remaining = existing.transports()
|
||||
if len(remaining) == 1:
|
||||
# Collapse: a session back down to one client is indistinguishable
|
||||
# from one that never fanned out.
|
||||
session["transport"] = remaining[0]
|
||||
return _session_has_live_transport(session, excluding=transport)
|
||||
viewers = session.get("viewers") or {}
|
||||
for viewer in list(viewers):
|
||||
if not existing.contains(viewer) or _transport_is_dead(viewer):
|
||||
viewers.pop(viewer, None)
|
||||
# Keep the surviving mailbox: collapsing to a bare transport lets
|
||||
# new frames overtake its already queued terminal/control events.
|
||||
return _session_has_live_transport(session, excluding=transport)
|
||||
|
||||
|
||||
def _detach_transport_from_sessions(transport) -> list[tuple[str, dict]]:
|
||||
"""Detach *transport* from every session holding it.
|
||||
|
||||
Returns the ``(sid, session)`` pairs left with NO live client — the ones the
|
||||
disconnect path must park or reap. Sessions that retain another attached
|
||||
client keep streaming and are not returned.
|
||||
"""
|
||||
"""Remove even closed/pruned peers' viewer entries; return clientless slots."""
|
||||
with _sessions_lock:
|
||||
attached = [
|
||||
(sid, s)
|
||||
for sid, s in _sessions.items()
|
||||
if _session_transport_contains(s, transport)
|
||||
]
|
||||
return [
|
||||
(sid, session)
|
||||
for sid, session in attached
|
||||
if not _detach_session_transport(session, transport)
|
||||
]
|
||||
|
||||
attached = []
|
||||
for sid, session in _sessions.items():
|
||||
existing = session.get("transport")
|
||||
if (existing is transport
|
||||
or isinstance(existing, FanoutTransport) and existing.contains(transport)
|
||||
or transport in (session.get("viewers") or {})):
|
||||
attached.append((sid, session))
|
||||
return [(sid, session) for sid, session in attached
|
||||
if not _detach_session_transport(session, transport)]
|
||||
|
||||
|
||||
def register(server) -> None:
|
||||
|
||||
+105
-232
@@ -9,8 +9,8 @@ A :class:`Transport` forwards a JSON-serialisable dict to its peer, so one dispa
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import concurrent.futures
|
||||
from collections import deque
|
||||
from dataclasses import dataclass, field
|
||||
import contextlib
|
||||
import contextvars
|
||||
import errno
|
||||
@@ -18,7 +18,6 @@ import json
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
import time
|
||||
from typing import Any, Callable, Optional, Protocol, runtime_checkable
|
||||
|
||||
# Errno values that mean "the peer is gone" rather than "the host has a real I/O problem". Anything
|
||||
@@ -35,65 +34,6 @@ logger = logging.getLogger(__name__)
|
||||
# fully buffered on a pipe, so this ONLY makes sense with ``-u``/``PYTHONUNBUFFERED=1``; otherwise the TUI hangs.
|
||||
_DISABLE_FLUSH = (os.environ.get("HERMES_TUI_GATEWAY_NO_FLUSH", "") or "").strip().lower() in {"1", "true", "yes", "on"}
|
||||
|
||||
# Worker pool behind FanoutTransport.write. A non-streaming WebSocket write
|
||||
# blocks the calling thread while the owning event loop flushes the frame — up
|
||||
# to ``tui_gateway.ws._WS_WRITE_TIMEOUT_S`` (10s) when that loop is stalled —
|
||||
# so walking a session's peers one at a time lets one wedged client hold the
|
||||
# frame back from every healthy client behind it. The pool is created lazily
|
||||
# and only ever used when a session has more than one peer, so the ordinary
|
||||
# single-client gateway never starts a thread.
|
||||
#
|
||||
# Sizing: these workers are idle except while a peer is actually wedged, so the
|
||||
# cap only has to cover the wedged peers of all sessions at once. If it were
|
||||
# ever saturated the excess writes queue and run in submission order, which is
|
||||
# the pre-pool behaviour — degraded, not broken.
|
||||
_FANOUT_POOL_MAX_WORKERS = 32
|
||||
|
||||
# Wall-clock bound on ONE fan-out write, measured from before the first peer is
|
||||
# dispatched. Deliberately larger than ``tui_gateway.ws._WS_WRITE_TIMEOUT_S``
|
||||
# (10.0) so a merely-slow peer reaches its own timeout and reports for itself;
|
||||
# this deadline only exists so a peer that never returns at all cannot pin the
|
||||
# emitting thread forever. Not imported from ``tui_gateway.ws``: that module
|
||||
# imports ``tui_gateway.server``, which imports this one, so the import would
|
||||
# be circular. Drift is harmless — a peer that misses this deadline is treated
|
||||
# as in-flight, not as dead (see ``FanoutTransport.write``).
|
||||
_FANOUT_WRITE_DEADLINE_S = 12.0
|
||||
|
||||
_fanout_pool: "concurrent.futures.ThreadPoolExecutor | None" = None
|
||||
_fanout_pool_lock = threading.Lock()
|
||||
|
||||
|
||||
def _fanout_write_pool() -> "concurrent.futures.ThreadPoolExecutor":
|
||||
"""The shared fan-out pool, created on first multi-peer write."""
|
||||
global _fanout_pool
|
||||
pool = _fanout_pool
|
||||
if pool is not None:
|
||||
return pool
|
||||
with _fanout_pool_lock:
|
||||
if _fanout_pool is None:
|
||||
_fanout_pool = concurrent.futures.ThreadPoolExecutor(
|
||||
max_workers=_FANOUT_POOL_MAX_WORKERS,
|
||||
thread_name_prefix="tui-fanout",
|
||||
)
|
||||
return _fanout_pool
|
||||
|
||||
|
||||
def _caller_is_on_event_loop() -> bool:
|
||||
"""True when the CALLING thread is running an asyncio event loop.
|
||||
|
||||
Mirrors the probe in ``tui_gateway.ws.WSTransport.write``, which takes a
|
||||
fire-and-forget path when it can see it is on its own loop and a BLOCKING
|
||||
``fut.result`` path when it cannot. ``FanoutTransport`` owns no loop and
|
||||
its peers may sit on different ones, so the test here is the conservative
|
||||
one: if this thread is running any loop at all, keep the peer writes on it.
|
||||
"""
|
||||
try:
|
||||
asyncio.get_running_loop()
|
||||
except RuntimeError:
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
class Transport(Protocol):
|
||||
"""Minimal interface every transport implements."""
|
||||
@@ -176,206 +116,139 @@ class StdioTransport:
|
||||
return None
|
||||
|
||||
|
||||
@dataclass(eq=False)
|
||||
class _FanoutPeer:
|
||||
transport: Transport
|
||||
pending: deque = field(default_factory=deque)
|
||||
pending_bytes: int = 0
|
||||
writing: bool = False
|
||||
attached: bool = True
|
||||
generation: int = 0
|
||||
|
||||
|
||||
class FanoutTransport:
|
||||
"""Delivers one JSON frame to every client attached to a session.
|
||||
"""Ordered, bounded session-event mailboxes; RPC replies remain request-local.
|
||||
|
||||
A session's ``transport`` slot used to hold exactly one client, so every
|
||||
``prompt.submit`` / ``session.resume`` / ``session.activate`` / queued-prompt
|
||||
drain REBOUND it and whichever client was there before went silent. This
|
||||
object goes in the same slot and satisfies the same :class:`Transport`
|
||||
protocol, so ``server.write_json`` — and every other reader of that slot —
|
||||
is unchanged; the only difference is that N clients receive the frame
|
||||
instead of the most recent one.
|
||||
|
||||
Scope: this carries a session's ASYNC EVENT stream only. Request/response
|
||||
RPCs still answer on the request's context-bound transport
|
||||
(``dispatch`` → ``current_transport()``), so a client only ever sees replies
|
||||
to its own calls.
|
||||
|
||||
Peers are pruned on failure: a transport that returns ``False`` (peer gone)
|
||||
or raises is dropped from the fan-out instead of being written to forever.
|
||||
Pruning is about dead peers, not slow ones — a client that is merely wedged
|
||||
is kept, and :meth:`write` runs the peer writes concurrently so its stall
|
||||
stays its own.
|
||||
One slow socket must not stop the emitting turn or any healthy subscriber.
|
||||
Each peer has at most one daemon writer and a bounded backlog. On overflow
|
||||
it loses its subscription (history/replay is the recovery path), not other
|
||||
sessions sharing its socket. A write already in the OS cannot be revoked.
|
||||
"""
|
||||
|
||||
__slots__ = ("_lock", "_transports")
|
||||
_MAX_PENDING_FRAMES = 256
|
||||
_MAX_PENDING_BYTES = 4 * 1024 * 1024
|
||||
|
||||
def __init__(self, *transports: "Transport") -> None:
|
||||
def __init__(self, *transports: Transport) -> None:
|
||||
self._lock = threading.Lock()
|
||||
self._transports: list["Transport"] = []
|
||||
self._peers: list[_FanoutPeer] = []
|
||||
for transport in transports:
|
||||
self.attach(transport)
|
||||
|
||||
def attach(self, transport: "Transport") -> bool:
|
||||
"""Add *transport*. Returns ``True`` when it was not already attached."""
|
||||
def attach(self, transport: Transport) -> bool:
|
||||
if transport is None or transport is self:
|
||||
return False
|
||||
with self._lock:
|
||||
for existing in self._transports:
|
||||
if existing is transport:
|
||||
return False
|
||||
self._transports.append(transport)
|
||||
for peer in self._peers:
|
||||
if peer.transport is transport:
|
||||
if peer.attached:
|
||||
return False
|
||||
# Reuse the in-flight writer: reconnect cannot spawn more
|
||||
# threads or overtake a write already inside this socket.
|
||||
peer.attached = True
|
||||
peer.generation += 1
|
||||
return True
|
||||
self._peers.append(_FanoutPeer(transport))
|
||||
return True
|
||||
|
||||
def detach(self, transport: "Transport") -> bool:
|
||||
"""Remove *transport*. Returns ``True`` when it was attached."""
|
||||
def _remove(self, peer: _FanoutPeer) -> None:
|
||||
# Membership lock held; identity fences a stale writer from removing
|
||||
# a later attachment of the same transport.
|
||||
peer.attached = False
|
||||
peer.pending.clear()
|
||||
peer.pending_bytes = 0
|
||||
if not peer.writing and peer in self._peers:
|
||||
self._peers.remove(peer)
|
||||
|
||||
def detach(self, transport: Transport) -> bool:
|
||||
with self._lock:
|
||||
for idx, existing in enumerate(self._transports):
|
||||
if existing is transport:
|
||||
del self._transports[idx]
|
||||
for peer in self._peers:
|
||||
if peer.attached and peer.transport is transport:
|
||||
self._remove(peer)
|
||||
return True
|
||||
return False
|
||||
|
||||
def contains(self, transport: "Transport") -> bool:
|
||||
def contains(self, transport: Transport) -> bool:
|
||||
with self._lock:
|
||||
return any(existing is transport for existing in self._transports)
|
||||
return any(peer.attached and peer.transport is transport for peer in self._peers)
|
||||
|
||||
def transports(self) -> list["Transport"]:
|
||||
"""A snapshot of the attached transports (safe to iterate)."""
|
||||
def transports(self) -> list[Transport]:
|
||||
with self._lock:
|
||||
return list(self._transports)
|
||||
return [peer.transport for peer in self._peers if peer.attached]
|
||||
|
||||
def has_transports(self, *, excluding: "Transport | None" = None) -> bool:
|
||||
"""True when any transport other than *excluding* is still attached."""
|
||||
with self._lock:
|
||||
return any(existing is not excluding for existing in self._transports)
|
||||
def has_transports(self, *, excluding: Transport | None = None) -> bool:
|
||||
return any(peer is not excluding for peer in self.transports())
|
||||
|
||||
@staticmethod
|
||||
def _deliver(transport: "Transport", obj: dict, dead: list) -> bool:
|
||||
"""Write *obj* to one peer inline. Records a dead peer in *dead*.
|
||||
|
||||
Returns ``True`` when the peer took the frame.
|
||||
"""
|
||||
try:
|
||||
ok = transport.write(obj)
|
||||
except Exception:
|
||||
logger.debug("fanout write failed; pruning peer", exc_info=True)
|
||||
dead.append(transport)
|
||||
return False
|
||||
if ok:
|
||||
return True
|
||||
dead.append(transport)
|
||||
return False
|
||||
|
||||
def write(self, obj: dict) -> bool:
|
||||
"""Deliver *obj* to every attached transport.
|
||||
|
||||
Iterates a SNAPSHOT and never holds the lock across a peer write, so a
|
||||
peer write can never block an attach or a detach.
|
||||
|
||||
With more than one peer the writes run CONCURRENTLY: every peer but one
|
||||
is handed to a shared worker pool and the last runs on the calling
|
||||
thread, which is going to block here anyway. Serial delivery would let
|
||||
a client wedged for the full WS write timeout hold the frame back from
|
||||
every healthy client behind it in the list; this way each peer's stall
|
||||
is its own. Two cases stay inline:
|
||||
|
||||
* exactly one peer, so a single-client session is untouched;
|
||||
* a caller already running on an event loop. ``WSTransport.write``
|
||||
fires and forgets when it sees its own loop and BLOCKS on
|
||||
``fut.result`` when it does not, so moving a loop-thread write onto a
|
||||
pool thread would make that thread wait on the loop it just left
|
||||
while the loop waits on it — a full write-timeout freeze.
|
||||
|
||||
Pruning is unchanged: a peer that returns ``False`` or raises is
|
||||
detached, whichever path it took. A peer that has not answered by
|
||||
``_FANOUT_WRITE_DEADLINE_S`` is NOT pruned — slow is not dead — and
|
||||
counts as delivered, which is what ``WSTransport.write`` itself reports
|
||||
when its own write times out: the frame is queued on that peer's loop
|
||||
and flushes when the loop breathes. If that peer really is gone, its
|
||||
next write returns ``False`` promptly and prunes it then.
|
||||
|
||||
The collect happens before the return, so frame A is on every peer that
|
||||
answered within the deadline before frame B is dispatched to any of
|
||||
them. A peer whose write was still queued on the pool when the deadline
|
||||
passed is the exception: the next frame can reach that peer's own
|
||||
``_token_lock`` while the previous one is still waiting for it, and the
|
||||
lock — not this method — decides the order they land in. Within a peer,
|
||||
ordering is that peer's own business (``WSTransport`` queues under its
|
||||
token lock, ``StdioTransport`` writes under its stream lock);
|
||||
``server.write_json`` takes no lock, so concurrent emitters interleave
|
||||
here exactly as they already did.
|
||||
|
||||
Returns ``True`` when at least one client accepted the frame or still
|
||||
has it in flight — an all-dead fan-out reports peer-gone exactly like a
|
||||
single dead transport would.
|
||||
"""
|
||||
targets = self.transports()
|
||||
if not targets:
|
||||
return False
|
||||
|
||||
dead: list["Transport"] = []
|
||||
delivered = False
|
||||
|
||||
if len(targets) == 1 or _caller_is_on_event_loop():
|
||||
for transport in targets:
|
||||
if self._deliver(transport, obj, dead):
|
||||
delivered = True
|
||||
for transport in dead:
|
||||
self.detach(transport)
|
||||
return delivered
|
||||
|
||||
deadline = time.monotonic() + _FANOUT_WRITE_DEADLINE_S
|
||||
pool = _fanout_write_pool()
|
||||
dispatched: list[tuple["Transport", "concurrent.futures.Future | None"]] = []
|
||||
for transport in targets[:-1]:
|
||||
def _drain(self, peer: _FanoutPeer) -> None:
|
||||
while True:
|
||||
with self._lock:
|
||||
if not peer.attached or not peer.pending:
|
||||
peer.writing = False
|
||||
if not peer.attached:
|
||||
self._remove(peer)
|
||||
return
|
||||
generation = peer.generation
|
||||
frame, size = peer.pending.popleft()
|
||||
peer.pending_bytes -= size
|
||||
try:
|
||||
dispatched.append((transport, pool.submit(transport.write, obj)))
|
||||
except RuntimeError:
|
||||
# Pool refused the work (interpreter shutting down). Deliver
|
||||
# inline below rather than dropping the frame.
|
||||
dispatched.append((transport, None))
|
||||
|
||||
# The last peer runs here: one fewer worker per write, and a two-client
|
||||
# session costs a single pool thread.
|
||||
if self._deliver(targets[-1], obj, dead):
|
||||
delivered = True
|
||||
|
||||
# One deadline for the whole batch, then read only the futures that
|
||||
# finished. Waiting first and calling ``result()`` on a settled future
|
||||
# keeps our deadline distinguishable from a peer that raised: on 3.11+
|
||||
# ``concurrent.futures.TimeoutError`` IS the builtin ``TimeoutError``,
|
||||
# so a ``result(timeout=...)`` cannot tell the two apart.
|
||||
futures = [fut for _, fut in dispatched if fut is not None]
|
||||
if futures:
|
||||
concurrent.futures.wait(
|
||||
futures, timeout=max(0.0, deadline - time.monotonic())
|
||||
)
|
||||
|
||||
for transport, fut in dispatched:
|
||||
if fut is None:
|
||||
if self._deliver(transport, obj, dead):
|
||||
delivered = True
|
||||
continue
|
||||
if not fut.done():
|
||||
# Slow, not dead: leave it attached and count the frame as in
|
||||
# flight. Cancelling would not stop a write already running.
|
||||
logger.warning(
|
||||
"fanout write still pending after %ss; peer left attached",
|
||||
_FANOUT_WRITE_DEADLINE_S,
|
||||
)
|
||||
delivered = True
|
||||
continue
|
||||
try:
|
||||
ok = fut.result()
|
||||
from tui_gateway.ws import WSTransport
|
||||
if isinstance(peer.transport, WSTransport):
|
||||
# write() acknowledges buffered tokens/timeouts, not socket
|
||||
# progress. Await the real send so WS cannot move an
|
||||
# unbounded backlog underneath this bounded mailbox.
|
||||
from agent.async_utils import safe_schedule_threadsafe
|
||||
future = safe_schedule_threadsafe(
|
||||
peer.transport.write_async(frame), peer.transport._loop)
|
||||
ok = future is not None and future.result()
|
||||
else:
|
||||
ok = peer.transport.write(frame)
|
||||
except Exception:
|
||||
logger.debug("fanout write failed; pruning peer", exc_info=True)
|
||||
dead.append(transport)
|
||||
continue
|
||||
if ok:
|
||||
delivered = True
|
||||
else:
|
||||
dead.append(transport)
|
||||
ok = False
|
||||
if not ok:
|
||||
with self._lock:
|
||||
if peer.generation != generation:
|
||||
continue
|
||||
peer.writing = False
|
||||
self._remove(peer)
|
||||
return
|
||||
|
||||
for transport in dead:
|
||||
self.detach(transport)
|
||||
return delivered
|
||||
def write(self, obj: dict) -> bool:
|
||||
# Freeze the queued frame so a caller cannot mutate it after admission.
|
||||
encoded = json.dumps(obj, ensure_ascii=False)
|
||||
size = len(encoded.encode("utf-8", errors="surrogatepass"))
|
||||
frame = json.loads(encoded)
|
||||
with self._lock:
|
||||
for peer in list(self._peers):
|
||||
if not peer.attached:
|
||||
continue
|
||||
if (len(peer.pending) >= self._MAX_PENDING_FRAMES
|
||||
or peer.pending_bytes + size > self._MAX_PENDING_BYTES):
|
||||
logger.warning("fanout subscriber backlog full; detaching peer")
|
||||
self._remove(peer)
|
||||
continue
|
||||
peer.pending.append((frame, size))
|
||||
peer.pending_bytes += size
|
||||
if not peer.writing:
|
||||
peer.writing = True
|
||||
threading.Thread(target=self._drain, args=(peer,),
|
||||
name="tui-fanout", daemon=True).start()
|
||||
return any(peer.attached for peer in self._peers)
|
||||
|
||||
def close(self) -> None:
|
||||
"""Detach every peer. Does NOT close them — each owns its own socket."""
|
||||
"""Detach without closing sockets owned by the connection handlers."""
|
||||
with self._lock:
|
||||
self._transports = []
|
||||
for peer in list(self._peers):
|
||||
self._remove(peer)
|
||||
|
||||
|
||||
class TeeTransport:
|
||||
|
||||
Reference in New Issue
Block a user