Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions backend/packages/app/src/windup_app/bootstrap/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,13 @@
from windup_app.server.character.model import Character # noqa: F401
from windup_app.server.project.model import Project # noqa: F401
from windup_app.server.user.model import User # noqa: F401
from windup_app.server.workflow_run.model import WorkflowRun # noqa: F401
from windup_app.web.api.auth import router as auth_router
from windup_app.web.api.character import router as character_router
from windup_app.web.api.generation import router as generation_router
from windup_app.web.api.media import router as media_router
from windup_app.web.api.project import router as project_router
from windup_app.web.api.workflow_run import router as workflow_run_router
from windup_app.web.handler.exception_handlers import register_exception_handlers
from windup_app.web.middleware.auth import AuthMiddleware
from windup_app.web.middleware.ratelimit import RateLimitMiddleware
Expand Down Expand Up @@ -85,6 +87,7 @@ def create_app() -> FastAPI:
app.include_router(auth_router)
app.include_router(project_router)
app.include_router(character_router)
app.include_router(workflow_run_router)
app.include_router(media_router)
app.include_router(generation_router)
register_exception_handlers(app)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,20 +18,18 @@

from abc import ABC, abstractmethod

from windup_app.server.workflow_run.model import (
RunStatus,
WorkflowRun,
)
from sqlalchemy.orm import Session

from windup_app.server.workflow_run.model import RunStatus, WorkflowRun


class WorkflowRunService(ABC):
"""执行记录用例的抽象边界。"""

# -- 执行记录 CRUD --------------------------------------------------------

@abstractmethod
def create_run(
self,
session: Session,
*,
project_id: int,
nodes: list | None = None,
Expand All @@ -42,22 +40,35 @@ def create_run(
"""

@abstractmethod
def get_run(self, run_id: int) -> WorkflowRun | None:
def get_run(self, session: Session, run_id: int) -> WorkflowRun | None:
"""获取执行记录详情(含 nodes JSONB)。"""

@abstractmethod
def list_runs(
self,
session: Session,
*,
project_id: int,
page: int = 1,
page_size: int = 20,
) -> tuple[list[WorkflowRun], int]:
"""分页查询项目下的执行记录,返回 (当前页数据, 总数)。"""

@abstractmethod
def update_run(
self,
session: Session,
run_id: int,
*,
nodes: list | None = None,
status: RunStatus | None = None,
) -> WorkflowRun:
"""全量更新执行记录
) -> WorkflowRun | None:
"""更新执行记录

前端维护节点树后,通过此接口全量写回。
返回更新后的记录;不存在时返回 None。
"""

@abstractmethod
def delete_run(self, run_id: int) -> None:
"""软删除执行记录。"""
def delete_run(self, session: Session, run_id: int) -> bool:
"""软删除执行记录。返回是否找到。"""
50 changes: 39 additions & 11 deletions backend/packages/app/src/windup_app/server/workflow_run/model.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,15 @@

from __future__ import annotations

from dataclasses import dataclass, field
from datetime import datetime, timezone
from enum import StrEnum

from sqlalchemy import BigInteger, DateTime, Integer, JSON, String
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.orm import Mapped, mapped_column

from windup_framework.db import Base


# -- 枚举 ----------------------------------------------------------------

Expand All @@ -21,19 +26,42 @@ class RunStatus(StrEnum):
SOFT_DELETED = "soft_deleted"


# -- 执行记录 -------------------------------------------------------------
# -- ORM -----------------------------------------------------------------


@dataclass
class WorkflowRun:
"""执行记录——前端维护的节点树的持久化容器。
class WorkflowRun(Base):
"""执行记录表——前端维护的节点树的持久化容器。

后端不校验 nodes 内部结构,仅做全量读写。
"""

id: int | None = None
project_id: int = 0
nodes: list = field(default_factory=list) # 节点树(前端自定义结构,后端不校验)
status: RunStatus = RunStatus.ACTIVE
version: int = 1
created_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc))
__tablename__ = "windup_workflow_run"

id: Mapped[int] = mapped_column(
BigInteger().with_variant(Integer, "sqlite"),
primary_key=True,
autoincrement=True,
)

project_id: Mapped[int] = mapped_column(BigInteger, nullable=False)

# 节点树(前端自定义结构,后端不校验);Postgres 上 JSONB,SQLite 上 JSON。
nodes: Mapped[list] = mapped_column(
JSON().with_variant(JSONB, "postgresql"),
nullable=False,
default=list,
)

status: Mapped[str] = mapped_column(
String(20), nullable=False, default=RunStatus.ACTIVE.value,
)

version: Mapped[int] = mapped_column(
Integer, nullable=False, default=1,
)

created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True),
nullable=False,
default=lambda: datetime.now(timezone.utc),
)
98 changes: 98 additions & 0 deletions backend/packages/app/src/windup_app/server/workflow_run/service.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
"""工作流执行记录领域服务的 SQLAlchemy 实现。

:class:`SqlAlchemyWorkflowRunService` 继承 :class:`WorkflowRunService` 接口,用同步
SQLAlchemy session 落库。无状态:``session`` 由调用方按请求传入,本对象可作
模块级单例(:data:`service`)。

事务边界由 ``windup_framework.db.get_session`` 依赖负责——成功 commit、异常
rollback,故本实现只 ``flush``(把变更发到当前事务、取回生成的主键),不 commit。
"""

from sqlalchemy import func, select
from sqlalchemy.orm import Session

from windup_app.server.workflow_run.interface import WorkflowRunService
from windup_app.server.workflow_run.model import RunStatus, WorkflowRun


class SqlAlchemyWorkflowRunService(WorkflowRunService):
"""基于 SQLAlchemy session 的执行记录 CRUD 实现。"""

def create_run(
self,
session: Session,
*,
project_id: int,
nodes: list | None = None,
) -> WorkflowRun:
run = WorkflowRun(
project_id=project_id,
nodes=nodes or [],
)
session.add(run)
session.flush()
return run

def get_run(self, session: Session, run_id: int) -> WorkflowRun | None:
return session.get(WorkflowRun, run_id)

def list_runs(
self,
session: Session,
*,
project_id: int,
page: int = 1,
page_size: int = 20,
) -> tuple[list[WorkflowRun], int]:
"""分页查询项目下的执行记录,返回 (当前页数据, 总数)。"""
count_stmt = (
select(func.count())
.select_from(WorkflowRun)
.where(
WorkflowRun.project_id == project_id,
WorkflowRun.status != RunStatus.SOFT_DELETED.value,
)
)
stmt = (
select(WorkflowRun)
.where(
WorkflowRun.project_id == project_id,
WorkflowRun.status != RunStatus.SOFT_DELETED.value,
)
.order_by(WorkflowRun.id.desc())
.offset((page - 1) * page_size)
.limit(page_size)
)
total = session.scalar(count_stmt) or 0
items = list(session.scalars(stmt))
return items, total

def update_run(
self,
session: Session,
run_id: int,
*,
nodes: list | None = None,
status: RunStatus | None = None,
) -> WorkflowRun | None:
run = session.get(WorkflowRun, run_id)
if run is None:
return None
if nodes is not None:
run.nodes = nodes
if status is not None:
run.status = status.value
run.version += 1
Comment thread
huyanxius marked this conversation as resolved.
session.flush()
return run

def delete_run(self, session: Session, run_id: int) -> bool:
run = session.get(WorkflowRun, run_id)
if run is None:
return False
run.status = RunStatus.SOFT_DELETED.value
session.flush()
return True


service = SqlAlchemyWorkflowRunService()
Loading