Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions src/openai/resources/beta/responses/responses.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions src/openai/resources/realtime/realtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions src/openai/resources/responses/responses.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
38 changes: 38 additions & 0 deletions tests/test_websocket_send_queue.py
Original file line number Diff line number Diff line change
@@ -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")