| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386 |
- 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
- reviewtopicsdb = control_db.teacher_review_topics
- 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 = {
- "questions": [
- {
- "question_id": "overall_experience",
- "type": "single_choice",
- "label": "这次服务体验怎么样?",
- "description": "选择最符合实际体验的一项。",
- "required": True,
- "weight": 100,
- "options": [
- {"value": "excellent", "label": "非常满意", "score": 5},
- {"value": "good", "label": "满意", "score": 4},
- {"value": "average", "label": "一般", "score": 3},
- {"value": "poor", "label": "不满意", "score": 2},
- {"value": "bad", "label": "非常不满意", "score": 1},
- ],
- },
- {
- "question_id": "service_highlights",
- "type": "multiple_choice",
- "label": "哪些方面值得肯定?",
- "description": "可多选,也可以跳过。",
- "required": False,
- "options": [
- {"value": "professional", "label": "专业可靠"},
- {"value": "communication", "label": "沟通顺畅"},
- {"value": "careful", "label": "服务细致"},
- {"value": "transparent", "label": "价格透明"},
- {"value": "punctual", "label": "准时守约"},
- ],
- },
- {
- "question_id": "comment",
- "type": "text",
- "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": "顾客到指定区域接受一次标准服务,具体地点和服务细节通过 Telegram 私聊沟通。",
- "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": "quoted",
- "amount": "0",
- "per_km": "0",
- "description": "无额外交通费;其他临时费用请提前沟通。",
- },
- "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": "技师提供一次标准上门服务,具体地点和服务细节通过 Telegram 私聊沟通。",
- "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",
- "description": "上门车费根据距离另行沟通。",
- },
- "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",
- "description": "上门车费、耗材费和超时费用另行沟通。",
- },
- "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") or ["at_store", "onsite"])
- )
- )
- if 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 公里之间。")
- out_of_range_policy = _clean_text(
- value.get("out_of_range_policy"),
- max_length=200,
- )
- 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)
- venue_address = _clean_text(venue.get("address_hint"), max_length=120)
- onsite_description = _clean_text(onsite.get("description"), max_length=300)
- 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)
- or "请通过 Telegram 私聊沟通具体服务信息",
- "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)
- await reviewtopicsdb.create_index([("technician_id", ASCENDING)], unique=True)
- await reviewtopicsdb.create_index([("forum_chat_id", ASCENDING), ("message_thread_id", ASCENDING)])
- _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, Any]:
- stored = await settingsdb.find_one({"settings_id": "global"}) or {}
- forum_chat_id = str(
- stored.get("service_review_forum_chat_id")
- or getattr(wbb, "SERVICE_REVIEW_FORUM_CHAT_ID", "")
- or ""
- ).strip()
- forum_username = str(
- stored.get("service_review_forum_username")
- or getattr(wbb, "SERVICE_REVIEW_FORUM_USERNAME", "")
- or ""
- ).strip().lstrip("@")
- 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
- ),
- ),
- ),
- "review_forum_chat_id": forum_chat_id,
- "review_forum_username": forum_username,
- }
- async def set_service_settings(values: dict[str, Any]) -> dict[str, Any]:
- 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"]))),
- ),
- "review_forum_chat_id": _clean_text(
- values.get("review_forum_chat_id", current["review_forum_chat_id"]),
- max_length=40,
- ),
- "review_forum_username": _clean_text(
- values.get(
- "review_forum_username",
- current["review_forum_username"],
- ),
- max_length=64,
- ).lstrip("@"),
- }
- 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"],
- "service_review_forum_chat_id": normalized["review_forum_chat_id"],
- "service_review_forum_username": normalized["review_forum_username"],
- "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 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", "未找到技师申请资料。")
- 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 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,
- "review_topic": await get_technician_review_topic(int(user_id)),
- }
- 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 | None = None,
- latitude: float | None = None,
- max_distance_meters: float | None = None,
- query: str = "",
- page: int = 1,
- page_size: int = 10,
- ) -> tuple[list[dict[str, Any]], int]:
- use_distance = longitude is not None and latitude is not None
- lon = float(longitude) if longitude is not None else 0.0
- lat = float(latitude) if latitude is not None else 0.0
- if use_distance and (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]] = []
- locations = {}
- if profiles and use_distance:
- locations = {
- int(location["user_id"]): location
- async for location in locationsdb.find(
- {"user_id": {"$in": list(profiles)}}
- )
- }
- for technician_id, profile in profiles.items():
- service_profile = profile.get("service_profile") or {}
- value = {
- "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": None,
- "distance_band": "",
- }
- if use_distance:
- location = locations.get(technician_id)
- if not location:
- continue
- 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
- value.update(
- {
- "distance_meters": round(distance, 2),
- "distance_band": distance_band(distance),
- }
- )
- values.append(value)
- if use_distance:
- values.sort(key=lambda item: (item["distance_meters"], item["user_id"]))
- else:
- values.sort(
- key=lambda item: (
- str(item.get("public_area_text") or "").casefold(),
- str(item.get("display_name") or "").casefold(),
- 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 _legacy_review_questions(values: dict[str, Any]) -> list[dict[str, Any]]:
- questions = []
- for item in values.get("rating_questions", []):
- questions.append(
- {
- **item,
- "type": "single_choice",
- "options": [
- {"value": str(score), "label": f"{score} 分", "score": score}
- for score in range(5, 0, -1)
- ],
- }
- )
- for item in values.get("text_questions", []):
- questions.append({**item, "type": "text"})
- return questions
- def normalize_review_template(values: dict[str, Any]) -> dict[str, Any]:
- raw_questions = values.get("questions")
- if raw_questions is None:
- raw_questions = _legacy_review_questions(values)
- if not isinstance(raw_questions, list) or not 1 <= len(raw_questions) <= 8:
- raise ServiceDataError(
- "invalid_review_questions",
- "评价模板必须包含 1 到 8 个问题。",
- )
- questions = []
- question_ids: set[str] = set()
- scoring_weight = 0
- for raw in raw_questions:
- if not isinstance(raw, dict):
- raise ServiceDataError("invalid_review_question", "评价问题格式无效。")
- question_type = str(raw.get("type") or "").strip()
- if question_type == "input":
- question_type = "text"
- if question_type not in {"single_choice", "multiple_choice", "text"}:
- raise ServiceDataError(
- "invalid_review_question_type",
- "评价问题仅支持单选、多选和输入框。",
- )
- question_id = _clean_text(
- raw.get("question_id") or uuid4().hex[:12],
- max_length=32,
- required=True,
- )
- if question_id in question_ids:
- raise ServiceDataError("duplicate_review_question", "评价问题 ID 不能重复。")
- question_ids.add(question_id)
- question = {
- "question_id": question_id,
- "type": question_type,
- "label": _clean_text(raw.get("label"), max_length=50, required=True),
- "description": _clean_text(raw.get("description"), max_length=160),
- "required": bool(raw.get("required")),
- }
- if question_type == "text":
- question["max_length"] = max(
- 50,
- min(1000, int(raw.get("max_length") or 500)),
- )
- else:
- raw_options = raw.get("options")
- if not isinstance(raw_options, list) or not 2 <= len(raw_options) <= 8:
- raise ServiceDataError(
- "invalid_review_options",
- "单选或多选问题必须设置 2 到 8 个选项。",
- )
- options = []
- option_values: set[str] = set()
- for index, raw_option in enumerate(raw_options):
- if isinstance(raw_option, str):
- raw_option = {"value": str(index + 1), "label": raw_option}
- if not isinstance(raw_option, dict):
- raise ServiceDataError(
- "invalid_review_options",
- "评价选项格式无效。",
- )
- option_value = _clean_text(
- raw_option.get("value") or str(index + 1),
- max_length=24,
- required=True,
- )
- if option_value in option_values:
- raise ServiceDataError(
- "duplicate_review_option",
- "同一问题的选项值不能重复。",
- )
- option_values.add(option_value)
- option = {
- "value": option_value,
- "label": _clean_text(
- raw_option.get("label"),
- max_length=24,
- required=True,
- ),
- }
- if question_type == "single_choice":
- score = int(raw_option.get("score") or 0)
- if not 1 <= score <= 5:
- raise ServiceDataError(
- "invalid_review_option_score",
- "单选项评分必须在 1 到 5 之间。",
- )
- option["score"] = score
- options.append(option)
- question["options"] = options
- if question_type == "single_choice":
- weight = int(raw.get("weight") or 0)
- if weight <= 0:
- raise ServiceDataError(
- "invalid_review_weight",
- "单选评分问题的权重必须大于 0。",
- )
- question["weight"] = weight
- scoring_weight += weight
- questions.append(question)
- if scoring_weight != 100:
- raise ServiceDataError(
- "invalid_review_weight",
- "全部单选评分问题的权重合计必须为 100%。",
- )
- return {"questions": questions}
- 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, **normalize_review_template(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
- ):
- snapshot = current.get("template_snapshot") or {}
- if "questions" not in snapshot:
- normalized_snapshot = {
- "template_id": snapshot.get("template_id") or template["template_id"],
- "version": snapshot.get("version") or template["version"],
- "questions": normalize_review_template(snapshot)["questions"],
- }
- question_index = int(current.get("rating_index") or 0) + int(
- current.get("text_index") or 0
- )
- await reviewdraftsdb.update_one(
- {"draft_id": current["draft_id"]},
- {
- "$set": {
- "template_snapshot": normalized_snapshot,
- "question_index": question_index,
- "updated_at": now,
- },
- "$unset": {"rating_index": "", "text_index": ""},
- },
- )
- current.update(
- {
- "template_snapshot": normalized_snapshot,
- "question_index": question_index,
- }
- )
- return current
- normalized_template = normalize_review_template(template)
- document = {
- "draft_id": uuid4().hex,
- "customer_id": int(customer_id),
- **identity,
- "template_snapshot": {
- "template_id": template["template_id"],
- "version": template["version"],
- "questions": normalized_template["questions"],
- },
- "question_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],
- question_index: int | None = None,
- rating_index: int = 0,
- text_index: int = 0,
- ) -> None:
- current_index = (
- max(0, int(question_index))
- if question_index is not None
- else max(0, int(rating_index) + int(text_index))
- )
- await reviewdraftsdb.update_one(
- {"customer_id": int(customer_id)},
- {
- "$set": {
- "answers": dict(answers),
- "question_index": current_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, Any], dict[str, str]]:
- questions = normalize_review_template(template)["questions"]
- ratings: dict[str, int] = {}
- choices: dict[str, Any] = {}
- texts: dict[str, str] = {}
- total = 0.0
- for question in questions:
- raw = answers.get(question["question_id"])
- if question["type"] == "text":
- text = _clean_text(
- raw,
- max_length=question["max_length"],
- required=question["required"],
- )
- if text:
- texts[question["question_id"]] = text
- continue
- options = {
- str(option["value"]): option for option in question["options"]
- }
- if question["type"] == "single_choice":
- if raw in (None, "") and not question["required"]:
- continue
- option = options.get(str(raw))
- if not option:
- raise ServiceDataError("invalid_review_answer", "请选择有效的单选项。")
- score = int(option["score"])
- ratings[question["question_id"]] = score
- choices[question["question_id"]] = option["label"]
- total += score * question["weight"] / 100
- continue
- selected = raw if isinstance(raw, list) else ([] if raw in (None, "") else [raw])
- selected_values = list(dict.fromkeys(str(item) for item in selected))
- if question["required"] and not selected_values:
- raise ServiceDataError("required_field", f"{question['label']} 为必选项。")
- if set(selected_values) - set(options):
- raise ServiceDataError("invalid_review_answer", "请选择有效的多选项。")
- if selected_values:
- choices[question["question_id"]] = [
- options[value]["label"] for value in selected_values
- ]
- return round(total, 4), ratings, choices, 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", "未找到所选套餐。")
- if not package_snapshot:
- packages = service_profile.get("packages") or []
- primary_package = packages[0] if packages else {}
- package_snapshot = {
- "package_id": None,
- "name": "技师服务",
- "category": primary_package.get("category") or "其他",
- }
- template = template_snapshot or await get_active_review_template()
- normalized_template = normalize_review_template(template)
- score, rating_answers, choice_answers, text_answers = _score_review(
- normalized_template,
- answers,
- )
- normalized_content = "\n".join(
- str(value).casefold()
- for _, value in sorted(
- {**choice_answers, **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": {
- "template_id": template["template_id"],
- "version": template["version"],
- "questions": normalized_template["questions"],
- },
- "rating_answers": rating_answers,
- "choice_answers": choice_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 >= 1,
- "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] = {"review_count": {"$gte": 1}}
- 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", {}),
- "choice_answers": item.get("choice_answers", {}),
- "text_answers": item.get("text_answers", {}),
- "template_snapshot": item.get("template_snapshot", {}),
- "approved_at": item.get("moderated_at"),
- "topic_url": item.get("topic_url"),
- }
- )
- return public_items, total
- async def get_technician_review_topic(
- technician_id: int,
- ) -> dict[str, Any] | None:
- await ensure_service_indexes()
- return await reviewtopicsdb.find_one(
- {"technician_id": int(technician_id), "status": "active"}
- )
- async def save_technician_review_topic(
- *,
- technician_id: int,
- forum_chat_id: str,
- forum_username: str,
- message_thread_id: int,
- topic_name: str,
- topic_url: str,
- ) -> dict[str, Any]:
- await ensure_service_indexes()
- now = utc_now()
- await reviewtopicsdb.update_one(
- {"technician_id": int(technician_id)},
- {
- "$set": {
- "forum_chat_id": str(forum_chat_id),
- "forum_username": str(forum_username),
- "message_thread_id": int(message_thread_id),
- "topic_name": _clean_text(
- topic_name,
- max_length=128,
- required=True,
- ),
- "topic_url": _clean_text(
- topic_url,
- max_length=300,
- required=True,
- ),
- "status": "active",
- "updated_at": now,
- },
- "$setOnInsert": {
- "technician_id": int(technician_id),
- "created_at": now,
- },
- },
- upsert=True,
- )
- return await get_technician_review_topic(int(technician_id)) or {}
- async def mark_review_topic_published(
- review_id: str,
- *,
- topic: dict[str, Any],
- message_id: int,
- ) -> dict[str, Any]:
- updated = await reviewsdb.find_one_and_update(
- {"review_id": str(review_id), "status": "approved"},
- {
- "$set": {
- "topic_url": topic["topic_url"],
- "topic_message_id": int(message_id),
- "topic_published_at": utc_now(),
- "updated_at": utc_now(),
- },
- "$unset": {"topic_publish_error": ""},
- },
- return_document=ReturnDocument.AFTER,
- )
- return updated or {}
- async def get_review(review_id: str) -> dict[str, Any] | None:
- return await reviewsdb.find_one({"review_id": str(review_id)})
- async def mark_review_topic_publish_error(
- review_id: str,
- error: str,
- ) -> None:
- await reviewsdb.update_one(
- {"review_id": str(review_id), "status": "approved"},
- {
- "$set": {
- "topic_publish_error": _clean_text(error, max_length=300),
- "updated_at": utc_now(),
- }
- },
- )
- 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 私聊成交不在统计范围。",
- }
|