Creating push_error_frame
This commit is contained in:
@@ -634,7 +634,7 @@ class FrameProcessor(BaseObject):
|
|||||||
|
|
||||||
async def push_error(
|
async def push_error(
|
||||||
self,
|
self,
|
||||||
error_message: Optional[str] = None,
|
error_msg: Optional[str] = None,
|
||||||
exception: Optional[Exception] = None,
|
exception: Optional[Exception] = None,
|
||||||
fatal: bool = False,
|
fatal: bool = False,
|
||||||
):
|
):
|
||||||
@@ -645,7 +645,7 @@ class FrameProcessor(BaseObject):
|
|||||||
which processor generated the error.
|
which processor generated the error.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
error_message: Optional descriptive message explaining the error condition.
|
error_msg: Optional descriptive message explaining the error condition.
|
||||||
exception: Optional exception object that caused the error, if available.
|
exception: Optional exception object that caused the error, if available.
|
||||||
This provides additional context for debugging and error handling.
|
This provides additional context for debugging and error handling.
|
||||||
fatal: Whether this error should be considered fatal to the pipeline.
|
fatal: Whether this error should be considered fatal to the pipeline.
|
||||||
@@ -664,14 +664,23 @@ class FrameProcessor(BaseObject):
|
|||||||
await self.push_error("Critical operation failed", exception=e, fatal=True)
|
await self.push_error("Critical operation failed", exception=e, fatal=True)
|
||||||
```
|
```
|
||||||
"""
|
"""
|
||||||
error_message = error_message or f"{self} exception: {exception}"
|
error_message = error_msg or f"{self} exception: {exception}"
|
||||||
logger.error(error_message)
|
logger.error(error_message)
|
||||||
error_frame = ErrorFrame(error=error_message, fatal=fatal, exception=exception)
|
error_frame = ErrorFrame(
|
||||||
await self._call_event_handler("on_error", error_frame)
|
error=error_message, fatal=fatal, exception=exception, processor=self
|
||||||
await self.push_frame(
|
|
||||||
error_frame,
|
|
||||||
FrameDirection.UPSTREAM,
|
|
||||||
)
|
)
|
||||||
|
await self.push_error_frame(error=error_frame)
|
||||||
|
|
||||||
|
async def push_error_frame(self, error: ErrorFrame):
|
||||||
|
"""Push an error frame upstream.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
error: The error frame to push.
|
||||||
|
"""
|
||||||
|
if not error.processor:
|
||||||
|
error.processor = self
|
||||||
|
await self._call_event_handler("on_error", error)
|
||||||
|
await self.push_frame(error, FrameDirection.UPSTREAM)
|
||||||
|
|
||||||
async def push_frame(self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM):
|
async def push_frame(self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM):
|
||||||
"""Push a frame to the next processor in the pipeline.
|
"""Push a frame to the next processor in the pipeline.
|
||||||
@@ -792,8 +801,7 @@ class FrameProcessor(BaseObject):
|
|||||||
await self.__cancel_process_task()
|
await self.__cancel_process_task()
|
||||||
self.__create_process_task()
|
self.__create_process_task()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Uncaught exception in {self} when handling _start_interruption: {e}")
|
await self.push_error(error_msg=f"Uncaught exception in {self} when handling _start_interruption: {e}", exception=e)
|
||||||
await self.push_error(ErrorFrame(str(e)))
|
|
||||||
|
|
||||||
async def __internal_push_frame(self, frame: Frame, direction: FrameDirection):
|
async def __internal_push_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
"""Internal method to push frames to adjacent processors.
|
"""Internal method to push frames to adjacent processors.
|
||||||
@@ -830,8 +838,7 @@ class FrameProcessor(BaseObject):
|
|||||||
await self._observer.on_push_frame(data)
|
await self._observer.on_push_frame(data)
|
||||||
await self._prev.queue_frame(frame, direction)
|
await self._prev.queue_frame(frame, direction)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Uncaught exception in {self}: {e}")
|
await self.push_error(exception=e)
|
||||||
await self.push_error(ErrorFrame(str(e)))
|
|
||||||
|
|
||||||
def _check_started(self, frame: Frame):
|
def _check_started(self, frame: Frame):
|
||||||
"""Check if the processor has been started.
|
"""Check if the processor has been started.
|
||||||
@@ -907,8 +914,7 @@ class FrameProcessor(BaseObject):
|
|||||||
|
|
||||||
await self._call_event_handler("on_after_process_frame", frame)
|
await self._call_event_handler("on_after_process_frame", frame)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"{self}: error processing frame: {e}")
|
await self.push_error(error_msg=f"{self}: error processing frame: {e}", exception=e)
|
||||||
await self.push_error(ErrorFrame(str(e)))
|
|
||||||
|
|
||||||
async def __input_frame_task_handler(self):
|
async def __input_frame_task_handler(self):
|
||||||
"""Handle frames from the input queue.
|
"""Handle frames from the input queue.
|
||||||
|
|||||||
Reference in New Issue
Block a user