Files
server-core/app/services/task_service.py
T
Pine ea904d3a50 feat(task): 完善任务系统 —— 分类字典/字段/四种接单方式/独占限额/交付时限
- Task 加 published_at/delivery_days/headcount/exclusive/publisher_id/category_id;新增 TaskCategory 分类字典表(35 类)。
- mode 归一 grab/bid/assign/recommend(旧 designated→assign、dispatch→recommend);TaskService 加 assign/recommend/select_recommend,claim 支持独占/限额、published+claimed 均可抢(多人)。
- TaskClaim 支持 source=assign/recommend、status=assigned/recommended/withdrawn + count_active_by_task;TaskCategoryRepository;Database 装配 task_categories。
- 端口:operator GET /admin/task-categories + POST tasks/{id}/assign|recommend|select;opc 同;park /api/tasks 返回 mode/deadline/delivery_days/headcount/exclusive/publisher_name。
- 迁移 0009(tasks 加列+task_categories)已应用+stamp;seed 35 分类字典 + task 新字段(含独占/推荐演示) + 补种已初始化库。
- 测试 test_task_system.py(独占/限额/指派/推荐/list_published/publish/category)+补齐 test_task_claim 语义, 15 passed。

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-25 19:30:16 +08:00

167 lines
8.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
"""业务层 · 任务状态机(抢单 / 投标 / 交付 / 中标 / 验收)。
业务层只处理业务逻辑;不感知 HTTP 之外的框架细节。依赖基础设施层
Repository(经 Database 门面),不反向依赖接口层。
"""
from __future__ import annotations
from fastapi import HTTPException
from ..infrastructure.repositories import Database
class TaskService:
"""统一任务状态流转:grab/bid/deliver/win/review。"""
def __init__(self, db: Database):
self.db = db
async def grab(self, task_id: str, actor: dict) -> dict:
"""抢单(兼容旧入口):published + grab → claimed,记录接单流水。"""
return await self.claim_by_id(task_id, actor, source="grab")
async def claim(self, task_code: str, actor: dict, source: str = "scan") -> dict:
"""扫码/按短码接单:解析 task_code 后领单。"""
task = await self.db.tasks.get_by_code(task_code)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
return await self.claim_by_id(task["id"], actor, source=source)
async def claim_by_id(self, task_id: str, actor: dict, source: str) -> dict:
"""published + grab 模式 → claimedclaimed_by/claimed_at + TaskClaim 流水)。
校验接单限制:独占(exclusive) 仅一人可接;headcount>0 时达上限拒绝。
"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if task["status"] not in ("published", "claimed"):
raise HTTPException(status_code=400, detail="任务不可接单")
if task["mode"] != "grab":
raise HTTPException(status_code=400, detail="任务不支持扫码接单")
active = await self.db.task_claims.count_active_by_task(task_id)
if task.get("exclusive") and active >= 1:
raise HTTPException(status_code=400, detail="该任务为独占,已被接单")
if (task.get("headcount") or 0) > 0 and active >= task["headcount"]:
raise HTTPException(status_code=400, detail="该任务接单人数已达上限")
updated = await self.db.tasks.claim(task_id, actor["id"])
await self.db.task_claims.create(
task_id, actor["id"],
actor.get("nickname") or actor.get("username", ""), source,
)
return updated
async def assign(self, task_id: str, taker_user_id: str, actor: dict) -> dict:
"""指派:发包方直接指派某人 → claimed(claimed_by=taker)。仅 published + assign 模式。"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if task["status"] != "published":
raise HTTPException(status_code=400, detail="任务不可指派")
if task["mode"] != "assign":
raise HTTPException(status_code=400, detail="任务不支持指派")
taker = await self.db.users.get_by_id(taker_user_id)
taker_name = (taker or {}).get("nickname") or (taker or {}).get("username") or taker_user_id
updated = await self.db.tasks.claim(task_id, taker_user_id)
await self.db.task_claims.create(
task_id, taker_user_id, taker_name, source="assign", status="assigned",
)
return updated
async def recommend(self, task_id: str, candidates: list[str], actor: dict) -> list[dict]:
"""推荐:发包方/系统推荐候选人,逐个写 TaskClaim(status=recommended)。仅 published + recommend。"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if task["status"] != "published":
raise HTTPException(status_code=400, detail="任务不可推荐")
if task["mode"] != "recommend":
raise HTTPException(status_code=400, detail="任务不支持推荐")
out: list[dict] = []
for uid in candidates or []:
u = await self.db.users.get_by_id(uid)
name = (u or {}).get("nickname") or (u or {}).get("username") or uid
record = await self.db.task_claims.create(
task_id, uid, name, source="recommend", status="recommended",
)
out.append(record)
return out
async def select_recommend(self, task_id: str, taker_user_id: str, actor: dict) -> dict:
"""选定推荐人:published → claimed(claimed_by=所选);该推荐 assigned、其余 withdrawn。"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if task["status"] != "published":
raise HTTPException(status_code=400, detail="任务不可选定")
taker = await self.db.users.get_by_id(taker_user_id)
taker_name = (taker or {}).get("nickname") or (taker or {}).get("username") or taker_user_id
for rec in await self.db.task_claims.list_by_task(task_id):
target = "assigned" if rec["claimer_user_id"] == taker_user_id else "withdrawn"
await self.db.task_claims.set_status(rec["id"], target)
updated = await self.db.tasks.claim(task_id, taker_user_id)
if not any(t["claimer_user_id"] == taker_user_id
for t in await self.db.task_claims.list_by_task(task_id)):
await self.db.task_claims.create(
task_id, taker_user_id, taker_name, source="recommend", status="assigned",
)
return updated
async def start_doing(self, task_id: str, actor: dict) -> dict:
"""claimed → doing。"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if task["status"] != "claimed":
raise HTTPException(status_code=400, detail="任务未接单,无法开始")
return await self.db.tasks.set_doing(task_id)
async def complete(self, task_id: str, actor: dict) -> dict:
"""doing → completed。"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if task["status"] != "doing":
raise HTTPException(status_code=400, detail="任务未在进行中,无法完成")
return await self.db.tasks.set_status(task_id, "completed")
async def bid(self, task_id: str, actor: dict, quote: int, plan: str) -> dict:
"""投标:仅 published + bid 模式可投。"""
task = await self.db.tasks.get(task_id)
if task is None or task["status"] != "published" or task["mode"] != "bid":
raise HTTPException(status_code=400, detail="任务不可投标")
return await self.db.bids.create(
task_id, actor["id"], actor.get("nickname") or actor["username"],
quote, plan,
)
async def deliver(self, task_id: str, actor: dict) -> dict:
"""交付:任务存在且处于进行中(in_progress)才可交付 → delivered。
补齐原实现「无任何状态预检」的缺陷,防随意交付。
"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="Task not found")
if task["status"] != "in_progress":
raise HTTPException(status_code=400, detail="任务未在进行中,无法交付")
return await self.db.tasks.set_status(task_id, "delivered")
async def enterprise_submit(self, task_id: str) -> dict:
"""企业提交审核:pending 任务 → published。"""
return await self.db.tasks.set_status(task_id, "published")
async def win_bid(self, task_id: str, bid_id: str) -> dict:
"""企业评标中标:bid → wintask → in_progress(双状态联动)。"""
bid = await self.db.bids.get(bid_id)
if bid is None or bid["task_id"] != task_id:
raise HTTPException(status_code=404, detail="竞标不存在")
await self.db.bids.set_status(bid_id, "win")
return await self.db.tasks.set_status(task_id, "in_progress")
async def review(self, task_id: str, accept: bool) -> dict:
"""企业验收:accept → completedreject → in_progress。"""
target = "completed" if accept else "in_progress"
return await self.db.tasks.set_status(task_id, target)