processors(rtvi): fix task cleanup
This commit is contained in:
@@ -59,5 +59,6 @@ class AsyncFrameProcessor(FrameProcessor):
|
|||||||
(frame, direction) = await self._push_queue.get()
|
(frame, direction) = await self._push_queue.get()
|
||||||
await self.push_frame(frame, direction)
|
await self.push_frame(frame, direction)
|
||||||
running = not isinstance(frame, EndFrame)
|
running = not isinstance(frame, EndFrame)
|
||||||
|
self._push_queue.task_done()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
|
|||||||
@@ -316,6 +316,7 @@ class RTVIProcessor(FrameProcessor):
|
|||||||
try:
|
try:
|
||||||
(frame, direction) = await self._frame_queue.get()
|
(frame, direction) = await self._frame_queue.get()
|
||||||
await self._handle_frame(frame, direction)
|
await self._handle_frame(frame, direction)
|
||||||
|
self._frame_queue.task_done()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
|
|
||||||
|
|||||||
@@ -114,9 +114,11 @@ class CartesiaTTSService(TTSService):
|
|||||||
try:
|
try:
|
||||||
if self._context_appending_task:
|
if self._context_appending_task:
|
||||||
self._context_appending_task.cancel()
|
self._context_appending_task.cancel()
|
||||||
|
await self._context_appending_task
|
||||||
self._context_appending_task = None
|
self._context_appending_task = None
|
||||||
if self._receive_task:
|
if self._receive_task:
|
||||||
self._receive_task.cancel()
|
self._receive_task.cancel()
|
||||||
|
await self._receive_task
|
||||||
self._receive_task = None
|
self._receive_task = None
|
||||||
if self._websocket:
|
if self._websocket:
|
||||||
ws = self._websocket
|
ws = self._websocket
|
||||||
|
|||||||
@@ -104,6 +104,7 @@ class BaseInputTransport(FrameProcessor):
|
|||||||
try:
|
try:
|
||||||
(frame, direction) = await self._push_queue.get()
|
(frame, direction) = await self._push_queue.get()
|
||||||
await self.push_frame(frame, direction)
|
await self.push_frame(frame, direction)
|
||||||
|
self._push_queue.task_done()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
|
|
||||||
@@ -185,6 +186,8 @@ class BaseInputTransport(FrameProcessor):
|
|||||||
# Push audio downstream if passthrough.
|
# Push audio downstream if passthrough.
|
||||||
if audio_passthrough:
|
if audio_passthrough:
|
||||||
await self._internal_push_frame(frame)
|
await self._internal_push_frame(frame)
|
||||||
|
|
||||||
|
self._audio_in_queue.task_done()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|||||||
@@ -204,6 +204,7 @@ class BaseOutputTransport(FrameProcessor):
|
|||||||
try:
|
try:
|
||||||
(frame, direction) = await self._push_queue.get()
|
(frame, direction) = await self._push_queue.get()
|
||||||
await self.push_frame(frame, direction)
|
await self.push_frame(frame, direction)
|
||||||
|
self._push_queue.task_done()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user