第一阶段:模型 → 判定 → /auth/me → 列表。统计接口与前端管理页随后。
1) 模型与迁移(head 从 k1l2m3n4o5p6 推进到 l1m2n3o4p5q6)
· business_groups 组定义,parent_id 表达「大组 > 小组」
· business_group_phases 可见范围,独立成表以支持多选 —— 需求要求
「范围可配置、不要写死」,单列存不下多个 phase
· business_group_members 成员,一人可属多组(这是「同时看生产+维修」的实现)
只建表、不写种子数据,所以可以先部署代码再建组。
2) DataScope 判定模块(app/services/data_scope_service.py)
全仓库唯一的权限谓词来源,业务代码里不准再出现 lifecycle_phase 过滤。
两条红线照抄部门隔离的教训:
· None(不限) 与 frozenset()(空) 语义相反,绝不共用一个哨兵值
· 空集合必须显式 false() —— SQLAlchemy 对 in_(()) 生 成 IN (NULL),
一旦退化成不过滤就是全量泄漏
解析优先级:SUPER_ADMIN 硬放行(不可被分组覆盖)
> 显式分组(分组优先于角色)
> 未分组 SUPERVISOR 默认全厂
> 未分组普通用户
3) 过渡期开关 DATA_SCOPE_UNGROUPED(默认 ALL)
直接上严格模式会让所有未分组工人当场看不到自己的任务、现场停摆。
默认 ALL 先放行并打 WARNING 记录「谁还没分组」,配好组后再改 NONE。
4) /auth/me 返回 scope,phase 中文标签由服务端下发
—— 前端已有两份 phase 词表副本,不再加第三份。
5) 列表接口接入
· get_all_products:过滤加在 offset/limit 之前(其下有 6 段基于 product_ids
的批量预计算,过滤晚了等于算完再丢)
· get_all_tasks:谓词进【共享 filters】,保证 count 与 select 两条独立语句
同时生效,否则 total 与实际页不一致、移动端 hasMore 判断跟着错
· 扫码 get_product_by_serial 刻意不过滤,理由写死在 docstring 里
实测:空 scope 生成 false、受限 scope 生成 JOIN + IN 谓词;
/products 返回 2 条、/tasks 的 total 与 returned 一致。
331 lines
12 KiB
Python
331 lines
12 KiB
Python
"""任务 API 端点 — 核心业务:接收、驳回返工、裂变转交、无限嵌套子任务"""
|
||
from __future__ import annotations
|
||
import uuid
|
||
from fastapi import APIRouter, Depends, Query, Body
|
||
from pydantic import BaseModel, Field
|
||
from app.services.auth_service import get_current_user
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.core.database import get_db
|
||
from app.core.deps import get_data_scope
|
||
from app.services.data_scope_service import DataScope
|
||
from app.schemas.task import (
|
||
TaskCreate,
|
||
TaskUpdate,
|
||
TaskCompleteRequest,
|
||
TaskRejectRequest,
|
||
TaskTransferRequest,
|
||
TaskTransferBranch,
|
||
SubtaskCreate,
|
||
TaskRecordCreate,
|
||
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筛选"),
|
||
assignee_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),
|
||
current_user: dict = Depends(get_current_user),
|
||
scope: DataScope = Depends(get_data_scope),
|
||
):
|
||
"""获取任务列表,可按产品/负责人筛选(只返回顶层任务)
|
||
|
||
结果受业务分组数据范围约束(见 app/services/data_scope_service.py)。
|
||
"""
|
||
pid = uuid.UUID(product_id) if product_id else None
|
||
return await task_service.get_all_tasks(
|
||
db, scope=scope, product_id=pid, assignee_id=assignee_id, skip=skip, limit=limit,
|
||
)
|
||
|
||
|
||
@router.get("/{task_id}", response_model=TaskResponse)
|
||
async def get_task(
|
||
task_id: str,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
获取任务详情 — 递归包含所有层级的子任务。
|
||
前端可根据此结果渲染完整的任务树。
|
||
"""
|
||
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),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""创建任务"""
|
||
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),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""更新任务"""
|
||
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),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**核心接口:完成任务 + 可选创建下一步任务(转交)**
|
||
|
||
卡点逻辑:
|
||
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,
|
||
operator_role=current_user.get("role"),
|
||
)
|
||
|
||
|
||
# ============================================================
|
||
# 核心业务 0:结束分支(终止当前节点,不创建下游)
|
||
# ============================================================
|
||
|
||
@router.post("/{task_id}/end", response_model=TaskResponse)
|
||
async def end_task_endpoint(
|
||
task_id: str,
|
||
operator_id: str | None = Query(None),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**结束当前分支:标记任务为 COMPLETED,不创建下游任务。**
|
||
|
||
用于工人认为工序已完结、无需转交下一人的场景。
|
||
"""
|
||
return await task_service.end_task(
|
||
db, uuid.UUID(task_id),
|
||
operator_id or current_user.get("username", "") or None,
|
||
operator_role=current_user.get("role"),
|
||
)
|
||
|
||
|
||
# ============================================================
|
||
# 核心业务 0.3:撤回转交 (PENDING → 删除 + 恢复父任务)
|
||
# ============================================================
|
||
|
||
@router.post("/{task_id}/recall", response_model=TaskResponse)
|
||
async def recall_task_endpoint(
|
||
task_id: str,
|
||
operator_id: str | None = Query(None),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**撤回转交:删除 PENDING 子任务,恢复父任务为 WIP。**
|
||
适用场景:转交后发现选错人,在对方接收前撤回。
|
||
"""
|
||
return await task_service.recall_task(
|
||
db, uuid.UUID(task_id),
|
||
operator_id or current_user.get("username", "") or None,
|
||
operator_role=current_user.get("role"),
|
||
)
|
||
|
||
|
||
# ============================================================
|
||
# 核心业务 0.5:并发派发协助分支 (WIP → 不改变状态,创建子任务)
|
||
# ============================================================
|
||
|
||
class SpawnRequest(BaseModel):
|
||
task_name: str = Field(..., max_length=200, description="工序名称")
|
||
assignee_id: str | None = Field(None, max_length=64, description="负责人ID")
|
||
remark: str | None = Field(None, max_length=2000, description="派发备注")
|
||
|
||
|
||
@router.post("/{task_id}/spawn", response_model=TaskResponse, status_code=201)
|
||
async def spawn_subtask_endpoint(
|
||
task_id: str,
|
||
data: SpawnRequest,
|
||
operator_id: str | None = Query(None),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**派发协助分支:在当前任务下创建并行子任务,父任务状态保持不变。**
|
||
用于 WIP 期间工人需要其他人协助协同的场景。
|
||
"""
|
||
return await task_service.spawn_subtask(
|
||
db, uuid.UUID(task_id), data,
|
||
operator_id or current_user.get("username", "") or None)
|
||
|
||
|
||
# ============================================================
|
||
# 核心业务 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"),
|
||
remark: str | None = Body(None, description="接收备注", embed=True),
|
||
task_name: str | None = Body(None, description="接收人选定的工序名称", embed=True),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**确认接收任务。工人选定工序名称后接收。**
|
||
|
||
校验:只有状态为 PENDING 的任务可接收。
|
||
动作:状态改为 WIP,记录 received_at,更新 task_name,同步产品宏观状态。
|
||
"""
|
||
return await task_service.receive_task(
|
||
db, uuid.UUID(task_id),
|
||
operator_id or current_user.get("username", "") or None,
|
||
remark, task_name,
|
||
operator_role=current_user.get("role"),
|
||
)
|
||
|
||
|
||
# ============================================================
|
||
# 核心业务 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),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**品质驳回:将任务标记为 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 or current_user.get("username", "") or None,
|
||
operator_role=current_user.get("role"),
|
||
)
|
||
|
||
|
||
# ============================================================
|
||
# 核心业务 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),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**完工并裂变转交:完成当前任务,批量创建下一道工序任务。**
|
||
|
||
动作 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 or current_user.get("username", "") or None,
|
||
operator_role=current_user.get("role"),
|
||
)
|
||
|
||
|
||
# ============================================================
|
||
# 无限层级子任务
|
||
# ============================================================
|
||
|
||
@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),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""
|
||
**创建子任务:支持无限层级嵌套。**
|
||
|
||
新子任务将自动继承父任务的 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),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""获取指定产品的顶层任务列表(不含子任务嵌套)"""
|
||
return await task_service.get_top_level_tasks(db, uuid.UUID(product_id))
|
||
|
||
|
||
# ============================================================
|
||
# 任务进度记录 — 备注/传图
|
||
# ============================================================
|
||
|
||
@router.patch("/{task_id}/records", response_model=TaskResponse)
|
||
async def add_task_record_endpoint(
|
||
task_id: str,
|
||
data: TaskRecordCreate,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: dict = Depends(get_current_user),
|
||
):
|
||
"""追加进度记录(备注+图片),不改变任务状态"""
|
||
return await task_service.add_task_record(db, uuid.UUID(task_id), data, current_user)
|