"""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]