Merge pull request #1880 from pipecat-ai/aleix/daily-resize-event-loop-fix
BaseOutputTransport: don't block event loop during image resize
This commit is contained in:
@@ -127,6 +127,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
### Fixed
|
### Fixed
|
||||||
|
|
||||||
|
- Fixed a `DailyTransport` issue that would cause images needing resize to block
|
||||||
|
the event loop.
|
||||||
|
|
||||||
- Fixed an issue with `ElevenLabsTTSService` where changing the model or voice
|
- Fixed an issue with `ElevenLabsTTSService` where changing the model or voice
|
||||||
while the service is running wasn't working.
|
while the service is running wasn't working.
|
||||||
|
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import asyncio
|
|||||||
import itertools
|
import itertools
|
||||||
import sys
|
import sys
|
||||||
import time
|
import time
|
||||||
|
from concurrent.futures import ThreadPoolExecutor
|
||||||
from typing import Any, AsyncGenerator, Dict, List, Mapping, Optional
|
from typing import Any, AsyncGenerator, Dict, List, Mapping, Optional
|
||||||
|
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
@@ -234,6 +235,9 @@ class BaseOutputTransport(FrameProcessor):
|
|||||||
self._audio_chunk_size = audio_chunk_size
|
self._audio_chunk_size = audio_chunk_size
|
||||||
self._params = params
|
self._params = params
|
||||||
|
|
||||||
|
# This is to resize images. We only need to resize one image at a time.
|
||||||
|
self._executor = ThreadPoolExecutor(max_workers=1)
|
||||||
|
|
||||||
# Buffer to keep track of incoming audio.
|
# Buffer to keep track of incoming audio.
|
||||||
self._audio_buffer = bytearray()
|
self._audio_buffer = bytearray()
|
||||||
|
|
||||||
@@ -558,18 +562,25 @@ class BaseOutputTransport(FrameProcessor):
|
|||||||
self._video_queue.task_done()
|
self._video_queue.task_done()
|
||||||
|
|
||||||
async def _draw_image(self, frame: OutputImageRawFrame):
|
async def _draw_image(self, frame: OutputImageRawFrame):
|
||||||
desired_size = (self._params.video_out_width, self._params.video_out_height)
|
def resize_frame(frame: OutputImageRawFrame) -> OutputImageRawFrame:
|
||||||
|
desired_size = (self._params.video_out_width, self._params.video_out_height)
|
||||||
|
|
||||||
# TODO: we should refactor in the future to support dynamic resolutions
|
# TODO: we should refactor in the future to support dynamic resolutions
|
||||||
# which is kind of what happens in P2P connections.
|
# which is kind of what happens in P2P connections.
|
||||||
# We need to add support for that inside the DailyTransport
|
# We need to add support for that inside the DailyTransport
|
||||||
if frame.size != desired_size:
|
if frame.size != desired_size:
|
||||||
image = Image.frombytes(frame.format, frame.size, frame.image)
|
image = Image.frombytes(frame.format, frame.size, frame.image)
|
||||||
resized_image = image.resize(desired_size)
|
resized_image = image.resize(desired_size)
|
||||||
# logger.warning(f"{frame} does not have the expected size {desired_size}, resizing")
|
# logger.warning(f"{frame} does not have the expected size {desired_size}, resizing")
|
||||||
frame = OutputImageRawFrame(
|
frame = OutputImageRawFrame(
|
||||||
resized_image.tobytes(), resized_image.size, resized_image.format
|
resized_image.tobytes(), resized_image.size, resized_image.format
|
||||||
)
|
)
|
||||||
|
|
||||||
|
return frame
|
||||||
|
|
||||||
|
frame = await self._transport.get_event_loop().run_in_executor(
|
||||||
|
self._executor, resize_frame, frame
|
||||||
|
)
|
||||||
|
|
||||||
await self._transport.write_raw_video_frame(frame, self._destination)
|
await self._transport.write_raw_video_frame(frame, self._destination)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user