Underground-Digital/Workflow-Engine
0
1from datetime import datetime2from typing import Any, Optional3 4from pydantic import BaseModel, Field5 6from core.workflow.graph_engine.entities.runtime_route_state import RouteNodeState7from core.workflow.nodes import NodeType8from core.workflow.nodes.base import BaseNodeData9 10 11class GraphEngineEvent(BaseModel):12 pass13 14 15###########################################16# Graph Events17###########################################18 19 20class BaseGraphEvent(GraphEngineEvent):21 pass22 23 24class GraphRunStartedEvent(BaseGraphEvent):25 pass26 27 28class GraphRunSucceededEvent(BaseGraphEvent):29 outputs: Optional[dict[str, Any]] = None30 """outputs"""31 32 33class GraphRunFailedEvent(BaseGraphEvent):34 error: str = Field(..., description="failed reason")35 36 37###########################################38# Node Events39###########################################40 41 42class BaseNodeEvent(GraphEngineEvent):43 id: str = Field(..., description="node execution id")44 node_id: str = Field(..., description="node id")45 node_type: NodeType = Field(..., description="node type")46 node_data: BaseNodeData = Field(..., description="node data")47 route_node_state: RouteNodeState = Field(..., description="route node state")48 parallel_id: Optional[str] = None49 """parallel id if node is in parallel"""50 parallel_start_node_id: Optional[str] = None51 """parallel start node id if node is in parallel"""52 parent_parallel_id: Optional[str] = None53 """parent parallel id if node is in parallel"""54 parent_parallel_start_node_id: Optional[str] = None55 """parent parallel start node id if node is in parallel"""56 in_iteration_id: Optional[str] = None57 """iteration id if node is in iteration"""58 59 60class NodeRunStartedEvent(BaseNodeEvent):61 predecessor_node_id: Optional[str] = None62 parallel_mode_run_id: Optional[str] = None63 """predecessor node id"""64 65 66class NodeRunStreamChunkEvent(BaseNodeEvent):67 chunk_content: str = Field(..., description="chunk content")68 from_variable_selector: Optional[list[str]] = None69 """from variable selector"""70 71 72class NodeRunRetrieverResourceEvent(BaseNodeEvent):73 retriever_resources: list[dict] = Field(..., description="retriever resources")74 context: str = Field(..., description="context")75 76 77class NodeRunSucceededEvent(BaseNodeEvent):78 pass79 80 81class NodeRunFailedEvent(BaseNodeEvent):82 error: str = Field(..., description="error")83 84 85class NodeInIterationFailedEvent(BaseNodeEvent):86 error: str = Field(..., description="error")87 88 89###########################################90# Parallel Branch Events91###########################################92 93 94class BaseParallelBranchEvent(GraphEngineEvent):95 parallel_id: str = Field(..., description="parallel id")96 """parallel id"""97 parallel_start_node_id: str = Field(..., description="parallel start node id")98 """parallel start node id"""99 parent_parallel_id: Optional[str] = None100 """parent parallel id if node is in parallel"""101 parent_parallel_start_node_id: Optional[str] = None102 """parent parallel start node id if node is in parallel"""103 in_iteration_id: Optional[str] = None104 """iteration id if node is in iteration"""105 106 107class ParallelBranchRunStartedEvent(BaseParallelBranchEvent):108 pass109 110 111class ParallelBranchRunSucceededEvent(BaseParallelBranchEvent):112 pass113 114 115class ParallelBranchRunFailedEvent(BaseParallelBranchEvent):116 error: str = Field(..., description="failed reason")117 118 119###########################################120# Iteration Events121###########################################122 123 124class BaseIterationEvent(GraphEngineEvent):125 iteration_id: str = Field(..., description="iteration node execution id")126 iteration_node_id: str = Field(..., description="iteration node id")127 iteration_node_type: NodeType = Field(..., description="node type, iteration or loop")128 iteration_node_data: BaseNodeData = Field(..., description="node data")129 parallel_id: Optional[str] = None130 """parallel id if node is in parallel"""131 parallel_start_node_id: Optional[str] = None132 """parallel start node id if node is in parallel"""133 parent_parallel_id: Optional[str] = None134 """parent parallel id if node is in parallel"""135 parent_parallel_start_node_id: Optional[str] = None136 """parent parallel start node id if node is in parallel"""137 parallel_mode_run_id: Optional[str] = None138 """iteratoin run in parallel mode run id"""139 140 141class IterationRunStartedEvent(BaseIterationEvent):142 start_at: datetime = Field(..., description="start at")143 inputs: Optional[dict[str, Any]] = None144 metadata: Optional[dict[str, Any]] = None145 predecessor_node_id: Optional[str] = None146 147 148class IterationRunNextEvent(BaseIterationEvent):149 index: int = Field(..., description="index")150 pre_iteration_output: Optional[Any] = Field(None, description="pre iteration output")151 152 153class IterationRunSucceededEvent(BaseIterationEvent):154 start_at: datetime = Field(..., description="start at")155 inputs: Optional[dict[str, Any]] = None156 outputs: Optional[dict[str, Any]] = None157 metadata: Optional[dict[str, Any]] = None158 steps: int = 0159 160 161class IterationRunFailedEvent(BaseIterationEvent):162 start_at: datetime = Field(..., description="start at")163 inputs: Optional[dict[str, Any]] = None164 outputs: Optional[dict[str, Any]] = None165 metadata: Optional[dict[str, Any]] = None166 steps: int = 0167 error: str = Field(..., description="failed reason")168 169 170InNodeEvent = BaseNodeEvent | BaseParallelBranchEvent | BaseIterationEvent171 