From c8a1d85af5bff617f99d7dd99b1484148f601a54 Mon Sep 17 00:00:00 2001 From: HughhhhCoder Date: Thu, 27 Aug 2026 14:22:23 +0800 Subject: [PATCH] Preserve explicitly provided send queues --- .../resources/beta/responses/responses.py | 4 +- src/openai/resources/realtime/realtime.py | 4 +- src/openai/resources/responses/responses.py | 4 +- tests/test_websocket_send_queue.py | 38 +++++++++++++++++++ 4 files changed, 44 insertions(+), 6 deletions(-) create mode 100644 tests/test_websocket_send_queue.py diff --git a/src/openai/resources/beta/responses/responses.py b/src/openai/resources/beta/responses/responses.py index c9d1120fcc..1a9c31ee51 100644 --- a/src/openai/resources/beta/responses/responses.py +++ b/src/openai/resources/beta/responses/responses.py @@ -4131,7 +4131,7 @@ def __init__( self._extra_headers = extra_headers self._intentionally_closed = False self._is_reconnecting = False - self._send_queue = send_queue or SendQueue() + self._send_queue = send_queue if send_queue is not None else SendQueue() self._event_handler_registry = EventHandlerRegistry(use_lock=False) self.response = AsyncResponsesResponseResource(self) @@ -4588,7 +4588,7 @@ def __init__( self._extra_headers = extra_headers self._intentionally_closed = False self._is_reconnecting = False - self._send_queue = send_queue or SendQueue() + self._send_queue = send_queue if send_queue is not None else SendQueue() self._event_handler_registry = EventHandlerRegistry(use_lock=True) self.response = ResponsesResponseResource(self) diff --git a/src/openai/resources/realtime/realtime.py b/src/openai/resources/realtime/realtime.py index 3cb815bddf..6ff204f2a5 100644 --- a/src/openai/resources/realtime/realtime.py +++ b/src/openai/resources/realtime/realtime.py @@ -292,7 +292,7 @@ def __init__( self._extra_headers = extra_headers self._intentionally_closed = False self._is_reconnecting = False - self._send_queue = send_queue or SendQueue() + self._send_queue = send_queue if send_queue is not None else SendQueue() self._event_handler_registry = EventHandlerRegistry(use_lock=False) self.session = AsyncRealtimeSessionResource(self) @@ -774,7 +774,7 @@ def __init__( self._extra_headers = extra_headers self._intentionally_closed = False self._is_reconnecting = False - self._send_queue = send_queue or SendQueue() + self._send_queue = send_queue if send_queue is not None else SendQueue() self._event_handler_registry = EventHandlerRegistry(use_lock=True) self.session = RealtimeSessionResource(self) diff --git a/src/openai/resources/responses/responses.py b/src/openai/resources/responses/responses.py index 2232b1a05d..3683545d49 100644 --- a/src/openai/resources/responses/responses.py +++ b/src/openai/resources/responses/responses.py @@ -4026,7 +4026,7 @@ def __init__( self._extra_headers = extra_headers self._intentionally_closed = False self._is_reconnecting = False - self._send_queue = send_queue or SendQueue() + self._send_queue = send_queue if send_queue is not None else SendQueue() self._event_handler_registry = EventHandlerRegistry(use_lock=False) self.response = AsyncResponsesResponseResource(self) @@ -4483,7 +4483,7 @@ def __init__( self._extra_headers = extra_headers self._intentionally_closed = False self._is_reconnecting = False - self._send_queue = send_queue or SendQueue() + self._send_queue = send_queue if send_queue is not None else SendQueue() self._event_handler_registry = EventHandlerRegistry(use_lock=True) self.response = ResponsesResponseResource(self) diff --git a/tests/test_websocket_send_queue.py b/tests/test_websocket_send_queue.py new file mode 100644 index 0000000000..3e2f39ec48 --- /dev/null +++ b/tests/test_websocket_send_queue.py @@ -0,0 +1,38 @@ +from __future__ import annotations + +from typing import Any, cast +from unittest.mock import Mock + +import pytest + +from openai._exceptions import WebSocketQueueFullError +from openai._send_queue import SendQueue +from openai.resources.realtime.realtime import RealtimeConnection, AsyncRealtimeConnection +from openai.resources.responses.responses import ( + ResponsesConnection, + AsyncResponsesConnection, +) +from openai.resources.beta.responses.responses import ( + ResponsesConnection as BetaResponsesConnection, + AsyncResponsesConnection as AsyncBetaResponsesConnection, +) + + +@pytest.mark.parametrize( + "connection_type", + [ + AsyncRealtimeConnection, + RealtimeConnection, + AsyncResponsesConnection, + ResponsesConnection, + AsyncBetaResponsesConnection, + BetaResponsesConnection, + ], +) +def test_connections_preserve_an_explicit_empty_send_queue(connection_type: Any) -> None: + send_queue = SendQueue(max_bytes=0) + connection = connection_type(cast(Any, Mock()), send_queue=send_queue) + + assert connection._send_queue is send_queue + with pytest.raises(WebSocketQueueFullError): + connection._send_queue.enqueue("message")