2024-04-08 18:51:46 +08:00
|
|
|
from enum import Enum
|
|
|
|
from typing import Any, Optional
|
|
|
|
|
2024-06-14 01:05:37 +08:00
|
|
|
from pydantic import BaseModel, field_validator
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
from core.model_runtime.entities.llm_entities import LLMResult, LLMResultChunk
|
|
|
|
from core.workflow.entities.base_node_data_entities import BaseNodeData
|
|
|
|
from core.workflow.entities.node_entities import NodeType
|
|
|
|
|
|
|
|
|
2024-06-14 01:05:37 +08:00
|
|
|
class QueueEvent(str, Enum):
|
2024-04-08 18:51:46 +08:00
|
|
|
"""
|
|
|
|
QueueEvent enum
|
|
|
|
"""
|
|
|
|
LLM_CHUNK = "llm_chunk"
|
|
|
|
TEXT_CHUNK = "text_chunk"
|
|
|
|
AGENT_MESSAGE = "agent_message"
|
|
|
|
MESSAGE_REPLACE = "message_replace"
|
|
|
|
MESSAGE_END = "message_end"
|
|
|
|
ADVANCED_CHAT_MESSAGE_END = "advanced_chat_message_end"
|
|
|
|
WORKFLOW_STARTED = "workflow_started"
|
|
|
|
WORKFLOW_SUCCEEDED = "workflow_succeeded"
|
|
|
|
WORKFLOW_FAILED = "workflow_failed"
|
2024-05-27 22:01:11 +08:00
|
|
|
ITERATION_START = "iteration_start"
|
|
|
|
ITERATION_NEXT = "iteration_next"
|
|
|
|
ITERATION_COMPLETED = "iteration_completed"
|
2024-04-08 18:51:46 +08:00
|
|
|
NODE_STARTED = "node_started"
|
|
|
|
NODE_SUCCEEDED = "node_succeeded"
|
|
|
|
NODE_FAILED = "node_failed"
|
|
|
|
RETRIEVER_RESOURCES = "retriever_resources"
|
|
|
|
ANNOTATION_REPLY = "annotation_reply"
|
|
|
|
AGENT_THOUGHT = "agent_thought"
|
|
|
|
MESSAGE_FILE = "message_file"
|
|
|
|
ERROR = "error"
|
|
|
|
PING = "ping"
|
|
|
|
STOP = "stop"
|
|
|
|
|
|
|
|
|
|
|
|
class AppQueueEvent(BaseModel):
|
|
|
|
"""
|
|
|
|
QueueEvent entity
|
|
|
|
"""
|
|
|
|
event: QueueEvent
|
|
|
|
|
|
|
|
|
|
|
|
class QueueLLMChunkEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueLLMChunkEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.LLM_CHUNK
|
2024-04-08 18:51:46 +08:00
|
|
|
chunk: LLMResultChunk
|
|
|
|
|
2024-05-27 22:01:11 +08:00
|
|
|
class QueueIterationStartEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueIterationStartEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.ITERATION_START
|
2024-05-27 22:01:11 +08:00
|
|
|
node_id: str
|
|
|
|
node_type: NodeType
|
|
|
|
node_data: BaseNodeData
|
|
|
|
|
|
|
|
node_run_index: int
|
|
|
|
inputs: dict = None
|
|
|
|
predecessor_node_id: Optional[str] = None
|
|
|
|
metadata: Optional[dict] = None
|
|
|
|
|
|
|
|
class QueueIterationNextEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueIterationNextEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.ITERATION_NEXT
|
2024-05-27 22:01:11 +08:00
|
|
|
|
|
|
|
index: int
|
|
|
|
node_id: str
|
|
|
|
node_type: NodeType
|
|
|
|
|
|
|
|
node_run_index: int
|
2024-06-14 01:05:37 +08:00
|
|
|
output: Optional[Any] = None # output for the current iteration
|
2024-05-27 22:01:11 +08:00
|
|
|
|
2024-06-14 01:05:37 +08:00
|
|
|
@field_validator('output', mode='before')
|
2024-06-17 10:04:28 +08:00
|
|
|
@classmethod
|
2024-05-27 22:01:11 +08:00
|
|
|
def set_output(cls, v):
|
|
|
|
"""
|
|
|
|
Set output
|
|
|
|
"""
|
|
|
|
if v is None:
|
|
|
|
return None
|
|
|
|
if isinstance(v, int | float | str | bool | dict | list):
|
|
|
|
return v
|
|
|
|
raise ValueError('output must be a valid type')
|
|
|
|
|
|
|
|
class QueueIterationCompletedEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueIterationCompletedEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event:QueueEvent = QueueEvent.ITERATION_COMPLETED
|
2024-05-27 22:01:11 +08:00
|
|
|
|
|
|
|
node_id: str
|
|
|
|
node_type: NodeType
|
|
|
|
|
|
|
|
node_run_index: int
|
|
|
|
outputs: dict
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
class QueueTextChunkEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueTextChunkEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.TEXT_CHUNK
|
2024-04-08 18:51:46 +08:00
|
|
|
text: str
|
|
|
|
metadata: Optional[dict] = None
|
|
|
|
|
|
|
|
|
|
|
|
class QueueAgentMessageEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueMessageEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.AGENT_MESSAGE
|
2024-04-08 18:51:46 +08:00
|
|
|
chunk: LLMResultChunk
|
|
|
|
|
|
|
|
|
|
|
|
class QueueMessageReplaceEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueMessageReplaceEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.MESSAGE_REPLACE
|
2024-04-08 18:51:46 +08:00
|
|
|
text: str
|
|
|
|
|
|
|
|
|
|
|
|
class QueueRetrieverResourcesEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueRetrieverResourcesEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.RETRIEVER_RESOURCES
|
2024-04-08 18:51:46 +08:00
|
|
|
retriever_resources: list[dict]
|
|
|
|
|
|
|
|
|
|
|
|
class QueueAnnotationReplyEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueAnnotationReplyEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.ANNOTATION_REPLY
|
2024-04-08 18:51:46 +08:00
|
|
|
message_annotation_id: str
|
|
|
|
|
|
|
|
|
|
|
|
class QueueMessageEndEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueMessageEndEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.MESSAGE_END
|
2024-04-08 18:51:46 +08:00
|
|
|
llm_result: Optional[LLMResult] = None
|
|
|
|
|
|
|
|
|
|
|
|
class QueueAdvancedChatMessageEndEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueAdvancedChatMessageEndEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.ADVANCED_CHAT_MESSAGE_END
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
|
|
|
|
class QueueWorkflowStartedEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueWorkflowStartedEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.WORKFLOW_STARTED
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
|
|
|
|
class QueueWorkflowSucceededEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueWorkflowSucceededEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.WORKFLOW_SUCCEEDED
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
|
|
|
|
class QueueWorkflowFailedEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueWorkflowFailedEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.WORKFLOW_FAILED
|
2024-04-08 18:51:46 +08:00
|
|
|
error: str
|
|
|
|
|
|
|
|
|
|
|
|
class QueueNodeStartedEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueNodeStartedEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.NODE_STARTED
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
node_id: str
|
|
|
|
node_type: NodeType
|
|
|
|
node_data: BaseNodeData
|
|
|
|
node_run_index: int = 1
|
|
|
|
predecessor_node_id: Optional[str] = None
|
|
|
|
|
|
|
|
|
|
|
|
class QueueNodeSucceededEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueNodeSucceededEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.NODE_SUCCEEDED
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
node_id: str
|
|
|
|
node_type: NodeType
|
|
|
|
node_data: BaseNodeData
|
|
|
|
|
|
|
|
inputs: Optional[dict] = None
|
|
|
|
process_data: Optional[dict] = None
|
|
|
|
outputs: Optional[dict] = None
|
|
|
|
execution_metadata: Optional[dict] = None
|
|
|
|
|
|
|
|
error: Optional[str] = None
|
|
|
|
|
|
|
|
|
|
|
|
class QueueNodeFailedEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueNodeFailedEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.NODE_FAILED
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
node_id: str
|
|
|
|
node_type: NodeType
|
|
|
|
node_data: BaseNodeData
|
|
|
|
|
|
|
|
inputs: Optional[dict] = None
|
|
|
|
outputs: Optional[dict] = None
|
|
|
|
process_data: Optional[dict] = None
|
|
|
|
|
|
|
|
error: str
|
|
|
|
|
|
|
|
|
|
|
|
class QueueAgentThoughtEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueAgentThoughtEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.AGENT_THOUGHT
|
2024-04-08 18:51:46 +08:00
|
|
|
agent_thought_id: str
|
|
|
|
|
|
|
|
|
|
|
|
class QueueMessageFileEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueAgentThoughtEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.MESSAGE_FILE
|
2024-04-08 18:51:46 +08:00
|
|
|
message_file_id: str
|
|
|
|
|
|
|
|
|
|
|
|
class QueueErrorEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueErrorEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.ERROR
|
|
|
|
error: Any = None
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
|
|
|
|
class QueuePingEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueuePingEvent entity
|
|
|
|
"""
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.PING
|
2024-04-08 18:51:46 +08:00
|
|
|
|
|
|
|
|
|
|
|
class QueueStopEvent(AppQueueEvent):
|
|
|
|
"""
|
|
|
|
QueueStopEvent entity
|
|
|
|
"""
|
|
|
|
class StopBy(Enum):
|
|
|
|
"""
|
|
|
|
Stop by enum
|
|
|
|
"""
|
|
|
|
USER_MANUAL = "user-manual"
|
|
|
|
ANNOTATION_REPLY = "annotation-reply"
|
|
|
|
OUTPUT_MODERATION = "output-moderation"
|
|
|
|
INPUT_MODERATION = "input-moderation"
|
|
|
|
|
2024-06-14 01:05:37 +08:00
|
|
|
event: QueueEvent = QueueEvent.STOP
|
2024-04-08 18:51:46 +08:00
|
|
|
stopped_by: StopBy
|
|
|
|
|
|
|
|
|
|
|
|
class QueueMessage(BaseModel):
|
|
|
|
"""
|
|
|
|
QueueMessage entity
|
|
|
|
"""
|
|
|
|
task_id: str
|
|
|
|
app_mode: str
|
|
|
|
event: AppQueueEvent
|
|
|
|
|
|
|
|
|
|
|
|
class MessageQueueMessage(QueueMessage):
|
|
|
|
"""
|
|
|
|
MessageQueueMessage entity
|
|
|
|
"""
|
|
|
|
message_id: str
|
|
|
|
conversation_id: str
|
|
|
|
|
|
|
|
|
|
|
|
class WorkflowQueueMessage(QueueMessage):
|
|
|
|
"""
|
|
|
|
WorkflowQueueMessage entity
|
|
|
|
"""
|
|
|
|
pass
|