Merge pull request #319 from pipecat-ai/aleix/daily-completion-callbacks-0.0.39-fix
transports(daily): fix completion callbacks handling
This commit is contained in:
@@ -5,6 +5,13 @@ All notable changes to **pipecat** will be documented in this file.
|
|||||||
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
|
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
|
||||||
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
|
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
|
||||||
|
|
||||||
|
## [0.0.39] - 2024-07-23
|
||||||
|
|
||||||
|
### Fixed
|
||||||
|
|
||||||
|
- Fixed a regression introduced in 0.0.38 that would cause Daily transcription
|
||||||
|
to stop the Pipeline.
|
||||||
|
|
||||||
## [0.0.38] - 2024-07-23
|
## [0.0.38] - 2024-07-23
|
||||||
|
|
||||||
### Added
|
### Added
|
||||||
|
|||||||
@@ -126,7 +126,10 @@ class DailyCallbacks(BaseModel):
|
|||||||
def completion_callback(future):
|
def completion_callback(future):
|
||||||
def _callback(*args):
|
def _callback(*args):
|
||||||
if not future.cancelled():
|
if not future.cancelled():
|
||||||
future.get_loop().call_soon_threadsafe(future.set_result, *args)
|
if len(args) > 1:
|
||||||
|
future.get_loop().call_soon_threadsafe(future.set_result, args)
|
||||||
|
else:
|
||||||
|
future.get_loop().call_soon_threadsafe(future.set_result, *args)
|
||||||
return _callback
|
return _callback
|
||||||
|
|
||||||
|
|
||||||
@@ -282,11 +285,12 @@ class DailyTransportClient(EventHandler):
|
|||||||
await self._callbacks.on_error(error_msg)
|
await self._callbacks.on_error(error_msg)
|
||||||
|
|
||||||
async def _start_transcription(self):
|
async def _start_transcription(self):
|
||||||
future = self._loop.create_future()
|
|
||||||
logger.info(f"Enabling transcription with settings {self._params.transcription_settings}")
|
logger.info(f"Enabling transcription with settings {self._params.transcription_settings}")
|
||||||
|
|
||||||
|
future = self._loop.create_future()
|
||||||
self._client.start_transcription(
|
self._client.start_transcription(
|
||||||
settings=self._params.transcription_settings.model_dump(exclude_none=True),
|
settings=self._params.transcription_settings.model_dump(exclude_none=True),
|
||||||
completion=lambda error: future.set_result(error)
|
completion=completion_callback(future)
|
||||||
)
|
)
|
||||||
error = await future
|
error = await future
|
||||||
if error:
|
if error:
|
||||||
@@ -295,14 +299,10 @@ class DailyTransportClient(EventHandler):
|
|||||||
async def _join(self):
|
async def _join(self):
|
||||||
future = self._loop.create_future()
|
future = self._loop.create_future()
|
||||||
|
|
||||||
def handle_join_response(data, error):
|
|
||||||
if not future.cancelled():
|
|
||||||
future.get_loop().call_soon_threadsafe(future.set_result, (data, error))
|
|
||||||
|
|
||||||
self._client.join(
|
self._client.join(
|
||||||
self._room_url,
|
self._room_url,
|
||||||
self._token,
|
self._token,
|
||||||
completion=handle_join_response,
|
completion=completion_callback(future),
|
||||||
client_settings={
|
client_settings={
|
||||||
"inputs": {
|
"inputs": {
|
||||||
"camera": {
|
"camera": {
|
||||||
@@ -370,20 +370,14 @@ class DailyTransportClient(EventHandler):
|
|||||||
|
|
||||||
async def _stop_transcription(self):
|
async def _stop_transcription(self):
|
||||||
future = self._loop.create_future()
|
future = self._loop.create_future()
|
||||||
self._client.stop_transcription(completion=lambda error: future.set_result(error))
|
self._client.stop_transcription(completion=completion_callback(future))
|
||||||
error = await future
|
error = await future
|
||||||
if error:
|
if error:
|
||||||
logger.error(f"Unable to stop transcription: {error}")
|
logger.error(f"Unable to stop transcription: {error}")
|
||||||
|
|
||||||
async def _leave(self):
|
async def _leave(self):
|
||||||
future = self._loop.create_future()
|
future = self._loop.create_future()
|
||||||
|
self._client.leave(completion=completion_callback(future))
|
||||||
def handle_leave_response(error):
|
|
||||||
if not future.cancelled():
|
|
||||||
future.get_loop().call_soon_threadsafe(future.set_result, error)
|
|
||||||
|
|
||||||
self._client.leave(completion=handle_leave_response)
|
|
||||||
|
|
||||||
return await asyncio.wait_for(future, timeout=10)
|
return await asyncio.wait_for(future, timeout=10)
|
||||||
|
|
||||||
async def cleanup(self):
|
async def cleanup(self):
|
||||||
|
|||||||
Reference in New Issue
Block a user