Skip to content

Commit 80e9686

Browse files
fix(live): preserve caller queues and prevent uncertain WebSocket replay (openai#3980)
If a Live connection opened with an empty manager queue, iterator recovery lost the caller's queued messages and byte budget. A typed command could also be sent again after the first socket actually received it but its send raised. Keep the passed queue even when empty and drop only the attempted message on a failed send/flush, leaving unattempted messages ordered. This matches the merged Realtime handling. Primary, fork and sideband fixes are separate commits for review. No startup, auth, default reconnect, imports or public method signatures change. Direct send errors still raise; queued flush failures still log the existing warning. The Live helper notes describe the recovery and recording-error policy. Validation: - Baseline: 337 Live/redirect tests pass under both Pydantic 1 and 2. - Reproduced all 18 queue/replay failures using actual loopback wire delivery, then 42 role tests pass, including per-role error/unknown-event/close ordering and helper disposal. - Final 399 Live, protected dial, Realtime and send-queue tests pass in each Pydantic lane; Ruff and full SDK Mypy (1,865 source files) pass. - Trusted main budget/isolation pass: 7,826/10,000 with unchanged generated snapshot. Independent review found only the queued-error doc clarification, which is included.
1 parent 78074d7 commit 80e9686

5 files changed

Lines changed: 341 additions & 63 deletions

File tree

‎src/openai/lib/live/README.md‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
# Live WebSocket lifecycle
2+
3+
`client.live.connect()` and `client.live.forks.connect(session_id=...)` leave
4+
startup to the caller: send `connection.session.start(...)` and wait for
5+
`session.started`. A sideband connection created by
6+
`client.live.sideband.connect(session_id=...)` attaches to an existing session.
7+
It does not send a start or require a new `session.started`; an attachment's
8+
short replay is not a complete session snapshot.
9+
10+
Direct `recv()` and `recv_bytes()` calls report transport errors to their
11+
caller. Iterator reconnection is available only when an `on_reconnecting`
12+
callback was explicitly supplied. The callback controls retry and can update
13+
credentials or query parameters. A new socket is not proof that the previous
14+
Live session, recording, or application state was restored. Applications own
15+
their recovery decision and any necessary startup or state reconstruction.
16+
Don't use multiple physical readers: a dispatcher owns the read loop while it
17+
runs. Detaching or closing one transcript grouper doesn't cancel another
18+
observer or the connection.
19+
20+
All modes preserve the manager's caller-configured `max_queue_size`, including
21+
if the queue was empty at connection time. Only messages that haven't been
22+
attempted on a socket remain eligible for flushing after a retry. A direct send
23+
exception is raised to its caller and cannot prove the server didn't receive
24+
that message. A failure while flushing an already queued message still logs a
25+
warning, as before; it cannot be raised at the original queueing call. Neither
26+
failed attempt is automatically retried. Unattempted pre-open messages and
27+
messages explicitly queued during recovery remain in order within their
28+
existing budget. No automatic session reopening or restoration is implied.
29+
30+
Treat each wire error as an event to handle, not as a successful session result.
31+
In particular, preserve `session_storage_failed` if `session.closed` follows;
32+
the close doesn't mean that recording storage succeeded. Unknown events and
33+
fields are preserved for callers that need them.

‎src/openai/resources/live/forks.py‎

Lines changed: 6 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -139,7 +139,7 @@ def __init__(
139139
self._extra_headers = extra_headers
140140
self._intentionally_closed = False
141141
self._is_reconnecting = False
142-
self._send_queue = send_queue or SendQueue()
142+
self._send_queue = send_queue if send_queue is not None else SendQueue()
143143
self._event_handler_registry = EventHandlerRegistry(use_lock=False)
144144

145145
self.session = AsyncForksSessionResource(self)
@@ -215,11 +215,7 @@ async def send(self, event: ForkClientEvent | ForkClientEventParam) -> None:
215215
if self._is_reconnecting:
216216
self._send_queue.enqueue(data)
217217
return
218-
try:
219-
await self._connection.send(data)
220-
except Exception:
221-
self._send_queue.enqueue(data)
222-
raise
218+
await self._connection.send(data)
223219

224220
async def send_raw(self, data: bytes | str) -> None:
225221
if self._is_reconnecting:
@@ -324,7 +320,7 @@ async def _send(data: str) -> None:
324320
await self._connection.send(data)
325321

326322
try:
327-
await self._send_queue.flush_async(_send)
323+
await self._send_queue.flush_async(_send, requeue_failed=False)
328324
except Exception:
329325
log.warning("Failed to flush send queue after reconnect")
330326

@@ -638,7 +634,7 @@ def __init__(
638634
self._extra_headers = extra_headers
639635
self._intentionally_closed = False
640636
self._is_reconnecting = False
641-
self._send_queue = send_queue or SendQueue()
637+
self._send_queue = send_queue if send_queue is not None else SendQueue()
642638
self._event_handler_registry = EventHandlerRegistry(use_lock=True)
643639

644640
self.session = ForksSessionResource(self)
@@ -714,11 +710,7 @@ def send(self, event: ForkClientEvent | ForkClientEventParam) -> None:
714710
if self._is_reconnecting:
715711
self._send_queue.enqueue(data)
716712
return
717-
try:
718-
self._connection.send(data)
719-
except Exception:
720-
self._send_queue.enqueue(data)
721-
raise
713+
self._connection.send(data)
722714

723715
def send_raw(self, data: bytes | str) -> None:
724716
if self._is_reconnecting:
@@ -817,7 +809,7 @@ def _reconnect(self, exc: Exception) -> bool:
817809
def _flush_send_queue(self) -> None:
818810
"""Send all queued messages over the current connection."""
819811
try:
820-
self._send_queue.flush_sync(lambda data: self._connection.send(data))
812+
self._send_queue.flush_sync(lambda data: self._connection.send(data), requeue_failed=False)
821813
except Exception:
822814
log.warning("Failed to flush send queue after reconnect")
823815

‎src/openai/resources/live/live.py‎

Lines changed: 6 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -360,7 +360,7 @@ def __init__(
360360
self._extra_headers = extra_headers
361361
self._intentionally_closed = False
362362
self._is_reconnecting = False
363-
self._send_queue = send_queue or SendQueue()
363+
self._send_queue = send_queue if send_queue is not None else SendQueue()
364364
self._event_handler_registry = EventHandlerRegistry(use_lock=False)
365365

366366
self.session = AsyncLiveSessionResource(self)
@@ -436,11 +436,7 @@ async def send(self, event: ClientEvent | ClientEventParam) -> None:
436436
if self._is_reconnecting:
437437
self._send_queue.enqueue(data)
438438
return
439-
try:
440-
await self._connection.send(data)
441-
except Exception:
442-
self._send_queue.enqueue(data)
443-
raise
439+
await self._connection.send(data)
444440

445441
async def send_raw(self, data: bytes | str) -> None:
446442
if self._is_reconnecting:
@@ -545,7 +541,7 @@ async def _send(data: str) -> None:
545541
await self._connection.send(data)
546542

547543
try:
548-
await self._send_queue.flush_async(_send)
544+
await self._send_queue.flush_async(_send, requeue_failed=False)
549545
except Exception:
550546
log.warning("Failed to flush send queue after reconnect")
551547

@@ -856,7 +852,7 @@ def __init__(
856852
self._extra_headers = extra_headers
857853
self._intentionally_closed = False
858854
self._is_reconnecting = False
859-
self._send_queue = send_queue or SendQueue()
855+
self._send_queue = send_queue if send_queue is not None else SendQueue()
860856
self._event_handler_registry = EventHandlerRegistry(use_lock=True)
861857

862858
self.session = LiveSessionResource(self)
@@ -932,11 +928,7 @@ def send(self, event: ClientEvent | ClientEventParam) -> None:
932928
if self._is_reconnecting:
933929
self._send_queue.enqueue(data)
934930
return
935-
try:
936-
self._connection.send(data)
937-
except Exception:
938-
self._send_queue.enqueue(data)
939-
raise
931+
self._connection.send(data)
940932

941933
def send_raw(self, data: bytes | str) -> None:
942934
if self._is_reconnecting:
@@ -1035,7 +1027,7 @@ def _reconnect(self, exc: Exception) -> bool:
10351027
def _flush_send_queue(self) -> None:
10361028
"""Send all queued messages over the current connection."""
10371029
try:
1038-
self._send_queue.flush_sync(lambda data: self._connection.send(data))
1030+
self._send_queue.flush_sync(lambda data: self._connection.send(data), requeue_failed=False)
10391031
except Exception:
10401032
log.warning("Failed to flush send queue after reconnect")
10411033

‎src/openai/resources/live/sideband.py‎

Lines changed: 6 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -142,7 +142,7 @@ def __init__(
142142
self._extra_headers = extra_headers
143143
self._intentionally_closed = False
144144
self._is_reconnecting = False
145-
self._send_queue = send_queue or SendQueue()
145+
self._send_queue = send_queue if send_queue is not None else SendQueue()
146146
self._event_handler_registry = EventHandlerRegistry(use_lock=False)
147147

148148
self.session = AsyncSidebandSessionResource(self)
@@ -218,11 +218,7 @@ async def send(self, event: ConnectClientEvent | ConnectClientEventParam) -> Non
218218
if self._is_reconnecting:
219219
self._send_queue.enqueue(data)
220220
return
221-
try:
222-
await self._connection.send(data)
223-
except Exception:
224-
self._send_queue.enqueue(data)
225-
raise
221+
await self._connection.send(data)
226222

227223
async def send_raw(self, data: bytes | str) -> None:
228224
if self._is_reconnecting:
@@ -329,7 +325,7 @@ async def _send(data: str) -> None:
329325
await self._connection.send(data)
330326

331327
try:
332-
await self._send_queue.flush_async(_send)
328+
await self._send_queue.flush_async(_send, requeue_failed=False)
333329
except Exception:
334330
log.warning("Failed to flush send queue after reconnect")
335331

@@ -646,7 +642,7 @@ def __init__(
646642
self._extra_headers = extra_headers
647643
self._intentionally_closed = False
648644
self._is_reconnecting = False
649-
self._send_queue = send_queue or SendQueue()
645+
self._send_queue = send_queue if send_queue is not None else SendQueue()
650646
self._event_handler_registry = EventHandlerRegistry(use_lock=True)
651647

652648
self.session = SidebandSessionResource(self)
@@ -722,11 +718,7 @@ def send(self, event: ConnectClientEvent | ConnectClientEventParam) -> None:
722718
if self._is_reconnecting:
723719
self._send_queue.enqueue(data)
724720
return
725-
try:
726-
self._connection.send(data)
727-
except Exception:
728-
self._send_queue.enqueue(data)
729-
raise
721+
self._connection.send(data)
730722

731723
def send_raw(self, data: bytes | str) -> None:
732724
if self._is_reconnecting:
@@ -827,7 +819,7 @@ def _reconnect(self, exc: Exception) -> bool:
827819
def _flush_send_queue(self) -> None:
828820
"""Send all queued messages over the current connection."""
829821
try:
830-
self._send_queue.flush_sync(lambda data: self._connection.send(data))
822+
self._send_queue.flush_sync(lambda data: self._connection.send(data), requeue_failed=False)
831823
except Exception:
832824
log.warning("Failed to flush send queue after reconnect")
833825

0 commit comments

Comments
 (0)