transports: disconnect client first
This commit is contained in:
@@ -324,28 +324,19 @@ class LiveKitInputTransport(BaseInputTransport):
|
||||
logger.info("LiveKitInputTransport started")
|
||||
|
||||
async def stop(self, frame: EndFrame):
|
||||
await self._client.disconnect()
|
||||
if self._audio_in_task:
|
||||
self._audio_in_task.cancel()
|
||||
try:
|
||||
await self._audio_in_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
await self._audio_in_task
|
||||
await super().stop(frame)
|
||||
await self._client.disconnect()
|
||||
logger.info("LiveKitInputTransport stopped")
|
||||
|
||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||
if isinstance(frame, EndFrame):
|
||||
await self.stop(frame)
|
||||
else:
|
||||
await super().process_frame(frame, direction)
|
||||
|
||||
async def cancel(self, frame: CancelFrame):
|
||||
await super().cancel(frame)
|
||||
await self._client.disconnect()
|
||||
if self._audio_in_task and (self._params.audio_in_enabled or self._params.vad_enabled):
|
||||
self._audio_in_task.cancel()
|
||||
await self._audio_in_task
|
||||
await super().cancel(frame)
|
||||
|
||||
def vad_analyzer(self) -> VADAnalyzer | None:
|
||||
return self._vad_analyzer
|
||||
@@ -407,19 +398,13 @@ class LiveKitOutputTransport(BaseOutputTransport):
|
||||
logger.info("LiveKitOutputTransport started")
|
||||
|
||||
async def stop(self, frame: EndFrame):
|
||||
await super().stop(frame)
|
||||
await self._client.disconnect()
|
||||
await super().stop(frame)
|
||||
logger.info("LiveKitOutputTransport stopped")
|
||||
|
||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||
if isinstance(frame, EndFrame):
|
||||
await self.stop(frame)
|
||||
else:
|
||||
await super().process_frame(frame, direction)
|
||||
|
||||
async def cancel(self, frame: CancelFrame):
|
||||
await super().cancel(frame)
|
||||
await self._client.disconnect()
|
||||
await super().cancel(frame)
|
||||
|
||||
async def send_message(self, frame: TransportMessageFrame | TransportMessageUrgentFrame):
|
||||
if isinstance(frame, (LiveKitTransportMessageFrame, LiveKitTransportMessageUrgentFrame)):
|
||||
@@ -526,12 +511,6 @@ class LiveKitTransport(BaseTransport):
|
||||
|
||||
async def _on_disconnected(self):
|
||||
await self._call_event_handler("on_disconnected")
|
||||
# Attempt to reconnect
|
||||
try:
|
||||
await self._client.connect()
|
||||
await self._call_event_handler("on_connected")
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to reconnect: {e}")
|
||||
|
||||
async def _on_participant_connected(self, participant_id: str):
|
||||
await self._call_event_handler("on_participant_connected", participant_id)
|
||||
|
||||
Reference in New Issue
Block a user