processor: add support for setting a processor parent
This commit is contained in:
@@ -91,5 +91,7 @@ class Pipeline(BasePipeline):
|
|||||||
def _link_processors(self):
|
def _link_processors(self):
|
||||||
prev = self._processors[0]
|
prev = self._processors[0]
|
||||||
for curr in self._processors[1:]:
|
for curr in self._processors[1:]:
|
||||||
|
prev.set_parent(self)
|
||||||
prev.link(curr)
|
prev.link(curr)
|
||||||
prev = curr
|
prev = curr
|
||||||
|
prev.set_parent(self)
|
||||||
|
|||||||
@@ -72,6 +72,7 @@ class FrameProcessor:
|
|||||||
**kwargs):
|
**kwargs):
|
||||||
self.id: int = obj_id()
|
self.id: int = obj_id()
|
||||||
self.name = name or f"{self.__class__.__name__}#{obj_count(self)}"
|
self.name = name or f"{self.__class__.__name__}#{obj_count(self)}"
|
||||||
|
self._parent: "FrameProcessor" | None = None
|
||||||
self._prev: "FrameProcessor" | None = None
|
self._prev: "FrameProcessor" | None = None
|
||||||
self._next: "FrameProcessor" | None = None
|
self._next: "FrameProcessor" | None = None
|
||||||
self._loop: asyncio.AbstractEventLoop = loop or asyncio.get_running_loop()
|
self._loop: asyncio.AbstractEventLoop = loop or asyncio.get_running_loop()
|
||||||
@@ -126,7 +127,7 @@ class FrameProcessor:
|
|||||||
async def cleanup(self):
|
async def cleanup(self):
|
||||||
pass
|
pass
|
||||||
|
|
||||||
def link(self, processor: 'FrameProcessor'):
|
def link(self, processor: "FrameProcessor"):
|
||||||
self._next = processor
|
self._next = processor
|
||||||
processor._prev = self
|
processor._prev = self
|
||||||
logger.debug(f"Linking {self} -> {self._next}")
|
logger.debug(f"Linking {self} -> {self._next}")
|
||||||
@@ -134,6 +135,12 @@ class FrameProcessor:
|
|||||||
def get_event_loop(self) -> asyncio.AbstractEventLoop:
|
def get_event_loop(self) -> asyncio.AbstractEventLoop:
|
||||||
return self._loop
|
return self._loop
|
||||||
|
|
||||||
|
def set_parent(self, parent: "FrameProcessor"):
|
||||||
|
self._parent = parent
|
||||||
|
|
||||||
|
def get_parent(self) -> "FrameProcessor":
|
||||||
|
return self._parent
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
if isinstance(frame, StartFrame):
|
if isinstance(frame, StartFrame):
|
||||||
self._allow_interruptions = frame.allow_interruptions
|
self._allow_interruptions = frame.allow_interruptions
|
||||||
|
|||||||
Reference in New Issue
Block a user