mirror of
https://github.com/langgenius/dify.git
synced 2026-09-28 14:23:33 +08:00
230 lines
7.7 KiB
Python
230 lines
7.7 KiB
Python
"""Workflow run response schemas for console APIs.
|
|
|
|
Workflow-run endpoints should document and serialize responses with the
|
|
Pydantic models in this module.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import datetime
|
|
from enum import Enum
|
|
from typing import Any
|
|
|
|
from pydantic import AliasChoices, Field, field_validator
|
|
from sqlalchemy.orm import Session
|
|
|
|
from fields.base import ResponseModel, SessionResponseSource
|
|
from fields.end_user_fields import SimpleEndUser
|
|
from fields.member_fields import SimpleAccount
|
|
from libs.helper import to_timestamp
|
|
|
|
|
|
class WorkflowRunForLogResponse(ResponseModel):
|
|
id: str
|
|
version: str | None = None
|
|
status: str | None = None
|
|
triggered_from: str | None = None
|
|
error: str | None = None
|
|
elapsed_time: float | None = None
|
|
total_tokens: int | None = None
|
|
total_steps: int | None = None
|
|
created_at: int | None = None
|
|
finished_at: int | None = None
|
|
exceptions_count: int | None = None
|
|
|
|
@field_validator("status", mode="before")
|
|
@classmethod
|
|
def _normalize_status(cls, value: Any) -> str | None:
|
|
if value is None or isinstance(value, str):
|
|
return value
|
|
return str(getattr(value, "value", value))
|
|
|
|
@field_validator("created_at", "finished_at", mode="before")
|
|
@classmethod
|
|
def _normalize_timestamp(cls, value: datetime | int | None) -> int | None:
|
|
return to_timestamp(value)
|
|
|
|
|
|
class WorkflowRunForArchivedLogResponse(ResponseModel):
|
|
id: str
|
|
status: str | None = None
|
|
triggered_from: str | None = None
|
|
elapsed_time: float | None = None
|
|
total_tokens: int | None = None
|
|
|
|
@field_validator("status", mode="before")
|
|
@classmethod
|
|
def _normalize_status(cls, value: Any) -> str | None:
|
|
if value is None or isinstance(value, str):
|
|
return value
|
|
return str(value.value if isinstance(value, Enum) else value)
|
|
|
|
|
|
class WorkflowRunForListResponse(ResponseModel):
|
|
id: str
|
|
version: str | None = None
|
|
status: str | None = None
|
|
elapsed_time: float | None = None
|
|
total_tokens: int | None = None
|
|
total_steps: int | None = None
|
|
created_by_account: SimpleAccount | None = None
|
|
created_at: int | None = None
|
|
finished_at: int | None = None
|
|
exceptions_count: int | None = None
|
|
retry_index: int | None = None
|
|
|
|
@field_validator("status", mode="before")
|
|
@classmethod
|
|
def _normalize_status(cls, value: Any) -> str | None:
|
|
if value is None or isinstance(value, str):
|
|
return value
|
|
return str(getattr(value, "value", value))
|
|
|
|
@field_validator("created_at", "finished_at", mode="before")
|
|
@classmethod
|
|
def _normalize_timestamp(cls, value: datetime | int | None) -> int | None:
|
|
return to_timestamp(value)
|
|
|
|
|
|
class AdvancedChatWorkflowRunForListResponse(WorkflowRunForListResponse):
|
|
conversation_id: str | None = None
|
|
message_id: str | None = None
|
|
|
|
|
|
class AdvancedChatWorkflowRunPaginationResponse(ResponseModel):
|
|
limit: int
|
|
has_more: bool
|
|
data: list[AdvancedChatWorkflowRunForListResponse]
|
|
|
|
|
|
class WorkflowRunPaginationResponse(ResponseModel):
|
|
limit: int
|
|
has_more: bool
|
|
data: list[WorkflowRunForListResponse]
|
|
|
|
|
|
class WorkflowRunCountResponse(ResponseModel):
|
|
total: int
|
|
running: int
|
|
succeeded: int
|
|
failed: int
|
|
stopped: int
|
|
partial_succeeded: int = Field(
|
|
alias="partial_succeeded",
|
|
validation_alias=AliasChoices("partial_succeeded", "partial-succeeded"),
|
|
)
|
|
|
|
|
|
class WorkflowRunDetailResponse(ResponseModel):
|
|
id: str
|
|
version: str | None = None
|
|
graph: Any = Field(validation_alias="graph_dict")
|
|
inputs: Any = Field(validation_alias="inputs_dict")
|
|
status: str | None = None
|
|
outputs: Any = Field(validation_alias="outputs_dict")
|
|
error: str | None = None
|
|
elapsed_time: float | None = None
|
|
total_tokens: int | None = None
|
|
total_steps: int | None = None
|
|
created_by_role: str | None = None
|
|
created_by_account: SimpleAccount | None = None
|
|
created_by_end_user: SimpleEndUser | None = None
|
|
created_at: int | None = None
|
|
finished_at: int | None = None
|
|
exceptions_count: int | None = None
|
|
|
|
@field_validator("status", mode="before")
|
|
@classmethod
|
|
def _normalize_status(cls, value: Any) -> str | None:
|
|
if value is None or isinstance(value, str):
|
|
return value
|
|
return str(getattr(value, "value", value))
|
|
|
|
@field_validator("created_at", "finished_at", mode="before")
|
|
@classmethod
|
|
def _normalize_timestamp(cls, value: datetime | int | None) -> int | None:
|
|
return to_timestamp(value)
|
|
|
|
|
|
class WorkflowRunNodeExecutionResponse(ResponseModel):
|
|
id: str
|
|
index: int | None = None
|
|
predecessor_node_id: str | None = None
|
|
node_id: str | None = None
|
|
node_type: str | None = None
|
|
title: str | None = None
|
|
inputs: Any = Field(default=None, validation_alias="inputs_dict")
|
|
process_data: Any = Field(default=None, validation_alias="process_data_dict")
|
|
outputs: Any = Field(default=None, validation_alias="outputs_dict")
|
|
status: str | None = None
|
|
error: str | None = None
|
|
elapsed_time: float | None = None
|
|
execution_metadata: Any = Field(default=None, validation_alias="execution_metadata_dict")
|
|
extras: Any = None
|
|
created_at: int | None = None
|
|
created_by_role: str | None = None
|
|
created_by_account: SimpleAccount | None = None
|
|
created_by_end_user: SimpleEndUser | None = None
|
|
finished_at: int | None = None
|
|
inputs_truncated: bool | None = None
|
|
outputs_truncated: bool | None = None
|
|
process_data_truncated: bool | None = None
|
|
retry_index: int | None = Field(default=None, exclude_if=lambda value: value is None)
|
|
|
|
@field_validator("status", mode="before")
|
|
@classmethod
|
|
def _normalize_status(cls, value: Any) -> str | None:
|
|
if value is None or isinstance(value, str):
|
|
return value
|
|
return str(getattr(value, "value", value))
|
|
|
|
@field_validator("created_at", "finished_at", mode="before")
|
|
@classmethod
|
|
def _normalize_timestamp(cls, value: datetime | int | None) -> int | None:
|
|
return to_timestamp(value)
|
|
|
|
|
|
class WorkflowRunResponseSource(SessionResponseSource[Any]):
|
|
"""Expose session-backed workflow-run accessors during response validation."""
|
|
|
|
@property
|
|
def created_by_account(self) -> Any:
|
|
return self._source.created_by_account(self._session)
|
|
|
|
@property
|
|
def created_by_end_user(self) -> Any:
|
|
return self._source.created_by_end_user(self._session)
|
|
|
|
|
|
def workflow_run_response_source(workflow_run: Any, *, session: Session) -> WorkflowRunResponseSource:
|
|
return WorkflowRunResponseSource(workflow_run, session=session)
|
|
|
|
|
|
def workflow_run_pagination_response_source(pagination: Any, *, session: Session) -> dict[str, Any]:
|
|
"""Wrap each run in a pagination payload so list responses resolve accessors via the session."""
|
|
return {
|
|
"limit": pagination.limit,
|
|
"has_more": pagination.has_more,
|
|
"data": [workflow_run_response_source(run, session=session) for run in pagination.data],
|
|
}
|
|
|
|
|
|
class WorkflowNodeExecutionResponseSource(SessionResponseSource[Any]):
|
|
"""Expose session-backed node-execution accessors during response validation."""
|
|
|
|
@property
|
|
def created_by_account(self) -> Any:
|
|
return self._source.created_by_account(self._session)
|
|
|
|
@property
|
|
def created_by_end_user(self) -> Any:
|
|
return self._source.created_by_end_user(self._session)
|
|
|
|
|
|
def node_execution_response_source(node_execution: Any, *, session: Session) -> WorkflowNodeExecutionResponseSource:
|
|
return WorkflowNodeExecutionResponseSource(node_execution, session=session)
|
|
|
|
|
|
class WorkflowRunNodeExecutionListResponse(ResponseModel):
|
|
data: list[WorkflowRunNodeExecutionResponse]
|