fix: preserve full root messages inbox

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
This commit is contained in:
John Kennedy
2026-07-09 04:29:18 +00:00
co-authored by open-swe[bot] <open-swe@users.noreply.github.com>
parent d0524a4473
commit a7cd539c94
4 changed files with 79 additions and 27 deletions
+17 -11
View File
@@ -912,9 +912,13 @@ class _SubgraphsProjection:
return self._subgraphs_iter()
@staticmethod
def _put_root_message(
root_inbox: asyncio.Queue[Event | None], item: Event
) -> None:
def _put_root_message(root_inbox: asyncio.Queue[Event | None], item: Event) -> None:
if root_inbox.maxsize > 0 and root_inbox.qsize() >= root_inbox.maxsize - 1:
raise RuntimeError(
"Root messages inbox exceeded max_queue_size while buffering "
"root-scope messages. Iterate thread.messages concurrently "
"or increase max_queue_size."
)
try:
root_inbox.put_nowait(item)
except asyncio.QueueFull as exc:
@@ -926,13 +930,14 @@ class _SubgraphsProjection:
@staticmethod
def _signal_root_inbox_closed(root_inbox: asyncio.Queue[Event | None]) -> None:
while True:
try:
root_inbox.put_nowait(None)
return
except asyncio.QueueFull:
with contextlib.suppress(asyncio.QueueEmpty):
root_inbox.get_nowait()
try:
root_inbox.put_nowait(None)
except asyncio.QueueFull as exc:
raise RuntimeError(
"Root messages inbox exceeded max_queue_size while closing "
"root-scope messages. Iterate thread.messages concurrently "
"or increase max_queue_size."
) from exc
async def _subgraphs_iter(self) -> AsyncGenerator[ScopedStreamHandle, None]:
if self._thread._transport is None:
@@ -1358,7 +1363,8 @@ class AsyncThreadStream:
that arrive at namespace `[]` before `thread.messages` has subscribed.
"""
if self._root_messages_inbox is None:
self._root_messages_inbox = asyncio.Queue(maxsize=self._max_queue_size)
maxsize = self._max_queue_size + 1 if self._max_queue_size > 0 else 0
self._root_messages_inbox = asyncio.Queue(maxsize=maxsize)
return self._root_messages_inbox
def _register_active_message_stream(self, stream: AsyncChatModelStream) -> None:
+15 -8
View File
@@ -956,6 +956,12 @@ class _SyncSubgraphsProjection:
@staticmethod
def _put_root_message(root_inbox: queue.Queue[Event | None], item: Event) -> None:
if root_inbox.maxsize > 0 and root_inbox.qsize() >= root_inbox.maxsize - 1:
raise RuntimeError(
"Root messages inbox exceeded max_queue_size while buffering "
"root-scope messages. Iterate thread.messages concurrently "
"or increase max_queue_size."
)
try:
root_inbox.put_nowait(item)
except queue.Full as exc:
@@ -967,13 +973,14 @@ class _SyncSubgraphsProjection:
@staticmethod
def _signal_root_inbox_closed(root_inbox: queue.Queue[Event | None]) -> None:
while True:
try:
root_inbox.put_nowait(None)
return
except queue.Full:
with contextlib.suppress(queue.Empty):
root_inbox.get_nowait()
try:
root_inbox.put_nowait(None)
except queue.Full as exc:
raise RuntimeError(
"Root messages inbox exceeded max_queue_size while closing "
"root-scope messages. Iterate thread.messages concurrently "
"or increase max_queue_size."
) from exc
def _subgraphs_iter(self) -> Iterator[SyncScopedStreamHandle]:
if self._thread._transport is None:
@@ -1243,7 +1250,7 @@ class SyncThreadStream:
def _activate_root_messages_inbox(self) -> queue.Queue[Event | None]:
if self._root_messages_inbox is None:
self._root_messages_inbox = queue.Queue(maxsize=1024)
self._root_messages_inbox = queue.Queue(maxsize=1025)
return self._root_messages_inbox
def _register_active_message_stream(self, stream: ChatModelStream) -> None:
@@ -3,10 +3,12 @@
from __future__ import annotations
import asyncio
from typing import cast
from unittest.mock import MagicMock
import httpx
import pytest
from langchain_protocol import Event
from langgraph_sdk._async.http import HttpClient
from langgraph_sdk._async.threads import ThreadsClient
@@ -452,7 +454,7 @@ def test_scoped_handle_inboxes_bounded_by_max_queue_size():
def test_root_messages_inbox_bounded_by_max_queue_size():
"""Root messages inbox must use the stream queue bound."""
"""Root messages inbox must reserve one terminal sentinel slot."""
from langgraph_sdk._async.stream import AsyncThreadStream
thread = AsyncThreadStream(
@@ -464,23 +466,43 @@ def test_root_messages_inbox_bounded_by_max_queue_size():
inbox = thread._activate_root_messages_inbox()
assert inbox.maxsize == 16
assert inbox.maxsize == 17
def test_subgraphs_root_message_overflow_raises_runtime_error():
"""Overflowing the root messages inbox must fail explicitly."""
from langgraph_sdk._async.stream import _SubgraphsProjection
inbox = asyncio.Queue(maxsize=1)
inbox.put_nowait(message_start_event(seq=1, message_id="msg-1"))
inbox: asyncio.Queue[Event | None] = asyncio.Queue(maxsize=2)
inbox.put_nowait(cast(Event, message_start_event(seq=1, message_id="msg-1")))
with pytest.raises(RuntimeError, match="Root messages inbox exceeded"):
_SubgraphsProjection._put_root_message(
inbox,
message_text_delta_event(seq=2, text="overflow", message_id="msg-1"),
cast(
Event,
message_text_delta_event(seq=2, text="overflow", message_id="msg-1"),
),
)
def test_subgraphs_root_message_close_preserves_full_inbox():
"""Closing a full root messages inbox must not evict buffered events."""
from langgraph_sdk._async.stream import _SubgraphsProjection
inbox: asyncio.Queue[Event | None] = asyncio.Queue(maxsize=3)
first = cast(Event, message_start_event(seq=1, message_id="msg-1"))
second = cast(Event, message_finish_event(seq=2, message_id="msg-1"))
inbox.put_nowait(first)
inbox.put_nowait(second)
_SubgraphsProjection._signal_root_inbox_closed(inbox)
assert inbox.get_nowait() is first
assert inbox.get_nowait() is second
assert inbox.get_nowait() is None
async def test_child_handle_inherits_max_queue_size_from_parent():
"""Grandchild ScopedStreamHandles created by _HandleSubgraphsProjection
inherit the parent's max_queue_size so all queues are consistently bounded."""
@@ -400,7 +400,7 @@ def test_sync_tool_calls_run_error_fails_active_handle():
def test_sync_root_messages_inbox_is_bounded():
"""Root messages inbox must have an explicit maximum size."""
"""Root messages inbox must reserve one terminal sentinel slot."""
fake = SyncFakeServer()
fake.script([lifecycle_completed_event(seq=1)])
fake.set_state({})
@@ -409,14 +409,14 @@ def test_sync_root_messages_inbox_is_bounded():
with threads.stream(thread_id="t-1", assistant_id="agent") as thread:
inbox = thread._activate_root_messages_inbox()
assert inbox.maxsize == 1024
assert inbox.maxsize == 1025
def test_sync_subgraphs_root_message_overflow_raises_runtime_error():
"""Overflowing the root messages inbox must fail explicitly."""
from langgraph_sdk._sync.stream import _SyncSubgraphsProjection
inbox: queue.Queue[Event | None] = queue.Queue(maxsize=1)
inbox: queue.Queue[Event | None] = queue.Queue(maxsize=2)
inbox.put_nowait(cast(Event, message_start_event(seq=1, message_id="msg-1")))
with pytest.raises(RuntimeError, match="Root messages inbox exceeded"):
@@ -429,6 +429,23 @@ def test_sync_subgraphs_root_message_overflow_raises_runtime_error():
)
def test_sync_subgraphs_root_message_close_preserves_full_inbox():
"""Closing a full root messages inbox must not evict buffered events."""
from langgraph_sdk._sync.stream import _SyncSubgraphsProjection
inbox: queue.Queue[Event | None] = queue.Queue(maxsize=3)
first = cast(Event, message_start_event(seq=1, message_id="msg-1"))
second = cast(Event, message_finish_event(seq=2, message_id="msg-1"))
inbox.put_nowait(first)
inbox.put_nowait(second)
_SyncSubgraphsProjection._signal_root_inbox_closed(inbox)
assert inbox.get_nowait() is first
assert inbox.get_nowait() is second
assert inbox.get_nowait() is None
def test_sync_drain_messages_inbox_pre_dispatches_before_yield():
"""When draining the root inbox, str(message.text) must work immediately on yield."""
fake = SyncFakeServer()