diff --git a/libs/sdk-py/langgraph_sdk/_async/stream.py b/libs/sdk-py/langgraph_sdk/_async/stream.py index b05ab40b5..b823df7f4 100644 --- a/libs/sdk-py/langgraph_sdk/_async/stream.py +++ b/libs/sdk-py/langgraph_sdk/_async/stream.py @@ -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: diff --git a/libs/sdk-py/langgraph_sdk/_sync/stream.py b/libs/sdk-py/langgraph_sdk/_sync/stream.py index 141cf01e8..8b43cbd7b 100644 --- a/libs/sdk-py/langgraph_sdk/_sync/stream.py +++ b/libs/sdk-py/langgraph_sdk/_sync/stream.py @@ -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: diff --git a/libs/sdk-py/tests/streaming/test_scoped_handles.py b/libs/sdk-py/tests/streaming/test_scoped_handles.py index 18841f9de..475c8a553 100644 --- a/libs/sdk-py/tests/streaming/test_scoped_handles.py +++ b/libs/sdk-py/tests/streaming/test_scoped_handles.py @@ -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.""" diff --git a/libs/sdk-py/tests/streaming/test_sync_projections.py b/libs/sdk-py/tests/streaming/test_sync_projections.py index 38749efdf..9961da7ba 100644 --- a/libs/sdk-py/tests/streaming/test_sync_projections.py +++ b/libs/sdk-py/tests/streaming/test_sync_projections.py @@ -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()