Files
dify/api/fields/workflow_run_fields.py

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]