From 8ead679d83a715f34bd4593a2e63c42a1000be3f Mon Sep 17 00:00:00 2001 From: Sankalp Gilda Date: Fri, 21 Aug 2026 10:19:52 -0400 Subject: [PATCH] docs(server): say what the evict-on-full check actually tests The eviction predicate is queue fullness at delivery time. The docstrings and the warning described it as an abandoned subscriber whose consumer stopped draining, which is a claim about the consumer that the check never makes: a consumer that is reading steadily but has fallen max_queue_size events behind satisfies the same predicate and is evicted the same way. Name the queue bound instead, and say explicitly that eviction does not imply the consumer stopped reading, so an operator reading the warning does not go hunting for a dead consumer and find a live one. No behavior change. Signed-off-by: Sankalp Gilda --- src/a2a/server/agent_execution/active_task.py | 10 +++-- src/a2a/server/events/event_queue_v2.py | 43 +++++++++++-------- 2 files changed, 32 insertions(+), 21 deletions(-) diff --git a/src/a2a/server/agent_execution/active_task.py b/src/a2a/server/agent_execution/active_task.py index e8e6e3333..a34ac20dd 100644 --- a/src/a2a/server/agent_execution/active_task.py +++ b/src/a2a/server/agent_execution/active_task.py @@ -633,10 +633,12 @@ async def subscribe( self._reference_count, ) - # Subscriber sinks belong to remote consumers that can be abandoned - # (e.g. a non-blocking send whose HTTP response already returned). - # evict_on_full keeps one undrained subscriber from wedging dispatch - # for every other subscriber and, transitively, the producer. + # Subscriber sinks belong to remote consumers that can fall behind + # the dispatcher, either because the consumer went away (e.g. a + # non-blocking send whose HTTP response already returned) or because + # it is merely slow. evict_on_full keeps one sink whose queue has + # filled from wedging dispatch for every other subscriber and, + # transitively, the producer. tapped_queue = await self._event_queue_subscribers.tap( evict_on_full=True ) diff --git a/src/a2a/server/events/event_queue_v2.py b/src/a2a/server/events/event_queue_v2.py index e0828714c..d760952cf 100644 --- a/src/a2a/server/events/event_queue_v2.py +++ b/src/a2a/server/events/event_queue_v2.py @@ -72,16 +72,20 @@ async def _deliver_to_sink( ) -> None: """Puts one event into one sink, evicting an evict-on-full sink if full. - A sink whose consumer stopped draining (an abandoned subscriber) fills - up and, with a blocking put, would stall this dispatcher's ``gather`` - forever: the incoming queue then fills, producers wedge in - ``enqueue_event``, and the task's event flow never recovers. Sinks - tapped with ``evict_on_full=True`` are therefore treated as broadcast - observers: when full, the sink is force-closed and detached so - dispatch to the remaining sinks continues and blocked producers - recover. Sinks with ``evict_on_full=False`` (including the default - sink) keep the blocking put, which is the documented flow-control - contract. + A sink whose queue has filled up would, with a blocking put, stall + this dispatcher's ``gather`` until the sink drains: if it never does, + the incoming queue then fills, producers wedge in ``enqueue_event``, + and the task's event flow never recovers. Sinks tapped with + ``evict_on_full=True`` are therefore treated as broadcast observers: + when full, the sink is force-closed and detached so dispatch to the + remaining sinks continues and blocked producers recover. Sinks with + ``evict_on_full=False`` (including the default sink) keep the blocking + put, which is the documented flow-control contract. + + The eviction test is queue fullness at delivery time and nothing else. + It makes no claim about the consumer: a consumer that has gone away and + a consumer that is reading steadily but has fallen ``max_queue_size`` + events behind the dispatcher both satisfy it and are both evicted. The ``full()`` pre-check is race-free: this dispatcher is the only writer to a sink queue and consumers only read, so the queue cannot @@ -89,8 +93,8 @@ async def _deliver_to_sink( A sink that closed after the dispatch snapshot was taken (e.g. a graceful ``close(immediate=False)`` whose consumer is still - draining) is skipped: treating its full queue as an abandoned - consumer would upgrade the close to immediate and discard the + draining) is skipped: treating its full queue as an eviction + candidate would upgrade the close to immediate and discard the events the consumer is draining. The check is best-effort — a close landing after it is harmless, because the put below tolerates a shut-down queue. @@ -99,9 +103,12 @@ async def _deliver_to_sink( return if sink._evict_on_full and sink.queue.full(): # noqa: SLF001 logger.warning( - 'Evicting event queue sink %r: its queue is full and the ' - 'consumer is not draining it. Closing the sink so dispatch ' - 'to the remaining sinks continues.', + 'Evicting event queue sink %r: its queue was full at delivery ' + 'time, so this sink is at least max_queue_size events behind ' + 'the dispatcher. Closing the sink so dispatch to the remaining ' + 'sinks continues. This does not mean the consumer stopped ' + 'reading: a consumer that is still reading but lagging by more ' + 'than the bound is evicted the same way.', sink, ) await sink.close(immediate=True) @@ -221,8 +228,10 @@ async def tap( sink is force-closed and detached instead of blocking dispatch (and, transitively, every producer) behind it. When False (the default), a full sink back-pressures dispatch, preserving the - documented flow-control behavior. Use True for sinks whose - consumer may be abandoned, such as remote subscribers. + documented flow-control behavior. Use True for sinks that may + fall more than ``max_queue_size`` events behind the dispatcher, + such as remote subscribers, whether because their consumer went + away or because it is merely slow. """ async with self._lock: if self._is_closed: