Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
queue_entities.py518 linesDownload Raw Back to entities
1from datetime import datetime2from enum import Enum3from typing import Any, Optional4 5from pydantic import BaseModel, field_validator6 7from core.model_runtime.entities.llm_entities import LLMResult, LLMResultChunk8from core.workflow.entities.node_entities import NodeRunMetadataKey9from core.workflow.graph_engine.entities.graph_runtime_state import GraphRuntimeState10from core.workflow.nodes import NodeType11from core.workflow.nodes.base import BaseNodeData12 13 14class QueueEvent(str, Enum):15    """16    QueueEvent enum17    """18 19    LLM_CHUNK = "llm_chunk"20    TEXT_CHUNK = "text_chunk"21    AGENT_MESSAGE = "agent_message"22    MESSAGE_REPLACE = "message_replace"23    MESSAGE_END = "message_end"24    ADVANCED_CHAT_MESSAGE_END = "advanced_chat_message_end"25    WORKFLOW_STARTED = "workflow_started"26    WORKFLOW_SUCCEEDED = "workflow_succeeded"27    WORKFLOW_FAILED = "workflow_failed"28    ITERATION_START = "iteration_start"29    ITERATION_NEXT = "iteration_next"30    ITERATION_COMPLETED = "iteration_completed"31    NODE_STARTED = "node_started"32    NODE_SUCCEEDED = "node_succeeded"33    NODE_FAILED = "node_failed"34    RETRIEVER_RESOURCES = "retriever_resources"35    ANNOTATION_REPLY = "annotation_reply"36    AGENT_THOUGHT = "agent_thought"37    MESSAGE_FILE = "message_file"38    PARALLEL_BRANCH_RUN_STARTED = "parallel_branch_run_started"39    PARALLEL_BRANCH_RUN_SUCCEEDED = "parallel_branch_run_succeeded"40    PARALLEL_BRANCH_RUN_FAILED = "parallel_branch_run_failed"41    ERROR = "error"42    PING = "ping"43    STOP = "stop"44 45 46class AppQueueEvent(BaseModel):47    """48    QueueEvent abstract entity49    """50 51    event: QueueEvent52 53 54class QueueLLMChunkEvent(AppQueueEvent):55    """56    QueueLLMChunkEvent entity57    Only for basic mode apps58    """59 60    event: QueueEvent = QueueEvent.LLM_CHUNK61    chunk: LLMResultChunk62 63 64class QueueIterationStartEvent(AppQueueEvent):65    """66    QueueIterationStartEvent entity67    """68 69    event: QueueEvent = QueueEvent.ITERATION_START70    node_execution_id: str71    node_id: str72    node_type: NodeType73    node_data: BaseNodeData74    parallel_id: Optional[str] = None75    """parallel id if node is in parallel"""76    parallel_start_node_id: Optional[str] = None77    """parallel start node id if node is in parallel"""78    parent_parallel_id: Optional[str] = None79    """parent parallel id if node is in parallel"""80    parent_parallel_start_node_id: Optional[str] = None81    """parent parallel start node id if node is in parallel"""82    start_at: datetime83 84    node_run_index: int85    inputs: Optional[dict[str, Any]] = None86    predecessor_node_id: Optional[str] = None87    metadata: Optional[dict[str, Any]] = None88 89 90class QueueIterationNextEvent(AppQueueEvent):91    """92    QueueIterationNextEvent entity93    """94 95    event: QueueEvent = QueueEvent.ITERATION_NEXT96 97    index: int98    node_execution_id: str99    node_id: str100    node_type: NodeType101    node_data: BaseNodeData102    parallel_id: Optional[str] = None103    """parallel id if node is in parallel"""104    parallel_start_node_id: Optional[str] = None105    """parallel start node id if node is in parallel"""106    parent_parallel_id: Optional[str] = None107    """parent parallel id if node is in parallel"""108    parent_parallel_start_node_id: Optional[str] = None109    """parent parallel start node id if node is in parallel"""110    parallel_mode_run_id: Optional[str] = None111    """iteratoin run in parallel mode run id"""112    node_run_index: int113    output: Optional[Any] = None  # output for the current iteration114 115    @field_validator("output", mode="before")116    @classmethod117    def set_output(cls, v):118        """119        Set output120        """121        if v is None:122            return None123        if isinstance(v, int | float | str | bool | dict | list):124            return v125        raise ValueError("output must be a valid type")126 127 128class QueueIterationCompletedEvent(AppQueueEvent):129    """130    QueueIterationCompletedEvent entity131    """132 133    event: QueueEvent = QueueEvent.ITERATION_COMPLETED134 135    node_execution_id: str136    node_id: str137    node_type: NodeType138    node_data: BaseNodeData139    parallel_id: Optional[str] = None140    """parallel id if node is in parallel"""141    parallel_start_node_id: Optional[str] = None142    """parallel start node id if node is in parallel"""143    parent_parallel_id: Optional[str] = None144    """parent parallel id if node is in parallel"""145    parent_parallel_start_node_id: Optional[str] = None146    """parent parallel start node id if node is in parallel"""147    start_at: datetime148 149    node_run_index: int150    inputs: Optional[dict[str, Any]] = None151    outputs: Optional[dict[str, Any]] = None152    metadata: Optional[dict[str, Any]] = None153    steps: int = 0154 155    error: Optional[str] = None156 157 158class QueueTextChunkEvent(AppQueueEvent):159    """160    QueueTextChunkEvent entity161    """162 163    event: QueueEvent = QueueEvent.TEXT_CHUNK164    text: str165    from_variable_selector: Optional[list[str]] = None166    """from variable selector"""167    in_iteration_id: Optional[str] = None168    """iteration id if node is in iteration"""169 170 171class QueueAgentMessageEvent(AppQueueEvent):172    """173    QueueMessageEvent entity174    """175 176    event: QueueEvent = QueueEvent.AGENT_MESSAGE177    chunk: LLMResultChunk178 179 180class QueueMessageReplaceEvent(AppQueueEvent):181    """182    QueueMessageReplaceEvent entity183    """184 185    event: QueueEvent = QueueEvent.MESSAGE_REPLACE186    text: str187 188 189class QueueRetrieverResourcesEvent(AppQueueEvent):190    """191    QueueRetrieverResourcesEvent entity192    """193 194    event: QueueEvent = QueueEvent.RETRIEVER_RESOURCES195    retriever_resources: list[dict]196    in_iteration_id: Optional[str] = None197    """iteration id if node is in iteration"""198 199 200class QueueAnnotationReplyEvent(AppQueueEvent):201    """202    QueueAnnotationReplyEvent entity203    """204 205    event: QueueEvent = QueueEvent.ANNOTATION_REPLY206    message_annotation_id: str207 208 209class QueueMessageEndEvent(AppQueueEvent):210    """211    QueueMessageEndEvent entity212    """213 214    event: QueueEvent = QueueEvent.MESSAGE_END215    llm_result: Optional[LLMResult] = None216 217 218class QueueAdvancedChatMessageEndEvent(AppQueueEvent):219    """220    QueueAdvancedChatMessageEndEvent entity221    """222 223    event: QueueEvent = QueueEvent.ADVANCED_CHAT_MESSAGE_END224 225 226class QueueWorkflowStartedEvent(AppQueueEvent):227    """228    QueueWorkflowStartedEvent entity229    """230 231    event: QueueEvent = QueueEvent.WORKFLOW_STARTED232    graph_runtime_state: GraphRuntimeState233 234 235class QueueWorkflowSucceededEvent(AppQueueEvent):236    """237    QueueWorkflowSucceededEvent entity238    """239 240    event: QueueEvent = QueueEvent.WORKFLOW_SUCCEEDED241    outputs: Optional[dict[str, Any]] = None242 243 244class QueueWorkflowFailedEvent(AppQueueEvent):245    """246    QueueWorkflowFailedEvent entity247    """248 249    event: QueueEvent = QueueEvent.WORKFLOW_FAILED250    error: str251 252 253class QueueNodeStartedEvent(AppQueueEvent):254    """255    QueueNodeStartedEvent entity256    """257 258    event: QueueEvent = QueueEvent.NODE_STARTED259 260    node_execution_id: str261    node_id: str262    node_type: NodeType263    node_data: BaseNodeData264    node_run_index: int = 1265    predecessor_node_id: Optional[str] = None266    parallel_id: Optional[str] = None267    """parallel id if node is in parallel"""268    parallel_start_node_id: Optional[str] = None269    """parallel start node id if node is in parallel"""270    parent_parallel_id: Optional[str] = None271    """parent parallel id if node is in parallel"""272    parent_parallel_start_node_id: Optional[str] = None273    """parent parallel start node id if node is in parallel"""274    in_iteration_id: Optional[str] = None275    """iteration id if node is in iteration"""276    start_at: datetime277    parallel_mode_run_id: Optional[str] = None278    """iteratoin run in parallel mode run id"""279 280 281class QueueNodeSucceededEvent(AppQueueEvent):282    """283    QueueNodeSucceededEvent entity284    """285 286    event: QueueEvent = QueueEvent.NODE_SUCCEEDED287 288    node_execution_id: str289    node_id: str290    node_type: NodeType291    node_data: BaseNodeData292    parallel_id: Optional[str] = None293    """parallel id if node is in parallel"""294    parallel_start_node_id: Optional[str] = None295    """parallel start node id if node is in parallel"""296    parent_parallel_id: Optional[str] = None297    """parent parallel id if node is in parallel"""298    parent_parallel_start_node_id: Optional[str] = None299    """parent parallel start node id if node is in parallel"""300    in_iteration_id: Optional[str] = None301    """iteration id if node is in iteration"""302    start_at: datetime303 304    inputs: Optional[dict[str, Any]] = None305    process_data: Optional[dict[str, Any]] = None306    outputs: Optional[dict[str, Any]] = None307    execution_metadata: Optional[dict[NodeRunMetadataKey, Any]] = None308 309    error: Optional[str] = None310 311 312class QueueNodeInIterationFailedEvent(AppQueueEvent):313    """314    QueueNodeInIterationFailedEvent entity315    """316 317    event: QueueEvent = QueueEvent.NODE_FAILED318 319    node_execution_id: str320    node_id: str321    node_type: NodeType322    node_data: BaseNodeData323    parallel_id: Optional[str] = None324    """parallel id if node is in parallel"""325    parallel_start_node_id: Optional[str] = None326    """parallel start node id if node is in parallel"""327    parent_parallel_id: Optional[str] = None328    """parent parallel id if node is in parallel"""329    parent_parallel_start_node_id: Optional[str] = None330    """parent parallel start node id if node is in parallel"""331    in_iteration_id: Optional[str] = None332    """iteration id if node is in iteration"""333    start_at: datetime334 335    inputs: Optional[dict[str, Any]] = None336    process_data: Optional[dict[str, Any]] = None337    outputs: Optional[dict[str, Any]] = None338    execution_metadata: Optional[dict[NodeRunMetadataKey, Any]] = None339 340    error: str341 342 343class QueueNodeFailedEvent(AppQueueEvent):344    """345    QueueNodeFailedEvent entity346    """347 348    event: QueueEvent = QueueEvent.NODE_FAILED349 350    node_execution_id: str351    node_id: str352    node_type: NodeType353    node_data: BaseNodeData354    parallel_id: Optional[str] = None355    """parallel id if node is in parallel"""356    parallel_start_node_id: Optional[str] = None357    """parallel start node id if node is in parallel"""358    parent_parallel_id: Optional[str] = None359    """parent parallel id if node is in parallel"""360    parent_parallel_start_node_id: Optional[str] = None361    """parent parallel start node id if node is in parallel"""362    in_iteration_id: Optional[str] = None363    """iteration id if node is in iteration"""364    start_at: datetime365 366    inputs: Optional[dict[str, Any]] = None367    process_data: Optional[dict[str, Any]] = None368    outputs: Optional[dict[str, Any]] = None369    execution_metadata: Optional[dict[NodeRunMetadataKey, Any]] = None370 371    error: str372 373 374class QueueAgentThoughtEvent(AppQueueEvent):375    """376    QueueAgentThoughtEvent entity377    """378 379    event: QueueEvent = QueueEvent.AGENT_THOUGHT380    agent_thought_id: str381 382 383class QueueMessageFileEvent(AppQueueEvent):384    """385    QueueAgentThoughtEvent entity386    """387 388    event: QueueEvent = QueueEvent.MESSAGE_FILE389    message_file_id: str390 391 392class QueueErrorEvent(AppQueueEvent):393    """394    QueueErrorEvent entity395    """396 397    event: QueueEvent = QueueEvent.ERROR398    error: Any = None399 400 401class QueuePingEvent(AppQueueEvent):402    """403    QueuePingEvent entity404    """405 406    event: QueueEvent = QueueEvent.PING407 408 409class QueueStopEvent(AppQueueEvent):410    """411    QueueStopEvent entity412    """413 414    class StopBy(Enum):415        """416        Stop by enum417        """418 419        USER_MANUAL = "user-manual"420        ANNOTATION_REPLY = "annotation-reply"421        OUTPUT_MODERATION = "output-moderation"422        INPUT_MODERATION = "input-moderation"423 424    event: QueueEvent = QueueEvent.STOP425    stopped_by: StopBy426 427    def get_stop_reason(self) -> str:428        """429        To stop reason430        """431        reason_mapping = {432            QueueStopEvent.StopBy.USER_MANUAL: "Stopped by user.",433            QueueStopEvent.StopBy.ANNOTATION_REPLY: "Stopped by annotation reply.",434            QueueStopEvent.StopBy.OUTPUT_MODERATION: "Stopped by output moderation.",435            QueueStopEvent.StopBy.INPUT_MODERATION: "Stopped by input moderation.",436        }437 438        return reason_mapping.get(self.stopped_by, "Stopped by unknown reason.")439 440 441class QueueMessage(BaseModel):442    """443    QueueMessage abstract entity444    """445 446    task_id: str447    app_mode: str448    event: AppQueueEvent449 450 451class MessageQueueMessage(QueueMessage):452    """453    MessageQueueMessage entity454    """455 456    message_id: str457    conversation_id: str458 459 460class WorkflowQueueMessage(QueueMessage):461    """462    WorkflowQueueMessage entity463    """464 465    pass466 467 468class QueueParallelBranchRunStartedEvent(AppQueueEvent):469    """470    QueueParallelBranchRunStartedEvent entity471    """472 473    event: QueueEvent = QueueEvent.PARALLEL_BRANCH_RUN_STARTED474 475    parallel_id: str476    parallel_start_node_id: str477    parent_parallel_id: Optional[str] = None478    """parent parallel id if node is in parallel"""479    parent_parallel_start_node_id: Optional[str] = None480    """parent parallel start node id if node is in parallel"""481    in_iteration_id: Optional[str] = None482    """iteration id if node is in iteration"""483 484 485class QueueParallelBranchRunSucceededEvent(AppQueueEvent):486    """487    QueueParallelBranchRunSucceededEvent entity488    """489 490    event: QueueEvent = QueueEvent.PARALLEL_BRANCH_RUN_SUCCEEDED491 492    parallel_id: str493    parallel_start_node_id: str494    parent_parallel_id: Optional[str] = None495    """parent parallel id if node is in parallel"""496    parent_parallel_start_node_id: Optional[str] = None497    """parent parallel start node id if node is in parallel"""498    in_iteration_id: Optional[str] = None499    """iteration id if node is in iteration"""500 501 502class QueueParallelBranchRunFailedEvent(AppQueueEvent):503    """504    QueueParallelBranchRunFailedEvent entity505    """506 507    event: QueueEvent = QueueEvent.PARALLEL_BRANCH_RUN_FAILED508 509    parallel_id: str510    parallel_start_node_id: str511    parent_parallel_id: Optional[str] = None512    """parent parallel id if node is in parallel"""513    parent_parallel_start_node_id: Optional[str] = None514    """parent parallel start node id if node is in parallel"""515    in_iteration_id: Optional[str] = None516    """iteration id if node is in iteration"""517    error: str518