| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052 |
- from __future__ import annotations
- import hashlib
- import json
- import math
- import re
- from collections import defaultdict
- from datetime import UTC, datetime, timedelta
- from decimal import Decimal, InvalidOperation
- from secrets import token_urlsafe
- from typing import Any
- from uuid import uuid4
- from cryptography.fernet import Fernet, InvalidToken
- from pymongo import ASCENDING, DESCENDING, ReturnDocument
- from pymongo.errors import DuplicateKeyError
- import wbb
- from wbb import control_db
- templatesdb = control_db.service_package_templates
- ordersdb = control_db.service_orders
- quotesdb = control_db.service_order_quotes
- addressesdb = control_db.service_order_addresses
- eventsdb = control_db.service_order_events
- customersdb = control_db.customer_profiles
- blocksdb = control_db.customer_blocks
- reportsdb = control_db.service_reports
- qrdb = control_db.teacher_review_qr_records
- reviewtemplatesdb = control_db.teacher_review_templates
- reviewdraftsdb = control_db.teacher_review_drafts
- reviewsdb = control_db.teacher_reviews
- statsdb = control_db.teacher_review_stats
- profilesdb = control_db.directory_profiles
- locationsdb = control_db.directory_locations
- membershipsdb = control_db.directory_memberships
- settingsdb = control_db.directory_settings
- SERVICE_MODES = {"at_store", "onsite"}
- PRICE_MODES = {"fixed", "starting_at", "range", "negotiable"}
- PRICE_UNITS = {"per_service", "per_hour", "per_item", "per_visit"}
- TRAVEL_FEE_MODES = {"included", "fixed", "per_km", "quoted"}
- ORDER_ACTIVE_STATUSES = {
- "requested",
- "quoted",
- "confirmed",
- "in_progress",
- "completion_pending",
- "disputed",
- }
- ORDER_TERMINAL_STATUSES = {
- "completed",
- "rejected",
- "canceled_customer",
- "canceled_technician",
- "voided",
- }
- ORDER_STATUSES = ORDER_ACTIVE_STATUSES | ORDER_TERMINAL_STATUSES
- REVIEW_SOURCES = {"qr_verified", "student_initiated"}
- REVIEW_STATUSES = {"pending", "approved", "rejected", "withdrawn", "voided"}
- DEFAULT_REVIEW_TEMPLATE = {
- "rating_questions": [
- {
- "question_id": "service_effect",
- "label": "服务效果",
- "description": "本次服务是否达到预期。",
- "required": True,
- "weight": 40,
- },
- {
- "question_id": "professionalism",
- "label": "专业程度",
- "description": "技师的专业能力与服务规范。",
- "required": True,
- "weight": 30,
- },
- {
- "question_id": "communication",
- "label": "沟通体验",
- "description": "沟通是否清晰、友好。",
- "required": True,
- "weight": 30,
- },
- ],
- "text_questions": [
- {
- "question_id": "comment",
- "label": "评价内容",
- "description": "分享真实体验,避免包含隐私信息。",
- "required": False,
- "max_length": 500,
- }
- ],
- }
- BUILTIN_PACKAGE_TEMPLATES = (
- {
- "template_id": "builtin_quick_at_store",
- "name": "快速到店服务",
- "sort_order": -300,
- "package": {
- "package_id": "builtin_package_at_store",
- "name": "标准到店服务",
- "category": "通用服务",
- "description": "顾客到店完成一次标准服务,具体内容可在接单后沟通确认。",
- "tags": ["到店", "标准服务"],
- "service_modes": ["at_store"],
- "price_mode": "fixed",
- "currency": "CNY",
- "price_unit": "per_service",
- "min_price": "200",
- "max_price": "200",
- "duration_minutes": 60,
- "included_items": "一次标准服务",
- "excluded_items": "额外耗材和临时加项",
- "preparation": "请提前说明具体需求",
- "addons": [],
- "travel_fee": {"mode": "included", "amount": "0", "per_km": "0"},
- "service_radius_km": None,
- "out_of_range_policy": "",
- },
- },
- {
- "template_id": "builtin_quick_onsite",
- "name": "快速上门服务",
- "sort_order": -200,
- "package": {
- "package_id": "builtin_package_onsite",
- "name": "标准上门服务",
- "category": "通用服务",
- "description": "技师上门完成一次标准服务,最终费用以确认报价为准。",
- "tags": ["上门", "标准服务"],
- "service_modes": ["onsite"],
- "price_mode": "starting_at",
- "currency": "CNY",
- "price_unit": "per_visit",
- "min_price": "300",
- "max_price": "300",
- "duration_minutes": 90,
- "included_items": "一次标准上门服务",
- "excluded_items": "交通费、额外耗材和临时加项",
- "preparation": "请提前发送位置并说明具体需求",
- "addons": [],
- "travel_fee": {"mode": "quoted", "amount": "0", "per_km": "0"},
- "service_radius_km": 20,
- "out_of_range_policy": "超出服务半径时,请先沟通交通费和是否可以接单。",
- },
- },
- {
- "template_id": "builtin_quick_hourly",
- "name": "快速按小时服务",
- "sort_order": -100,
- "package": {
- "package_id": "builtin_package_hourly",
- "name": "按小时专业服务",
- "category": "通用服务",
- "description": "按服务时长计价,支持到店或上门,具体工作范围提前确认。",
- "tags": ["按小时", "到店", "上门"],
- "service_modes": ["at_store", "onsite"],
- "price_mode": "starting_at",
- "currency": "CNY",
- "price_unit": "per_hour",
- "min_price": "200",
- "max_price": "200",
- "duration_minutes": 60,
- "included_items": "一小时专业服务",
- "excluded_items": "交通费、额外耗材和超时服务",
- "preparation": "请提前说明服务内容和预计时长",
- "addons": [],
- "travel_fee": {"mode": "quoted", "amount": "0", "per_km": "0"},
- "service_radius_km": 20,
- "out_of_range_policy": "超出服务半径时,请先沟通交通费和是否可以接单。",
- },
- },
- )
- _indexes_ready = False
- class ServiceDataError(ValueError):
- def __init__(self, code: str, message: str):
- super().__init__(message)
- self.code = code
- def utc_now() -> datetime:
- return datetime.now(UTC)
- def _as_utc(value: datetime) -> datetime:
- return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
- def _clean_text(value: Any, *, max_length: int, required: bool = False) -> str:
- text = re.sub(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]", "", str(value or ""))
- text = " ".join(text.strip().split())
- if required and not text:
- raise ServiceDataError("required_field", "必填内容不能为空。")
- if len(text) > max_length:
- raise ServiceDataError("text_too_long", f"内容不能超过 {max_length} 个字符。")
- return text
- def _positive_decimal(value: Any, field: str, *, allow_zero: bool = True) -> str:
- try:
- number = Decimal(str(value or 0)).quantize(Decimal("0.01"))
- except (InvalidOperation, ValueError) as exc:
- raise ServiceDataError("invalid_price", f"{field} 必须是有效金额。") from exc
- if not number.is_finite():
- raise ServiceDataError("invalid_price", f"{field} 必须是有效金额。")
- if number < 0 or (not allow_zero and number == 0):
- raise ServiceDataError("invalid_price", f"{field} 不能小于 0。")
- return format(number, "f")
- def _normalize_tags(values: Any) -> list[str]:
- if isinstance(values, str):
- values = values.replace(",", ",").split(",")
- tags = []
- for value in values or []:
- tag = _clean_text(value, max_length=20)
- if tag and tag.casefold() not in {item.casefold() for item in tags}:
- tags.append(tag)
- if len(tags) > 8:
- raise ServiceDataError("too_many_tags", "服务标签最多 8 个。")
- return tags
- def normalize_package(value: dict[str, Any], *, package_id: str | None = None) -> dict[str, Any]:
- modes = list(dict.fromkeys(str(item) for item in value.get("service_modes", [])))
- if not modes or set(modes) - SERVICE_MODES:
- raise ServiceDataError("invalid_service_mode", "套餐必须选择到店或上门服务方式。")
- price_mode = str(value.get("price_mode") or "negotiable")
- if price_mode not in PRICE_MODES:
- raise ServiceDataError("invalid_price_mode", "不支持该价格模式。")
- price_unit = str(value.get("price_unit") or "per_service")
- if price_unit not in PRICE_UNITS:
- raise ServiceDataError("invalid_price_unit", "不支持该计价单位。")
- currency = str(value.get("currency") or "CNY").upper()
- if not re.fullmatch(r"[A-Z]{3}", currency):
- raise ServiceDataError("invalid_currency", "币种必须使用三位 ISO 代码。")
- min_price = _positive_decimal(value.get("min_price"), "最低价格")
- max_price = _positive_decimal(value.get("max_price"), "最高价格")
- if Decimal(max_price) and Decimal(max_price) < Decimal(min_price):
- raise ServiceDataError("invalid_price_range", "最高价格不能低于最低价格。")
- if price_mode == "fixed" and Decimal(min_price) != Decimal(max_price):
- raise ServiceDataError("invalid_fixed_price", "固定价的最低价格和最高价格必须一致。")
- duration = value.get("duration_minutes")
- if duration in (None, "", 0, "0"):
- duration = None
- else:
- try:
- duration = int(duration)
- except (TypeError, ValueError) as exc:
- raise ServiceDataError("invalid_duration", "服务时长必须是分钟数。") from exc
- if not 15 <= duration <= 480:
- raise ServiceDataError("invalid_duration", "服务时长必须在 15 到 480 分钟之间。")
- travel = value.get("travel_fee") if isinstance(value.get("travel_fee"), dict) else {}
- travel_mode = str(travel.get("mode") or "quoted")
- if travel_mode not in TRAVEL_FEE_MODES:
- raise ServiceDataError("invalid_travel_fee", "不支持该交通费模式。")
- addons = []
- for item in value.get("addons", []) or []:
- if len(addons) >= 10:
- raise ServiceDataError("too_many_addons", "标准加项最多 10 个。")
- addons.append(
- {
- "addon_id": str(item.get("addon_id") or uuid4().hex),
- "name": _clean_text(item.get("name"), max_length=40, required=True),
- "unit": _clean_text(item.get("unit"), max_length=20) or "项",
- "price": _positive_decimal(item.get("price"), "加项价格"),
- "extra_minutes": max(0, min(480, int(item.get("extra_minutes") or 0))),
- }
- )
- radius = value.get("service_radius_km")
- if radius in (None, ""):
- radius = None
- else:
- radius = float(radius)
- if not 0.5 <= radius <= 200:
- raise ServiceDataError("invalid_service_radius", "上门半径必须在 0.5 到 200 公里之间。")
- if "onsite" in modes and radius is None:
- raise ServiceDataError("service_radius_required", "上门套餐必须填写服务半径。")
- out_of_range_policy = _clean_text(
- value.get("out_of_range_policy"),
- max_length=200,
- required="onsite" in modes,
- )
- source_type = str(value.get("source_type") or "custom")
- if source_type not in {"template", "custom"}:
- raise ServiceDataError("invalid_package_source", "套餐来源无效。")
- return {
- "package_id": str(package_id or value.get("package_id") or uuid4().hex),
- "name": _clean_text(value.get("name"), max_length=50, required=True),
- "category": _clean_text(value.get("category"), max_length=30, required=True),
- "description": _clean_text(value.get("description"), max_length=800, required=True),
- "tags": _normalize_tags(value.get("tags", [])),
- "service_modes": modes,
- "price_mode": price_mode,
- "currency": currency,
- "price_unit": price_unit,
- "min_price": min_price,
- "max_price": max_price,
- "price_display": _clean_text(value.get("price_display"), max_length=80),
- "duration_minutes": duration,
- "included_items": _clean_text(value.get("included_items"), max_length=500),
- "excluded_items": _clean_text(value.get("excluded_items"), max_length=500),
- "preparation": _clean_text(value.get("preparation"), max_length=500),
- "addons": addons,
- "travel_fee": {
- "mode": travel_mode,
- "amount": _positive_decimal(travel.get("amount"), "交通费"),
- "per_km": _positive_decimal(travel.get("per_km"), "每公里交通费"),
- "description": _clean_text(travel.get("description"), max_length=200),
- },
- "service_radius_km": radius,
- "out_of_range_policy": out_of_range_policy,
- "source_type": source_type,
- "source_template_id": value.get("source_template_id"),
- "source_template_version": value.get("source_template_version"),
- "customized_from_template": bool(value.get("customized_from_template")),
- }
- def _normalize_profile(value: dict[str, Any]) -> dict[str, Any]:
- packages = [normalize_package(item) for item in value.get("packages", [])]
- if not 1 <= len(packages) <= 5:
- raise ServiceDataError("invalid_package_count", "技师必须发布 1 到 5 个套餐。")
- modes = sorted({mode for item in packages for mode in item["service_modes"]})
- venue = value.get("venue") if isinstance(value.get("venue"), dict) else {}
- onsite = value.get("onsite_policy") if isinstance(value.get("onsite_policy"), dict) else {}
- public_area_text = _clean_text(
- value.get("public_area_text"), max_length=80, required=True
- )
- venue_name = _clean_text(
- venue.get("name"), max_length=80, required="at_store" in modes
- )
- venue_address = _clean_text(
- venue.get("address_hint"), max_length=120, required="at_store" in modes
- )
- onsite_description = _clean_text(
- onsite.get("description"), max_length=300, required="onsite" in modes
- )
- return {
- "headline": _clean_text(value.get("headline"), max_length=80, required=True),
- "bio": _clean_text(value.get("bio"), max_length=800, required=True),
- "tags": _normalize_tags(value.get("tags", [])),
- "contact_hours": _clean_text(value.get("contact_hours"), max_length=120, required=True),
- "service_modes": modes,
- "accepting_requests": bool(value.get("accepting_requests", True)),
- "public_area_text": public_area_text,
- "venue": {
- "name": venue_name,
- "address_hint": venue_address,
- },
- "onsite_policy": {
- "description": onsite_description,
- },
- "packages": packages,
- "is_complete": True,
- }
- async def ensure_service_indexes() -> None:
- global _indexes_ready
- if _indexes_ready:
- return
- await templatesdb.create_index([("template_id", ASCENDING)], unique=True)
- await templatesdb.create_index([("status", ASCENDING), ("sort_order", ASCENDING)])
- await ordersdb.create_index([("order_id", ASCENDING)], unique=True)
- await ordersdb.create_index([("customer_id", ASCENDING), ("created_at", DESCENDING)])
- await ordersdb.create_index([("technician_id", ASCENDING), ("created_at", DESCENDING)])
- await ordersdb.create_index([("status", ASCENDING), ("updated_at", DESCENDING)])
- await quotesdb.create_index([("quote_id", ASCENDING)], unique=True)
- await quotesdb.create_index([("order_id", ASCENDING), ("version", DESCENDING)], unique=True)
- await addressesdb.create_index([("order_id", ASCENDING)], unique=True)
- await addressesdb.create_index([("redact_after", ASCENDING)])
- await eventsdb.create_index([("event_id", ASCENDING)], unique=True)
- await eventsdb.create_index([("order_id", ASCENDING), ("created_at", ASCENDING)])
- await customersdb.create_index([("user_id", ASCENDING)], unique=True)
- await blocksdb.create_index([("scope", ASCENDING), ("technician_id", ASCENDING), ("customer_id", ASCENDING)], unique=True)
- await reportsdb.create_index([("report_id", ASCENDING)], unique=True)
- await qrdb.create_index([("qr_id", ASCENDING)], unique=True)
- await qrdb.create_index([("token_hash", ASCENDING)], unique=True)
- await qrdb.create_index([("order_id", ASCENDING), ("created_at", DESCENDING)])
- await reviewtemplatesdb.create_index([("version", DESCENDING)], unique=True)
- await reviewdraftsdb.create_index([("customer_id", ASCENDING)], unique=True)
- await reviewdraftsdb.create_index([("expires_at", ASCENDING)], expireAfterSeconds=0)
- await reviewsdb.create_index([("review_id", ASCENDING)], unique=True)
- await reviewsdb.create_index([("qr_id", ASCENDING)], unique=True, sparse=True)
- await reviewsdb.create_index([("technician_id", ASCENDING), ("status", ASCENDING), ("created_at", DESCENDING)])
- await statsdb.create_index([("technician_id", ASCENDING), ("category", ASCENDING)], unique=True)
- _indexes_ready = True
- async def record_service_event(
- order_id: str,
- event_type: str,
- *,
- actor_id: int | str,
- reason: str = "",
- metadata: dict[str, Any] | None = None,
- ) -> dict[str, Any]:
- await ensure_service_indexes()
- document = {
- "event_id": uuid4().hex,
- "order_id": str(order_id),
- "event_type": str(event_type),
- "actor_id": actor_id,
- "reason": _clean_text(reason, max_length=500),
- "metadata": metadata or {},
- "created_at": utc_now(),
- }
- await eventsdb.insert_one(document)
- return document
- async def record_service_funnel_event(user_id: int, event_type: str) -> dict[str, Any]:
- if event_type not in {"directory_viewed", "technician_viewed", "request_started"}:
- raise ServiceDataError("invalid_funnel_event", "不支持该履约漏斗事件。")
- return await record_service_event(
- "",
- event_type,
- actor_id=int(user_id),
- metadata={"funnel": True},
- )
- async def observe_customer(user: Any, *, accepted_terms: bool = False) -> dict[str, Any]:
- await ensure_service_indexes()
- now = utc_now()
- display_name = " ".join(
- value for value in (str(getattr(user, "first_name", "") or "").strip(), str(getattr(user, "last_name", "") or "").strip()) if value
- ) or f"用户 {user.id}"
- updates: dict[str, Any] = {
- "username": getattr(user, "username", None),
- "display_name": display_name,
- "updated_at": now,
- }
- if accepted_terms:
- updates["terms_accepted_at"] = now
- await customersdb.update_one(
- {"user_id": int(user.id)},
- {"$set": updates, "$setOnInsert": {"created_at": now, "blocked": False}},
- upsert=True,
- )
- return await customersdb.find_one({"user_id": int(user.id)}) or {}
- async def require_customer_allowed(user_id: int) -> dict[str, Any]:
- await ensure_service_indexes()
- customer = await customersdb.find_one({"user_id": int(user_id)})
- if not customer or not customer.get("terms_accepted_at"):
- raise ServiceDataError("terms_required", "请先接受服务规则和隐私说明。")
- if customer.get("blocked"):
- raise ServiceDataError("customer_blocked", "你的服务功能已被暂停。")
- global_block = await blocksdb.find_one({"scope": "global", "customer_id": int(user_id), "active": True})
- if global_block:
- raise ServiceDataError("customer_blocked", "你的服务功能已被暂停。")
- return customer
- async def get_service_settings() -> dict[str, int]:
- stored = await settingsdb.find_one({"settings_id": "global"}) or {}
- return {
- "customer_max_open_orders": max(
- 1,
- min(
- 20,
- int(
- stored.get("service_customer_max_open_orders")
- or getattr(wbb, "SERVICE_CUSTOMER_MAX_OPEN_ORDERS", 3)
- or 3
- ),
- ),
- ),
- "customer_daily_request_limit": max(
- 1,
- min(
- 100,
- int(
- stored.get("service_customer_daily_request_limit")
- or getattr(wbb, "SERVICE_CUSTOMER_DAILY_REQUEST_LIMIT", 10)
- or 10
- ),
- ),
- ),
- "review_qr_expiry_hours": max(
- 1,
- min(
- 168,
- int(
- stored.get("service_review_qr_expiry_hours")
- or getattr(wbb, "SERVICE_REVIEW_QR_EXPIRY_HOURS", 24)
- or 24
- ),
- ),
- ),
- }
- async def set_service_settings(values: dict[str, Any]) -> dict[str, int]:
- current = await get_service_settings()
- try:
- normalized = {
- "customer_max_open_orders": max(
- 1,
- min(20, int(values.get("customer_max_open_orders", current["customer_max_open_orders"]))),
- ),
- "customer_daily_request_limit": max(
- 1,
- min(
- 100,
- int(
- values.get(
- "customer_daily_request_limit",
- current["customer_daily_request_limit"],
- )
- ),
- ),
- ),
- "review_qr_expiry_hours": max(
- 1,
- min(168, int(values.get("review_qr_expiry_hours", current["review_qr_expiry_hours"]))),
- ),
- }
- except (TypeError, ValueError) as exc:
- raise ServiceDataError("invalid_service_settings", "履约设置必须是有效整数。") from exc
- await settingsdb.update_one(
- {"settings_id": "global"},
- {
- "$set": {
- "service_customer_max_open_orders": normalized["customer_max_open_orders"],
- "service_customer_daily_request_limit": normalized["customer_daily_request_limit"],
- "service_review_qr_expiry_hours": normalized["review_qr_expiry_hours"],
- "updated_at": utc_now(),
- },
- "$setOnInsert": {"created_at": utc_now()},
- },
- upsert=True,
- )
- return normalized
- async def list_package_templates(*, include_archived: bool = False) -> list[dict[str, Any]]:
- await ensure_service_indexes()
- filters = {} if include_archived else {"status": {"$ne": "archived"}}
- return await templatesdb.find(filters).sort([("sort_order", ASCENDING), ("created_at", ASCENDING)]).to_list(length=500)
- async def ensure_builtin_package_templates() -> list[dict[str, Any]]:
- """Seed editable starter templates once without overwriting administrator changes."""
- await ensure_service_indexes()
- now = utc_now()
- template_ids = []
- for item in BUILTIN_PACKAGE_TEMPLATES:
- template_id = str(item["template_id"])
- template_ids.append(template_id)
- package = normalize_package(
- item["package"],
- package_id=str(item["package"]["package_id"]),
- )
- package["source_type"] = "template"
- await templatesdb.update_one(
- {"template_id": template_id},
- {
- "$setOnInsert": {
- "template_id": template_id,
- "version": 1,
- "name": item["name"],
- "admin_note": "系统内置快速模板,可在后台编辑、停用或归档。",
- "package": package,
- "status": "enabled",
- "sort_order": int(item["sort_order"]),
- "builtin": True,
- "created_by": "system",
- "created_at": now,
- "updated_at": now,
- }
- },
- upsert=True,
- )
- return await templatesdb.find({"template_id": {"$in": template_ids}}).sort(
- [("sort_order", ASCENDING)]
- ).to_list(length=len(template_ids))
- async def create_package_template(values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
- await ensure_service_indexes()
- now = utc_now()
- package = normalize_package(values.get("package") or values)
- package["source_type"] = "template"
- document = {
- "template_id": uuid4().hex,
- "version": 1,
- "name": _clean_text(values.get("template_name") or package["name"], max_length=80, required=True),
- "admin_note": _clean_text(values.get("admin_note"), max_length=500),
- "package": package,
- "status": "enabled" if values.get("enabled", True) else "disabled",
- "sort_order": int(values.get("sort_order") or 0),
- "created_by": str(actor_id),
- "created_at": now,
- "updated_at": now,
- }
- await templatesdb.insert_one(document)
- return document
- async def update_package_template(template_id: str, values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
- current = await templatesdb.find_one({"template_id": str(template_id)})
- if not current:
- raise ServiceDataError("template_not_found", "未找到套餐模板。")
- package = normalize_package(values.get("package") or values, package_id=current["package"]["package_id"])
- package["source_type"] = "template"
- updated = await templatesdb.find_one_and_update(
- {"template_id": str(template_id)},
- {
- "$set": {
- "name": _clean_text(values.get("template_name") or package["name"], max_length=80, required=True),
- "admin_note": _clean_text(values.get("admin_note"), max_length=500),
- "package": package,
- "sort_order": int(values.get("sort_order") or 0),
- "updated_by": str(actor_id),
- "updated_at": utc_now(),
- },
- "$inc": {"version": 1},
- },
- return_document=ReturnDocument.AFTER,
- )
- return updated or {}
- async def set_package_template_status(template_id: str, action: str, *, actor_id: str) -> dict[str, Any]:
- status_map = {"enable": "enabled", "disable": "disabled", "archive": "archived"}
- if action not in status_map:
- raise ServiceDataError("invalid_template_action", "不支持该模板操作。")
- updated = await templatesdb.find_one_and_update(
- {"template_id": str(template_id)},
- {"$set": {"status": status_map[action], "updated_by": str(actor_id), "updated_at": utc_now()}},
- return_document=ReturnDocument.AFTER,
- )
- if not updated:
- raise ServiceDataError("template_not_found", "未找到套餐模板。")
- return updated
- async def duplicate_package_template(template_id: str, *, actor_id: str) -> dict[str, Any]:
- current = await templatesdb.find_one({"template_id": str(template_id)})
- if not current:
- raise ServiceDataError("template_not_found", "未找到套餐模板。")
- values = {
- "template_name": f"{current['name']}(副本)",
- "admin_note": current.get("admin_note", ""),
- "sort_order": current.get("sort_order", 0),
- "enabled": False,
- "package": current["package"],
- }
- return await create_package_template(values, actor_id=actor_id)
- async def publish_technician_profile(user_id: int, values: dict[str, Any], *, actor_id: int | str) -> dict[str, Any]:
- await ensure_service_indexes()
- profile = await profilesdb.find_one({"user_id": int(user_id)})
- if not profile or profile.get("application_status") != "approved":
- raise ServiceDataError("technician_required", "只有已认证技师可以发布服务资料。")
- if not profile.get("username"):
- raise ServiceDataError("username_required", "技师必须设置有效的 Telegram 用户名。")
- if not await locationsdb.find_one({"user_id": int(user_id)}):
- raise ServiceDataError("technician_location_required", "技师必须设置有效服务位置。")
- if not await membershipsdb.find_one({"user_id": int(user_id), "active": True}):
- raise ServiceDataError("technician_membership_required", "技师必须仍是受管群当前成员。")
- normalized = _normalize_profile(values)
- now = utc_now()
- normalized.update({"updated_at": now, "updated_by": actor_id, "version": int((profile.get("service_profile") or {}).get("version") or 0) + 1})
- await profilesdb.update_one(
- {"user_id": int(user_id)},
- {"$set": {"service_profile": normalized, "updated_at": now}},
- )
- return await profilesdb.find_one({"user_id": int(user_id)}) or {}
- async def get_technician_service_profile(user_id: int) -> dict[str, Any]:
- profile = await profilesdb.find_one({"user_id": int(user_id)})
- if not profile:
- raise ServiceDataError("technician_not_found", "未找到技师资料。")
- return profile
- async def get_technician_self_service_context(user_id: int) -> dict[str, Any]:
- profile = await profilesdb.find_one({"user_id": int(user_id)})
- if not profile:
- raise ServiceDataError("technician_not_found", "未找到技师申请资料。")
- location_configured = bool(
- await locationsdb.find_one({"user_id": int(user_id)})
- )
- membership_active = bool(
- await membershipsdb.find_one({"user_id": int(user_id), "active": True})
- )
- issues = []
- if profile.get("application_status") != "approved":
- issues.append({"code": "approval_required", "message": "技师认证尚未通过。"})
- if not profile.get("username"):
- issues.append(
- {
- "code": "username_required",
- "message": "请先在 Telegram 设置用户名。",
- }
- )
- if not location_configured:
- issues.append(
- {
- "code": "location_required",
- "message": "请先在 Bot 中设置服务基准位置。",
- }
- )
- if not membership_active:
- issues.append(
- {
- "code": "membership_required",
- "message": "当前不在受管群成员名单中。",
- }
- )
- return {
- "user_id": int(user_id),
- "display_name": profile.get("display_name") or f"技师 {user_id}",
- "username": profile.get("username"),
- "eligibility": {"ready": not issues, "issues": issues},
- "service_profile": profile.get("service_profile") or None,
- }
- async def set_technician_accepting_requests(
- user_id: int,
- accepting_requests: bool,
- ) -> dict[str, Any]:
- profile = await get_technician_service_profile(user_id)
- values = dict(profile.get("service_profile") or {})
- if not values.get("is_complete"):
- raise ServiceDataError("technician_profile_required", "请先发布至少一个服务套餐。")
- values["accepting_requests"] = bool(accepting_requests)
- return await publish_technician_profile(
- user_id,
- values,
- actor_id=user_id,
- )
- async def instantiate_package_template(template_id: str) -> dict[str, Any]:
- template = await templatesdb.find_one(
- {"template_id": str(template_id), "status": "enabled"}
- )
- if not template:
- raise ServiceDataError("template_not_found", "套餐模板不存在或已停用。")
- package = dict(template["package"])
- package.update(
- {
- "package_id": uuid4().hex,
- "source_type": "template",
- "source_template_id": template["template_id"],
- "source_template_version": template["version"],
- "customized_from_template": False,
- }
- )
- return package
- def _fernet() -> Fernet:
- raw = str(getattr(wbb, "SERVICE_ADDRESS_ENCRYPTION_KEY", "") or "").strip()
- if not raw:
- raise ServiceDataError("address_encryption_required", "服务地址加密密钥尚未配置。")
- try:
- return Fernet(raw.encode())
- except (ValueError, TypeError) as exc:
- raise ServiceDataError("invalid_address_encryption_key", "服务地址加密密钥无效。") from exc
- def _encrypt_address(payload: dict[str, Any]) -> str:
- return _fernet().encrypt(json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode()).decode()
- def _decrypt_address(token: str) -> dict[str, Any]:
- try:
- return json.loads(_fernet().decrypt(str(token).encode()).decode())
- except (InvalidToken, ValueError, json.JSONDecodeError) as exc:
- raise ServiceDataError("address_unavailable", "精确地址无法解密。") from exc
- async def _schedule_address_redaction(order_id: str, closed_at: datetime) -> None:
- retention = max(
- 1,
- min(365, int(getattr(wbb, "SERVICE_ADDRESS_RETENTION_DAYS", 7) or 7)),
- )
- await addressesdb.update_one(
- {"order_id": str(order_id), "redacted": False},
- {
- "$set": {
- "redact_after": closed_at + timedelta(days=retention),
- "updated_at": closed_at,
- }
- },
- )
- def _distance_meters(lon1: float, lat1: float, lon2: float, lat2: float) -> float:
- radius = 6_371_000.0
- phi1, phi2 = math.radians(lat1), math.radians(lat2)
- delta_phi = math.radians(lat2 - lat1)
- delta_lambda = math.radians(lon2 - lon1)
- a = math.sin(delta_phi / 2) ** 2 + math.cos(phi1) * math.cos(phi2) * math.sin(delta_lambda / 2) ** 2
- return radius * 2 * math.atan2(math.sqrt(a), math.sqrt(1 - a))
- def distance_band(distance_meters: float) -> str:
- km = max(0.0, distance_meters / 1000)
- for threshold in (1, 3, 5, 10, 20, 50):
- if km <= threshold:
- return f"{threshold} 公里内"
- return "50 公里以上"
- async def _get_package(technician_id: int, package_id: str) -> tuple[dict[str, Any], dict[str, Any]]:
- profile = await profilesdb.find_one({"user_id": int(technician_id)})
- if not profile or profile.get("application_status") != "approved":
- raise ServiceDataError("technician_not_found", "未找到已认证技师。")
- if not profile.get("listed") or not profile.get("username"):
- raise ServiceDataError("technician_unavailable", "该技师暂未公开接单。")
- if not await membershipsdb.find_one({"user_id": int(technician_id), "active": True}):
- raise ServiceDataError("technician_unavailable", "该技师当前不满足接单资格。")
- service_profile = profile.get("service_profile") or {}
- if not service_profile.get("is_complete") or not service_profile.get("accepting_requests", True):
- raise ServiceDataError("technician_unavailable", "该技师暂未接单。")
- package = next((item for item in service_profile.get("packages", []) if str(item.get("package_id")) == str(package_id)), None)
- if not package:
- raise ServiceDataError("package_not_found", "未找到该服务套餐。")
- return profile, package
- async def list_available_technicians(
- *,
- longitude: float,
- latitude: float,
- max_distance_meters: float | None = None,
- query: str = "",
- page: int = 1,
- page_size: int = 10,
- ) -> tuple[list[dict[str, Any]], int]:
- lon, lat = float(longitude), float(latitude)
- if not -180 <= lon <= 180 or not -90 <= lat <= 90:
- raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
- filters: dict[str, Any] = {
- "application_status": "approved",
- "listed": True,
- "username": {"$nin": [None, ""]},
- "service_profile.is_complete": True,
- "service_profile.accepting_requests": True,
- }
- if query:
- pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
- filters["$or"] = [
- {"username": pattern},
- {"display_name": pattern},
- {"service_profile.tags": pattern},
- {"service_profile.packages.category": pattern},
- ]
- active_members = {
- int(item["user_id"])
- async for item in membershipsdb.find({"active": True}, {"user_id": 1})
- }
- profiles = {
- int(item["user_id"]): item
- async for item in profilesdb.find(filters)
- if int(item["user_id"]) in active_members
- }
- values: list[dict[str, Any]] = []
- if profiles:
- async for location in locationsdb.find({"user_id": {"$in": list(profiles)}}):
- point = location.get("point", {}).get("coordinates") or [
- location.get("longitude"),
- location.get("latitude"),
- ]
- if len(point) != 2 or point[0] is None or point[1] is None:
- continue
- distance = _distance_meters(lon, lat, float(point[0]), float(point[1]))
- if max_distance_meters is not None and distance > float(max_distance_meters):
- continue
- profile = profiles[int(location["user_id"])]
- service_profile = profile.get("service_profile") or {}
- values.append(
- {
- "user_id": int(profile["user_id"]),
- "username": profile.get("username"),
- "display_name": profile.get("display_name"),
- "headline": service_profile.get("headline"),
- "bio": service_profile.get("bio"),
- "tags": service_profile.get("tags", []),
- "public_area_text": service_profile.get("public_area_text"),
- "service_modes": service_profile.get("service_modes", []),
- "packages": service_profile.get("packages", []),
- "distance_meters": round(distance, 2),
- "distance_band": distance_band(distance),
- }
- )
- values.sort(key=lambda item: (item["distance_meters"], item["user_id"]))
- total = len(values)
- start = (max(1, page) - 1) * max(1, page_size)
- return values[start : start + max(1, page_size)], total
- async def create_service_order(
- *,
- customer: Any,
- technician_id: int,
- package_id: str,
- service_mode: str,
- requirements: str,
- longitude: float | None = None,
- latitude: float | None = None,
- address_text: str = "",
- ) -> dict[str, Any]:
- await ensure_service_indexes()
- await observe_customer(customer)
- await require_customer_allowed(int(customer.id))
- if int(customer.id) == int(technician_id):
- raise ServiceDataError("self_service_not_allowed", "不能向自己发起服务请求。")
- personal_block = await blocksdb.find_one({"scope": "technician", "technician_id": int(technician_id), "customer_id": int(customer.id), "active": True})
- if personal_block:
- raise ServiceDataError("customer_blocked", "该技师暂不接受你的服务请求。")
- service_settings = await get_service_settings()
- max_open = service_settings["customer_max_open_orders"]
- open_count = await ordersdb.count_documents({"customer_id": int(customer.id), "status": {"$in": list(ORDER_ACTIVE_STATUSES)}})
- if open_count >= max_open:
- raise ServiceDataError("too_many_open_orders", f"最多同时保留 {max_open} 个未关闭服务单。")
- day_start = utc_now().replace(hour=0, minute=0, second=0, microsecond=0)
- daily_limit = service_settings["customer_daily_request_limit"]
- if await ordersdb.count_documents({"customer_id": int(customer.id), "created_at": {"$gte": day_start}}) >= daily_limit:
- raise ServiceDataError("daily_request_limit", "今天发起的服务请求已达到上限。")
- profile, package = await _get_package(int(technician_id), str(package_id))
- if service_mode not in package["service_modes"]:
- raise ServiceDataError("service_mode_unavailable", "该套餐不支持所选服务方式。")
- technician_location = await locationsdb.find_one({"user_id": int(technician_id)})
- if not technician_location:
- raise ServiceDataError("technician_location_required", "技师尚未配置服务位置。")
- address_payload: dict[str, Any] | None = None
- distance = 0.0
- if service_mode == "onsite":
- if longitude is None or latitude is None:
- raise ServiceDataError("customer_location_required", "上门服务必须提供位置。")
- lon, lat = float(longitude), float(latitude)
- if not -180 <= lon <= 180 or not -90 <= lat <= 90:
- raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
- address = _clean_text(address_text, max_length=300, required=True)
- point = technician_location.get("point", {}).get("coordinates") or [technician_location.get("longitude"), technician_location.get("latitude")]
- distance = _distance_meters(float(point[0]), float(point[1]), lon, lat)
- radius = float(package.get("service_radius_km") or 0)
- if radius and distance > radius * 1000:
- raise ServiceDataError("outside_service_radius", "顾客位置超出该套餐的上门服务范围。")
- address_payload = {"kind": "customer", "longitude": lon, "latitude": lat, "address_text": address}
- else:
- if longitude is None or latitude is None:
- raise ServiceDataError(
- "customer_location_required", "到店服务必须先通过附近查找提供位置。"
- )
- lon, lat = float(longitude), float(latitude)
- if not -180 <= lon <= 180 or not -90 <= lat <= 90:
- raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
- point = technician_location.get("point", {}).get("coordinates") or [
- technician_location.get("longitude"),
- technician_location.get("latitude"),
- ]
- distance = _distance_meters(float(point[0]), float(point[1]), lon, lat)
- venue = (profile.get("service_profile") or {}).get("venue") or {}
- address_payload = {
- "kind": "venue",
- "longitude": technician_location.get("longitude"),
- "latitude": technician_location.get("latitude"),
- "address_text": _clean_text(venue.get("address_hint"), max_length=300, required=True),
- }
- now = utc_now()
- order_id = uuid4().hex
- order = {
- "order_id": order_id,
- "customer_id": int(customer.id),
- "customer_name": _clean_text(getattr(customer, "first_name", ""), max_length=80) or f"用户 {customer.id}",
- "technician_id": int(technician_id),
- "technician_name": profile.get("display_name") or f"技师 {technician_id}",
- "package_id": str(package_id),
- "package_snapshot": package,
- "category": package["category"],
- "service_mode": service_mode,
- "requirements": _clean_text(requirements, max_length=800, required=True),
- "distance_meters": round(distance, 2),
- "distance_band": distance_band(distance),
- "status": "requested",
- "version": 1,
- "created_at": now,
- "updated_at": now,
- }
- encrypted_address = _encrypt_address(address_payload)
- await ordersdb.insert_one(order)
- try:
- await addressesdb.insert_one(
- {
- "order_id": order_id,
- "encrypted_payload": encrypted_address,
- "redacted": False,
- "redact_after": None,
- "created_at": now,
- "updated_at": now,
- }
- )
- except Exception:
- await ordersdb.delete_one({"order_id": order_id, "status": "requested"})
- raise
- await record_service_event(order_id, "order_requested", actor_id=int(customer.id))
- return order
- def normalize_quote(value: dict[str, Any], *, package: dict[str, Any]) -> dict[str, Any]:
- currency = str(value.get("currency") or package.get("currency") or "CNY").upper()
- if currency != str(package.get("currency") or currency):
- raise ServiceDataError("currency_mismatch", "最终报价币种必须与套餐一致。")
- base = _positive_decimal(value.get("base_amount"), "服务金额")
- travel = _positive_decimal(value.get("travel_fee"), "交通费")
- discount = _positive_decimal(value.get("discount"), "优惠金额")
- addons = []
- addon_total = Decimal("0")
- for item in value.get("addons", []) or []:
- quantity = max(1, min(999, int(item.get("quantity") or 1)))
- unit_price = Decimal(_positive_decimal(item.get("unit_price"), "加项单价"))
- addons.append({"name": _clean_text(item.get("name"), max_length=60, required=True), "quantity": quantity, "unit_price": format(unit_price, "f")})
- addon_total += unit_price * quantity
- total = Decimal(base) + Decimal(travel) + addon_total - Decimal(discount)
- if total < 0:
- raise ServiceDataError("invalid_total", "最终报价总额不能小于 0。")
- provided = value.get("total_amount")
- if provided not in (None, "") and Decimal(_positive_decimal(provided, "总价")) != total:
- raise ServiceDataError("quote_total_mismatch", "报价明细与总价不一致。")
- return {
- "currency": currency,
- "base_amount": base,
- "travel_fee": travel,
- "addons": addons,
- "discount": discount,
- "total_amount": format(total.quantize(Decimal("0.01")), "f"),
- "note": _clean_text(value.get("note"), max_length=500),
- }
- async def submit_service_quote(order_id: str, technician_id: int, values: dict[str, Any]) -> tuple[dict[str, Any], dict[str, Any]]:
- await ensure_service_indexes()
- order = await ordersdb.find_one({"order_id": str(order_id)})
- if not order or int(order["technician_id"]) != int(technician_id):
- raise ServiceDataError("order_not_found", "未找到可报价服务单。")
- if order["status"] not in {"requested", "quoted"}:
- raise ServiceDataError("invalid_order_state", "当前服务单不能修改报价。")
- quote = normalize_quote(values, package=order["package_snapshot"])
- version = int(order.get("quote_version") or 0) + 1
- quote.update({"quote_id": uuid4().hex, "order_id": str(order_id), "version": version, "created_by": int(technician_id), "created_at": utc_now()})
- try:
- await quotesdb.insert_one(quote)
- except DuplicateKeyError as exc:
- raise ServiceDataError(
- "concurrent_order_update", "服务单状态已变化,请刷新后重试。"
- ) from exc
- updated = await ordersdb.find_one_and_update(
- {
- "order_id": str(order_id),
- "status": {"$in": ["requested", "quoted"]},
- "version": order["version"],
- },
- {
- "$set": {
- "status": "quoted",
- "current_quote_id": quote["quote_id"],
- "quote_snapshot": quote,
- "quote_version": version,
- "updated_at": utc_now(),
- },
- "$inc": {"version": 1},
- },
- return_document=ReturnDocument.AFTER,
- )
- if not updated:
- await quotesdb.delete_one({"quote_id": quote["quote_id"]})
- raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
- await record_service_event(order_id, "quote_submitted", actor_id=int(technician_id), metadata={"quote_id": quote["quote_id"], "version": version})
- return updated, quote
- async def confirm_service_quote(order_id: str, customer_id: int) -> dict[str, Any]:
- order = await ordersdb.find_one_and_update(
- {"order_id": str(order_id), "customer_id": int(customer_id), "status": "quoted"},
- {"$set": {"status": "confirmed", "confirmed_at": utc_now(), "updated_at": utc_now()}, "$inc": {"version": 1}},
- return_document=ReturnDocument.AFTER,
- )
- if not order:
- raise ServiceDataError("invalid_order_state", "当前报价无法确认。")
- await record_service_event(order_id, "quote_confirmed", actor_id=int(customer_id))
- return order
- async def reject_service_order(order_id: str, technician_id: int, reason: str) -> dict[str, Any]:
- now = utc_now()
- order = await ordersdb.find_one_and_update(
- {"order_id": str(order_id), "technician_id": int(technician_id), "status": {"$in": ["requested", "quoted"]}},
- {"$set": {"status": "rejected", "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
- return_document=ReturnDocument.AFTER,
- )
- if not order:
- raise ServiceDataError("invalid_order_state", "当前服务单无法拒绝。")
- await _schedule_address_redaction(order_id, now)
- await record_service_event(order_id, "order_rejected", actor_id=int(technician_id), reason=reason)
- return order
- async def start_service_order(order_id: str, technician_id: int) -> dict[str, Any]:
- order = await ordersdb.find_one_and_update(
- {"order_id": str(order_id), "technician_id": int(technician_id), "status": "confirmed"},
- {"$set": {"status": "in_progress", "started_at": utc_now(), "updated_at": utc_now()}, "$inc": {"version": 1}},
- return_document=ReturnDocument.AFTER,
- )
- if not order:
- raise ServiceDataError("invalid_order_state", "只有已确认服务单可以开始。")
- await record_service_event(order_id, "service_started", actor_id=int(technician_id))
- return order
- async def cancel_service_order(order_id: str, actor_id: int, reason: str) -> dict[str, Any]:
- order = await ordersdb.find_one({"order_id": str(order_id)})
- if not order or int(actor_id) not in {int(order["customer_id"]), int(order["technician_id"])}:
- raise ServiceDataError("order_not_found", "未找到该服务单。")
- if order["status"] not in ORDER_ACTIVE_STATUSES - {"disputed"}:
- raise ServiceDataError("invalid_order_state", "当前服务单不能取消。")
- status = "canceled_customer" if int(actor_id) == int(order["customer_id"]) else "canceled_technician"
- now = utc_now()
- updated = await ordersdb.find_one_and_update(
- {"order_id": str(order_id), "status": order["status"], "version": order["version"]},
- {"$set": {"status": status, "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
- return_document=ReturnDocument.AFTER,
- )
- if not updated:
- raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
- await _schedule_address_redaction(order_id, now)
- await qrdb.update_many({"order_id": str(order_id), "status": "issued"}, {"$set": {"status": "revoked", "revoked_at": now, "revoke_reason": "服务单已取消"}})
- await record_service_event(order_id, "order_canceled", actor_id=int(actor_id), reason=reason, metadata={"status": status})
- return updated
- async def dispute_service_order(order_id: str, actor_id: int, reason: str) -> dict[str, Any]:
- current = await ordersdb.find_one(
- {
- "order_id": str(order_id),
- "status": {"$in": list(ORDER_ACTIVE_STATUSES - {"disputed"})},
- "$or": [{"customer_id": int(actor_id)}, {"technician_id": int(actor_id)}],
- }
- )
- if not current:
- raise ServiceDataError("invalid_order_state", "当前服务单不能发起争议。")
- order = await ordersdb.find_one_and_update(
- {"order_id": str(order_id), "status": current["status"], "version": current["version"]},
- {"$set": {"status": "disputed", "status_before_dispute": current["status"], "dispute_reason": _clean_text(reason, max_length=500, required=True), "updated_at": utc_now()}, "$inc": {"version": 1}},
- return_document=ReturnDocument.AFTER,
- )
- if not order:
- raise ServiceDataError("invalid_order_state", "当前服务单不能发起争议。")
- await qrdb.update_many({"order_id": str(order_id), "status": "issued"}, {"$set": {"status": "revoked", "revoked_at": utc_now(), "revoke_reason": "服务单争议"}})
- await record_service_event(order_id, "order_disputed", actor_id=int(actor_id), reason=reason)
- return order
- async def resolve_service_dispute(
- order_id: str,
- *,
- actor_id: str,
- action: str,
- reason: str,
- ) -> dict[str, Any]:
- if action == "void":
- return await admin_void_service_order(order_id, actor_id=actor_id, reason=reason)
- if action != "resume":
- raise ServiceDataError("invalid_dispute_action", "不支持该争议处理操作。")
- current = await ordersdb.find_one({"order_id": str(order_id), "status": "disputed"})
- if not current:
- raise ServiceDataError("invalid_order_state", "该服务单当前不在争议中。")
- resume_status = str(current.get("status_before_dispute") or "confirmed")
- if resume_status == "completion_pending":
- resume_status = "in_progress"
- if resume_status not in {"requested", "quoted", "confirmed", "in_progress"}:
- resume_status = "confirmed"
- updated = await ordersdb.find_one_and_update(
- {"order_id": str(order_id), "status": "disputed", "version": current["version"]},
- {
- "$set": {
- "status": resume_status,
- "dispute_resolved_at": utc_now(),
- "dispute_resolution_reason": _clean_text(reason, max_length=500, required=True),
- "updated_at": utc_now(),
- },
- "$unset": {"status_before_dispute": ""},
- "$inc": {"version": 1},
- },
- return_document=ReturnDocument.AFTER,
- )
- if not updated:
- raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
- await record_service_event(order_id, "dispute_resolved", actor_id=actor_id, reason=reason, metadata={"status": resume_status})
- return updated
- async def admin_void_service_order(order_id: str, *, actor_id: str, reason: str) -> dict[str, Any]:
- now = utc_now()
- order = await ordersdb.find_one_and_update(
- {"order_id": str(order_id), "status": {"$ne": "voided"}},
- {"$set": {"status": "voided", "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
- return_document=ReturnDocument.AFTER,
- )
- if not order:
- raise ServiceDataError("order_not_found", "未找到可作废服务单。")
- await _schedule_address_redaction(order_id, now)
- await qrdb.update_many({"order_id": str(order_id), "status": {"$ne": "voided"}}, {"$set": {"status": "voided", "revoked_at": now}})
- await reviewsdb.update_many({"order_id": str(order_id), "status": "approved"}, {"$set": {"status": "voided", "moderated_at": now, "moderation_reason": reason, "updated_at": now}})
- await recompute_review_stats()
- await record_service_event(order_id, "order_voided", actor_id=actor_id, reason=reason)
- return order
- async def issue_review_qr(order_id: str, technician_id: int) -> tuple[dict[str, Any], str]:
- await ensure_service_indexes()
- order = await ordersdb.find_one({"order_id": str(order_id), "technician_id": int(technician_id)})
- if order and order.get("status") == "completion_pending":
- previous = await qrdb.find_one(
- {"order_id": str(order_id), "status": "issued"},
- sort=[("created_at", DESCENDING)],
- )
- if previous and _as_utc(previous["expires_at"]) <= utc_now():
- now = utc_now()
- expired = await qrdb.find_one_and_update(
- {"qr_id": previous["qr_id"], "status": "issued"},
- {"$set": {"status": "expired", "expired_at": now}},
- return_document=ReturnDocument.AFTER,
- )
- if expired:
- await ordersdb.update_one(
- {
- "order_id": str(order_id),
- "status": "completion_pending",
- "current_review_qr_id": previous["qr_id"],
- },
- {
- "$set": {"status": "in_progress", "updated_at": now},
- "$unset": {"current_review_qr_id": ""},
- "$inc": {"version": 1},
- },
- )
- order = await ordersdb.find_one(
- {"order_id": str(order_id), "technician_id": int(technician_id)}
- )
- if not order or order.get("status") != "in_progress":
- raise ServiceDataError("invalid_order_state", "只有服务中的订单可以生成评价二维码。")
- if await qrdb.find_one({"order_id": str(order_id), "status": {"$in": ["issued", "claimed"]}}):
- raise ServiceDataError("review_qr_exists", "该服务单已经生成评价二维码。")
- token = token_urlsafe(18)
- now = utc_now()
- hours = (await get_service_settings())["review_qr_expiry_hours"]
- document = {
- "qr_id": uuid4().hex,
- "order_id": str(order_id),
- "technician_id": int(technician_id),
- "customer_id": int(order["customer_id"]),
- "package_snapshot": order["package_snapshot"],
- "quote_snapshot": order.get("quote_snapshot"),
- "token_hash": hashlib.sha256(token.encode()).hexdigest(),
- "status": "issued",
- "expires_at": now + timedelta(hours=hours),
- "created_at": now,
- }
- try:
- await qrdb.insert_one(document)
- except DuplicateKeyError as exc:
- raise ServiceDataError("review_qr_exists", "该服务单已经生成评价二维码。") from exc
- updated = await ordersdb.find_one_and_update(
- {
- "order_id": str(order_id),
- "status": "in_progress",
- "version": order["version"],
- },
- {
- "$set": {
- "status": "completion_pending",
- "current_review_qr_id": document["qr_id"],
- "completion_requested_at": now,
- "updated_at": now,
- },
- "$inc": {"version": 1},
- },
- return_document=ReturnDocument.AFTER,
- )
- if not updated:
- await qrdb.update_one(
- {"qr_id": document["qr_id"], "status": "issued"},
- {
- "$set": {
- "status": "voided",
- "revoked_at": utc_now(),
- "revoke_reason": "服务单状态已变化",
- }
- },
- )
- raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
- await record_service_event(order_id, "review_qr_issued", actor_id=int(technician_id), metadata={"qr_id": document["qr_id"]})
- return document, token
- async def preview_review_qr(token: str, customer_id: int) -> dict[str, Any]:
- document = await qrdb.find_one({"token_hash": hashlib.sha256(str(token).encode()).hexdigest()})
- if not document:
- raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
- if int(document["customer_id"]) != int(customer_id):
- raise ServiceDataError("review_qr_customer_mismatch", "该二维码不属于当前顾客。")
- if document.get("status") == "claimed":
- if await reviewsdb.find_one({"qr_id": document["qr_id"]}):
- raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
- return document
- if document.get("status") != "issued":
- raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
- if _as_utc(document["expires_at"]) <= utc_now():
- now = utc_now()
- expired = await qrdb.find_one_and_update(
- {"qr_id": document["qr_id"], "status": "issued"},
- {"$set": {"status": "expired", "expired_at": now}},
- return_document=ReturnDocument.AFTER,
- )
- if expired:
- await ordersdb.update_one(
- {
- "order_id": document["order_id"],
- "status": "completion_pending",
- "current_review_qr_id": document["qr_id"],
- },
- {
- "$set": {"status": "in_progress", "updated_at": now},
- "$unset": {"current_review_qr_id": ""},
- "$inc": {"version": 1},
- },
- )
- raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
- return document
- async def claim_review_qr(token: str, customer_id: int) -> tuple[dict[str, Any], dict[str, Any]]:
- token_hash = hashlib.sha256(str(token).encode()).hexdigest()
- now = utc_now()
- await preview_review_qr(token, customer_id)
- document = await qrdb.find_one_and_update(
- {"token_hash": token_hash, "customer_id": int(customer_id), "status": "issued"},
- {"$set": {"status": "claimed", "claimed_at": now}},
- return_document=ReturnDocument.AFTER,
- )
- if not document:
- raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
- order = await ordersdb.find_one_and_update(
- {"order_id": document["order_id"], "status": "completion_pending", "customer_id": int(customer_id)},
- {"$set": {"status": "completed", "completed_at": now, "closed_at": now, "updated_at": now}, "$inc": {"version": 1}},
- return_document=ReturnDocument.AFTER,
- )
- if not order:
- await qrdb.update_one({"qr_id": document["qr_id"]}, {"$set": {"status": "voided"}})
- raise ServiceDataError("invalid_order_state", "服务单状态已变化,无法完成评价。")
- await _schedule_address_redaction(order["order_id"], now)
- await record_service_event(order["order_id"], "service_completed", actor_id=int(customer_id), metadata={"qr_id": document["qr_id"]})
- return document, order
- def normalize_review_template(values: dict[str, Any]) -> dict[str, Any]:
- ratings = []
- for item in values.get("rating_questions", []):
- if len(ratings) >= 5:
- raise ServiceDataError("too_many_review_questions", "评分维度最多 5 个。")
- weight = int(item.get("weight") or 0)
- if weight <= 0:
- raise ServiceDataError("invalid_review_weight", "评分权重必须大于 0。")
- ratings.append({"question_id": str(item.get("question_id") or uuid4().hex[:12]), "label": _clean_text(item.get("label"), max_length=40, required=True), "description": _clean_text(item.get("description"), max_length=160), "required": bool(item.get("required", True)), "weight": weight})
- if not ratings or sum(item["weight"] for item in ratings) != 100:
- raise ServiceDataError("invalid_review_weight", "评分维度权重合计必须为 100%。")
- texts = []
- for item in values.get("text_questions", []):
- if len(texts) >= 2:
- raise ServiceDataError("too_many_review_questions", "文字问题最多 2 个。")
- texts.append({"question_id": str(item.get("question_id") or uuid4().hex[:12]), "label": _clean_text(item.get("label"), max_length=40, required=True), "description": _clean_text(item.get("description"), max_length=160), "required": bool(item.get("required")), "max_length": max(50, min(1000, int(item.get("max_length") or 500)))})
- return {"rating_questions": ratings, "text_questions": texts}
- async def get_active_review_template() -> dict[str, Any]:
- await ensure_service_indexes()
- current = await reviewtemplatesdb.find_one({"status": "active"}, sort=[("version", DESCENDING)])
- if current:
- return current
- now = utc_now()
- document = {"template_id": uuid4().hex, "version": 1, "status": "active", **normalize_review_template(DEFAULT_REVIEW_TEMPLATE), "created_at": now, "published_at": now}
- try:
- await reviewtemplatesdb.insert_one(document)
- except DuplicateKeyError:
- return await reviewtemplatesdb.find_one({"status": "active"}, sort=[("version", DESCENDING)]) or document
- return document
- async def publish_review_template(values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
- normalized = normalize_review_template(values)
- await ensure_service_indexes()
- while True:
- current = await reviewtemplatesdb.find_one(sort=[("version", DESCENDING)])
- version = int((current or {}).get("version") or 0) + 1
- now = utc_now()
- document = {
- "template_id": uuid4().hex,
- "version": version,
- "status": "active",
- **normalized,
- "published_by": str(actor_id),
- "created_at": now,
- "published_at": now,
- }
- try:
- await reviewtemplatesdb.insert_one(document)
- break
- except DuplicateKeyError:
- continue
- await reviewtemplatesdb.update_many(
- {"status": "active", "template_id": {"$ne": document["template_id"]}},
- {"$set": {"status": "retired", "retired_at": now}},
- )
- return document
- async def begin_review_draft(
- *,
- customer_id: int,
- technician_id: int,
- source: str,
- template: dict[str, Any],
- package_id: str | None = None,
- qr_id: str | None = None,
- ) -> dict[str, Any]:
- await ensure_service_indexes()
- now = utc_now()
- current = await reviewdraftsdb.find_one({"customer_id": int(customer_id)})
- identity = {
- "technician_id": int(technician_id),
- "source": str(source),
- "package_id": str(package_id) if package_id else None,
- "qr_id": str(qr_id) if qr_id else None,
- }
- if (
- current
- and all(current.get(key) == value for key, value in identity.items())
- and _as_utc(current["expires_at"]) > now
- ):
- return current
- document = {
- "draft_id": uuid4().hex,
- "customer_id": int(customer_id),
- **identity,
- "template_snapshot": {
- key: template[key]
- for key in (
- "template_id",
- "version",
- "rating_questions",
- "text_questions",
- )
- },
- "rating_index": 0,
- "text_index": 0,
- "answers": {},
- "created_at": now,
- "updated_at": now,
- "expires_at": now + timedelta(days=30),
- }
- await reviewdraftsdb.update_one(
- {"customer_id": int(customer_id)}, {"$set": document}, upsert=True
- )
- return document
- async def save_review_draft_progress(
- customer_id: int,
- *,
- answers: dict[str, Any],
- rating_index: int,
- text_index: int,
- ) -> None:
- await reviewdraftsdb.update_one(
- {"customer_id": int(customer_id)},
- {
- "$set": {
- "answers": dict(answers),
- "rating_index": max(0, int(rating_index)),
- "text_index": max(0, int(text_index)),
- "updated_at": utc_now(),
- }
- },
- )
- async def delete_review_draft(customer_id: int) -> None:
- await reviewdraftsdb.delete_one({"customer_id": int(customer_id)})
- def _score_review(template: dict[str, Any], answers: dict[str, Any]) -> tuple[float, dict[str, int], dict[str, str]]:
- ratings: dict[str, int] = {}
- texts: dict[str, str] = {}
- total = 0.0
- for question in template["rating_questions"]:
- raw = answers.get(question["question_id"])
- if raw in (None, "") and not question["required"]:
- continue
- try:
- score = int(raw)
- except (TypeError, ValueError) as exc:
- raise ServiceDataError("invalid_review_answer", "评分必须是 1 到 5 的整数。") from exc
- if not 1 <= score <= 5:
- raise ServiceDataError("invalid_review_answer", "评分必须在 1 到 5 之间。")
- ratings[question["question_id"]] = score
- total += score * question["weight"] / 100
- for question in template["text_questions"]:
- text = _clean_text(answers.get(question["question_id"]), max_length=question["max_length"], required=question["required"])
- if text:
- texts[question["question_id"]] = text
- return round(total, 4), ratings, texts
- async def submit_review(
- *,
- customer: Any,
- technician_id: int,
- source: str,
- answers: dict[str, Any],
- anonymous: bool,
- package_id: str | None = None,
- qr_id: str | None = None,
- template_snapshot: dict[str, Any] | None = None,
- ) -> dict[str, Any]:
- await observe_customer(customer)
- await require_customer_allowed(int(customer.id))
- if source not in REVIEW_SOURCES:
- raise ServiceDataError("invalid_review_source", "评价来源无效。")
- profile = await profilesdb.find_one({"user_id": int(technician_id), "application_status": "approved"})
- if not profile:
- raise ServiceDataError("technician_not_found", "未找到已认证技师。")
- if int(customer.id) == int(technician_id):
- raise ServiceDataError("self_review_not_allowed", "不能评价自己。")
- order_id = None
- package_snapshot = None
- if source == "qr_verified":
- qr = await qrdb.find_one({"qr_id": str(qr_id), "status": "claimed", "customer_id": int(customer.id), "technician_id": int(technician_id)})
- if not qr:
- raise ServiceDataError("review_qr_required", "缺少已完成服务单的评价资格。")
- if await reviewsdb.find_one({"qr_id": str(qr_id)}):
- raise ServiceDataError("review_already_submitted", "该服务单已经提交评价。")
- order_id = qr["order_id"]
- package_snapshot = qr["package_snapshot"]
- else:
- service_profile = profile.get("service_profile") or {}
- package_snapshot = next((item for item in service_profile.get("packages", []) if str(item.get("package_id")) == str(package_id)), None)
- if package_id and not package_snapshot:
- raise ServiceDataError("package_not_found", "未找到所选套餐。")
- package_snapshot = package_snapshot or {"package_id": None, "name": "其他服务", "category": "其他"}
- template = template_snapshot or await get_active_review_template()
- score, rating_answers, text_answers = _score_review(template, answers)
- normalized_content = "\n".join(
- value.casefold() for _, value in sorted(text_answers.items()) if value
- )
- now = utc_now()
- document = {
- "review_id": uuid4().hex,
- "technician_id": int(technician_id),
- "technician_name": profile.get("display_name") or f"技师 {technician_id}",
- "customer_id": int(customer.id),
- "customer_name": _clean_text(getattr(customer, "first_name", ""), max_length=80) or f"用户 {customer.id}",
- "source": source,
- "package_snapshot": package_snapshot,
- "category": package_snapshot.get("category") or "其他",
- "template_snapshot": {key: template[key] for key in ("template_id", "version", "rating_questions", "text_questions")},
- "rating_answers": rating_answers,
- "text_answers": text_answers,
- "content_fingerprint": (
- hashlib.sha256(normalized_content.encode()).hexdigest()
- if normalized_content
- else None
- ),
- "score": score,
- "anonymous": bool(anonymous),
- "status": "pending",
- "created_at": now,
- "updated_at": now,
- }
- if qr_id:
- document["qr_id"] = str(qr_id)
- document["order_id"] = order_id
- try:
- await reviewsdb.insert_one(document)
- except DuplicateKeyError as exc:
- if qr_id:
- raise ServiceDataError(
- "review_already_submitted", "该服务单已经提交评价。"
- ) from exc
- raise
- return document
- async def moderate_review(review_id: str, action: str, *, actor_id: str, reason: str = "") -> dict[str, Any]:
- status_map = {"approve": "approved", "reject": "rejected", "void": "voided"}
- source_status = {"approve": "pending", "reject": "pending", "void": "approved"}
- if action not in status_map:
- raise ServiceDataError("invalid_review_action", "不支持该评价操作。")
- current = await reviewsdb.find_one({"review_id": str(review_id)})
- if not current:
- raise ServiceDataError("review_not_found", "未找到评价。")
- if current["status"] != source_status[action]:
- raise ServiceDataError("invalid_review_state", "当前评价状态不能执行该审核操作。")
- if action in {"reject", "void"} and not _clean_text(reason, max_length=500):
- raise ServiceDataError("reason_required", "拒绝或作废必须填写原因。")
- updated = await reviewsdb.find_one_and_update(
- {"review_id": str(review_id), "status": source_status[action]},
- {"$set": {"status": status_map[action], "moderated_by": str(actor_id), "moderated_at": utc_now(), "moderation_reason": _clean_text(reason, max_length=500), "updated_at": utc_now()}},
- return_document=ReturnDocument.AFTER,
- )
- if not updated:
- raise ServiceDataError("concurrent_review_update", "评价状态已变化,请刷新后重试。")
- await recompute_review_stats()
- return updated
- async def recompute_review_stats() -> None:
- await ensure_service_indexes()
- approved = await reviewsdb.find({"status": "approved"}).to_list(length=100000)
- grouped: dict[tuple[int, str], list[dict[str, Any]]] = defaultdict(list)
- category_scores: dict[str, list[float]] = defaultdict(list)
- for review in approved:
- key = (int(review["technician_id"]), str(review.get("category") or "其他"))
- grouped[key].append(review)
- category_scores[key[1]].append(float(review["score"]))
- now = utc_now()
- await statsdb.delete_many({})
- documents = []
- for (technician_id, category), values in grouped.items():
- count = len(values)
- average = sum(float(item["score"]) for item in values) / count
- category_average = sum(category_scores[category]) / len(category_scores[category])
- rank_score = count / (count + 5) * average + 5 / (count + 5) * category_average
- dimensions: dict[str, list[int]] = defaultdict(list)
- for review in values:
- for key, score in review.get("rating_answers", {}).items():
- dimensions[key].append(int(score))
- documents.append(
- {
- "technician_id": technician_id,
- "category": category,
- "review_count": count,
- "average_score": round(average, 4),
- "category_average": round(category_average, 4),
- "rank_score": round(rank_score, 6),
- "eligible": count >= 3,
- "dimension_averages": {key: round(sum(items) / len(items), 4) for key, items in dimensions.items()},
- "last_review_at": max(item.get("moderated_at") or item["created_at"] for item in values),
- "updated_at": now,
- }
- )
- if documents:
- await statsdb.insert_many(documents)
- async def list_leaderboard(*, category: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
- await ensure_service_indexes()
- filters: dict[str, Any] = {"eligible": True}
- if category:
- filters["category"] = str(category)
- visible_ids = [
- int(item["user_id"])
- async for item in profilesdb.find(
- {
- "application_status": "approved",
- "listed": True,
- "username": {"$nin": [None, ""]},
- },
- {"user_id": 1},
- )
- ]
- filters["technician_id"] = {"$in": visible_ids}
- total = await statsdb.count_documents(filters)
- items = await statsdb.find(filters).sort([("rank_score", DESCENDING), ("review_count", DESCENDING), ("last_review_at", DESCENDING)]).skip((max(1, page) - 1) * page_size).limit(page_size).to_list(length=page_size)
- profiles = {int(item["user_id"]): item async for item in profilesdb.find({"user_id": {"$in": [item["technician_id"] for item in items]}, "application_status": "approved", "listed": True})}
- visible = [{**item, "display_name": profiles[item["technician_id"]].get("display_name"), "username": profiles[item["technician_id"]].get("username")} for item in items if item["technician_id"] in profiles]
- return visible, total
- async def list_public_reviews(
- technician_id: int,
- *,
- page: int = 1,
- page_size: int = 10,
- ) -> tuple[list[dict[str, Any]], int]:
- filters = {"technician_id": int(technician_id), "status": "approved"}
- total = await reviewsdb.count_documents(filters)
- items = await reviewsdb.find(filters).sort("moderated_at", DESCENDING).skip(
- (max(1, page) - 1) * page_size
- ).limit(page_size).to_list(length=page_size)
- public_items = []
- for item in items:
- public_items.append(
- {
- "review_id": item["review_id"],
- "source": item["source"],
- "source_label": (
- "完成服务单评价"
- if item["source"] == "qr_verified"
- else "用户主动评价"
- ),
- "customer_name": "匿名顾客" if item.get("anonymous", True) else item.get("customer_name"),
- "package_name": (item.get("package_snapshot") or {}).get("name"),
- "category": item.get("category"),
- "score": item.get("score"),
- "rating_answers": item.get("rating_answers", {}),
- "text_answers": item.get("text_answers", {}),
- "approved_at": item.get("moderated_at"),
- }
- )
- return public_items, total
- async def get_technician_review_summary(technician_id: int) -> dict[str, Any]:
- items = await reviewsdb.find(
- {"technician_id": int(technician_id), "status": "approved"},
- {"score": 1, "source": 1},
- ).to_list(length=100000)
- count = len(items)
- return {
- "review_count": count,
- "average_score": round(
- sum(float(item["score"]) for item in items) / count,
- 4,
- )
- if count
- else 0,
- "source_counts": {
- source: sum(1 for item in items if item.get("source") == source)
- for source in REVIEW_SOURCES
- },
- "ranking_title": "审核评价排行",
- }
- async def list_service_orders(*, query: str = "", status: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
- await ensure_service_indexes()
- filters: dict[str, Any] = {}
- if status:
- filters["status"] = status
- if query:
- if query.isdigit():
- filters["$or"] = [{"customer_id": int(query)}, {"technician_id": int(query)}]
- else:
- pattern = re.compile(re.escape(query), re.IGNORECASE)
- filters["$or"] = [{"order_id": pattern}, {"customer_name": pattern}, {"technician_name": pattern}]
- total = await ordersdb.count_documents(filters)
- items = await ordersdb.find(filters).sort("created_at", DESCENDING).skip((max(1, page) - 1) * page_size).limit(page_size).to_list(length=page_size)
- return items, total
- async def list_actor_service_orders(
- actor_id: int,
- *,
- role: str = "all",
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- if role == "customer":
- filters: dict[str, Any] = {"customer_id": int(actor_id)}
- elif role == "technician":
- filters = {"technician_id": int(actor_id)}
- else:
- filters = {
- "$or": [
- {"customer_id": int(actor_id)},
- {"technician_id": int(actor_id)},
- ]
- }
- total = await ordersdb.count_documents(filters)
- items = await ordersdb.find(filters).sort("created_at", DESCENDING).skip(
- (max(1, page) - 1) * page_size
- ).limit(page_size).to_list(length=page_size)
- return items, total
- async def get_admin_service_order(order_id: str) -> dict[str, Any]:
- order = await ordersdb.find_one({"order_id": str(order_id)})
- if not order:
- raise ServiceDataError("order_not_found", "未找到服务单。")
- result = dict(order)
- result["quotes"] = await quotesdb.find({"order_id": str(order_id)}).sort(
- "version", ASCENDING
- ).to_list(length=100)
- result["events"] = await eventsdb.find({"order_id": str(order_id)}).sort(
- "created_at", ASCENDING
- ).to_list(length=500)
- result["qr_records"] = await qrdb.find({"order_id": str(order_id)}).sort(
- "created_at", DESCENDING
- ).to_list(length=100)
- return result
- async def list_review_qr_records(
- *,
- status: str = "",
- query: str = "",
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- filters: dict[str, Any] = {}
- if status:
- filters["status"] = status
- if query:
- filters["$or"] = [
- {"order_id": re.compile(re.escape(query), re.IGNORECASE)},
- {"qr_id": re.compile(re.escape(query), re.IGNORECASE)},
- {"customer_id": int(query) if query.isdigit() else -1},
- {"technician_id": int(query) if query.isdigit() else -1},
- ]
- total = await qrdb.count_documents(filters)
- items = await qrdb.find(filters).sort("created_at", DESCENDING).skip(
- (max(1, page) - 1) * page_size
- ).limit(page_size).to_list(length=page_size)
- for item in items:
- item.pop("token_hash", None)
- return items, total
- async def list_reviews(*, status: str = "", source: str = "", query: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
- await ensure_service_indexes()
- filters: dict[str, Any] = {}
- if status:
- filters["status"] = status
- if source:
- filters["source"] = source
- if query:
- pattern = re.compile(re.escape(query), re.IGNORECASE)
- filters["$or"] = [{"technician_name": pattern}, {"customer_name": pattern}, {"review_id": pattern}]
- total = await reviewsdb.count_documents(filters)
- items = await reviewsdb.find(filters).sort("created_at", DESCENDING).skip((max(1, page) - 1) * page_size).limit(page_size).to_list(length=page_size)
- for item in items:
- item["customer_review_count"] = await reviewsdb.count_documents({"customer_id": item["customer_id"], "technician_id": item["technician_id"]})
- item["customer_review_count_30d"] = await reviewsdb.count_documents({"customer_id": item["customer_id"], "technician_id": item["technician_id"], "created_at": {"$gte": utc_now() - timedelta(days=30)}})
- item["customer_review_count_90d"] = await reviewsdb.count_documents({"customer_id": item["customer_id"], "technician_id": item["technician_id"], "created_at": {"$gte": utc_now() - timedelta(days=90)}})
- item["similar_content_count"] = (
- await reviewsdb.count_documents(
- {
- "customer_id": item["customer_id"],
- "technician_id": item["technician_id"],
- "content_fingerprint": item["content_fingerprint"],
- }
- )
- if item.get("content_fingerprint")
- else 0
- )
- return items, total
- async def get_order_for_actor(order_id: str, actor_id: int, *, reveal_address: bool = False) -> dict[str, Any]:
- order = await ordersdb.find_one({"order_id": str(order_id)})
- if not order or int(actor_id) not in {int(order["customer_id"]), int(order["technician_id"])}:
- raise ServiceDataError("order_not_found", "未找到服务单。")
- result = dict(order)
- can_reveal = order["status"] in {"confirmed", "in_progress", "completion_pending", "completed", "disputed"}
- if reveal_address and can_reveal:
- address = await addressesdb.find_one({"order_id": str(order_id), "redacted": False})
- if address:
- result["exact_address"] = _decrypt_address(address["encrypted_payload"])
- return result
- async def admin_reveal_order_address(order_id: str, *, actor_id: str, reason: str) -> dict[str, Any]:
- reason = _clean_text(reason, max_length=500, required=True)
- order = await ordersdb.find_one({"order_id": str(order_id)})
- address = await addressesdb.find_one({"order_id": str(order_id), "redacted": False})
- if not order or not address:
- raise ServiceDataError("address_unavailable", "精确地址已脱敏或不存在。")
- await record_service_event(order_id, "admin_address_revealed", actor_id=actor_id, reason=reason)
- return _decrypt_address(address["encrypted_payload"])
- async def redact_expired_addresses() -> int:
- now = utc_now()
- candidates = await addressesdb.find({"redacted": False, "redact_after": {"$lte": now}}).to_list(length=1000)
- count = 0
- for item in candidates:
- order = await ordersdb.find_one({"order_id": item["order_id"]})
- if not order or order.get("status") not in ORDER_TERMINAL_STATUSES:
- continue
- result = await addressesdb.update_one({"_id": item["_id"], "redacted": False}, {"$set": {"redacted": True, "redacted_at": now, "updated_at": now}, "$unset": {"encrypted_payload": ""}})
- if result.modified_count:
- await ordersdb.update_one(
- {"order_id": item["order_id"]},
- {
- "$set": {"address_redacted": True, "updated_at": now},
- "$unset": {"distance_meters": ""},
- },
- )
- count += result.modified_count
- return count
- async def set_customer_block(*, customer_id: int, scope: str, technician_id: int | None, active: bool, actor_id: int | str, reason: str) -> dict[str, Any]:
- if scope not in {"global", "technician"}:
- raise ServiceDataError("invalid_block_scope", "拉黑范围无效。")
- if scope == "technician" and technician_id is None:
- raise ServiceDataError("technician_required", "技师拉黑必须指定技师。")
- now = utc_now()
- await blocksdb.update_one(
- {"scope": scope, "technician_id": int(technician_id or 0), "customer_id": int(customer_id)},
- {"$set": {"active": bool(active), "reason": _clean_text(reason, max_length=500, required=active), "actor_id": actor_id, "updated_at": now}, "$setOnInsert": {"created_at": now}},
- upsert=True,
- )
- return await blocksdb.find_one({"scope": scope, "technician_id": int(technician_id or 0), "customer_id": int(customer_id)}) or {}
- async def set_technician_customer_block(
- *,
- technician_id: int,
- customer_id: int,
- active: bool,
- reason: str,
- ) -> dict[str, Any]:
- profile = await profilesdb.find_one(
- {"user_id": int(technician_id), "application_status": "approved"}
- )
- if not profile:
- raise ServiceDataError("technician_required", "只有已认证技师可以管理个人拉黑。")
- if not await ordersdb.find_one(
- {"technician_id": int(technician_id), "customer_id": int(customer_id)}
- ):
- raise ServiceDataError("customer_relationship_required", "只能拉黑曾向你发起服务请求的顾客。")
- return await set_customer_block(
- customer_id=int(customer_id),
- scope="technician",
- technician_id=int(technician_id),
- active=active,
- actor_id=int(technician_id),
- reason=reason,
- )
- async def list_customer_blocks(
- *,
- scope: str = "global",
- active: bool | None = None,
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- filters: dict[str, Any] = {"scope": scope}
- if active is not None:
- filters["active"] = active
- total = await blocksdb.count_documents(filters)
- items = await blocksdb.find(filters).sort("updated_at", DESCENDING).skip(
- (max(1, page) - 1) * page_size
- ).limit(page_size).to_list(length=page_size)
- return items, total
- async def create_service_report(*, reporter_id: int, target_type: str, target_id: str, reason: str) -> dict[str, Any]:
- if target_type not in {"technician", "customer", "order", "review"}:
- raise ServiceDataError("invalid_report_target", "举报对象无效。")
- document = {"report_id": uuid4().hex, "reporter_id": int(reporter_id), "target_type": target_type, "target_id": str(target_id), "reason": _clean_text(reason, max_length=800, required=True), "status": "open", "created_at": utc_now()}
- await reportsdb.insert_one(document)
- return document
- async def list_service_reports(
- *,
- status: str = "",
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- filters = {"status": status} if status else {}
- total = await reportsdb.count_documents(filters)
- items = await reportsdb.find(filters).sort("created_at", DESCENDING).skip(
- (max(1, page) - 1) * page_size
- ).limit(page_size).to_list(length=page_size)
- return items, total
- async def moderate_service_report(
- report_id: str,
- *,
- action: str,
- actor_id: str,
- reason: str,
- ) -> dict[str, Any]:
- status_map = {"resolve": "resolved", "dismiss": "dismissed"}
- if action not in status_map:
- raise ServiceDataError("invalid_report_action", "不支持该举报操作。")
- updated = await reportsdb.find_one_and_update(
- {"report_id": str(report_id), "status": "open"},
- {
- "$set": {
- "status": status_map[action],
- "handled_by": str(actor_id),
- "handled_at": utc_now(),
- "handling_reason": _clean_text(reason, max_length=500, required=True),
- }
- },
- return_document=ReturnDocument.AFTER,
- )
- if not updated:
- raise ServiceDataError("report_not_found", "举报不存在或已处理。")
- return updated
- async def fulfillment_metrics() -> dict[str, Any]:
- await ensure_service_indexes()
- counts = {status: await ordersdb.count_documents({"status": status}) for status in ORDER_STATUSES}
- valid_requests = sum(counts.values()) - counts["voided"]
- quoted = await ordersdb.count_documents(
- {"status": {"$ne": "voided"}, "current_quote_id": {"$exists": True}}
- )
- confirmed = await ordersdb.count_documents(
- {"status": {"$ne": "voided"}, "confirmed_at": {"$exists": True}}
- )
- completed = counts["completed"]
- canceled = counts["canceled_customer"] + counts["canceled_technician"]
- directory_views = await eventsdb.count_documents({"event_type": "directory_viewed"})
- technician_views = await eventsdb.count_documents({"event_type": "technician_viewed"})
- request_starts = await eventsdb.count_documents({"event_type": "request_started"})
- return {
- "counts": counts,
- "valid_requests": valid_requests,
- "quoted_requests": quoted,
- "confirmed_requests": confirmed,
- "directory_views": directory_views,
- "technician_views": technician_views,
- "request_starts": request_starts,
- "request_submission_rate": min(1.0, round(valid_requests / request_starts, 4))
- if request_starts
- else 0,
- "quote_rate": round(quoted / valid_requests, 4) if valid_requests else 0,
- "customer_confirmation_rate": round(confirmed / quoted, 4) if quoted else 0,
- "match_success_rate": round(confirmed / valid_requests, 4) if valid_requests else 0,
- "fulfillment_rate": round(completed / confirmed, 4) if confirmed else 0,
- "cancellation_rate": round(canceled / valid_requests, 4) if valid_requests else 0,
- "cancellation_distribution": {
- "customer": counts["canceled_customer"],
- "technician": counts["canceled_technician"],
- },
- "scope_note": "仅统计平台服务单;直接 Telegram 私聊成交不在统计范围。",
- }
|