From 9ec4f8efa764a3318f88721ec4879d2ca9ceb1b4 Mon Sep 17 00:00:00 2001 From: duxingchen Date: Tue, 4 Aug 2026 17:03:43 +0800 Subject: [PATCH] =?UTF-8?q?feat(api):=20=E6=96=B0=E5=A2=9E=E6=A0=B8?= =?UTF-8?q?=E5=BF=83=20API=20=E7=AB=AF=E7=82=B9=20+=20=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E5=BA=93=E8=BF=81=E7=A7=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增 3 个 API 端点: - POST /tasks/{id}/receive: 确认接收 (PENDING→WIP) - POST /tasks/{id}/reject: 品质驳回 + 返工闭环 - POST /tasks/{id}/transfer: 完工裂变转交 (多路分支+入库) 数据库迁移 (a1b2c3d4e5f6): - tasks 表新增 received_at, completed_at, reject_reason, is_rework - products 表新增 current_location_id - 数据迁移: 已有任务状态小写→大写 (pending→PENDING 等) --- ...f6_add_task_rework_and_product_location.py | 102 +++++++++ backend/app/api/v1/endpoints/tasks.py | 210 ++++++++++++++++++ 2 files changed, 312 insertions(+) create mode 100644 backend/alembic/versions/a1b2c3d4e5f6_add_task_rework_and_product_location.py create mode 100644 backend/app/api/v1/endpoints/tasks.py diff --git a/backend/alembic/versions/a1b2c3d4e5f6_add_task_rework_and_product_location.py b/backend/alembic/versions/a1b2c3d4e5f6_add_task_rework_and_product_location.py new file mode 100644 index 0000000..1941125 --- /dev/null +++ b/backend/alembic/versions/a1b2c3d4e5f6_add_task_rework_and_product_location.py @@ -0,0 +1,102 @@ +"""add task rework fields and product current_location + +Revision ID: a1b2c3d4e5f6 +Revises: 80f8cb64544a +Create Date: 2026-08-04 12:00:00.000000 + +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = 'a1b2c3d4e5f6' +down_revision: Union[str, Sequence[str], None] = '80f8cb64544a' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + """Upgrade schema.""" + # === 数据迁移:将已有任务状态从小写转为大写 === + op.execute(""" + UPDATE tasks SET status = 'PENDING' WHERE lower(status) = 'pending'; + """) + op.execute(""" + UPDATE tasks SET status = 'WIP' WHERE lower(status) = 'wip'; + """) + op.execute(""" + UPDATE tasks SET status = 'COMPLETED' WHERE lower(status) = 'completed'; + """) + op.execute(""" + UPDATE tasks SET status = 'REJECTED' WHERE lower(status) = 'rejected'; + """) + op.execute(""" + UPDATE tasks SET status = 'ARCHIVED' WHERE lower(status) = 'archived'; + """) + + # === tasks 表新增字段 === + op.add_column('tasks', sa.Column( + 'received_at', + sa.DateTime(timezone=True), + nullable=True, + comment='操作员确认接收时间', + )) + op.add_column('tasks', sa.Column( + 'completed_at', + sa.DateTime(timezone=True), + nullable=True, + comment='任务完工转交时间', + )) + op.add_column('tasks', sa.Column( + 'reject_reason', + sa.String(length=500), + nullable=True, + comment='驳回原因', + )) + op.add_column('tasks', sa.Column( + 'is_rework', + sa.Boolean(), + nullable=False, + server_default=sa.text('false'), + comment='是否为返工任务', + )) + + # === products 表新增字段 === + op.add_column('products', sa.Column( + 'current_location_id', + sa.String(length=64), + nullable=True, + comment="当前持有者ID 或 'virtual_warehouse'(仓库)", + )) + + +def downgrade() -> None: + """Downgrade schema.""" + # === products 表移除字段 === + op.drop_column('products', 'current_location_id') + + # === tasks 表移除字段 === + op.drop_column('tasks', 'is_rework') + op.drop_column('tasks', 'reject_reason') + op.drop_column('tasks', 'completed_at') + op.drop_column('tasks', 'received_at') + + # === 数据回迁:将任务状态从大写转回小写 === + op.execute(""" + UPDATE tasks SET status = 'pending' WHERE status = 'PENDING'; + """) + op.execute(""" + UPDATE tasks SET status = 'wip' WHERE status = 'WIP'; + """) + op.execute(""" + UPDATE tasks SET status = 'completed' WHERE status = 'COMPLETED'; + """) + op.execute(""" + UPDATE tasks SET status = 'rejected' WHERE status = 'REJECTED'; + """) + op.execute(""" + UPDATE tasks SET status = 'archived' WHERE status = 'ARCHIVED'; + """) diff --git a/backend/app/api/v1/endpoints/tasks.py b/backend/app/api/v1/endpoints/tasks.py new file mode 100644 index 0000000..5066f03 --- /dev/null +++ b/backend/app/api/v1/endpoints/tasks.py @@ -0,0 +1,210 @@ +"""任务 API 端点 — 核心业务:接收、驳回返工、裂变转交、无限嵌套子任务""" +from __future__ import annotations +import uuid +from fastapi import APIRouter, Depends, Query +from sqlalchemy.ext.asyncio import AsyncSession + +from app.core.database import get_db +from app.schemas.task import ( + TaskCreate, + TaskUpdate, + TaskCompleteRequest, + TaskRejectRequest, + TaskTransferRequest, + SubtaskCreate, + TaskResponse, + TaskCompleteResponse, + TaskTransferResponse, + TaskSummaryResponse, + TaskListResponse, +) +from app.services import task_service + +router = APIRouter(prefix="/tasks", tags=["任务管理"]) + + +# ============================================================ +# 任务 CRUD +# ============================================================ + +@router.get("/", response_model=TaskListResponse) +async def list_tasks( + product_id: str | None = Query(None, description="按产品ID筛选"), + skip: int = Query(0, ge=0), + limit: int = Query(50, ge=1, le=200), + db: AsyncSession = Depends(get_db), +): + """获取任务列表,可按产品筛选(只返回顶层任务)""" + pid = uuid.UUID(product_id) if product_id else None + return await task_service.get_all_tasks(db, product_id=pid, skip=skip, limit=limit) + + +@router.get("/{task_id}", response_model=TaskResponse) +async def get_task( + task_id: str, + db: AsyncSession = Depends(get_db), +): + """ + 获取任务详情 — 递归包含所有层级的子任务。 + 前端可根据此结果渲染完整的任务树。 + """ + return await task_service.get_task(db, uuid.UUID(task_id)) + + +@router.post("/", response_model=TaskResponse, status_code=201) +async def create_task_endpoint( + data: TaskCreate, + db: AsyncSession = Depends(get_db), +): + """创建任务""" + return await task_service.create_task(db, data) + + +@router.patch("/{task_id}", response_model=TaskResponse) +async def update_task_endpoint( + task_id: str, + data: TaskUpdate, + db: AsyncSession = Depends(get_db), +): + """更新任务""" + return await task_service.update_task(db, uuid.UUID(task_id), data) + + +# ============================================================ +# 核心卡点逻辑:任务完成 / 转交 +# ============================================================ + +@router.post("/{task_id}/complete", response_model=TaskCompleteResponse) +async def complete_task_endpoint( + task_id: str, + request: TaskCompleteRequest, + db: AsyncSession = Depends(get_db), +): + """ + **核心接口:完成任务 + 可选创建下一步任务(转交)** + + 卡点逻辑: + 1. 检查当前任务是否已完成(幂等保护) + 2. 查询所有 `notify_parent_on_complete=True` 的子任务 + → 如果存在未完成的,返回 HTTP 400:「请等待相关子任务完成」 + 3. 全部通过后,标记任务为 completed,写入操作日志 + 4. 若提供了 `next_task_name` + `next_assignee_id`,自动创建下一步任务 + + 典型场景: + - 某个加工步骤完成,需要检查所有必须的前置工序(子任务)是否已完成 + - 完成后自动创建下一步任务并指定负责人 + """ + return await task_service.complete_task( + db, uuid.UUID(task_id), request + ) + + +# ============================================================ +# 核心业务 1:确认接收 (PENDING → WIP) +# ============================================================ + +@router.post("/{task_id}/receive", response_model=TaskResponse) +async def receive_task_endpoint( + task_id: str, + operator_id: str | None = Query(None, description="操作人ID"), + db: AsyncSession = Depends(get_db), +): + """ + **确认接收任务。** + + 校验:只有状态为 PENDING 的任务可接收。 + 动作:将状态改为 WIP,记录 received_at 为当前时间。 + """ + return await task_service.receive_task( + db, uuid.UUID(task_id), operator_id + ) + + +# ============================================================ +# 核心业务 2:品质驳回 (→ REJECTED + 返工闭环) +# ============================================================ + +@router.post("/{task_id}/reject", response_model=TaskResponse) +async def reject_task_endpoint( + task_id: str, + request: TaskRejectRequest, + operator_id: str | None = Query(None, description="操作人ID"), + db: AsyncSession = Depends(get_db), +): + """ + **品质驳回:将任务标记为 REJECTED,自动创建返工任务。** + + 防呆闭环逻辑: + 1. 将当前任务状态改为 REJECTED,记录 reject_reason 和 completed_at。 + 2. 查找上一道工序的负责人(父任务的 assignee_id)。 + 3. 为该负责人新建返工任务(is_rework=True, status=PENDING)。 + """ + return await task_service.reject_task( + db, uuid.UUID(task_id), request, operator_id + ) + + +# ============================================================ +# 核心业务 3:完工并裂变转交 (→ COMPLETED + 多路裂变) +# ============================================================ + +@router.post("/{task_id}/transfer", response_model=TaskTransferResponse) +async def transfer_task_endpoint( + task_id: str, + request: TaskTransferRequest, + operator_id: str | None = Query(None, description="操作人ID"), + db: AsyncSession = Depends(get_db), +): + """ + **完工并裂变转交:完成当前任务,批量创建下一道工序任务。** + + 动作 1(闭环当前节点): + - 将当前任务状态改为 COMPLETED,记录 completed_at。 + + 动作 2(解析下家): + - 遍历 next_assignees 列表。 + - 如果包含 'virtual_warehouse',则将 Product 的 current_location_id 设为仓库。 + - 为每一个 assignee_id 新建 PENDING 任务。 + + 裂变逻辑: + - next_assignees > 1 → 多路裂变,新任务挂在当前任务下形成树状分支。 + - 当前任务是子任务 → 单路转交也保持在同一父任务下。 + - 否则 → 顶层同级转交。 + """ + return await task_service.transfer_task( + db, uuid.UUID(task_id), request, operator_id + ) + + +# ============================================================ +# 无限层级子任务 +# ============================================================ + +@router.post("/{task_id}/subtasks", response_model=TaskResponse, status_code=201) +async def create_subtask_endpoint( + task_id: str, + data: SubtaskCreate, + db: AsyncSession = Depends(get_db), +): + """ + **创建子任务:支持无限层级嵌套。** + + 新子任务将自动继承父任务的 product_id。 + 若父任务已完成,拒绝创建。 + """ + return await task_service.create_subtask( + db, uuid.UUID(task_id), data + ) + + +# ============================================================ +# 查询产品顶层任务(便捷接口) +# ============================================================ + +@router.get("/by-product/{product_id}", response_model=list[TaskSummaryResponse]) +async def get_tasks_by_product( + product_id: str, + db: AsyncSession = Depends(get_db), +): + """获取指定产品的顶层任务列表(不含子任务嵌套)""" + return await task_service.get_top_level_tasks(db, uuid.UUID(product_id))