some event loop parameter updates
This commit is contained in:
@@ -46,7 +46,7 @@ class PipelineRunner:
|
|||||||
return self._running
|
return self._running
|
||||||
|
|
||||||
def _setup_sigint(self):
|
def _setup_sigint(self):
|
||||||
loop = asyncio.get_event_loop()
|
loop = asyncio.get_running_loop()
|
||||||
loop.add_signal_handler(
|
loop.add_signal_handler(
|
||||||
signal.SIGINT,
|
signal.SIGINT,
|
||||||
lambda *args: asyncio.create_task(self._sigint_handler())
|
lambda *args: asyncio.create_task(self._sigint_handler())
|
||||||
|
|||||||
@@ -5,7 +5,7 @@
|
|||||||
#
|
#
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
from asyncio import AbstractEventLoop
|
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
|
|
||||||
from pipecat.frames.frames import ErrorFrame, Frame
|
from pipecat.frames.frames import ErrorFrame, Frame
|
||||||
@@ -21,12 +21,12 @@ class FrameDirection(Enum):
|
|||||||
|
|
||||||
class FrameProcessor:
|
class FrameProcessor:
|
||||||
|
|
||||||
def __init__(self):
|
def __init__(self, loop: asyncio.AbstractEventLoop | None = None):
|
||||||
self.id: int = obj_id()
|
self.id: int = obj_id()
|
||||||
self.name = f"{self.__class__.__name__}#{obj_count(self)}"
|
self.name = f"{self.__class__.__name__}#{obj_count(self)}"
|
||||||
self._prev: "FrameProcessor" | None = None
|
self._prev: "FrameProcessor" | None = None
|
||||||
self._next: "FrameProcessor" | None = None
|
self._next: "FrameProcessor" | None = None
|
||||||
self._loop: AbstractEventLoop = asyncio.get_event_loop()
|
self._loop: asyncio.AbstractEventLoop = loop or asyncio.get_running_loop()
|
||||||
|
|
||||||
async def cleanup(self):
|
async def cleanup(self):
|
||||||
pass
|
pass
|
||||||
@@ -36,7 +36,7 @@ class FrameProcessor:
|
|||||||
processor._prev = self
|
processor._prev = self
|
||||||
logger.debug(f"Linking {self} -> {self._next}")
|
logger.debug(f"Linking {self} -> {self._next}")
|
||||||
|
|
||||||
def get_event_loop(self) -> AbstractEventLoop:
|
def get_event_loop(self) -> asyncio.AbstractEventLoop:
|
||||||
return self._loop
|
return self._loop
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
|
|||||||
@@ -41,8 +41,8 @@ class TransportParams(BaseModel):
|
|||||||
|
|
||||||
class BaseTransport(ABC):
|
class BaseTransport(ABC):
|
||||||
|
|
||||||
def __init__(self, loop: asyncio.AbstractEventLoop):
|
def __init__(self, loop: asyncio.AbstractEventLoop | None):
|
||||||
self._loop = loop
|
self._loop = loop or asyncio.get_running_loop()
|
||||||
self._event_handlers: dict = {}
|
self._event_handlers: dict = {}
|
||||||
|
|
||||||
@abstractmethod
|
@abstractmethod
|
||||||
|
|||||||
@@ -223,7 +223,7 @@ class WebsocketServerTransport(BaseTransport):
|
|||||||
host: str = "localhost",
|
host: str = "localhost",
|
||||||
port: int = 8765,
|
port: int = 8765,
|
||||||
params: WebsocketServerParams = WebsocketServerParams(),
|
params: WebsocketServerParams = WebsocketServerParams(),
|
||||||
loop: asyncio.AbstractEventLoop = asyncio.get_event_loop()):
|
loop: asyncio.AbstractEventLoop | None = None):
|
||||||
super().__init__(loop)
|
super().__init__(loop)
|
||||||
self._host = host
|
self._host = host
|
||||||
self._port = port
|
self._port = port
|
||||||
|
|||||||
@@ -643,7 +643,7 @@ class DailyTransport(BaseTransport):
|
|||||||
token: str | None,
|
token: str | None,
|
||||||
bot_name: str,
|
bot_name: str,
|
||||||
params: DailyParams,
|
params: DailyParams,
|
||||||
loop: asyncio.AbstractEventLoop = asyncio.get_event_loop()):
|
loop: asyncio.AbstractEventLoop | None = None):
|
||||||
super().__init__(loop)
|
super().__init__(loop)
|
||||||
|
|
||||||
callbacks = DailyCallbacks(
|
callbacks = DailyCallbacks(
|
||||||
@@ -663,7 +663,8 @@ class DailyTransport(BaseTransport):
|
|||||||
)
|
)
|
||||||
self._params = params
|
self._params = params
|
||||||
|
|
||||||
self._client = DailyTransportClient(room_url, token, bot_name, params, callbacks, loop)
|
self._client = DailyTransportClient(
|
||||||
|
room_url, token, bot_name, params, callbacks, self._loop)
|
||||||
self._input: DailyInputTransport | None = None
|
self._input: DailyInputTransport | None = None
|
||||||
self._output: DailyOutputTransport | None = None
|
self._output: DailyOutputTransport | None = None
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user