Files
server-core/app/services/task_service.py
T

293 lines
15 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
import asyncio
import json
from fastapi import HTTPException
from ..infrastructure.repositories import Database
def _sync_task_group(task_id: str) -> None:
"""任务群后台同步(fire-and-forget,不阻断主流程;IM 服务不可用时静默降级)。"""
try:
from ..im import client as im_client
asyncio.create_task(im_client.sync_task_group(task_id))
except Exception: # noqa: BLE001
pass
class TaskService:
"""统一任务状态流转:grab/bid/deliver/win/review。"""
def __init__(self, db: Database):
self.db = db
async def can_accept(self, task: dict | None, actor: dict | None) -> dict:
"""接单资格判定:可见性(仅被指派) + 归属园区(专属) + eligibility 条件集。
返回 {ok, reasons[]}reasons 为空即通过。
"""
if not task:
return {"ok": False, "reasons": ["任务不存在"]}
reasons: list[str] = []
actor = actor or {}
full = ((await self.db.users.get_by_id(actor["id"])) or {}) if actor.get("id") else {}
user_park = full.get("park_company_id") or full.get("park_id") or actor.get("park_id")
uid = actor.get("id")
# 仅被指派可见 → 只有被指派对象(OPC 或园区)可接
if uid and task.get("visibility") == "assigned_only":
target = task.get("assign_opc_id") or task.get("assign_park_id")
if target and target != uid and target != user_park:
reasons.append("仅被指派方可接")
# 园区专属任务 → 仅本园
if task.get("park_id") and user_park != task.get("park_id"):
reasons.append("任务为本园区专属")
# 指派给园区但尚未发单 → 需园区端先操作
if task.get("assign_type") == "park" and not task.get("park_released") and task.get("assign_park_id") != user_park:
reasons.append("由园区分派")
# eligibility 条件集
try:
cond = json.loads(task.get("eligibility") or "{}")
except Exception:
cond = {}
if not isinstance(cond, dict):
cond = {}
if cond.get("opc_certified") and not full.get("opc_certified") and not full.get("opc_cert"):
reasons.append("需 OPC 认证")
if cond.get("region") and cond["region"] != "any" and (full.get("region_id") or "") != cond["region"]:
reasons.append("地域不符")
if cond.get("gender") and full.get("gender") != cond["gender"]:
reasons.append("不符合性别要求")
if cond.get("years_min") and (int(full.get("exp_years") or 0) < int(cond["years_min"])):
reasons.append(f"需从业 {cond['years_min']} 年以上")
if cond.get("field") and cond["field"] not in (full.get("fields") or []) and cond["field"] not in (full.get("field") or ""):
reasons.append("领域不符")
if cond.get("skill") and cond["skill"] not in (full.get("skills") or []) and cond["skill"] not in (full.get("skill") or ""):
reasons.append("专长不符")
if cond.get("case_req") and (int(full.get("case_count") or 0) < int(cond["case_req"])):
reasons.append(f"{cond['case_req']} 例服务案例")
reasons = [r for r in reasons if r]
return {"ok": not reasons, "reasons": reasons}
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="该任务接单人数已达上限")
acc = await self.can_accept(task, actor)
if not acc["ok"]:
raise HTTPException(status_code=403, detail="不满足接单条件:" + "".join(acc["reasons"]))
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,
)
_sync_task_group(task_id)
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",
)
_sync_task_group(task_id)
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",
)
_sync_task_group(task_id)
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="任务不可投标")
acc = await self.can_accept(task, actor)
if not acc["ok"]:
raise HTTPException(status_code=403, detail="不符合报名条件:" + "".join(acc["reasons"]))
if task.get("bid_quota"):
existing = await self.db.bids.list_for_task(task_id)
if len(existing) >= task["bid_quota"]:
raise HTTPException(status_code=400, detail="该竞标报名人数已达上限")
for b in await self.db.bids.list_for_task(task_id):
if b.get("opc_id") == actor["id"]:
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 register(self, task_id: str, actor: dict) -> dict:
"""报名(大厅 register 模式,0035):报名即占位、可多人、无报价。
与 bid 的区别——不报价、不评标,先到先得;register_quota>0 时达上限即拒绝。
资格判定与 grab 共用 can_accept;落 TaskClaim(status=registered, source=register)。
任务保持 published(多报名不翻转任务状态,由发包方/运营端选定后转 claimed)。
"""
task = await self.db.tasks.get(task_id)
if task is None or task["status"] != "published" or task["mode"] != "register":
raise HTTPException(status_code=400, detail="任务不可报名")
for c in await self.db.task_claims.list_by_task(task_id):
if c.get("claimer_user_id") == actor["id"] and c.get("status") not in ("withdrawn",):
raise HTTPException(status_code=400, detail="你已报名该任务")
quota = task.get("register_quota") or 0
if quota > 0:
active = sum(1 for c in await self.db.task_claims.list_by_task(task_id)
if c.get("claim_source") == "register" and c.get("status") == "registered")
if active >= quota:
raise HTTPException(status_code=400, detail="该任务报名人数已达上限")
acc = await self.can_accept(task, actor)
if not acc["ok"]:
raise HTTPException(status_code=403, detail="不满足报名条件:" + "".join(acc["reasons"]))
return await self.db.task_claims.create(
task_id, actor["id"], actor.get("nickname") or actor.get("username", ""),
source="register", status="registered",
)
async def register_select(self, task_id: str, claim_id: str, actor: dict) -> dict:
"""报名模式选定:所选 registered → assigned(任务 → claimed),其余 → withdrawn。"""
task = await self.db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
for c in await self.db.task_claims.list_by_task(task_id):
target = "assigned" if c["id"] == claim_id else "withdrawn"
if c.get("claim_source") == "register" and c.get("status") == "registered":
await self.db.task_claims.set_status(c["id"], target)
chosen = await self.db.task_claims.get(claim_id)
if chosen is None:
raise HTTPException(status_code=404, detail="报名记录不存在")
result = await self.db.tasks.claim(task_id, chosen.get("claimer_user_id") or "")
_sync_task_group(task_id)
return result
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")
result = await self.db.tasks.set_status(task_id, "in_progress")
_sync_task_group(task_id)
return result
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)