feat(api): 新增核心 API 端点 + 数据库迁移
新增 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 等)
This commit is contained in:
@ -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';
|
||||||
|
""")
|
||||||
210
backend/app/api/v1/endpoints/tasks.py
Normal file
210
backend/app/api/v1/endpoints/tasks.py
Normal file
@ -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))
|
||||||
Reference in New Issue
Block a user