Compare commits
1 Commits
vp-trace-c
...
fix/fastap
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ae8b9f0756 |
@@ -278,6 +278,13 @@ class FastAPIWebsocketInputTransport(BaseInputTransport):
|
|||||||
|
|
||||||
async def _receive_messages(self):
|
async def _receive_messages(self):
|
||||||
"""Main message receiving loop for WebSocket messages."""
|
"""Main message receiving loop for WebSocket messages."""
|
||||||
|
|
||||||
|
async def trigger_disconnect_if_needed():
|
||||||
|
# Trigger `on_client_disconnected` if the client actually disconnects,
|
||||||
|
# that is, we are not the ones disconnecting.
|
||||||
|
if not self._client.is_closing:
|
||||||
|
await self._client.trigger_client_disconnected()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
async for message in self._client.receive():
|
async for message in self._client.receive():
|
||||||
if not self._params.serializer:
|
if not self._params.serializer:
|
||||||
@@ -294,11 +301,14 @@ class FastAPIWebsocketInputTransport(BaseInputTransport):
|
|||||||
await self.push_frame(frame)
|
await self.push_frame(frame)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} exception receiving data: {e.__class__.__name__} ({e})")
|
logger.error(f"{self} exception receiving data: {e.__class__.__name__} ({e})")
|
||||||
|
finally:
|
||||||
# Trigger `on_client_disconnected` if the client actually disconnects,
|
# Use shield to prevent cancellation from stopping the disconnect callback
|
||||||
# that is, we are not the ones disconnecting.
|
try:
|
||||||
if not self._client.is_closing:
|
await asyncio.shield(trigger_disconnect_if_needed())
|
||||||
await self._client.trigger_client_disconnected()
|
except asyncio.CancelledError:
|
||||||
|
# Even if we're cancelled, try to trigger the disconnect
|
||||||
|
await trigger_disconnect_if_needed()
|
||||||
|
raise
|
||||||
|
|
||||||
async def _monitor_websocket(self):
|
async def _monitor_websocket(self):
|
||||||
"""Wait for self._params.session_timeout seconds, if the websocket is still open, trigger timeout event."""
|
"""Wait for self._params.session_timeout seconds, if the websocket is still open, trigger timeout event."""
|
||||||
|
|||||||
Reference in New Issue
Block a user