Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
workflow_logging_callback.py222 linesDownload Raw Back to callbacks
1from typing import Optional2 3from core.model_runtime.utils.encoders import jsonable_encoder4from core.workflow.graph_engine.entities.event import (5    GraphEngineEvent,6    GraphRunFailedEvent,7    GraphRunStartedEvent,8    GraphRunSucceededEvent,9    IterationRunFailedEvent,10    IterationRunNextEvent,11    IterationRunStartedEvent,12    IterationRunSucceededEvent,13    NodeRunFailedEvent,14    NodeRunStartedEvent,15    NodeRunStreamChunkEvent,16    NodeRunSucceededEvent,17    ParallelBranchRunFailedEvent,18    ParallelBranchRunStartedEvent,19    ParallelBranchRunSucceededEvent,20)21 22from .base_workflow_callback import WorkflowCallback23 24_TEXT_COLOR_MAPPING = {25    "blue": "36;1",26    "yellow": "33;1",27    "pink": "38;5;200",28    "green": "32;1",29    "red": "31;1",30}31 32 33class WorkflowLoggingCallback(WorkflowCallback):34    def __init__(self) -> None:35        self.current_node_id = None36 37    def on_event(self, event: GraphEngineEvent) -> None:38        if isinstance(event, GraphRunStartedEvent):39            self.print_text("\n[GraphRunStartedEvent]", color="pink")40        elif isinstance(event, GraphRunSucceededEvent):41            self.print_text("\n[GraphRunSucceededEvent]", color="green")42        elif isinstance(event, GraphRunFailedEvent):43            self.print_text(f"\n[GraphRunFailedEvent] reason: {event.error}", color="red")44        elif isinstance(event, NodeRunStartedEvent):45            self.on_workflow_node_execute_started(event=event)46        elif isinstance(event, NodeRunSucceededEvent):47            self.on_workflow_node_execute_succeeded(event=event)48        elif isinstance(event, NodeRunFailedEvent):49            self.on_workflow_node_execute_failed(event=event)50        elif isinstance(event, NodeRunStreamChunkEvent):51            self.on_node_text_chunk(event=event)52        elif isinstance(event, ParallelBranchRunStartedEvent):53            self.on_workflow_parallel_started(event=event)54        elif isinstance(event, ParallelBranchRunSucceededEvent | ParallelBranchRunFailedEvent):55            self.on_workflow_parallel_completed(event=event)56        elif isinstance(event, IterationRunStartedEvent):57            self.on_workflow_iteration_started(event=event)58        elif isinstance(event, IterationRunNextEvent):59            self.on_workflow_iteration_next(event=event)60        elif isinstance(event, IterationRunSucceededEvent | IterationRunFailedEvent):61            self.on_workflow_iteration_completed(event=event)62        else:63            self.print_text(f"\n[{event.__class__.__name__}]", color="blue")64 65    def on_workflow_node_execute_started(self, event: NodeRunStartedEvent) -> None:66        """67        Workflow node execute started68        """69        self.print_text("\n[NodeRunStartedEvent]", color="yellow")70        self.print_text(f"Node ID: {event.node_id}", color="yellow")71        self.print_text(f"Node Title: {event.node_data.title}", color="yellow")72        self.print_text(f"Type: {event.node_type.value}", color="yellow")73 74    def on_workflow_node_execute_succeeded(self, event: NodeRunSucceededEvent) -> None:75        """76        Workflow node execute succeeded77        """78        route_node_state = event.route_node_state79 80        self.print_text("\n[NodeRunSucceededEvent]", color="green")81        self.print_text(f"Node ID: {event.node_id}", color="green")82        self.print_text(f"Node Title: {event.node_data.title}", color="green")83        self.print_text(f"Type: {event.node_type.value}", color="green")84 85        if route_node_state.node_run_result:86            node_run_result = route_node_state.node_run_result87            self.print_text(88                f"Inputs: {jsonable_encoder(node_run_result.inputs) if node_run_result.inputs else ''}",89                color="green",90            )91            self.print_text(92                f"Process Data: "93                f"{jsonable_encoder(node_run_result.process_data) if node_run_result.process_data else ''}",94                color="green",95            )96            self.print_text(97                f"Outputs: {jsonable_encoder(node_run_result.outputs) if node_run_result.outputs else ''}",98                color="green",99            )100            self.print_text(101                f"Metadata: {jsonable_encoder(node_run_result.metadata) if node_run_result.metadata else ''}",102                color="green",103            )104 105    def on_workflow_node_execute_failed(self, event: NodeRunFailedEvent) -> None:106        """107        Workflow node execute failed108        """109        route_node_state = event.route_node_state110 111        self.print_text("\n[NodeRunFailedEvent]", color="red")112        self.print_text(f"Node ID: {event.node_id}", color="red")113        self.print_text(f"Node Title: {event.node_data.title}", color="red")114        self.print_text(f"Type: {event.node_type.value}", color="red")115 116        if route_node_state.node_run_result:117            node_run_result = route_node_state.node_run_result118            self.print_text(f"Error: {node_run_result.error}", color="red")119            self.print_text(120                f"Inputs: {jsonable_encoder(node_run_result.inputs) if node_run_result.inputs else ''}",121                color="red",122            )123            self.print_text(124                f"Process Data: "125                f"{jsonable_encoder(node_run_result.process_data) if node_run_result.process_data else ''}",126                color="red",127            )128            self.print_text(129                f"Outputs: {jsonable_encoder(node_run_result.outputs) if node_run_result.outputs else ''}",130                color="red",131            )132 133    def on_node_text_chunk(self, event: NodeRunStreamChunkEvent) -> None:134        """135        Publish text chunk136        """137        route_node_state = event.route_node_state138        if not self.current_node_id or self.current_node_id != route_node_state.node_id:139            self.current_node_id = route_node_state.node_id140            self.print_text("\n[NodeRunStreamChunkEvent]")141            self.print_text(f"Node ID: {route_node_state.node_id}")142 143            node_run_result = route_node_state.node_run_result144            if node_run_result:145                self.print_text(146                    f"Metadata: {jsonable_encoder(node_run_result.metadata) if node_run_result.metadata else ''}"147                )148 149        self.print_text(event.chunk_content, color="pink", end="")150 151    def on_workflow_parallel_started(self, event: ParallelBranchRunStartedEvent) -> None:152        """153        Publish parallel started154        """155        self.print_text("\n[ParallelBranchRunStartedEvent]", color="blue")156        self.print_text(f"Parallel ID: {event.parallel_id}", color="blue")157        self.print_text(f"Branch ID: {event.parallel_start_node_id}", color="blue")158        if event.in_iteration_id:159            self.print_text(f"Iteration ID: {event.in_iteration_id}", color="blue")160 161    def on_workflow_parallel_completed(162        self, event: ParallelBranchRunSucceededEvent | ParallelBranchRunFailedEvent163    ) -> None:164        """165        Publish parallel completed166        """167        if isinstance(event, ParallelBranchRunSucceededEvent):168            color = "blue"169        elif isinstance(event, ParallelBranchRunFailedEvent):170            color = "red"171 172        self.print_text(173            "\n[ParallelBranchRunSucceededEvent]"174            if isinstance(event, ParallelBranchRunSucceededEvent)175            else "\n[ParallelBranchRunFailedEvent]",176            color=color,177        )178        self.print_text(f"Parallel ID: {event.parallel_id}", color=color)179        self.print_text(f"Branch ID: {event.parallel_start_node_id}", color=color)180        if event.in_iteration_id:181            self.print_text(f"Iteration ID: {event.in_iteration_id}", color=color)182 183        if isinstance(event, ParallelBranchRunFailedEvent):184            self.print_text(f"Error: {event.error}", color=color)185 186    def on_workflow_iteration_started(self, event: IterationRunStartedEvent) -> None:187        """188        Publish iteration started189        """190        self.print_text("\n[IterationRunStartedEvent]", color="blue")191        self.print_text(f"Iteration Node ID: {event.iteration_id}", color="blue")192 193    def on_workflow_iteration_next(self, event: IterationRunNextEvent) -> None:194        """195        Publish iteration next196        """197        self.print_text("\n[IterationRunNextEvent]", color="blue")198        self.print_text(f"Iteration Node ID: {event.iteration_id}", color="blue")199        self.print_text(f"Iteration Index: {event.index}", color="blue")200 201    def on_workflow_iteration_completed(self, event: IterationRunSucceededEvent | IterationRunFailedEvent) -> None:202        """203        Publish iteration completed204        """205        self.print_text(206            "\n[IterationRunSucceededEvent]"207            if isinstance(event, IterationRunSucceededEvent)208            else "\n[IterationRunFailedEvent]",209            color="blue",210        )211        self.print_text(f"Node ID: {event.iteration_id}", color="blue")212 213    def print_text(self, text: str, color: Optional[str] = None, end: str = "\n") -> None:214        """Print text with highlighting and no end characters."""215        text_to_print = self._get_colored_text(text, color) if color else text216        print(f"{text_to_print}", end=end)217 218    def _get_colored_text(self, text: str, color: str) -> str:219        """Get colored text."""220        color_str = _TEXT_COLOR_MAPPING[color]221        return f"\u001b[{color_str}m\033[1;3m{text}\u001b[0m"222