dbservice.py 83 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876
  1. from __future__ import annotations
  2. import hashlib
  3. import json
  4. import math
  5. import re
  6. from collections import defaultdict
  7. from datetime import UTC, datetime, timedelta
  8. from decimal import Decimal, InvalidOperation
  9. from secrets import token_urlsafe
  10. from typing import Any
  11. from uuid import uuid4
  12. from cryptography.fernet import Fernet, InvalidToken
  13. from pymongo import ASCENDING, DESCENDING, ReturnDocument
  14. from pymongo.errors import DuplicateKeyError
  15. import wbb
  16. from wbb import control_db
  17. templatesdb = control_db.service_package_templates
  18. ordersdb = control_db.service_orders
  19. quotesdb = control_db.service_order_quotes
  20. addressesdb = control_db.service_order_addresses
  21. eventsdb = control_db.service_order_events
  22. customersdb = control_db.customer_profiles
  23. blocksdb = control_db.customer_blocks
  24. reportsdb = control_db.service_reports
  25. qrdb = control_db.teacher_review_qr_records
  26. reviewtemplatesdb = control_db.teacher_review_templates
  27. reviewdraftsdb = control_db.teacher_review_drafts
  28. reviewsdb = control_db.teacher_reviews
  29. statsdb = control_db.teacher_review_stats
  30. profilesdb = control_db.directory_profiles
  31. locationsdb = control_db.directory_locations
  32. membershipsdb = control_db.directory_memberships
  33. settingsdb = control_db.directory_settings
  34. SERVICE_MODES = {"at_store", "onsite"}
  35. PRICE_MODES = {"fixed", "starting_at", "range", "negotiable"}
  36. PRICE_UNITS = {"per_service", "per_hour", "per_item", "per_visit"}
  37. TRAVEL_FEE_MODES = {"included", "fixed", "per_km", "quoted"}
  38. ORDER_ACTIVE_STATUSES = {
  39. "requested",
  40. "quoted",
  41. "confirmed",
  42. "in_progress",
  43. "completion_pending",
  44. "disputed",
  45. }
  46. ORDER_TERMINAL_STATUSES = {
  47. "completed",
  48. "rejected",
  49. "canceled_customer",
  50. "canceled_technician",
  51. "voided",
  52. }
  53. ORDER_STATUSES = ORDER_ACTIVE_STATUSES | ORDER_TERMINAL_STATUSES
  54. REVIEW_SOURCES = {"qr_verified", "student_initiated"}
  55. REVIEW_STATUSES = {"pending", "approved", "rejected", "withdrawn", "voided"}
  56. DEFAULT_REVIEW_TEMPLATE = {
  57. "rating_questions": [
  58. {
  59. "question_id": "service_effect",
  60. "label": "服务效果",
  61. "description": "本次服务是否达到预期。",
  62. "required": True,
  63. "weight": 40,
  64. },
  65. {
  66. "question_id": "professionalism",
  67. "label": "专业程度",
  68. "description": "技师的专业能力与服务规范。",
  69. "required": True,
  70. "weight": 30,
  71. },
  72. {
  73. "question_id": "communication",
  74. "label": "沟通体验",
  75. "description": "沟通是否清晰、友好。",
  76. "required": True,
  77. "weight": 30,
  78. },
  79. ],
  80. "text_questions": [
  81. {
  82. "question_id": "comment",
  83. "label": "评价内容",
  84. "description": "分享真实体验,避免包含隐私信息。",
  85. "required": False,
  86. "max_length": 500,
  87. }
  88. ],
  89. }
  90. _indexes_ready = False
  91. class ServiceDataError(ValueError):
  92. def __init__(self, code: str, message: str):
  93. super().__init__(message)
  94. self.code = code
  95. def utc_now() -> datetime:
  96. return datetime.now(UTC)
  97. def _as_utc(value: datetime) -> datetime:
  98. return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
  99. def _clean_text(value: Any, *, max_length: int, required: bool = False) -> str:
  100. text = re.sub(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]", "", str(value or ""))
  101. text = " ".join(text.strip().split())
  102. if required and not text:
  103. raise ServiceDataError("required_field", "必填内容不能为空。")
  104. if len(text) > max_length:
  105. raise ServiceDataError("text_too_long", f"内容不能超过 {max_length} 个字符。")
  106. return text
  107. def _positive_decimal(value: Any, field: str, *, allow_zero: bool = True) -> str:
  108. try:
  109. number = Decimal(str(value or 0)).quantize(Decimal("0.01"))
  110. except (InvalidOperation, ValueError) as exc:
  111. raise ServiceDataError("invalid_price", f"{field} 必须是有效金额。") from exc
  112. if not number.is_finite():
  113. raise ServiceDataError("invalid_price", f"{field} 必须是有效金额。")
  114. if number < 0 or (not allow_zero and number == 0):
  115. raise ServiceDataError("invalid_price", f"{field} 不能小于 0。")
  116. return format(number, "f")
  117. def _normalize_tags(values: Any) -> list[str]:
  118. if isinstance(values, str):
  119. values = values.replace(",", ",").split(",")
  120. tags = []
  121. for value in values or []:
  122. tag = _clean_text(value, max_length=20)
  123. if tag and tag.casefold() not in {item.casefold() for item in tags}:
  124. tags.append(tag)
  125. if len(tags) > 8:
  126. raise ServiceDataError("too_many_tags", "服务标签最多 8 个。")
  127. return tags
  128. def normalize_package(value: dict[str, Any], *, package_id: str | None = None) -> dict[str, Any]:
  129. modes = list(dict.fromkeys(str(item) for item in value.get("service_modes", [])))
  130. if not modes or set(modes) - SERVICE_MODES:
  131. raise ServiceDataError("invalid_service_mode", "套餐必须选择到店或上门服务方式。")
  132. price_mode = str(value.get("price_mode") or "negotiable")
  133. if price_mode not in PRICE_MODES:
  134. raise ServiceDataError("invalid_price_mode", "不支持该价格模式。")
  135. price_unit = str(value.get("price_unit") or "per_service")
  136. if price_unit not in PRICE_UNITS:
  137. raise ServiceDataError("invalid_price_unit", "不支持该计价单位。")
  138. currency = str(value.get("currency") or "CNY").upper()
  139. if not re.fullmatch(r"[A-Z]{3}", currency):
  140. raise ServiceDataError("invalid_currency", "币种必须使用三位 ISO 代码。")
  141. min_price = _positive_decimal(value.get("min_price"), "最低价格")
  142. max_price = _positive_decimal(value.get("max_price"), "最高价格")
  143. if Decimal(max_price) and Decimal(max_price) < Decimal(min_price):
  144. raise ServiceDataError("invalid_price_range", "最高价格不能低于最低价格。")
  145. if price_mode == "fixed" and Decimal(min_price) != Decimal(max_price):
  146. raise ServiceDataError("invalid_fixed_price", "固定价的最低价格和最高价格必须一致。")
  147. duration = value.get("duration_minutes")
  148. if duration in (None, "", 0, "0"):
  149. duration = None
  150. else:
  151. try:
  152. duration = int(duration)
  153. except (TypeError, ValueError) as exc:
  154. raise ServiceDataError("invalid_duration", "服务时长必须是分钟数。") from exc
  155. if not 15 <= duration <= 480:
  156. raise ServiceDataError("invalid_duration", "服务时长必须在 15 到 480 分钟之间。")
  157. travel = value.get("travel_fee") if isinstance(value.get("travel_fee"), dict) else {}
  158. travel_mode = str(travel.get("mode") or "quoted")
  159. if travel_mode not in TRAVEL_FEE_MODES:
  160. raise ServiceDataError("invalid_travel_fee", "不支持该交通费模式。")
  161. addons = []
  162. for item in value.get("addons", []) or []:
  163. if len(addons) >= 10:
  164. raise ServiceDataError("too_many_addons", "标准加项最多 10 个。")
  165. addons.append(
  166. {
  167. "addon_id": str(item.get("addon_id") or uuid4().hex),
  168. "name": _clean_text(item.get("name"), max_length=40, required=True),
  169. "unit": _clean_text(item.get("unit"), max_length=20) or "项",
  170. "price": _positive_decimal(item.get("price"), "加项价格"),
  171. "extra_minutes": max(0, min(480, int(item.get("extra_minutes") or 0))),
  172. }
  173. )
  174. radius = value.get("service_radius_km")
  175. if radius in (None, ""):
  176. radius = None
  177. else:
  178. radius = float(radius)
  179. if not 0.5 <= radius <= 200:
  180. raise ServiceDataError("invalid_service_radius", "上门半径必须在 0.5 到 200 公里之间。")
  181. if "onsite" in modes and radius is None:
  182. raise ServiceDataError("service_radius_required", "上门套餐必须填写服务半径。")
  183. out_of_range_policy = _clean_text(
  184. value.get("out_of_range_policy"),
  185. max_length=200,
  186. required="onsite" in modes,
  187. )
  188. source_type = str(value.get("source_type") or "custom")
  189. if source_type not in {"template", "custom"}:
  190. raise ServiceDataError("invalid_package_source", "套餐来源无效。")
  191. return {
  192. "package_id": str(package_id or value.get("package_id") or uuid4().hex),
  193. "name": _clean_text(value.get("name"), max_length=50, required=True),
  194. "category": _clean_text(value.get("category"), max_length=30, required=True),
  195. "description": _clean_text(value.get("description"), max_length=800, required=True),
  196. "tags": _normalize_tags(value.get("tags", [])),
  197. "service_modes": modes,
  198. "price_mode": price_mode,
  199. "currency": currency,
  200. "price_unit": price_unit,
  201. "min_price": min_price,
  202. "max_price": max_price,
  203. "price_display": _clean_text(value.get("price_display"), max_length=80),
  204. "duration_minutes": duration,
  205. "included_items": _clean_text(value.get("included_items"), max_length=500),
  206. "excluded_items": _clean_text(value.get("excluded_items"), max_length=500),
  207. "preparation": _clean_text(value.get("preparation"), max_length=500),
  208. "addons": addons,
  209. "travel_fee": {
  210. "mode": travel_mode,
  211. "amount": _positive_decimal(travel.get("amount"), "交通费"),
  212. "per_km": _positive_decimal(travel.get("per_km"), "每公里交通费"),
  213. "description": _clean_text(travel.get("description"), max_length=200),
  214. },
  215. "service_radius_km": radius,
  216. "out_of_range_policy": out_of_range_policy,
  217. "source_type": source_type,
  218. "source_template_id": value.get("source_template_id"),
  219. "source_template_version": value.get("source_template_version"),
  220. "customized_from_template": bool(value.get("customized_from_template")),
  221. }
  222. def _normalize_profile(value: dict[str, Any]) -> dict[str, Any]:
  223. packages = [normalize_package(item) for item in value.get("packages", [])]
  224. if not 1 <= len(packages) <= 5:
  225. raise ServiceDataError("invalid_package_count", "技师必须发布 1 到 5 个套餐。")
  226. modes = sorted({mode for item in packages for mode in item["service_modes"]})
  227. venue = value.get("venue") if isinstance(value.get("venue"), dict) else {}
  228. onsite = value.get("onsite_policy") if isinstance(value.get("onsite_policy"), dict) else {}
  229. public_area_text = _clean_text(
  230. value.get("public_area_text"), max_length=80, required=True
  231. )
  232. venue_name = _clean_text(
  233. venue.get("name"), max_length=80, required="at_store" in modes
  234. )
  235. venue_address = _clean_text(
  236. venue.get("address_hint"), max_length=120, required="at_store" in modes
  237. )
  238. onsite_description = _clean_text(
  239. onsite.get("description"), max_length=300, required="onsite" in modes
  240. )
  241. return {
  242. "headline": _clean_text(value.get("headline"), max_length=80, required=True),
  243. "bio": _clean_text(value.get("bio"), max_length=800, required=True),
  244. "tags": _normalize_tags(value.get("tags", [])),
  245. "contact_hours": _clean_text(value.get("contact_hours"), max_length=120, required=True),
  246. "service_modes": modes,
  247. "accepting_requests": bool(value.get("accepting_requests", True)),
  248. "public_area_text": public_area_text,
  249. "venue": {
  250. "name": venue_name,
  251. "address_hint": venue_address,
  252. },
  253. "onsite_policy": {
  254. "description": onsite_description,
  255. },
  256. "packages": packages,
  257. "is_complete": True,
  258. }
  259. async def ensure_service_indexes() -> None:
  260. global _indexes_ready
  261. if _indexes_ready:
  262. return
  263. await templatesdb.create_index([("template_id", ASCENDING)], unique=True)
  264. await templatesdb.create_index([("status", ASCENDING), ("sort_order", ASCENDING)])
  265. await ordersdb.create_index([("order_id", ASCENDING)], unique=True)
  266. await ordersdb.create_index([("customer_id", ASCENDING), ("created_at", DESCENDING)])
  267. await ordersdb.create_index([("technician_id", ASCENDING), ("created_at", DESCENDING)])
  268. await ordersdb.create_index([("status", ASCENDING), ("updated_at", DESCENDING)])
  269. await quotesdb.create_index([("quote_id", ASCENDING)], unique=True)
  270. await quotesdb.create_index([("order_id", ASCENDING), ("version", DESCENDING)], unique=True)
  271. await addressesdb.create_index([("order_id", ASCENDING)], unique=True)
  272. await addressesdb.create_index([("redact_after", ASCENDING)])
  273. await eventsdb.create_index([("event_id", ASCENDING)], unique=True)
  274. await eventsdb.create_index([("order_id", ASCENDING), ("created_at", ASCENDING)])
  275. await customersdb.create_index([("user_id", ASCENDING)], unique=True)
  276. await blocksdb.create_index([("scope", ASCENDING), ("technician_id", ASCENDING), ("customer_id", ASCENDING)], unique=True)
  277. await reportsdb.create_index([("report_id", ASCENDING)], unique=True)
  278. await qrdb.create_index([("qr_id", ASCENDING)], unique=True)
  279. await qrdb.create_index([("token_hash", ASCENDING)], unique=True)
  280. await qrdb.create_index([("order_id", ASCENDING), ("created_at", DESCENDING)])
  281. await reviewtemplatesdb.create_index([("version", DESCENDING)], unique=True)
  282. await reviewdraftsdb.create_index([("customer_id", ASCENDING)], unique=True)
  283. await reviewdraftsdb.create_index([("expires_at", ASCENDING)], expireAfterSeconds=0)
  284. await reviewsdb.create_index([("review_id", ASCENDING)], unique=True)
  285. await reviewsdb.create_index([("qr_id", ASCENDING)], unique=True, sparse=True)
  286. await reviewsdb.create_index([("technician_id", ASCENDING), ("status", ASCENDING), ("created_at", DESCENDING)])
  287. await statsdb.create_index([("technician_id", ASCENDING), ("category", ASCENDING)], unique=True)
  288. _indexes_ready = True
  289. async def record_service_event(
  290. order_id: str,
  291. event_type: str,
  292. *,
  293. actor_id: int | str,
  294. reason: str = "",
  295. metadata: dict[str, Any] | None = None,
  296. ) -> dict[str, Any]:
  297. await ensure_service_indexes()
  298. document = {
  299. "event_id": uuid4().hex,
  300. "order_id": str(order_id),
  301. "event_type": str(event_type),
  302. "actor_id": actor_id,
  303. "reason": _clean_text(reason, max_length=500),
  304. "metadata": metadata or {},
  305. "created_at": utc_now(),
  306. }
  307. await eventsdb.insert_one(document)
  308. return document
  309. async def record_service_funnel_event(user_id: int, event_type: str) -> dict[str, Any]:
  310. if event_type not in {"directory_viewed", "technician_viewed", "request_started"}:
  311. raise ServiceDataError("invalid_funnel_event", "不支持该履约漏斗事件。")
  312. return await record_service_event(
  313. "",
  314. event_type,
  315. actor_id=int(user_id),
  316. metadata={"funnel": True},
  317. )
  318. async def observe_customer(user: Any, *, accepted_terms: bool = False) -> dict[str, Any]:
  319. await ensure_service_indexes()
  320. now = utc_now()
  321. display_name = " ".join(
  322. value for value in (str(getattr(user, "first_name", "") or "").strip(), str(getattr(user, "last_name", "") or "").strip()) if value
  323. ) or f"用户 {user.id}"
  324. updates: dict[str, Any] = {
  325. "username": getattr(user, "username", None),
  326. "display_name": display_name,
  327. "updated_at": now,
  328. }
  329. if accepted_terms:
  330. updates["terms_accepted_at"] = now
  331. await customersdb.update_one(
  332. {"user_id": int(user.id)},
  333. {"$set": updates, "$setOnInsert": {"created_at": now, "blocked": False}},
  334. upsert=True,
  335. )
  336. return await customersdb.find_one({"user_id": int(user.id)}) or {}
  337. async def require_customer_allowed(user_id: int) -> dict[str, Any]:
  338. await ensure_service_indexes()
  339. customer = await customersdb.find_one({"user_id": int(user_id)})
  340. if not customer or not customer.get("terms_accepted_at"):
  341. raise ServiceDataError("terms_required", "请先接受服务规则和隐私说明。")
  342. if customer.get("blocked"):
  343. raise ServiceDataError("customer_blocked", "你的服务功能已被暂停。")
  344. global_block = await blocksdb.find_one({"scope": "global", "customer_id": int(user_id), "active": True})
  345. if global_block:
  346. raise ServiceDataError("customer_blocked", "你的服务功能已被暂停。")
  347. return customer
  348. async def get_service_settings() -> dict[str, int]:
  349. stored = await settingsdb.find_one({"settings_id": "global"}) or {}
  350. return {
  351. "customer_max_open_orders": max(
  352. 1,
  353. min(
  354. 20,
  355. int(
  356. stored.get("service_customer_max_open_orders")
  357. or getattr(wbb, "SERVICE_CUSTOMER_MAX_OPEN_ORDERS", 3)
  358. or 3
  359. ),
  360. ),
  361. ),
  362. "customer_daily_request_limit": max(
  363. 1,
  364. min(
  365. 100,
  366. int(
  367. stored.get("service_customer_daily_request_limit")
  368. or getattr(wbb, "SERVICE_CUSTOMER_DAILY_REQUEST_LIMIT", 10)
  369. or 10
  370. ),
  371. ),
  372. ),
  373. "review_qr_expiry_hours": max(
  374. 1,
  375. min(
  376. 168,
  377. int(
  378. stored.get("service_review_qr_expiry_hours")
  379. or getattr(wbb, "SERVICE_REVIEW_QR_EXPIRY_HOURS", 24)
  380. or 24
  381. ),
  382. ),
  383. ),
  384. }
  385. async def set_service_settings(values: dict[str, Any]) -> dict[str, int]:
  386. current = await get_service_settings()
  387. try:
  388. normalized = {
  389. "customer_max_open_orders": max(
  390. 1,
  391. min(20, int(values.get("customer_max_open_orders", current["customer_max_open_orders"]))),
  392. ),
  393. "customer_daily_request_limit": max(
  394. 1,
  395. min(
  396. 100,
  397. int(
  398. values.get(
  399. "customer_daily_request_limit",
  400. current["customer_daily_request_limit"],
  401. )
  402. ),
  403. ),
  404. ),
  405. "review_qr_expiry_hours": max(
  406. 1,
  407. min(168, int(values.get("review_qr_expiry_hours", current["review_qr_expiry_hours"]))),
  408. ),
  409. }
  410. except (TypeError, ValueError) as exc:
  411. raise ServiceDataError("invalid_service_settings", "履约设置必须是有效整数。") from exc
  412. await settingsdb.update_one(
  413. {"settings_id": "global"},
  414. {
  415. "$set": {
  416. "service_customer_max_open_orders": normalized["customer_max_open_orders"],
  417. "service_customer_daily_request_limit": normalized["customer_daily_request_limit"],
  418. "service_review_qr_expiry_hours": normalized["review_qr_expiry_hours"],
  419. "updated_at": utc_now(),
  420. },
  421. "$setOnInsert": {"created_at": utc_now()},
  422. },
  423. upsert=True,
  424. )
  425. return normalized
  426. async def list_package_templates(*, include_archived: bool = False) -> list[dict[str, Any]]:
  427. await ensure_service_indexes()
  428. filters = {} if include_archived else {"status": {"$ne": "archived"}}
  429. return await templatesdb.find(filters).sort([("sort_order", ASCENDING), ("created_at", ASCENDING)]).to_list(length=500)
  430. async def create_package_template(values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
  431. await ensure_service_indexes()
  432. now = utc_now()
  433. package = normalize_package(values.get("package") or values)
  434. package["source_type"] = "template"
  435. document = {
  436. "template_id": uuid4().hex,
  437. "version": 1,
  438. "name": _clean_text(values.get("template_name") or package["name"], max_length=80, required=True),
  439. "admin_note": _clean_text(values.get("admin_note"), max_length=500),
  440. "package": package,
  441. "status": "enabled" if values.get("enabled", True) else "disabled",
  442. "sort_order": int(values.get("sort_order") or 0),
  443. "created_by": str(actor_id),
  444. "created_at": now,
  445. "updated_at": now,
  446. }
  447. await templatesdb.insert_one(document)
  448. return document
  449. async def update_package_template(template_id: str, values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
  450. current = await templatesdb.find_one({"template_id": str(template_id)})
  451. if not current:
  452. raise ServiceDataError("template_not_found", "未找到套餐模板。")
  453. package = normalize_package(values.get("package") or values, package_id=current["package"]["package_id"])
  454. package["source_type"] = "template"
  455. updated = await templatesdb.find_one_and_update(
  456. {"template_id": str(template_id)},
  457. {
  458. "$set": {
  459. "name": _clean_text(values.get("template_name") or package["name"], max_length=80, required=True),
  460. "admin_note": _clean_text(values.get("admin_note"), max_length=500),
  461. "package": package,
  462. "sort_order": int(values.get("sort_order") or 0),
  463. "updated_by": str(actor_id),
  464. "updated_at": utc_now(),
  465. },
  466. "$inc": {"version": 1},
  467. },
  468. return_document=ReturnDocument.AFTER,
  469. )
  470. return updated or {}
  471. async def set_package_template_status(template_id: str, action: str, *, actor_id: str) -> dict[str, Any]:
  472. status_map = {"enable": "enabled", "disable": "disabled", "archive": "archived"}
  473. if action not in status_map:
  474. raise ServiceDataError("invalid_template_action", "不支持该模板操作。")
  475. updated = await templatesdb.find_one_and_update(
  476. {"template_id": str(template_id)},
  477. {"$set": {"status": status_map[action], "updated_by": str(actor_id), "updated_at": utc_now()}},
  478. return_document=ReturnDocument.AFTER,
  479. )
  480. if not updated:
  481. raise ServiceDataError("template_not_found", "未找到套餐模板。")
  482. return updated
  483. async def duplicate_package_template(template_id: str, *, actor_id: str) -> dict[str, Any]:
  484. current = await templatesdb.find_one({"template_id": str(template_id)})
  485. if not current:
  486. raise ServiceDataError("template_not_found", "未找到套餐模板。")
  487. values = {
  488. "template_name": f"{current['name']}(副本)",
  489. "admin_note": current.get("admin_note", ""),
  490. "sort_order": current.get("sort_order", 0),
  491. "enabled": False,
  492. "package": current["package"],
  493. }
  494. return await create_package_template(values, actor_id=actor_id)
  495. async def publish_technician_profile(user_id: int, values: dict[str, Any], *, actor_id: int | str) -> dict[str, Any]:
  496. await ensure_service_indexes()
  497. profile = await profilesdb.find_one({"user_id": int(user_id)})
  498. if not profile or profile.get("application_status") != "approved":
  499. raise ServiceDataError("technician_required", "只有已认证技师可以发布服务资料。")
  500. if not profile.get("username"):
  501. raise ServiceDataError("username_required", "技师必须设置有效的 Telegram 用户名。")
  502. if not await locationsdb.find_one({"user_id": int(user_id)}):
  503. raise ServiceDataError("technician_location_required", "技师必须设置有效服务位置。")
  504. if not await membershipsdb.find_one({"user_id": int(user_id), "active": True}):
  505. raise ServiceDataError("technician_membership_required", "技师必须仍是受管群当前成员。")
  506. normalized = _normalize_profile(values)
  507. now = utc_now()
  508. normalized.update({"updated_at": now, "updated_by": actor_id, "version": int((profile.get("service_profile") or {}).get("version") or 0) + 1})
  509. await profilesdb.update_one(
  510. {"user_id": int(user_id)},
  511. {"$set": {"service_profile": normalized, "updated_at": now}},
  512. )
  513. return await profilesdb.find_one({"user_id": int(user_id)}) or {}
  514. async def get_technician_service_profile(user_id: int) -> dict[str, Any]:
  515. profile = await profilesdb.find_one({"user_id": int(user_id)})
  516. if not profile:
  517. raise ServiceDataError("technician_not_found", "未找到技师资料。")
  518. return profile
  519. async def instantiate_package_template(template_id: str) -> dict[str, Any]:
  520. template = await templatesdb.find_one(
  521. {"template_id": str(template_id), "status": "enabled"}
  522. )
  523. if not template:
  524. raise ServiceDataError("template_not_found", "套餐模板不存在或已停用。")
  525. package = dict(template["package"])
  526. package.update(
  527. {
  528. "package_id": uuid4().hex,
  529. "source_type": "template",
  530. "source_template_id": template["template_id"],
  531. "source_template_version": template["version"],
  532. "customized_from_template": False,
  533. }
  534. )
  535. return package
  536. def _fernet() -> Fernet:
  537. raw = str(getattr(wbb, "SERVICE_ADDRESS_ENCRYPTION_KEY", "") or "").strip()
  538. if not raw:
  539. raise ServiceDataError("address_encryption_required", "服务地址加密密钥尚未配置。")
  540. try:
  541. return Fernet(raw.encode())
  542. except (ValueError, TypeError) as exc:
  543. raise ServiceDataError("invalid_address_encryption_key", "服务地址加密密钥无效。") from exc
  544. def _encrypt_address(payload: dict[str, Any]) -> str:
  545. return _fernet().encrypt(json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode()).decode()
  546. def _decrypt_address(token: str) -> dict[str, Any]:
  547. try:
  548. return json.loads(_fernet().decrypt(str(token).encode()).decode())
  549. except (InvalidToken, ValueError, json.JSONDecodeError) as exc:
  550. raise ServiceDataError("address_unavailable", "精确地址无法解密。") from exc
  551. async def _schedule_address_redaction(order_id: str, closed_at: datetime) -> None:
  552. retention = max(
  553. 1,
  554. min(365, int(getattr(wbb, "SERVICE_ADDRESS_RETENTION_DAYS", 7) or 7)),
  555. )
  556. await addressesdb.update_one(
  557. {"order_id": str(order_id), "redacted": False},
  558. {
  559. "$set": {
  560. "redact_after": closed_at + timedelta(days=retention),
  561. "updated_at": closed_at,
  562. }
  563. },
  564. )
  565. def _distance_meters(lon1: float, lat1: float, lon2: float, lat2: float) -> float:
  566. radius = 6_371_000.0
  567. phi1, phi2 = math.radians(lat1), math.radians(lat2)
  568. delta_phi = math.radians(lat2 - lat1)
  569. delta_lambda = math.radians(lon2 - lon1)
  570. a = math.sin(delta_phi / 2) ** 2 + math.cos(phi1) * math.cos(phi2) * math.sin(delta_lambda / 2) ** 2
  571. return radius * 2 * math.atan2(math.sqrt(a), math.sqrt(1 - a))
  572. def distance_band(distance_meters: float) -> str:
  573. km = max(0.0, distance_meters / 1000)
  574. for threshold in (1, 3, 5, 10, 20, 50):
  575. if km <= threshold:
  576. return f"{threshold} 公里内"
  577. return "50 公里以上"
  578. async def _get_package(technician_id: int, package_id: str) -> tuple[dict[str, Any], dict[str, Any]]:
  579. profile = await profilesdb.find_one({"user_id": int(technician_id)})
  580. if not profile or profile.get("application_status") != "approved":
  581. raise ServiceDataError("technician_not_found", "未找到已认证技师。")
  582. if not profile.get("listed") or not profile.get("username"):
  583. raise ServiceDataError("technician_unavailable", "该技师暂未公开接单。")
  584. if not await membershipsdb.find_one({"user_id": int(technician_id), "active": True}):
  585. raise ServiceDataError("technician_unavailable", "该技师当前不满足接单资格。")
  586. service_profile = profile.get("service_profile") or {}
  587. if not service_profile.get("is_complete") or not service_profile.get("accepting_requests", True):
  588. raise ServiceDataError("technician_unavailable", "该技师暂未接单。")
  589. package = next((item for item in service_profile.get("packages", []) if str(item.get("package_id")) == str(package_id)), None)
  590. if not package:
  591. raise ServiceDataError("package_not_found", "未找到该服务套餐。")
  592. return profile, package
  593. async def list_available_technicians(
  594. *,
  595. longitude: float,
  596. latitude: float,
  597. max_distance_meters: float | None = None,
  598. query: str = "",
  599. page: int = 1,
  600. page_size: int = 10,
  601. ) -> tuple[list[dict[str, Any]], int]:
  602. lon, lat = float(longitude), float(latitude)
  603. if not -180 <= lon <= 180 or not -90 <= lat <= 90:
  604. raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
  605. filters: dict[str, Any] = {
  606. "application_status": "approved",
  607. "listed": True,
  608. "username": {"$nin": [None, ""]},
  609. "service_profile.is_complete": True,
  610. "service_profile.accepting_requests": True,
  611. }
  612. if query:
  613. pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
  614. filters["$or"] = [
  615. {"username": pattern},
  616. {"display_name": pattern},
  617. {"service_profile.tags": pattern},
  618. {"service_profile.packages.category": pattern},
  619. ]
  620. active_members = {
  621. int(item["user_id"])
  622. async for item in membershipsdb.find({"active": True}, {"user_id": 1})
  623. }
  624. profiles = {
  625. int(item["user_id"]): item
  626. async for item in profilesdb.find(filters)
  627. if int(item["user_id"]) in active_members
  628. }
  629. values: list[dict[str, Any]] = []
  630. if profiles:
  631. async for location in locationsdb.find({"user_id": {"$in": list(profiles)}}):
  632. point = location.get("point", {}).get("coordinates") or [
  633. location.get("longitude"),
  634. location.get("latitude"),
  635. ]
  636. if len(point) != 2 or point[0] is None or point[1] is None:
  637. continue
  638. distance = _distance_meters(lon, lat, float(point[0]), float(point[1]))
  639. if max_distance_meters is not None and distance > float(max_distance_meters):
  640. continue
  641. profile = profiles[int(location["user_id"])]
  642. service_profile = profile.get("service_profile") or {}
  643. values.append(
  644. {
  645. "user_id": int(profile["user_id"]),
  646. "username": profile.get("username"),
  647. "display_name": profile.get("display_name"),
  648. "headline": service_profile.get("headline"),
  649. "bio": service_profile.get("bio"),
  650. "tags": service_profile.get("tags", []),
  651. "public_area_text": service_profile.get("public_area_text"),
  652. "service_modes": service_profile.get("service_modes", []),
  653. "packages": service_profile.get("packages", []),
  654. "distance_meters": round(distance, 2),
  655. "distance_band": distance_band(distance),
  656. }
  657. )
  658. values.sort(key=lambda item: (item["distance_meters"], item["user_id"]))
  659. total = len(values)
  660. start = (max(1, page) - 1) * max(1, page_size)
  661. return values[start : start + max(1, page_size)], total
  662. async def create_service_order(
  663. *,
  664. customer: Any,
  665. technician_id: int,
  666. package_id: str,
  667. service_mode: str,
  668. requirements: str,
  669. longitude: float | None = None,
  670. latitude: float | None = None,
  671. address_text: str = "",
  672. ) -> dict[str, Any]:
  673. await ensure_service_indexes()
  674. await observe_customer(customer)
  675. await require_customer_allowed(int(customer.id))
  676. if int(customer.id) == int(technician_id):
  677. raise ServiceDataError("self_service_not_allowed", "不能向自己发起服务请求。")
  678. personal_block = await blocksdb.find_one({"scope": "technician", "technician_id": int(technician_id), "customer_id": int(customer.id), "active": True})
  679. if personal_block:
  680. raise ServiceDataError("customer_blocked", "该技师暂不接受你的服务请求。")
  681. service_settings = await get_service_settings()
  682. max_open = service_settings["customer_max_open_orders"]
  683. open_count = await ordersdb.count_documents({"customer_id": int(customer.id), "status": {"$in": list(ORDER_ACTIVE_STATUSES)}})
  684. if open_count >= max_open:
  685. raise ServiceDataError("too_many_open_orders", f"最多同时保留 {max_open} 个未关闭服务单。")
  686. day_start = utc_now().replace(hour=0, minute=0, second=0, microsecond=0)
  687. daily_limit = service_settings["customer_daily_request_limit"]
  688. if await ordersdb.count_documents({"customer_id": int(customer.id), "created_at": {"$gte": day_start}}) >= daily_limit:
  689. raise ServiceDataError("daily_request_limit", "今天发起的服务请求已达到上限。")
  690. profile, package = await _get_package(int(technician_id), str(package_id))
  691. if service_mode not in package["service_modes"]:
  692. raise ServiceDataError("service_mode_unavailable", "该套餐不支持所选服务方式。")
  693. technician_location = await locationsdb.find_one({"user_id": int(technician_id)})
  694. if not technician_location:
  695. raise ServiceDataError("technician_location_required", "技师尚未配置服务位置。")
  696. address_payload: dict[str, Any] | None = None
  697. distance = 0.0
  698. if service_mode == "onsite":
  699. if longitude is None or latitude is None:
  700. raise ServiceDataError("customer_location_required", "上门服务必须提供位置。")
  701. lon, lat = float(longitude), float(latitude)
  702. if not -180 <= lon <= 180 or not -90 <= lat <= 90:
  703. raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
  704. address = _clean_text(address_text, max_length=300, required=True)
  705. point = technician_location.get("point", {}).get("coordinates") or [technician_location.get("longitude"), technician_location.get("latitude")]
  706. distance = _distance_meters(float(point[0]), float(point[1]), lon, lat)
  707. radius = float(package.get("service_radius_km") or 0)
  708. if radius and distance > radius * 1000:
  709. raise ServiceDataError("outside_service_radius", "顾客位置超出该套餐的上门服务范围。")
  710. address_payload = {"kind": "customer", "longitude": lon, "latitude": lat, "address_text": address}
  711. else:
  712. if longitude is None or latitude is None:
  713. raise ServiceDataError(
  714. "customer_location_required", "到店服务必须先通过附近查找提供位置。"
  715. )
  716. lon, lat = float(longitude), float(latitude)
  717. if not -180 <= lon <= 180 or not -90 <= lat <= 90:
  718. raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
  719. point = technician_location.get("point", {}).get("coordinates") or [
  720. technician_location.get("longitude"),
  721. technician_location.get("latitude"),
  722. ]
  723. distance = _distance_meters(float(point[0]), float(point[1]), lon, lat)
  724. venue = (profile.get("service_profile") or {}).get("venue") or {}
  725. address_payload = {
  726. "kind": "venue",
  727. "longitude": technician_location.get("longitude"),
  728. "latitude": technician_location.get("latitude"),
  729. "address_text": _clean_text(venue.get("address_hint"), max_length=300, required=True),
  730. }
  731. now = utc_now()
  732. order_id = uuid4().hex
  733. order = {
  734. "order_id": order_id,
  735. "customer_id": int(customer.id),
  736. "customer_name": _clean_text(getattr(customer, "first_name", ""), max_length=80) or f"用户 {customer.id}",
  737. "technician_id": int(technician_id),
  738. "technician_name": profile.get("display_name") or f"技师 {technician_id}",
  739. "package_id": str(package_id),
  740. "package_snapshot": package,
  741. "category": package["category"],
  742. "service_mode": service_mode,
  743. "requirements": _clean_text(requirements, max_length=800, required=True),
  744. "distance_meters": round(distance, 2),
  745. "distance_band": distance_band(distance),
  746. "status": "requested",
  747. "version": 1,
  748. "created_at": now,
  749. "updated_at": now,
  750. }
  751. encrypted_address = _encrypt_address(address_payload)
  752. await ordersdb.insert_one(order)
  753. try:
  754. await addressesdb.insert_one(
  755. {
  756. "order_id": order_id,
  757. "encrypted_payload": encrypted_address,
  758. "redacted": False,
  759. "redact_after": None,
  760. "created_at": now,
  761. "updated_at": now,
  762. }
  763. )
  764. except Exception:
  765. await ordersdb.delete_one({"order_id": order_id, "status": "requested"})
  766. raise
  767. await record_service_event(order_id, "order_requested", actor_id=int(customer.id))
  768. return order
  769. def normalize_quote(value: dict[str, Any], *, package: dict[str, Any]) -> dict[str, Any]:
  770. currency = str(value.get("currency") or package.get("currency") or "CNY").upper()
  771. if currency != str(package.get("currency") or currency):
  772. raise ServiceDataError("currency_mismatch", "最终报价币种必须与套餐一致。")
  773. base = _positive_decimal(value.get("base_amount"), "服务金额")
  774. travel = _positive_decimal(value.get("travel_fee"), "交通费")
  775. discount = _positive_decimal(value.get("discount"), "优惠金额")
  776. addons = []
  777. addon_total = Decimal("0")
  778. for item in value.get("addons", []) or []:
  779. quantity = max(1, min(999, int(item.get("quantity") or 1)))
  780. unit_price = Decimal(_positive_decimal(item.get("unit_price"), "加项单价"))
  781. addons.append({"name": _clean_text(item.get("name"), max_length=60, required=True), "quantity": quantity, "unit_price": format(unit_price, "f")})
  782. addon_total += unit_price * quantity
  783. total = Decimal(base) + Decimal(travel) + addon_total - Decimal(discount)
  784. if total < 0:
  785. raise ServiceDataError("invalid_total", "最终报价总额不能小于 0。")
  786. provided = value.get("total_amount")
  787. if provided not in (None, "") and Decimal(_positive_decimal(provided, "总价")) != total:
  788. raise ServiceDataError("quote_total_mismatch", "报价明细与总价不一致。")
  789. return {
  790. "currency": currency,
  791. "base_amount": base,
  792. "travel_fee": travel,
  793. "addons": addons,
  794. "discount": discount,
  795. "total_amount": format(total.quantize(Decimal("0.01")), "f"),
  796. "note": _clean_text(value.get("note"), max_length=500),
  797. }
  798. async def submit_service_quote(order_id: str, technician_id: int, values: dict[str, Any]) -> tuple[dict[str, Any], dict[str, Any]]:
  799. await ensure_service_indexes()
  800. order = await ordersdb.find_one({"order_id": str(order_id)})
  801. if not order or int(order["technician_id"]) != int(technician_id):
  802. raise ServiceDataError("order_not_found", "未找到可报价服务单。")
  803. if order["status"] not in {"requested", "quoted"}:
  804. raise ServiceDataError("invalid_order_state", "当前服务单不能修改报价。")
  805. quote = normalize_quote(values, package=order["package_snapshot"])
  806. version = int(order.get("quote_version") or 0) + 1
  807. quote.update({"quote_id": uuid4().hex, "order_id": str(order_id), "version": version, "created_by": int(technician_id), "created_at": utc_now()})
  808. try:
  809. await quotesdb.insert_one(quote)
  810. except DuplicateKeyError as exc:
  811. raise ServiceDataError(
  812. "concurrent_order_update", "服务单状态已变化,请刷新后重试。"
  813. ) from exc
  814. updated = await ordersdb.find_one_and_update(
  815. {
  816. "order_id": str(order_id),
  817. "status": {"$in": ["requested", "quoted"]},
  818. "version": order["version"],
  819. },
  820. {
  821. "$set": {
  822. "status": "quoted",
  823. "current_quote_id": quote["quote_id"],
  824. "quote_snapshot": quote,
  825. "quote_version": version,
  826. "updated_at": utc_now(),
  827. },
  828. "$inc": {"version": 1},
  829. },
  830. return_document=ReturnDocument.AFTER,
  831. )
  832. if not updated:
  833. await quotesdb.delete_one({"quote_id": quote["quote_id"]})
  834. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  835. await record_service_event(order_id, "quote_submitted", actor_id=int(technician_id), metadata={"quote_id": quote["quote_id"], "version": version})
  836. return updated, quote
  837. async def confirm_service_quote(order_id: str, customer_id: int) -> dict[str, Any]:
  838. order = await ordersdb.find_one_and_update(
  839. {"order_id": str(order_id), "customer_id": int(customer_id), "status": "quoted"},
  840. {"$set": {"status": "confirmed", "confirmed_at": utc_now(), "updated_at": utc_now()}, "$inc": {"version": 1}},
  841. return_document=ReturnDocument.AFTER,
  842. )
  843. if not order:
  844. raise ServiceDataError("invalid_order_state", "当前报价无法确认。")
  845. await record_service_event(order_id, "quote_confirmed", actor_id=int(customer_id))
  846. return order
  847. async def reject_service_order(order_id: str, technician_id: int, reason: str) -> dict[str, Any]:
  848. now = utc_now()
  849. order = await ordersdb.find_one_and_update(
  850. {"order_id": str(order_id), "technician_id": int(technician_id), "status": {"$in": ["requested", "quoted"]}},
  851. {"$set": {"status": "rejected", "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
  852. return_document=ReturnDocument.AFTER,
  853. )
  854. if not order:
  855. raise ServiceDataError("invalid_order_state", "当前服务单无法拒绝。")
  856. await _schedule_address_redaction(order_id, now)
  857. await record_service_event(order_id, "order_rejected", actor_id=int(technician_id), reason=reason)
  858. return order
  859. async def start_service_order(order_id: str, technician_id: int) -> dict[str, Any]:
  860. order = await ordersdb.find_one_and_update(
  861. {"order_id": str(order_id), "technician_id": int(technician_id), "status": "confirmed"},
  862. {"$set": {"status": "in_progress", "started_at": utc_now(), "updated_at": utc_now()}, "$inc": {"version": 1}},
  863. return_document=ReturnDocument.AFTER,
  864. )
  865. if not order:
  866. raise ServiceDataError("invalid_order_state", "只有已确认服务单可以开始。")
  867. await record_service_event(order_id, "service_started", actor_id=int(technician_id))
  868. return order
  869. async def cancel_service_order(order_id: str, actor_id: int, reason: str) -> dict[str, Any]:
  870. order = await ordersdb.find_one({"order_id": str(order_id)})
  871. if not order or int(actor_id) not in {int(order["customer_id"]), int(order["technician_id"])}:
  872. raise ServiceDataError("order_not_found", "未找到该服务单。")
  873. if order["status"] not in ORDER_ACTIVE_STATUSES - {"disputed"}:
  874. raise ServiceDataError("invalid_order_state", "当前服务单不能取消。")
  875. status = "canceled_customer" if int(actor_id) == int(order["customer_id"]) else "canceled_technician"
  876. now = utc_now()
  877. updated = await ordersdb.find_one_and_update(
  878. {"order_id": str(order_id), "status": order["status"], "version": order["version"]},
  879. {"$set": {"status": status, "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
  880. return_document=ReturnDocument.AFTER,
  881. )
  882. if not updated:
  883. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  884. await _schedule_address_redaction(order_id, now)
  885. await qrdb.update_many({"order_id": str(order_id), "status": "issued"}, {"$set": {"status": "revoked", "revoked_at": now, "revoke_reason": "服务单已取消"}})
  886. await record_service_event(order_id, "order_canceled", actor_id=int(actor_id), reason=reason, metadata={"status": status})
  887. return updated
  888. async def dispute_service_order(order_id: str, actor_id: int, reason: str) -> dict[str, Any]:
  889. current = await ordersdb.find_one(
  890. {
  891. "order_id": str(order_id),
  892. "status": {"$in": list(ORDER_ACTIVE_STATUSES - {"disputed"})},
  893. "$or": [{"customer_id": int(actor_id)}, {"technician_id": int(actor_id)}],
  894. }
  895. )
  896. if not current:
  897. raise ServiceDataError("invalid_order_state", "当前服务单不能发起争议。")
  898. order = await ordersdb.find_one_and_update(
  899. {"order_id": str(order_id), "status": current["status"], "version": current["version"]},
  900. {"$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}},
  901. return_document=ReturnDocument.AFTER,
  902. )
  903. if not order:
  904. raise ServiceDataError("invalid_order_state", "当前服务单不能发起争议。")
  905. await qrdb.update_many({"order_id": str(order_id), "status": "issued"}, {"$set": {"status": "revoked", "revoked_at": utc_now(), "revoke_reason": "服务单争议"}})
  906. await record_service_event(order_id, "order_disputed", actor_id=int(actor_id), reason=reason)
  907. return order
  908. async def resolve_service_dispute(
  909. order_id: str,
  910. *,
  911. actor_id: str,
  912. action: str,
  913. reason: str,
  914. ) -> dict[str, Any]:
  915. if action == "void":
  916. return await admin_void_service_order(order_id, actor_id=actor_id, reason=reason)
  917. if action != "resume":
  918. raise ServiceDataError("invalid_dispute_action", "不支持该争议处理操作。")
  919. current = await ordersdb.find_one({"order_id": str(order_id), "status": "disputed"})
  920. if not current:
  921. raise ServiceDataError("invalid_order_state", "该服务单当前不在争议中。")
  922. resume_status = str(current.get("status_before_dispute") or "confirmed")
  923. if resume_status == "completion_pending":
  924. resume_status = "in_progress"
  925. if resume_status not in {"requested", "quoted", "confirmed", "in_progress"}:
  926. resume_status = "confirmed"
  927. updated = await ordersdb.find_one_and_update(
  928. {"order_id": str(order_id), "status": "disputed", "version": current["version"]},
  929. {
  930. "$set": {
  931. "status": resume_status,
  932. "dispute_resolved_at": utc_now(),
  933. "dispute_resolution_reason": _clean_text(reason, max_length=500, required=True),
  934. "updated_at": utc_now(),
  935. },
  936. "$unset": {"status_before_dispute": ""},
  937. "$inc": {"version": 1},
  938. },
  939. return_document=ReturnDocument.AFTER,
  940. )
  941. if not updated:
  942. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  943. await record_service_event(order_id, "dispute_resolved", actor_id=actor_id, reason=reason, metadata={"status": resume_status})
  944. return updated
  945. async def admin_void_service_order(order_id: str, *, actor_id: str, reason: str) -> dict[str, Any]:
  946. now = utc_now()
  947. order = await ordersdb.find_one_and_update(
  948. {"order_id": str(order_id), "status": {"$ne": "voided"}},
  949. {"$set": {"status": "voided", "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
  950. return_document=ReturnDocument.AFTER,
  951. )
  952. if not order:
  953. raise ServiceDataError("order_not_found", "未找到可作废服务单。")
  954. await _schedule_address_redaction(order_id, now)
  955. await qrdb.update_many({"order_id": str(order_id), "status": {"$ne": "voided"}}, {"$set": {"status": "voided", "revoked_at": now}})
  956. await reviewsdb.update_many({"order_id": str(order_id), "status": "approved"}, {"$set": {"status": "voided", "moderated_at": now, "moderation_reason": reason, "updated_at": now}})
  957. await recompute_review_stats()
  958. await record_service_event(order_id, "order_voided", actor_id=actor_id, reason=reason)
  959. return order
  960. async def issue_review_qr(order_id: str, technician_id: int) -> tuple[dict[str, Any], str]:
  961. await ensure_service_indexes()
  962. order = await ordersdb.find_one({"order_id": str(order_id), "technician_id": int(technician_id)})
  963. if order and order.get("status") == "completion_pending":
  964. previous = await qrdb.find_one(
  965. {"order_id": str(order_id), "status": "issued"},
  966. sort=[("created_at", DESCENDING)],
  967. )
  968. if previous and _as_utc(previous["expires_at"]) <= utc_now():
  969. now = utc_now()
  970. expired = await qrdb.find_one_and_update(
  971. {"qr_id": previous["qr_id"], "status": "issued"},
  972. {"$set": {"status": "expired", "expired_at": now}},
  973. return_document=ReturnDocument.AFTER,
  974. )
  975. if expired:
  976. await ordersdb.update_one(
  977. {
  978. "order_id": str(order_id),
  979. "status": "completion_pending",
  980. "current_review_qr_id": previous["qr_id"],
  981. },
  982. {
  983. "$set": {"status": "in_progress", "updated_at": now},
  984. "$unset": {"current_review_qr_id": ""},
  985. "$inc": {"version": 1},
  986. },
  987. )
  988. order = await ordersdb.find_one(
  989. {"order_id": str(order_id), "technician_id": int(technician_id)}
  990. )
  991. if not order or order.get("status") != "in_progress":
  992. raise ServiceDataError("invalid_order_state", "只有服务中的订单可以生成评价二维码。")
  993. if await qrdb.find_one({"order_id": str(order_id), "status": {"$in": ["issued", "claimed"]}}):
  994. raise ServiceDataError("review_qr_exists", "该服务单已经生成评价二维码。")
  995. token = token_urlsafe(18)
  996. now = utc_now()
  997. hours = (await get_service_settings())["review_qr_expiry_hours"]
  998. document = {
  999. "qr_id": uuid4().hex,
  1000. "order_id": str(order_id),
  1001. "technician_id": int(technician_id),
  1002. "customer_id": int(order["customer_id"]),
  1003. "package_snapshot": order["package_snapshot"],
  1004. "quote_snapshot": order.get("quote_snapshot"),
  1005. "token_hash": hashlib.sha256(token.encode()).hexdigest(),
  1006. "status": "issued",
  1007. "expires_at": now + timedelta(hours=hours),
  1008. "created_at": now,
  1009. }
  1010. try:
  1011. await qrdb.insert_one(document)
  1012. except DuplicateKeyError as exc:
  1013. raise ServiceDataError("review_qr_exists", "该服务单已经生成评价二维码。") from exc
  1014. updated = await ordersdb.find_one_and_update(
  1015. {
  1016. "order_id": str(order_id),
  1017. "status": "in_progress",
  1018. "version": order["version"],
  1019. },
  1020. {
  1021. "$set": {
  1022. "status": "completion_pending",
  1023. "current_review_qr_id": document["qr_id"],
  1024. "completion_requested_at": now,
  1025. "updated_at": now,
  1026. },
  1027. "$inc": {"version": 1},
  1028. },
  1029. return_document=ReturnDocument.AFTER,
  1030. )
  1031. if not updated:
  1032. await qrdb.update_one(
  1033. {"qr_id": document["qr_id"], "status": "issued"},
  1034. {
  1035. "$set": {
  1036. "status": "voided",
  1037. "revoked_at": utc_now(),
  1038. "revoke_reason": "服务单状态已变化",
  1039. }
  1040. },
  1041. )
  1042. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  1043. await record_service_event(order_id, "review_qr_issued", actor_id=int(technician_id), metadata={"qr_id": document["qr_id"]})
  1044. return document, token
  1045. async def preview_review_qr(token: str, customer_id: int) -> dict[str, Any]:
  1046. document = await qrdb.find_one({"token_hash": hashlib.sha256(str(token).encode()).hexdigest()})
  1047. if not document:
  1048. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1049. if int(document["customer_id"]) != int(customer_id):
  1050. raise ServiceDataError("review_qr_customer_mismatch", "该二维码不属于当前顾客。")
  1051. if document.get("status") == "claimed":
  1052. if await reviewsdb.find_one({"qr_id": document["qr_id"]}):
  1053. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1054. return document
  1055. if document.get("status") != "issued":
  1056. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1057. if _as_utc(document["expires_at"]) <= utc_now():
  1058. now = utc_now()
  1059. expired = await qrdb.find_one_and_update(
  1060. {"qr_id": document["qr_id"], "status": "issued"},
  1061. {"$set": {"status": "expired", "expired_at": now}},
  1062. return_document=ReturnDocument.AFTER,
  1063. )
  1064. if expired:
  1065. await ordersdb.update_one(
  1066. {
  1067. "order_id": document["order_id"],
  1068. "status": "completion_pending",
  1069. "current_review_qr_id": document["qr_id"],
  1070. },
  1071. {
  1072. "$set": {"status": "in_progress", "updated_at": now},
  1073. "$unset": {"current_review_qr_id": ""},
  1074. "$inc": {"version": 1},
  1075. },
  1076. )
  1077. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1078. return document
  1079. async def claim_review_qr(token: str, customer_id: int) -> tuple[dict[str, Any], dict[str, Any]]:
  1080. token_hash = hashlib.sha256(str(token).encode()).hexdigest()
  1081. now = utc_now()
  1082. await preview_review_qr(token, customer_id)
  1083. document = await qrdb.find_one_and_update(
  1084. {"token_hash": token_hash, "customer_id": int(customer_id), "status": "issued"},
  1085. {"$set": {"status": "claimed", "claimed_at": now}},
  1086. return_document=ReturnDocument.AFTER,
  1087. )
  1088. if not document:
  1089. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1090. order = await ordersdb.find_one_and_update(
  1091. {"order_id": document["order_id"], "status": "completion_pending", "customer_id": int(customer_id)},
  1092. {"$set": {"status": "completed", "completed_at": now, "closed_at": now, "updated_at": now}, "$inc": {"version": 1}},
  1093. return_document=ReturnDocument.AFTER,
  1094. )
  1095. if not order:
  1096. await qrdb.update_one({"qr_id": document["qr_id"]}, {"$set": {"status": "voided"}})
  1097. raise ServiceDataError("invalid_order_state", "服务单状态已变化,无法完成评价。")
  1098. await _schedule_address_redaction(order["order_id"], now)
  1099. await record_service_event(order["order_id"], "service_completed", actor_id=int(customer_id), metadata={"qr_id": document["qr_id"]})
  1100. return document, order
  1101. def normalize_review_template(values: dict[str, Any]) -> dict[str, Any]:
  1102. ratings = []
  1103. for item in values.get("rating_questions", []):
  1104. if len(ratings) >= 5:
  1105. raise ServiceDataError("too_many_review_questions", "评分维度最多 5 个。")
  1106. weight = int(item.get("weight") or 0)
  1107. if weight <= 0:
  1108. raise ServiceDataError("invalid_review_weight", "评分权重必须大于 0。")
  1109. ratings.append({"question_id": str(item.get("question_id") or uuid4().hex[:12]), "label": _clean_text(item.get("label"), max_length=40, required=True), "description": _clean_text(item.get("description"), max_length=160), "required": bool(item.get("required", True)), "weight": weight})
  1110. if not ratings or sum(item["weight"] for item in ratings) != 100:
  1111. raise ServiceDataError("invalid_review_weight", "评分维度权重合计必须为 100%。")
  1112. texts = []
  1113. for item in values.get("text_questions", []):
  1114. if len(texts) >= 2:
  1115. raise ServiceDataError("too_many_review_questions", "文字问题最多 2 个。")
  1116. texts.append({"question_id": str(item.get("question_id") or uuid4().hex[:12]), "label": _clean_text(item.get("label"), max_length=40, required=True), "description": _clean_text(item.get("description"), max_length=160), "required": bool(item.get("required")), "max_length": max(50, min(1000, int(item.get("max_length") or 500)))})
  1117. return {"rating_questions": ratings, "text_questions": texts}
  1118. async def get_active_review_template() -> dict[str, Any]:
  1119. await ensure_service_indexes()
  1120. current = await reviewtemplatesdb.find_one({"status": "active"}, sort=[("version", DESCENDING)])
  1121. if current:
  1122. return current
  1123. now = utc_now()
  1124. document = {"template_id": uuid4().hex, "version": 1, "status": "active", **normalize_review_template(DEFAULT_REVIEW_TEMPLATE), "created_at": now, "published_at": now}
  1125. try:
  1126. await reviewtemplatesdb.insert_one(document)
  1127. except DuplicateKeyError:
  1128. return await reviewtemplatesdb.find_one({"status": "active"}, sort=[("version", DESCENDING)]) or document
  1129. return document
  1130. async def publish_review_template(values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
  1131. normalized = normalize_review_template(values)
  1132. await ensure_service_indexes()
  1133. while True:
  1134. current = await reviewtemplatesdb.find_one(sort=[("version", DESCENDING)])
  1135. version = int((current or {}).get("version") or 0) + 1
  1136. now = utc_now()
  1137. document = {
  1138. "template_id": uuid4().hex,
  1139. "version": version,
  1140. "status": "active",
  1141. **normalized,
  1142. "published_by": str(actor_id),
  1143. "created_at": now,
  1144. "published_at": now,
  1145. }
  1146. try:
  1147. await reviewtemplatesdb.insert_one(document)
  1148. break
  1149. except DuplicateKeyError:
  1150. continue
  1151. await reviewtemplatesdb.update_many(
  1152. {"status": "active", "template_id": {"$ne": document["template_id"]}},
  1153. {"$set": {"status": "retired", "retired_at": now}},
  1154. )
  1155. return document
  1156. async def begin_review_draft(
  1157. *,
  1158. customer_id: int,
  1159. technician_id: int,
  1160. source: str,
  1161. template: dict[str, Any],
  1162. package_id: str | None = None,
  1163. qr_id: str | None = None,
  1164. ) -> dict[str, Any]:
  1165. await ensure_service_indexes()
  1166. now = utc_now()
  1167. current = await reviewdraftsdb.find_one({"customer_id": int(customer_id)})
  1168. identity = {
  1169. "technician_id": int(technician_id),
  1170. "source": str(source),
  1171. "package_id": str(package_id) if package_id else None,
  1172. "qr_id": str(qr_id) if qr_id else None,
  1173. }
  1174. if (
  1175. current
  1176. and all(current.get(key) == value for key, value in identity.items())
  1177. and _as_utc(current["expires_at"]) > now
  1178. ):
  1179. return current
  1180. document = {
  1181. "draft_id": uuid4().hex,
  1182. "customer_id": int(customer_id),
  1183. **identity,
  1184. "template_snapshot": {
  1185. key: template[key]
  1186. for key in (
  1187. "template_id",
  1188. "version",
  1189. "rating_questions",
  1190. "text_questions",
  1191. )
  1192. },
  1193. "rating_index": 0,
  1194. "text_index": 0,
  1195. "answers": {},
  1196. "created_at": now,
  1197. "updated_at": now,
  1198. "expires_at": now + timedelta(days=30),
  1199. }
  1200. await reviewdraftsdb.update_one(
  1201. {"customer_id": int(customer_id)}, {"$set": document}, upsert=True
  1202. )
  1203. return document
  1204. async def save_review_draft_progress(
  1205. customer_id: int,
  1206. *,
  1207. answers: dict[str, Any],
  1208. rating_index: int,
  1209. text_index: int,
  1210. ) -> None:
  1211. await reviewdraftsdb.update_one(
  1212. {"customer_id": int(customer_id)},
  1213. {
  1214. "$set": {
  1215. "answers": dict(answers),
  1216. "rating_index": max(0, int(rating_index)),
  1217. "text_index": max(0, int(text_index)),
  1218. "updated_at": utc_now(),
  1219. }
  1220. },
  1221. )
  1222. async def delete_review_draft(customer_id: int) -> None:
  1223. await reviewdraftsdb.delete_one({"customer_id": int(customer_id)})
  1224. def _score_review(template: dict[str, Any], answers: dict[str, Any]) -> tuple[float, dict[str, int], dict[str, str]]:
  1225. ratings: dict[str, int] = {}
  1226. texts: dict[str, str] = {}
  1227. total = 0.0
  1228. for question in template["rating_questions"]:
  1229. raw = answers.get(question["question_id"])
  1230. if raw in (None, "") and not question["required"]:
  1231. continue
  1232. try:
  1233. score = int(raw)
  1234. except (TypeError, ValueError) as exc:
  1235. raise ServiceDataError("invalid_review_answer", "评分必须是 1 到 5 的整数。") from exc
  1236. if not 1 <= score <= 5:
  1237. raise ServiceDataError("invalid_review_answer", "评分必须在 1 到 5 之间。")
  1238. ratings[question["question_id"]] = score
  1239. total += score * question["weight"] / 100
  1240. for question in template["text_questions"]:
  1241. text = _clean_text(answers.get(question["question_id"]), max_length=question["max_length"], required=question["required"])
  1242. if text:
  1243. texts[question["question_id"]] = text
  1244. return round(total, 4), ratings, texts
  1245. async def submit_review(
  1246. *,
  1247. customer: Any,
  1248. technician_id: int,
  1249. source: str,
  1250. answers: dict[str, Any],
  1251. anonymous: bool,
  1252. package_id: str | None = None,
  1253. qr_id: str | None = None,
  1254. template_snapshot: dict[str, Any] | None = None,
  1255. ) -> dict[str, Any]:
  1256. await observe_customer(customer)
  1257. await require_customer_allowed(int(customer.id))
  1258. if source not in REVIEW_SOURCES:
  1259. raise ServiceDataError("invalid_review_source", "评价来源无效。")
  1260. profile = await profilesdb.find_one({"user_id": int(technician_id), "application_status": "approved"})
  1261. if not profile:
  1262. raise ServiceDataError("technician_not_found", "未找到已认证技师。")
  1263. if int(customer.id) == int(technician_id):
  1264. raise ServiceDataError("self_review_not_allowed", "不能评价自己。")
  1265. order_id = None
  1266. package_snapshot = None
  1267. if source == "qr_verified":
  1268. qr = await qrdb.find_one({"qr_id": str(qr_id), "status": "claimed", "customer_id": int(customer.id), "technician_id": int(technician_id)})
  1269. if not qr:
  1270. raise ServiceDataError("review_qr_required", "缺少已完成服务单的评价资格。")
  1271. if await reviewsdb.find_one({"qr_id": str(qr_id)}):
  1272. raise ServiceDataError("review_already_submitted", "该服务单已经提交评价。")
  1273. order_id = qr["order_id"]
  1274. package_snapshot = qr["package_snapshot"]
  1275. else:
  1276. service_profile = profile.get("service_profile") or {}
  1277. package_snapshot = next((item for item in service_profile.get("packages", []) if str(item.get("package_id")) == str(package_id)), None)
  1278. if package_id and not package_snapshot:
  1279. raise ServiceDataError("package_not_found", "未找到所选套餐。")
  1280. package_snapshot = package_snapshot or {"package_id": None, "name": "其他服务", "category": "其他"}
  1281. template = template_snapshot or await get_active_review_template()
  1282. score, rating_answers, text_answers = _score_review(template, answers)
  1283. normalized_content = "\n".join(
  1284. value.casefold() for _, value in sorted(text_answers.items()) if value
  1285. )
  1286. now = utc_now()
  1287. document = {
  1288. "review_id": uuid4().hex,
  1289. "technician_id": int(technician_id),
  1290. "technician_name": profile.get("display_name") or f"技师 {technician_id}",
  1291. "customer_id": int(customer.id),
  1292. "customer_name": _clean_text(getattr(customer, "first_name", ""), max_length=80) or f"用户 {customer.id}",
  1293. "source": source,
  1294. "package_snapshot": package_snapshot,
  1295. "category": package_snapshot.get("category") or "其他",
  1296. "template_snapshot": {key: template[key] for key in ("template_id", "version", "rating_questions", "text_questions")},
  1297. "rating_answers": rating_answers,
  1298. "text_answers": text_answers,
  1299. "content_fingerprint": (
  1300. hashlib.sha256(normalized_content.encode()).hexdigest()
  1301. if normalized_content
  1302. else None
  1303. ),
  1304. "score": score,
  1305. "anonymous": bool(anonymous),
  1306. "status": "pending",
  1307. "created_at": now,
  1308. "updated_at": now,
  1309. }
  1310. if qr_id:
  1311. document["qr_id"] = str(qr_id)
  1312. document["order_id"] = order_id
  1313. try:
  1314. await reviewsdb.insert_one(document)
  1315. except DuplicateKeyError as exc:
  1316. if qr_id:
  1317. raise ServiceDataError(
  1318. "review_already_submitted", "该服务单已经提交评价。"
  1319. ) from exc
  1320. raise
  1321. return document
  1322. async def moderate_review(review_id: str, action: str, *, actor_id: str, reason: str = "") -> dict[str, Any]:
  1323. status_map = {"approve": "approved", "reject": "rejected", "void": "voided"}
  1324. source_status = {"approve": "pending", "reject": "pending", "void": "approved"}
  1325. if action not in status_map:
  1326. raise ServiceDataError("invalid_review_action", "不支持该评价操作。")
  1327. current = await reviewsdb.find_one({"review_id": str(review_id)})
  1328. if not current:
  1329. raise ServiceDataError("review_not_found", "未找到评价。")
  1330. if current["status"] != source_status[action]:
  1331. raise ServiceDataError("invalid_review_state", "当前评价状态不能执行该审核操作。")
  1332. if action in {"reject", "void"} and not _clean_text(reason, max_length=500):
  1333. raise ServiceDataError("reason_required", "拒绝或作废必须填写原因。")
  1334. updated = await reviewsdb.find_one_and_update(
  1335. {"review_id": str(review_id), "status": source_status[action]},
  1336. {"$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()}},
  1337. return_document=ReturnDocument.AFTER,
  1338. )
  1339. if not updated:
  1340. raise ServiceDataError("concurrent_review_update", "评价状态已变化,请刷新后重试。")
  1341. await recompute_review_stats()
  1342. return updated
  1343. async def recompute_review_stats() -> None:
  1344. await ensure_service_indexes()
  1345. approved = await reviewsdb.find({"status": "approved"}).to_list(length=100000)
  1346. grouped: dict[tuple[int, str], list[dict[str, Any]]] = defaultdict(list)
  1347. category_scores: dict[str, list[float]] = defaultdict(list)
  1348. for review in approved:
  1349. key = (int(review["technician_id"]), str(review.get("category") or "其他"))
  1350. grouped[key].append(review)
  1351. category_scores[key[1]].append(float(review["score"]))
  1352. now = utc_now()
  1353. await statsdb.delete_many({})
  1354. documents = []
  1355. for (technician_id, category), values in grouped.items():
  1356. count = len(values)
  1357. average = sum(float(item["score"]) for item in values) / count
  1358. category_average = sum(category_scores[category]) / len(category_scores[category])
  1359. rank_score = count / (count + 5) * average + 5 / (count + 5) * category_average
  1360. dimensions: dict[str, list[int]] = defaultdict(list)
  1361. for review in values:
  1362. for key, score in review.get("rating_answers", {}).items():
  1363. dimensions[key].append(int(score))
  1364. documents.append(
  1365. {
  1366. "technician_id": technician_id,
  1367. "category": category,
  1368. "review_count": count,
  1369. "average_score": round(average, 4),
  1370. "category_average": round(category_average, 4),
  1371. "rank_score": round(rank_score, 6),
  1372. "eligible": count >= 3,
  1373. "dimension_averages": {key: round(sum(items) / len(items), 4) for key, items in dimensions.items()},
  1374. "last_review_at": max(item.get("moderated_at") or item["created_at"] for item in values),
  1375. "updated_at": now,
  1376. }
  1377. )
  1378. if documents:
  1379. await statsdb.insert_many(documents)
  1380. async def list_leaderboard(*, category: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
  1381. await ensure_service_indexes()
  1382. filters: dict[str, Any] = {"eligible": True}
  1383. if category:
  1384. filters["category"] = str(category)
  1385. visible_ids = [
  1386. int(item["user_id"])
  1387. async for item in profilesdb.find(
  1388. {
  1389. "application_status": "approved",
  1390. "listed": True,
  1391. "username": {"$nin": [None, ""]},
  1392. },
  1393. {"user_id": 1},
  1394. )
  1395. ]
  1396. filters["technician_id"] = {"$in": visible_ids}
  1397. total = await statsdb.count_documents(filters)
  1398. 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)
  1399. 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})}
  1400. 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]
  1401. return visible, total
  1402. async def list_public_reviews(
  1403. technician_id: int,
  1404. *,
  1405. page: int = 1,
  1406. page_size: int = 10,
  1407. ) -> tuple[list[dict[str, Any]], int]:
  1408. filters = {"technician_id": int(technician_id), "status": "approved"}
  1409. total = await reviewsdb.count_documents(filters)
  1410. items = await reviewsdb.find(filters).sort("moderated_at", DESCENDING).skip(
  1411. (max(1, page) - 1) * page_size
  1412. ).limit(page_size).to_list(length=page_size)
  1413. public_items = []
  1414. for item in items:
  1415. public_items.append(
  1416. {
  1417. "review_id": item["review_id"],
  1418. "source": item["source"],
  1419. "source_label": (
  1420. "完成服务单评价"
  1421. if item["source"] == "qr_verified"
  1422. else "用户主动评价"
  1423. ),
  1424. "customer_name": "匿名顾客" if item.get("anonymous", True) else item.get("customer_name"),
  1425. "package_name": (item.get("package_snapshot") or {}).get("name"),
  1426. "category": item.get("category"),
  1427. "score": item.get("score"),
  1428. "rating_answers": item.get("rating_answers", {}),
  1429. "text_answers": item.get("text_answers", {}),
  1430. "approved_at": item.get("moderated_at"),
  1431. }
  1432. )
  1433. return public_items, total
  1434. async def get_technician_review_summary(technician_id: int) -> dict[str, Any]:
  1435. items = await reviewsdb.find(
  1436. {"technician_id": int(technician_id), "status": "approved"},
  1437. {"score": 1, "source": 1},
  1438. ).to_list(length=100000)
  1439. count = len(items)
  1440. return {
  1441. "review_count": count,
  1442. "average_score": round(
  1443. sum(float(item["score"]) for item in items) / count,
  1444. 4,
  1445. )
  1446. if count
  1447. else 0,
  1448. "source_counts": {
  1449. source: sum(1 for item in items if item.get("source") == source)
  1450. for source in REVIEW_SOURCES
  1451. },
  1452. "ranking_title": "审核评价排行",
  1453. }
  1454. async def list_service_orders(*, query: str = "", status: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
  1455. await ensure_service_indexes()
  1456. filters: dict[str, Any] = {}
  1457. if status:
  1458. filters["status"] = status
  1459. if query:
  1460. if query.isdigit():
  1461. filters["$or"] = [{"customer_id": int(query)}, {"technician_id": int(query)}]
  1462. else:
  1463. pattern = re.compile(re.escape(query), re.IGNORECASE)
  1464. filters["$or"] = [{"order_id": pattern}, {"customer_name": pattern}, {"technician_name": pattern}]
  1465. total = await ordersdb.count_documents(filters)
  1466. items = await ordersdb.find(filters).sort("created_at", DESCENDING).skip((max(1, page) - 1) * page_size).limit(page_size).to_list(length=page_size)
  1467. return items, total
  1468. async def list_actor_service_orders(
  1469. actor_id: int,
  1470. *,
  1471. role: str = "all",
  1472. page: int = 1,
  1473. page_size: int = 20,
  1474. ) -> tuple[list[dict[str, Any]], int]:
  1475. if role == "customer":
  1476. filters: dict[str, Any] = {"customer_id": int(actor_id)}
  1477. elif role == "technician":
  1478. filters = {"technician_id": int(actor_id)}
  1479. else:
  1480. filters = {
  1481. "$or": [
  1482. {"customer_id": int(actor_id)},
  1483. {"technician_id": int(actor_id)},
  1484. ]
  1485. }
  1486. total = await ordersdb.count_documents(filters)
  1487. items = await ordersdb.find(filters).sort("created_at", DESCENDING).skip(
  1488. (max(1, page) - 1) * page_size
  1489. ).limit(page_size).to_list(length=page_size)
  1490. return items, total
  1491. async def get_admin_service_order(order_id: str) -> dict[str, Any]:
  1492. order = await ordersdb.find_one({"order_id": str(order_id)})
  1493. if not order:
  1494. raise ServiceDataError("order_not_found", "未找到服务单。")
  1495. result = dict(order)
  1496. result["quotes"] = await quotesdb.find({"order_id": str(order_id)}).sort(
  1497. "version", ASCENDING
  1498. ).to_list(length=100)
  1499. result["events"] = await eventsdb.find({"order_id": str(order_id)}).sort(
  1500. "created_at", ASCENDING
  1501. ).to_list(length=500)
  1502. result["qr_records"] = await qrdb.find({"order_id": str(order_id)}).sort(
  1503. "created_at", DESCENDING
  1504. ).to_list(length=100)
  1505. return result
  1506. async def list_review_qr_records(
  1507. *,
  1508. status: str = "",
  1509. query: str = "",
  1510. page: int = 1,
  1511. page_size: int = 20,
  1512. ) -> tuple[list[dict[str, Any]], int]:
  1513. filters: dict[str, Any] = {}
  1514. if status:
  1515. filters["status"] = status
  1516. if query:
  1517. filters["$or"] = [
  1518. {"order_id": re.compile(re.escape(query), re.IGNORECASE)},
  1519. {"qr_id": re.compile(re.escape(query), re.IGNORECASE)},
  1520. {"customer_id": int(query) if query.isdigit() else -1},
  1521. {"technician_id": int(query) if query.isdigit() else -1},
  1522. ]
  1523. total = await qrdb.count_documents(filters)
  1524. items = await qrdb.find(filters).sort("created_at", DESCENDING).skip(
  1525. (max(1, page) - 1) * page_size
  1526. ).limit(page_size).to_list(length=page_size)
  1527. for item in items:
  1528. item.pop("token_hash", None)
  1529. return items, total
  1530. async def list_reviews(*, status: str = "", source: str = "", query: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
  1531. await ensure_service_indexes()
  1532. filters: dict[str, Any] = {}
  1533. if status:
  1534. filters["status"] = status
  1535. if source:
  1536. filters["source"] = source
  1537. if query:
  1538. pattern = re.compile(re.escape(query), re.IGNORECASE)
  1539. filters["$or"] = [{"technician_name": pattern}, {"customer_name": pattern}, {"review_id": pattern}]
  1540. total = await reviewsdb.count_documents(filters)
  1541. items = await reviewsdb.find(filters).sort("created_at", DESCENDING).skip((max(1, page) - 1) * page_size).limit(page_size).to_list(length=page_size)
  1542. for item in items:
  1543. item["customer_review_count"] = await reviewsdb.count_documents({"customer_id": item["customer_id"], "technician_id": item["technician_id"]})
  1544. 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)}})
  1545. 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)}})
  1546. item["similar_content_count"] = (
  1547. await reviewsdb.count_documents(
  1548. {
  1549. "customer_id": item["customer_id"],
  1550. "technician_id": item["technician_id"],
  1551. "content_fingerprint": item["content_fingerprint"],
  1552. }
  1553. )
  1554. if item.get("content_fingerprint")
  1555. else 0
  1556. )
  1557. return items, total
  1558. async def get_order_for_actor(order_id: str, actor_id: int, *, reveal_address: bool = False) -> dict[str, Any]:
  1559. order = await ordersdb.find_one({"order_id": str(order_id)})
  1560. if not order or int(actor_id) not in {int(order["customer_id"]), int(order["technician_id"])}:
  1561. raise ServiceDataError("order_not_found", "未找到服务单。")
  1562. result = dict(order)
  1563. can_reveal = order["status"] in {"confirmed", "in_progress", "completion_pending", "completed", "disputed"}
  1564. if reveal_address and can_reveal:
  1565. address = await addressesdb.find_one({"order_id": str(order_id), "redacted": False})
  1566. if address:
  1567. result["exact_address"] = _decrypt_address(address["encrypted_payload"])
  1568. return result
  1569. async def admin_reveal_order_address(order_id: str, *, actor_id: str, reason: str) -> dict[str, Any]:
  1570. reason = _clean_text(reason, max_length=500, required=True)
  1571. order = await ordersdb.find_one({"order_id": str(order_id)})
  1572. address = await addressesdb.find_one({"order_id": str(order_id), "redacted": False})
  1573. if not order or not address:
  1574. raise ServiceDataError("address_unavailable", "精确地址已脱敏或不存在。")
  1575. await record_service_event(order_id, "admin_address_revealed", actor_id=actor_id, reason=reason)
  1576. return _decrypt_address(address["encrypted_payload"])
  1577. async def redact_expired_addresses() -> int:
  1578. now = utc_now()
  1579. candidates = await addressesdb.find({"redacted": False, "redact_after": {"$lte": now}}).to_list(length=1000)
  1580. count = 0
  1581. for item in candidates:
  1582. order = await ordersdb.find_one({"order_id": item["order_id"]})
  1583. if not order or order.get("status") not in ORDER_TERMINAL_STATUSES:
  1584. continue
  1585. result = await addressesdb.update_one({"_id": item["_id"], "redacted": False}, {"$set": {"redacted": True, "redacted_at": now, "updated_at": now}, "$unset": {"encrypted_payload": ""}})
  1586. if result.modified_count:
  1587. await ordersdb.update_one(
  1588. {"order_id": item["order_id"]},
  1589. {
  1590. "$set": {"address_redacted": True, "updated_at": now},
  1591. "$unset": {"distance_meters": ""},
  1592. },
  1593. )
  1594. count += result.modified_count
  1595. return count
  1596. 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]:
  1597. if scope not in {"global", "technician"}:
  1598. raise ServiceDataError("invalid_block_scope", "拉黑范围无效。")
  1599. if scope == "technician" and technician_id is None:
  1600. raise ServiceDataError("technician_required", "技师拉黑必须指定技师。")
  1601. now = utc_now()
  1602. await blocksdb.update_one(
  1603. {"scope": scope, "technician_id": int(technician_id or 0), "customer_id": int(customer_id)},
  1604. {"$set": {"active": bool(active), "reason": _clean_text(reason, max_length=500, required=active), "actor_id": actor_id, "updated_at": now}, "$setOnInsert": {"created_at": now}},
  1605. upsert=True,
  1606. )
  1607. return await blocksdb.find_one({"scope": scope, "technician_id": int(technician_id or 0), "customer_id": int(customer_id)}) or {}
  1608. async def set_technician_customer_block(
  1609. *,
  1610. technician_id: int,
  1611. customer_id: int,
  1612. active: bool,
  1613. reason: str,
  1614. ) -> dict[str, Any]:
  1615. profile = await profilesdb.find_one(
  1616. {"user_id": int(technician_id), "application_status": "approved"}
  1617. )
  1618. if not profile:
  1619. raise ServiceDataError("technician_required", "只有已认证技师可以管理个人拉黑。")
  1620. if not await ordersdb.find_one(
  1621. {"technician_id": int(technician_id), "customer_id": int(customer_id)}
  1622. ):
  1623. raise ServiceDataError("customer_relationship_required", "只能拉黑曾向你发起服务请求的顾客。")
  1624. return await set_customer_block(
  1625. customer_id=int(customer_id),
  1626. scope="technician",
  1627. technician_id=int(technician_id),
  1628. active=active,
  1629. actor_id=int(technician_id),
  1630. reason=reason,
  1631. )
  1632. async def list_customer_blocks(
  1633. *,
  1634. scope: str = "global",
  1635. active: bool | None = None,
  1636. page: int = 1,
  1637. page_size: int = 20,
  1638. ) -> tuple[list[dict[str, Any]], int]:
  1639. filters: dict[str, Any] = {"scope": scope}
  1640. if active is not None:
  1641. filters["active"] = active
  1642. total = await blocksdb.count_documents(filters)
  1643. items = await blocksdb.find(filters).sort("updated_at", DESCENDING).skip(
  1644. (max(1, page) - 1) * page_size
  1645. ).limit(page_size).to_list(length=page_size)
  1646. return items, total
  1647. async def create_service_report(*, reporter_id: int, target_type: str, target_id: str, reason: str) -> dict[str, Any]:
  1648. if target_type not in {"technician", "customer", "order", "review"}:
  1649. raise ServiceDataError("invalid_report_target", "举报对象无效。")
  1650. 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()}
  1651. await reportsdb.insert_one(document)
  1652. return document
  1653. async def list_service_reports(
  1654. *,
  1655. status: str = "",
  1656. page: int = 1,
  1657. page_size: int = 20,
  1658. ) -> tuple[list[dict[str, Any]], int]:
  1659. filters = {"status": status} if status else {}
  1660. total = await reportsdb.count_documents(filters)
  1661. items = await reportsdb.find(filters).sort("created_at", DESCENDING).skip(
  1662. (max(1, page) - 1) * page_size
  1663. ).limit(page_size).to_list(length=page_size)
  1664. return items, total
  1665. async def moderate_service_report(
  1666. report_id: str,
  1667. *,
  1668. action: str,
  1669. actor_id: str,
  1670. reason: str,
  1671. ) -> dict[str, Any]:
  1672. status_map = {"resolve": "resolved", "dismiss": "dismissed"}
  1673. if action not in status_map:
  1674. raise ServiceDataError("invalid_report_action", "不支持该举报操作。")
  1675. updated = await reportsdb.find_one_and_update(
  1676. {"report_id": str(report_id), "status": "open"},
  1677. {
  1678. "$set": {
  1679. "status": status_map[action],
  1680. "handled_by": str(actor_id),
  1681. "handled_at": utc_now(),
  1682. "handling_reason": _clean_text(reason, max_length=500, required=True),
  1683. }
  1684. },
  1685. return_document=ReturnDocument.AFTER,
  1686. )
  1687. if not updated:
  1688. raise ServiceDataError("report_not_found", "举报不存在或已处理。")
  1689. return updated
  1690. async def fulfillment_metrics() -> dict[str, Any]:
  1691. await ensure_service_indexes()
  1692. counts = {status: await ordersdb.count_documents({"status": status}) for status in ORDER_STATUSES}
  1693. valid_requests = sum(counts.values()) - counts["voided"]
  1694. quoted = await ordersdb.count_documents(
  1695. {"status": {"$ne": "voided"}, "current_quote_id": {"$exists": True}}
  1696. )
  1697. confirmed = await ordersdb.count_documents(
  1698. {"status": {"$ne": "voided"}, "confirmed_at": {"$exists": True}}
  1699. )
  1700. completed = counts["completed"]
  1701. canceled = counts["canceled_customer"] + counts["canceled_technician"]
  1702. directory_views = await eventsdb.count_documents({"event_type": "directory_viewed"})
  1703. technician_views = await eventsdb.count_documents({"event_type": "technician_viewed"})
  1704. request_starts = await eventsdb.count_documents({"event_type": "request_started"})
  1705. return {
  1706. "counts": counts,
  1707. "valid_requests": valid_requests,
  1708. "quoted_requests": quoted,
  1709. "confirmed_requests": confirmed,
  1710. "directory_views": directory_views,
  1711. "technician_views": technician_views,
  1712. "request_starts": request_starts,
  1713. "request_submission_rate": min(1.0, round(valid_requests / request_starts, 4))
  1714. if request_starts
  1715. else 0,
  1716. "quote_rate": round(quoted / valid_requests, 4) if valid_requests else 0,
  1717. "customer_confirmation_rate": round(confirmed / quoted, 4) if quoted else 0,
  1718. "match_success_rate": round(confirmed / valid_requests, 4) if valid_requests else 0,
  1719. "fulfillment_rate": round(completed / confirmed, 4) if confirmed else 0,
  1720. "cancellation_rate": round(canceled / valid_requests, 4) if valid_requests else 0,
  1721. "cancellation_distribution": {
  1722. "customer": counts["canceled_customer"],
  1723. "technician": counts["canceled_technician"],
  1724. },
  1725. "scope_note": "仅统计平台服务单;直接 Telegram 私聊成交不在统计范围。",
  1726. }