Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
event.py171 linesDownload Raw Back to entities
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