Underground-Digital/Workflow-Engine
0
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 