Underground-Digital/Workflow-Engine
0
1from collections.abc import Mapping, Sequence2from enum import Enum3from typing import Any, Optional4 5from pydantic import BaseModel, ConfigDict6 7from core.model_runtime.entities.llm_entities import LLMResult8from core.model_runtime.utils.encoders import jsonable_encoder9from models.workflow import WorkflowNodeExecutionStatus10 11 12class TaskState(BaseModel):13 """14 TaskState entity15 """16 17 metadata: dict = {}18 19 20class EasyUITaskState(TaskState):21 """22 EasyUITaskState entity23 """24 25 llm_result: LLMResult26 27 28class WorkflowTaskState(TaskState):29 """30 WorkflowTaskState entity31 """32 33 answer: str = ""34 35 36class StreamEvent(Enum):37 """38 Stream event39 """40 41 PING = "ping"42 ERROR = "error"43 MESSAGE = "message"44 MESSAGE_END = "message_end"45 TTS_MESSAGE = "tts_message"46 TTS_MESSAGE_END = "tts_message_end"47 MESSAGE_FILE = "message_file"48 MESSAGE_REPLACE = "message_replace"49 AGENT_THOUGHT = "agent_thought"50 AGENT_MESSAGE = "agent_message"51 WORKFLOW_STARTED = "workflow_started"52 WORKFLOW_FINISHED = "workflow_finished"53 NODE_STARTED = "node_started"54 NODE_FINISHED = "node_finished"55 PARALLEL_BRANCH_STARTED = "parallel_branch_started"56 PARALLEL_BRANCH_FINISHED = "parallel_branch_finished"57 ITERATION_STARTED = "iteration_started"58 ITERATION_NEXT = "iteration_next"59 ITERATION_COMPLETED = "iteration_completed"60 TEXT_CHUNK = "text_chunk"61 TEXT_REPLACE = "text_replace"62 63 64class StreamResponse(BaseModel):65 """66 StreamResponse entity67 """68 69 event: StreamEvent70 task_id: str71 72 def to_dict(self) -> dict:73 return jsonable_encoder(self)74 75 76class ErrorStreamResponse(StreamResponse):77 """78 ErrorStreamResponse entity79 """80 81 event: StreamEvent = StreamEvent.ERROR82 err: Exception83 model_config = ConfigDict(arbitrary_types_allowed=True)84 85 86class MessageStreamResponse(StreamResponse):87 """88 MessageStreamResponse entity89 """90 91 event: StreamEvent = StreamEvent.MESSAGE92 id: str93 answer: str94 from_variable_selector: Optional[list[str]] = None95 96 97class MessageAudioStreamResponse(StreamResponse):98 """99 MessageStreamResponse entity100 """101 102 event: StreamEvent = StreamEvent.TTS_MESSAGE103 audio: str104 105 106class MessageAudioEndStreamResponse(StreamResponse):107 """108 MessageStreamResponse entity109 """110 111 event: StreamEvent = StreamEvent.TTS_MESSAGE_END112 audio: str113 114 115class MessageEndStreamResponse(StreamResponse):116 """117 MessageEndStreamResponse entity118 """119 120 event: StreamEvent = StreamEvent.MESSAGE_END121 id: str122 metadata: dict = {}123 files: Optional[Sequence[Mapping[str, Any]]] = None124 125 126class MessageFileStreamResponse(StreamResponse):127 """128 MessageFileStreamResponse entity129 """130 131 event: StreamEvent = StreamEvent.MESSAGE_FILE132 id: str133 type: str134 belongs_to: str135 url: str136 137 138class MessageReplaceStreamResponse(StreamResponse):139 """140 MessageReplaceStreamResponse entity141 """142 143 event: StreamEvent = StreamEvent.MESSAGE_REPLACE144 answer: str145 146 147class AgentThoughtStreamResponse(StreamResponse):148 """149 AgentThoughtStreamResponse entity150 """151 152 event: StreamEvent = StreamEvent.AGENT_THOUGHT153 id: str154 position: int155 thought: Optional[str] = None156 observation: Optional[str] = None157 tool: Optional[str] = None158 tool_labels: Optional[dict] = None159 tool_input: Optional[str] = None160 message_files: Optional[list[str]] = None161 162 163class AgentMessageStreamResponse(StreamResponse):164 """165 AgentMessageStreamResponse entity166 """167 168 event: StreamEvent = StreamEvent.AGENT_MESSAGE169 id: str170 answer: str171 172 173class WorkflowStartStreamResponse(StreamResponse):174 """175 WorkflowStartStreamResponse entity176 """177 178 class Data(BaseModel):179 """180 Data entity181 """182 183 id: str184 workflow_id: str185 sequence_number: int186 inputs: dict187 created_at: int188 189 event: StreamEvent = StreamEvent.WORKFLOW_STARTED190 workflow_run_id: str191 data: Data192 193 194class WorkflowFinishStreamResponse(StreamResponse):195 """196 WorkflowFinishStreamResponse entity197 """198 199 class Data(BaseModel):200 """201 Data entity202 """203 204 id: str205 workflow_id: str206 sequence_number: int207 status: str208 outputs: Optional[dict] = None209 error: Optional[str] = None210 elapsed_time: float211 total_tokens: int212 total_steps: int213 created_by: Optional[dict] = None214 created_at: int215 finished_at: int216 files: Optional[Sequence[Mapping[str, Any]]] = []217 218 event: StreamEvent = StreamEvent.WORKFLOW_FINISHED219 workflow_run_id: str220 data: Data221 222 223class NodeStartStreamResponse(StreamResponse):224 """225 NodeStartStreamResponse entity226 """227 228 class Data(BaseModel):229 """230 Data entity231 """232 233 id: str234 node_id: str235 node_type: str236 title: str237 index: int238 predecessor_node_id: Optional[str] = None239 inputs: Optional[dict] = None240 created_at: int241 extras: dict = {}242 parallel_id: Optional[str] = None243 parallel_start_node_id: Optional[str] = None244 parent_parallel_id: Optional[str] = None245 parent_parallel_start_node_id: Optional[str] = None246 iteration_id: Optional[str] = None247 parallel_run_id: Optional[str] = None248 249 event: StreamEvent = StreamEvent.NODE_STARTED250 workflow_run_id: str251 data: Data252 253 def to_ignore_detail_dict(self):254 return {255 "event": self.event.value,256 "task_id": self.task_id,257 "workflow_run_id": self.workflow_run_id,258 "data": {259 "id": self.data.id,260 "node_id": self.data.node_id,261 "node_type": self.data.node_type,262 "title": self.data.title,263 "index": self.data.index,264 "predecessor_node_id": self.data.predecessor_node_id,265 "inputs": None,266 "created_at": self.data.created_at,267 "extras": {},268 "parallel_id": self.data.parallel_id,269 "parallel_start_node_id": self.data.parallel_start_node_id,270 "parent_parallel_id": self.data.parent_parallel_id,271 "parent_parallel_start_node_id": self.data.parent_parallel_start_node_id,272 "iteration_id": self.data.iteration_id,273 },274 }275 276 277class NodeFinishStreamResponse(StreamResponse):278 """279 NodeFinishStreamResponse entity280 """281 282 class Data(BaseModel):283 """284 Data entity285 """286 287 id: str288 node_id: str289 node_type: str290 title: str291 index: int292 predecessor_node_id: Optional[str] = None293 inputs: Optional[dict] = None294 process_data: Optional[dict] = None295 outputs: Optional[dict] = None296 status: str297 error: Optional[str] = None298 elapsed_time: float299 execution_metadata: Optional[dict] = None300 created_at: int301 finished_at: int302 files: Optional[Sequence[Mapping[str, Any]]] = []303 parallel_id: Optional[str] = None304 parallel_start_node_id: Optional[str] = None305 parent_parallel_id: Optional[str] = None306 parent_parallel_start_node_id: Optional[str] = None307 iteration_id: Optional[str] = None308 309 event: StreamEvent = StreamEvent.NODE_FINISHED310 workflow_run_id: str311 data: Data312 313 def to_ignore_detail_dict(self):314 return {315 "event": self.event.value,316 "task_id": self.task_id,317 "workflow_run_id": self.workflow_run_id,318 "data": {319 "id": self.data.id,320 "node_id": self.data.node_id,321 "node_type": self.data.node_type,322 "title": self.data.title,323 "index": self.data.index,324 "predecessor_node_id": self.data.predecessor_node_id,325 "inputs": None,326 "process_data": None,327 "outputs": None,328 "status": self.data.status,329 "error": None,330 "elapsed_time": self.data.elapsed_time,331 "execution_metadata": None,332 "created_at": self.data.created_at,333 "finished_at": self.data.finished_at,334 "files": [],335 "parallel_id": self.data.parallel_id,336 "parallel_start_node_id": self.data.parallel_start_node_id,337 "parent_parallel_id": self.data.parent_parallel_id,338 "parent_parallel_start_node_id": self.data.parent_parallel_start_node_id,339 "iteration_id": self.data.iteration_id,340 },341 }342 343 344class ParallelBranchStartStreamResponse(StreamResponse):345 """346 ParallelBranchStartStreamResponse entity347 """348 349 class Data(BaseModel):350 """351 Data entity352 """353 354 parallel_id: str355 parallel_branch_id: str356 parent_parallel_id: Optional[str] = None357 parent_parallel_start_node_id: Optional[str] = None358 iteration_id: Optional[str] = None359 created_at: int360 361 event: StreamEvent = StreamEvent.PARALLEL_BRANCH_STARTED362 workflow_run_id: str363 data: Data364 365 366class ParallelBranchFinishedStreamResponse(StreamResponse):367 """368 ParallelBranchFinishedStreamResponse entity369 """370 371 class Data(BaseModel):372 """373 Data entity374 """375 376 parallel_id: str377 parallel_branch_id: str378 parent_parallel_id: Optional[str] = None379 parent_parallel_start_node_id: Optional[str] = None380 iteration_id: Optional[str] = None381 status: str382 error: Optional[str] = None383 created_at: int384 385 event: StreamEvent = StreamEvent.PARALLEL_BRANCH_FINISHED386 workflow_run_id: str387 data: Data388 389 390class IterationNodeStartStreamResponse(StreamResponse):391 """392 NodeStartStreamResponse entity393 """394 395 class Data(BaseModel):396 """397 Data entity398 """399 400 id: str401 node_id: str402 node_type: str403 title: str404 created_at: int405 extras: dict = {}406 metadata: dict = {}407 inputs: dict = {}408 parallel_id: Optional[str] = None409 parallel_start_node_id: Optional[str] = None410 411 event: StreamEvent = StreamEvent.ITERATION_STARTED412 workflow_run_id: str413 data: Data414 415 416class IterationNodeNextStreamResponse(StreamResponse):417 """418 NodeStartStreamResponse entity419 """420 421 class Data(BaseModel):422 """423 Data entity424 """425 426 id: str427 node_id: str428 node_type: str429 title: str430 index: int431 created_at: int432 pre_iteration_output: Optional[Any] = None433 extras: dict = {}434 parallel_id: Optional[str] = None435 parallel_start_node_id: Optional[str] = None436 parallel_mode_run_id: Optional[str] = None437 438 event: StreamEvent = StreamEvent.ITERATION_NEXT439 workflow_run_id: str440 data: Data441 442 443class IterationNodeCompletedStreamResponse(StreamResponse):444 """445 NodeCompletedStreamResponse entity446 """447 448 class Data(BaseModel):449 """450 Data entity451 """452 453 id: str454 node_id: str455 node_type: str456 title: str457 outputs: Optional[dict] = None458 created_at: int459 extras: Optional[dict] = None460 inputs: Optional[dict] = None461 status: WorkflowNodeExecutionStatus462 error: Optional[str] = None463 elapsed_time: float464 total_tokens: int465 execution_metadata: Optional[dict] = None466 finished_at: int467 steps: int468 parallel_id: Optional[str] = None469 parallel_start_node_id: Optional[str] = None470 471 event: StreamEvent = StreamEvent.ITERATION_COMPLETED472 workflow_run_id: str473 data: Data474 475 476class TextChunkStreamResponse(StreamResponse):477 """478 TextChunkStreamResponse entity479 """480 481 class Data(BaseModel):482 """483 Data entity484 """485 486 text: str487 from_variable_selector: Optional[list[str]] = None488 489 event: StreamEvent = StreamEvent.TEXT_CHUNK490 data: Data491 492 493class TextReplaceStreamResponse(StreamResponse):494 """495 TextReplaceStreamResponse entity496 """497 498 class Data(BaseModel):499 """500 Data entity501 """502 503 text: str504 505 event: StreamEvent = StreamEvent.TEXT_REPLACE506 data: Data507 508 509class PingStreamResponse(StreamResponse):510 """511 PingStreamResponse entity512 """513 514 event: StreamEvent = StreamEvent.PING515 516 517class AppStreamResponse(BaseModel):518 """519 AppStreamResponse entity520 """521 522 stream_response: StreamResponse523 524 525class ChatbotAppStreamResponse(AppStreamResponse):526 """527 ChatbotAppStreamResponse entity528 """529 530 conversation_id: str531 message_id: str532 created_at: int533 534 535class CompletionAppStreamResponse(AppStreamResponse):536 """537 CompletionAppStreamResponse entity538 """539 540 message_id: str541 created_at: int542 543 544class WorkflowAppStreamResponse(AppStreamResponse):545 """546 WorkflowAppStreamResponse entity547 """548 549 workflow_run_id: Optional[str] = None550 551 552class AppBlockingResponse(BaseModel):553 """554 AppBlockingResponse entity555 """556 557 task_id: str558 559 def to_dict(self) -> dict:560 return jsonable_encoder(self)561 562 563class ChatbotAppBlockingResponse(AppBlockingResponse):564 """565 ChatbotAppBlockingResponse entity566 """567 568 class Data(BaseModel):569 """570 Data entity571 """572 573 id: str574 mode: str575 conversation_id: str576 message_id: str577 answer: str578 metadata: dict = {}579 created_at: int580 581 data: Data582 583 584class CompletionAppBlockingResponse(AppBlockingResponse):585 """586 CompletionAppBlockingResponse entity587 """588 589 class Data(BaseModel):590 """591 Data entity592 """593 594 id: str595 mode: str596 message_id: str597 answer: str598 metadata: dict = {}599 created_at: int600 601 data: Data602 603 604class WorkflowAppBlockingResponse(AppBlockingResponse):605 """606 WorkflowAppBlockingResponse entity607 """608 609 class Data(BaseModel):610 """611 Data entity612 """613 614 id: str615 workflow_id: str616 status: str617 outputs: Optional[dict] = None618 error: Optional[str] = None619 elapsed_time: float620 total_tokens: int621 total_steps: int622 created_at: int623 finished_at: int624 625 workflow_run_id: str626 data: Data627 