dbservice.py 101 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386
  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. reviewtopicsdb = control_db.teacher_review_topics
  31. profilesdb = control_db.directory_profiles
  32. locationsdb = control_db.directory_locations
  33. membershipsdb = control_db.directory_memberships
  34. settingsdb = control_db.directory_settings
  35. SERVICE_MODES = {"at_store", "onsite"}
  36. PRICE_MODES = {"fixed", "starting_at", "range", "negotiable"}
  37. PRICE_UNITS = {"per_service", "per_hour", "per_item", "per_visit"}
  38. TRAVEL_FEE_MODES = {"included", "fixed", "per_km", "quoted"}
  39. ORDER_ACTIVE_STATUSES = {
  40. "requested",
  41. "quoted",
  42. "confirmed",
  43. "in_progress",
  44. "completion_pending",
  45. "disputed",
  46. }
  47. ORDER_TERMINAL_STATUSES = {
  48. "completed",
  49. "rejected",
  50. "canceled_customer",
  51. "canceled_technician",
  52. "voided",
  53. }
  54. ORDER_STATUSES = ORDER_ACTIVE_STATUSES | ORDER_TERMINAL_STATUSES
  55. REVIEW_SOURCES = {"qr_verified", "student_initiated"}
  56. REVIEW_STATUSES = {"pending", "approved", "rejected", "withdrawn", "voided"}
  57. DEFAULT_REVIEW_TEMPLATE = {
  58. "questions": [
  59. {
  60. "question_id": "overall_experience",
  61. "type": "single_choice",
  62. "label": "这次服务体验怎么样?",
  63. "description": "选择最符合实际体验的一项。",
  64. "required": True,
  65. "weight": 100,
  66. "options": [
  67. {"value": "excellent", "label": "非常满意", "score": 5},
  68. {"value": "good", "label": "满意", "score": 4},
  69. {"value": "average", "label": "一般", "score": 3},
  70. {"value": "poor", "label": "不满意", "score": 2},
  71. {"value": "bad", "label": "非常不满意", "score": 1},
  72. ],
  73. },
  74. {
  75. "question_id": "service_highlights",
  76. "type": "multiple_choice",
  77. "label": "哪些方面值得肯定?",
  78. "description": "可多选,也可以跳过。",
  79. "required": False,
  80. "options": [
  81. {"value": "professional", "label": "专业可靠"},
  82. {"value": "communication", "label": "沟通顺畅"},
  83. {"value": "careful", "label": "服务细致"},
  84. {"value": "transparent", "label": "价格透明"},
  85. {"value": "punctual", "label": "准时守约"},
  86. ],
  87. },
  88. {
  89. "question_id": "comment",
  90. "type": "text",
  91. "label": "补充评价",
  92. "description": "可选,请勿填写联系方式或精确地址。",
  93. "required": False,
  94. "max_length": 500,
  95. },
  96. ]
  97. }
  98. BUILTIN_PACKAGE_TEMPLATES = (
  99. {
  100. "template_id": "builtin_quick_at_store",
  101. "name": "快速到店服务",
  102. "sort_order": -300,
  103. "package": {
  104. "package_id": "builtin_package_at_store",
  105. "name": "标准到店服务",
  106. "category": "通用服务",
  107. "description": "顾客到指定区域接受一次标准服务,具体地点和服务细节通过 Telegram 私聊沟通。",
  108. "tags": ["到店", "标准服务"],
  109. "service_modes": ["at_store"],
  110. "price_mode": "fixed",
  111. "currency": "CNY",
  112. "price_unit": "per_service",
  113. "min_price": "200",
  114. "max_price": "200",
  115. "duration_minutes": 60,
  116. "included_items": "一次标准服务",
  117. "excluded_items": "额外耗材和临时加项",
  118. "preparation": "请提前说明具体需求",
  119. "addons": [],
  120. "travel_fee": {
  121. "mode": "quoted",
  122. "amount": "0",
  123. "per_km": "0",
  124. "description": "无额外交通费;其他临时费用请提前沟通。",
  125. },
  126. "service_radius_km": None,
  127. "out_of_range_policy": "",
  128. },
  129. },
  130. {
  131. "template_id": "builtin_quick_onsite",
  132. "name": "快速上门服务",
  133. "sort_order": -200,
  134. "package": {
  135. "package_id": "builtin_package_onsite",
  136. "name": "标准上门服务",
  137. "category": "通用服务",
  138. "description": "技师提供一次标准上门服务,具体地点和服务细节通过 Telegram 私聊沟通。",
  139. "tags": ["上门", "标准服务"],
  140. "service_modes": ["onsite"],
  141. "price_mode": "starting_at",
  142. "currency": "CNY",
  143. "price_unit": "per_visit",
  144. "min_price": "300",
  145. "max_price": "300",
  146. "duration_minutes": 90,
  147. "included_items": "一次标准上门服务",
  148. "excluded_items": "交通费、额外耗材和临时加项",
  149. "preparation": "请提前发送位置并说明具体需求",
  150. "addons": [],
  151. "travel_fee": {
  152. "mode": "quoted",
  153. "amount": "0",
  154. "per_km": "0",
  155. "description": "上门车费根据距离另行沟通。",
  156. },
  157. "service_radius_km": 20,
  158. "out_of_range_policy": "超出服务半径时,请先沟通交通费和是否可以接单。",
  159. },
  160. },
  161. {
  162. "template_id": "builtin_quick_hourly",
  163. "name": "快速按小时服务",
  164. "sort_order": -100,
  165. "package": {
  166. "package_id": "builtin_package_hourly",
  167. "name": "按小时专业服务",
  168. "category": "通用服务",
  169. "description": "按服务时长计价,支持到店或上门,具体工作范围提前确认。",
  170. "tags": ["按小时", "到店", "上门"],
  171. "service_modes": ["at_store", "onsite"],
  172. "price_mode": "starting_at",
  173. "currency": "CNY",
  174. "price_unit": "per_hour",
  175. "min_price": "200",
  176. "max_price": "200",
  177. "duration_minutes": 60,
  178. "included_items": "一小时专业服务",
  179. "excluded_items": "交通费、额外耗材和超时服务",
  180. "preparation": "请提前说明服务内容和预计时长",
  181. "addons": [],
  182. "travel_fee": {
  183. "mode": "quoted",
  184. "amount": "0",
  185. "per_km": "0",
  186. "description": "上门车费、耗材费和超时费用另行沟通。",
  187. },
  188. "service_radius_km": 20,
  189. "out_of_range_policy": "超出服务半径时,请先沟通交通费和是否可以接单。",
  190. },
  191. },
  192. )
  193. _indexes_ready = False
  194. class ServiceDataError(ValueError):
  195. def __init__(self, code: str, message: str):
  196. super().__init__(message)
  197. self.code = code
  198. def utc_now() -> datetime:
  199. return datetime.now(UTC)
  200. def _as_utc(value: datetime) -> datetime:
  201. return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
  202. def _clean_text(value: Any, *, max_length: int, required: bool = False) -> str:
  203. text = re.sub(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]", "", str(value or ""))
  204. text = " ".join(text.strip().split())
  205. if required and not text:
  206. raise ServiceDataError("required_field", "必填内容不能为空。")
  207. if len(text) > max_length:
  208. raise ServiceDataError("text_too_long", f"内容不能超过 {max_length} 个字符。")
  209. return text
  210. def _positive_decimal(value: Any, field: str, *, allow_zero: bool = True) -> str:
  211. try:
  212. number = Decimal(str(value or 0)).quantize(Decimal("0.01"))
  213. except (InvalidOperation, ValueError) as exc:
  214. raise ServiceDataError("invalid_price", f"{field} 必须是有效金额。") from exc
  215. if not number.is_finite():
  216. raise ServiceDataError("invalid_price", f"{field} 必须是有效金额。")
  217. if number < 0 or (not allow_zero and number == 0):
  218. raise ServiceDataError("invalid_price", f"{field} 不能小于 0。")
  219. return format(number, "f")
  220. def _normalize_tags(values: Any) -> list[str]:
  221. if isinstance(values, str):
  222. values = values.replace(",", ",").split(",")
  223. tags = []
  224. for value in values or []:
  225. tag = _clean_text(value, max_length=20)
  226. if tag and tag.casefold() not in {item.casefold() for item in tags}:
  227. tags.append(tag)
  228. if len(tags) > 8:
  229. raise ServiceDataError("too_many_tags", "服务标签最多 8 个。")
  230. return tags
  231. def normalize_package(value: dict[str, Any], *, package_id: str | None = None) -> dict[str, Any]:
  232. modes = list(
  233. dict.fromkeys(
  234. str(item)
  235. for item in (value.get("service_modes") or ["at_store", "onsite"])
  236. )
  237. )
  238. if set(modes) - SERVICE_MODES:
  239. raise ServiceDataError("invalid_service_mode", "服务方式无效。")
  240. price_mode = str(value.get("price_mode") or "negotiable")
  241. if price_mode not in PRICE_MODES:
  242. raise ServiceDataError("invalid_price_mode", "不支持该价格模式。")
  243. price_unit = str(value.get("price_unit") or "per_service")
  244. if price_unit not in PRICE_UNITS:
  245. raise ServiceDataError("invalid_price_unit", "不支持该计价单位。")
  246. currency = str(value.get("currency") or "CNY").upper()
  247. if not re.fullmatch(r"[A-Z]{3}", currency):
  248. raise ServiceDataError("invalid_currency", "币种必须使用三位 ISO 代码。")
  249. min_price = _positive_decimal(value.get("min_price"), "最低价格")
  250. max_price = _positive_decimal(value.get("max_price"), "最高价格")
  251. if Decimal(max_price) and Decimal(max_price) < Decimal(min_price):
  252. raise ServiceDataError("invalid_price_range", "最高价格不能低于最低价格。")
  253. if price_mode == "fixed" and Decimal(min_price) != Decimal(max_price):
  254. raise ServiceDataError("invalid_fixed_price", "固定价的最低价格和最高价格必须一致。")
  255. duration = value.get("duration_minutes")
  256. if duration in (None, "", 0, "0"):
  257. duration = None
  258. else:
  259. try:
  260. duration = int(duration)
  261. except (TypeError, ValueError) as exc:
  262. raise ServiceDataError("invalid_duration", "服务时长必须是分钟数。") from exc
  263. if not 15 <= duration <= 480:
  264. raise ServiceDataError("invalid_duration", "服务时长必须在 15 到 480 分钟之间。")
  265. travel = value.get("travel_fee") if isinstance(value.get("travel_fee"), dict) else {}
  266. travel_mode = str(travel.get("mode") or "quoted")
  267. if travel_mode not in TRAVEL_FEE_MODES:
  268. raise ServiceDataError("invalid_travel_fee", "不支持该交通费模式。")
  269. addons = []
  270. for item in value.get("addons", []) or []:
  271. if len(addons) >= 10:
  272. raise ServiceDataError("too_many_addons", "标准加项最多 10 个。")
  273. addons.append(
  274. {
  275. "addon_id": str(item.get("addon_id") or uuid4().hex),
  276. "name": _clean_text(item.get("name"), max_length=40, required=True),
  277. "unit": _clean_text(item.get("unit"), max_length=20) or "项",
  278. "price": _positive_decimal(item.get("price"), "加项价格"),
  279. "extra_minutes": max(0, min(480, int(item.get("extra_minutes") or 0))),
  280. }
  281. )
  282. radius = value.get("service_radius_km")
  283. if radius in (None, ""):
  284. radius = None
  285. else:
  286. radius = float(radius)
  287. if not 0.5 <= radius <= 200:
  288. raise ServiceDataError("invalid_service_radius", "上门半径必须在 0.5 到 200 公里之间。")
  289. out_of_range_policy = _clean_text(
  290. value.get("out_of_range_policy"),
  291. max_length=200,
  292. )
  293. source_type = str(value.get("source_type") or "custom")
  294. if source_type not in {"template", "custom"}:
  295. raise ServiceDataError("invalid_package_source", "套餐来源无效。")
  296. return {
  297. "package_id": str(package_id or value.get("package_id") or uuid4().hex),
  298. "name": _clean_text(value.get("name"), max_length=50, required=True),
  299. "category": _clean_text(value.get("category"), max_length=30, required=True),
  300. "description": _clean_text(value.get("description"), max_length=800, required=True),
  301. "tags": _normalize_tags(value.get("tags", [])),
  302. "service_modes": modes,
  303. "price_mode": price_mode,
  304. "currency": currency,
  305. "price_unit": price_unit,
  306. "min_price": min_price,
  307. "max_price": max_price,
  308. "price_display": _clean_text(value.get("price_display"), max_length=80),
  309. "duration_minutes": duration,
  310. "included_items": _clean_text(value.get("included_items"), max_length=500),
  311. "excluded_items": _clean_text(value.get("excluded_items"), max_length=500),
  312. "preparation": _clean_text(value.get("preparation"), max_length=500),
  313. "addons": addons,
  314. "travel_fee": {
  315. "mode": travel_mode,
  316. "amount": _positive_decimal(travel.get("amount"), "交通费"),
  317. "per_km": _positive_decimal(travel.get("per_km"), "每公里交通费"),
  318. "description": _clean_text(travel.get("description"), max_length=200),
  319. },
  320. "service_radius_km": radius,
  321. "out_of_range_policy": out_of_range_policy,
  322. "source_type": source_type,
  323. "source_template_id": value.get("source_template_id"),
  324. "source_template_version": value.get("source_template_version"),
  325. "customized_from_template": bool(value.get("customized_from_template")),
  326. }
  327. def _normalize_profile(value: dict[str, Any]) -> dict[str, Any]:
  328. packages = [normalize_package(item) for item in value.get("packages", [])]
  329. if not 1 <= len(packages) <= 5:
  330. raise ServiceDataError("invalid_package_count", "技师必须发布 1 到 5 个套餐。")
  331. modes = sorted({mode for item in packages for mode in item["service_modes"]})
  332. venue = value.get("venue") if isinstance(value.get("venue"), dict) else {}
  333. onsite = value.get("onsite_policy") if isinstance(value.get("onsite_policy"), dict) else {}
  334. public_area_text = _clean_text(
  335. value.get("public_area_text"), max_length=80, required=True
  336. )
  337. venue_name = _clean_text(venue.get("name"), max_length=80)
  338. venue_address = _clean_text(venue.get("address_hint"), max_length=120)
  339. onsite_description = _clean_text(onsite.get("description"), max_length=300)
  340. return {
  341. "headline": _clean_text(value.get("headline"), max_length=80, required=True),
  342. "bio": _clean_text(value.get("bio"), max_length=800, required=True),
  343. "tags": _normalize_tags(value.get("tags", [])),
  344. "contact_hours": _clean_text(value.get("contact_hours"), max_length=120)
  345. or "请通过 Telegram 私聊沟通具体服务信息",
  346. "service_modes": modes,
  347. "accepting_requests": bool(value.get("accepting_requests", True)),
  348. "public_area_text": public_area_text,
  349. "venue": {
  350. "name": venue_name,
  351. "address_hint": venue_address,
  352. },
  353. "onsite_policy": {
  354. "description": onsite_description,
  355. },
  356. "packages": packages,
  357. "is_complete": True,
  358. }
  359. async def ensure_service_indexes() -> None:
  360. global _indexes_ready
  361. if _indexes_ready:
  362. return
  363. await templatesdb.create_index([("template_id", ASCENDING)], unique=True)
  364. await templatesdb.create_index([("status", ASCENDING), ("sort_order", ASCENDING)])
  365. await ordersdb.create_index([("order_id", ASCENDING)], unique=True)
  366. await ordersdb.create_index([("customer_id", ASCENDING), ("created_at", DESCENDING)])
  367. await ordersdb.create_index([("technician_id", ASCENDING), ("created_at", DESCENDING)])
  368. await ordersdb.create_index([("status", ASCENDING), ("updated_at", DESCENDING)])
  369. await quotesdb.create_index([("quote_id", ASCENDING)], unique=True)
  370. await quotesdb.create_index([("order_id", ASCENDING), ("version", DESCENDING)], unique=True)
  371. await addressesdb.create_index([("order_id", ASCENDING)], unique=True)
  372. await addressesdb.create_index([("redact_after", ASCENDING)])
  373. await eventsdb.create_index([("event_id", ASCENDING)], unique=True)
  374. await eventsdb.create_index([("order_id", ASCENDING), ("created_at", ASCENDING)])
  375. await customersdb.create_index([("user_id", ASCENDING)], unique=True)
  376. await blocksdb.create_index([("scope", ASCENDING), ("technician_id", ASCENDING), ("customer_id", ASCENDING)], unique=True)
  377. await reportsdb.create_index([("report_id", ASCENDING)], unique=True)
  378. await qrdb.create_index([("qr_id", ASCENDING)], unique=True)
  379. await qrdb.create_index([("token_hash", ASCENDING)], unique=True)
  380. await qrdb.create_index([("order_id", ASCENDING), ("created_at", DESCENDING)])
  381. await reviewtemplatesdb.create_index([("version", DESCENDING)], unique=True)
  382. await reviewdraftsdb.create_index([("customer_id", ASCENDING)], unique=True)
  383. await reviewdraftsdb.create_index([("expires_at", ASCENDING)], expireAfterSeconds=0)
  384. await reviewsdb.create_index([("review_id", ASCENDING)], unique=True)
  385. await reviewsdb.create_index([("qr_id", ASCENDING)], unique=True, sparse=True)
  386. await reviewsdb.create_index([("technician_id", ASCENDING), ("status", ASCENDING), ("created_at", DESCENDING)])
  387. await statsdb.create_index([("technician_id", ASCENDING), ("category", ASCENDING)], unique=True)
  388. await reviewtopicsdb.create_index([("technician_id", ASCENDING)], unique=True)
  389. await reviewtopicsdb.create_index([("forum_chat_id", ASCENDING), ("message_thread_id", ASCENDING)])
  390. _indexes_ready = True
  391. async def record_service_event(
  392. order_id: str,
  393. event_type: str,
  394. *,
  395. actor_id: int | str,
  396. reason: str = "",
  397. metadata: dict[str, Any] | None = None,
  398. ) -> dict[str, Any]:
  399. await ensure_service_indexes()
  400. document = {
  401. "event_id": uuid4().hex,
  402. "order_id": str(order_id),
  403. "event_type": str(event_type),
  404. "actor_id": actor_id,
  405. "reason": _clean_text(reason, max_length=500),
  406. "metadata": metadata or {},
  407. "created_at": utc_now(),
  408. }
  409. await eventsdb.insert_one(document)
  410. return document
  411. async def record_service_funnel_event(user_id: int, event_type: str) -> dict[str, Any]:
  412. if event_type not in {"directory_viewed", "technician_viewed", "request_started"}:
  413. raise ServiceDataError("invalid_funnel_event", "不支持该履约漏斗事件。")
  414. return await record_service_event(
  415. "",
  416. event_type,
  417. actor_id=int(user_id),
  418. metadata={"funnel": True},
  419. )
  420. async def observe_customer(user: Any, *, accepted_terms: bool = False) -> dict[str, Any]:
  421. await ensure_service_indexes()
  422. now = utc_now()
  423. display_name = " ".join(
  424. value for value in (str(getattr(user, "first_name", "") or "").strip(), str(getattr(user, "last_name", "") or "").strip()) if value
  425. ) or f"用户 {user.id}"
  426. updates: dict[str, Any] = {
  427. "username": getattr(user, "username", None),
  428. "display_name": display_name,
  429. "updated_at": now,
  430. }
  431. if accepted_terms:
  432. updates["terms_accepted_at"] = now
  433. await customersdb.update_one(
  434. {"user_id": int(user.id)},
  435. {"$set": updates, "$setOnInsert": {"created_at": now, "blocked": False}},
  436. upsert=True,
  437. )
  438. return await customersdb.find_one({"user_id": int(user.id)}) or {}
  439. async def require_customer_allowed(user_id: int) -> dict[str, Any]:
  440. await ensure_service_indexes()
  441. customer = await customersdb.find_one({"user_id": int(user_id)})
  442. if not customer or not customer.get("terms_accepted_at"):
  443. raise ServiceDataError("terms_required", "请先接受服务规则和隐私说明。")
  444. if customer.get("blocked"):
  445. raise ServiceDataError("customer_blocked", "你的服务功能已被暂停。")
  446. global_block = await blocksdb.find_one({"scope": "global", "customer_id": int(user_id), "active": True})
  447. if global_block:
  448. raise ServiceDataError("customer_blocked", "你的服务功能已被暂停。")
  449. return customer
  450. async def get_service_settings() -> dict[str, Any]:
  451. stored = await settingsdb.find_one({"settings_id": "global"}) or {}
  452. forum_chat_id = str(
  453. stored.get("service_review_forum_chat_id")
  454. or getattr(wbb, "SERVICE_REVIEW_FORUM_CHAT_ID", "")
  455. or ""
  456. ).strip()
  457. forum_username = str(
  458. stored.get("service_review_forum_username")
  459. or getattr(wbb, "SERVICE_REVIEW_FORUM_USERNAME", "")
  460. or ""
  461. ).strip().lstrip("@")
  462. return {
  463. "customer_max_open_orders": max(
  464. 1,
  465. min(
  466. 20,
  467. int(
  468. stored.get("service_customer_max_open_orders")
  469. or getattr(wbb, "SERVICE_CUSTOMER_MAX_OPEN_ORDERS", 3)
  470. or 3
  471. ),
  472. ),
  473. ),
  474. "customer_daily_request_limit": max(
  475. 1,
  476. min(
  477. 100,
  478. int(
  479. stored.get("service_customer_daily_request_limit")
  480. or getattr(wbb, "SERVICE_CUSTOMER_DAILY_REQUEST_LIMIT", 10)
  481. or 10
  482. ),
  483. ),
  484. ),
  485. "review_qr_expiry_hours": max(
  486. 1,
  487. min(
  488. 168,
  489. int(
  490. stored.get("service_review_qr_expiry_hours")
  491. or getattr(wbb, "SERVICE_REVIEW_QR_EXPIRY_HOURS", 24)
  492. or 24
  493. ),
  494. ),
  495. ),
  496. "review_forum_chat_id": forum_chat_id,
  497. "review_forum_username": forum_username,
  498. }
  499. async def set_service_settings(values: dict[str, Any]) -> dict[str, Any]:
  500. current = await get_service_settings()
  501. try:
  502. normalized = {
  503. "customer_max_open_orders": max(
  504. 1,
  505. min(20, int(values.get("customer_max_open_orders", current["customer_max_open_orders"]))),
  506. ),
  507. "customer_daily_request_limit": max(
  508. 1,
  509. min(
  510. 100,
  511. int(
  512. values.get(
  513. "customer_daily_request_limit",
  514. current["customer_daily_request_limit"],
  515. )
  516. ),
  517. ),
  518. ),
  519. "review_qr_expiry_hours": max(
  520. 1,
  521. min(168, int(values.get("review_qr_expiry_hours", current["review_qr_expiry_hours"]))),
  522. ),
  523. "review_forum_chat_id": _clean_text(
  524. values.get("review_forum_chat_id", current["review_forum_chat_id"]),
  525. max_length=40,
  526. ),
  527. "review_forum_username": _clean_text(
  528. values.get(
  529. "review_forum_username",
  530. current["review_forum_username"],
  531. ),
  532. max_length=64,
  533. ).lstrip("@"),
  534. }
  535. except (TypeError, ValueError) as exc:
  536. raise ServiceDataError("invalid_service_settings", "服务设置格式无效。") from exc
  537. await settingsdb.update_one(
  538. {"settings_id": "global"},
  539. {
  540. "$set": {
  541. "service_customer_max_open_orders": normalized["customer_max_open_orders"],
  542. "service_customer_daily_request_limit": normalized["customer_daily_request_limit"],
  543. "service_review_qr_expiry_hours": normalized["review_qr_expiry_hours"],
  544. "service_review_forum_chat_id": normalized["review_forum_chat_id"],
  545. "service_review_forum_username": normalized["review_forum_username"],
  546. "updated_at": utc_now(),
  547. },
  548. "$setOnInsert": {"created_at": utc_now()},
  549. },
  550. upsert=True,
  551. )
  552. return normalized
  553. async def list_package_templates(*, include_archived: bool = False) -> list[dict[str, Any]]:
  554. await ensure_service_indexes()
  555. filters = {} if include_archived else {"status": {"$ne": "archived"}}
  556. return await templatesdb.find(filters).sort([("sort_order", ASCENDING), ("created_at", ASCENDING)]).to_list(length=500)
  557. async def ensure_builtin_package_templates() -> list[dict[str, Any]]:
  558. """Seed editable starter templates once without overwriting administrator changes."""
  559. await ensure_service_indexes()
  560. now = utc_now()
  561. template_ids = []
  562. for item in BUILTIN_PACKAGE_TEMPLATES:
  563. template_id = str(item["template_id"])
  564. template_ids.append(template_id)
  565. package = normalize_package(
  566. item["package"],
  567. package_id=str(item["package"]["package_id"]),
  568. )
  569. package["source_type"] = "template"
  570. await templatesdb.update_one(
  571. {"template_id": template_id},
  572. {
  573. "$setOnInsert": {
  574. "template_id": template_id,
  575. "version": 1,
  576. "name": item["name"],
  577. "admin_note": "系统内置快速模板,可在后台编辑、停用或归档。",
  578. "package": package,
  579. "status": "enabled",
  580. "sort_order": int(item["sort_order"]),
  581. "builtin": True,
  582. "created_by": "system",
  583. "created_at": now,
  584. "updated_at": now,
  585. }
  586. },
  587. upsert=True,
  588. )
  589. return await templatesdb.find({"template_id": {"$in": template_ids}}).sort(
  590. [("sort_order", ASCENDING)]
  591. ).to_list(length=len(template_ids))
  592. async def create_package_template(values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
  593. await ensure_service_indexes()
  594. now = utc_now()
  595. package = normalize_package(values.get("package") or values)
  596. package["source_type"] = "template"
  597. document = {
  598. "template_id": uuid4().hex,
  599. "version": 1,
  600. "name": _clean_text(values.get("template_name") or package["name"], max_length=80, required=True),
  601. "admin_note": _clean_text(values.get("admin_note"), max_length=500),
  602. "package": package,
  603. "status": "enabled" if values.get("enabled", True) else "disabled",
  604. "sort_order": int(values.get("sort_order") or 0),
  605. "created_by": str(actor_id),
  606. "created_at": now,
  607. "updated_at": now,
  608. }
  609. await templatesdb.insert_one(document)
  610. return document
  611. async def update_package_template(template_id: str, values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
  612. current = await templatesdb.find_one({"template_id": str(template_id)})
  613. if not current:
  614. raise ServiceDataError("template_not_found", "未找到套餐模板。")
  615. package = normalize_package(values.get("package") or values, package_id=current["package"]["package_id"])
  616. package["source_type"] = "template"
  617. updated = await templatesdb.find_one_and_update(
  618. {"template_id": str(template_id)},
  619. {
  620. "$set": {
  621. "name": _clean_text(values.get("template_name") or package["name"], max_length=80, required=True),
  622. "admin_note": _clean_text(values.get("admin_note"), max_length=500),
  623. "package": package,
  624. "sort_order": int(values.get("sort_order") or 0),
  625. "updated_by": str(actor_id),
  626. "updated_at": utc_now(),
  627. },
  628. "$inc": {"version": 1},
  629. },
  630. return_document=ReturnDocument.AFTER,
  631. )
  632. return updated or {}
  633. async def set_package_template_status(template_id: str, action: str, *, actor_id: str) -> dict[str, Any]:
  634. status_map = {"enable": "enabled", "disable": "disabled", "archive": "archived"}
  635. if action not in status_map:
  636. raise ServiceDataError("invalid_template_action", "不支持该模板操作。")
  637. updated = await templatesdb.find_one_and_update(
  638. {"template_id": str(template_id)},
  639. {"$set": {"status": status_map[action], "updated_by": str(actor_id), "updated_at": utc_now()}},
  640. return_document=ReturnDocument.AFTER,
  641. )
  642. if not updated:
  643. raise ServiceDataError("template_not_found", "未找到套餐模板。")
  644. return updated
  645. async def duplicate_package_template(template_id: str, *, actor_id: str) -> dict[str, Any]:
  646. current = await templatesdb.find_one({"template_id": str(template_id)})
  647. if not current:
  648. raise ServiceDataError("template_not_found", "未找到套餐模板。")
  649. values = {
  650. "template_name": f"{current['name']}(副本)",
  651. "admin_note": current.get("admin_note", ""),
  652. "sort_order": current.get("sort_order", 0),
  653. "enabled": False,
  654. "package": current["package"],
  655. }
  656. return await create_package_template(values, actor_id=actor_id)
  657. async def publish_technician_profile(user_id: int, values: dict[str, Any], *, actor_id: int | str) -> dict[str, Any]:
  658. await ensure_service_indexes()
  659. profile = await profilesdb.find_one({"user_id": int(user_id)})
  660. if not profile or profile.get("application_status") != "approved":
  661. raise ServiceDataError("technician_required", "只有已认证技师可以发布服务资料。")
  662. if not profile.get("username"):
  663. raise ServiceDataError("username_required", "技师必须设置有效的 Telegram 用户名。")
  664. if not await membershipsdb.find_one({"user_id": int(user_id), "active": True}):
  665. raise ServiceDataError("technician_membership_required", "技师必须仍是受管群当前成员。")
  666. normalized = _normalize_profile(values)
  667. now = utc_now()
  668. normalized.update({"updated_at": now, "updated_by": actor_id, "version": int((profile.get("service_profile") or {}).get("version") or 0) + 1})
  669. await profilesdb.update_one(
  670. {"user_id": int(user_id)},
  671. {"$set": {"service_profile": normalized, "updated_at": now}},
  672. )
  673. return await profilesdb.find_one({"user_id": int(user_id)}) or {}
  674. async def get_technician_service_profile(user_id: int) -> dict[str, Any]:
  675. profile = await profilesdb.find_one({"user_id": int(user_id)})
  676. if not profile:
  677. raise ServiceDataError("technician_not_found", "未找到技师资料。")
  678. return profile
  679. async def get_technician_self_service_context(user_id: int) -> dict[str, Any]:
  680. profile = await profilesdb.find_one({"user_id": int(user_id)})
  681. if not profile:
  682. raise ServiceDataError("technician_not_found", "未找到技师申请资料。")
  683. membership_active = bool(
  684. await membershipsdb.find_one({"user_id": int(user_id), "active": True})
  685. )
  686. issues = []
  687. if profile.get("application_status") != "approved":
  688. issues.append({"code": "approval_required", "message": "技师认证尚未通过。"})
  689. if not profile.get("username"):
  690. issues.append(
  691. {
  692. "code": "username_required",
  693. "message": "请先在 Telegram 设置用户名。",
  694. }
  695. )
  696. if not membership_active:
  697. issues.append(
  698. {
  699. "code": "membership_required",
  700. "message": "当前不在受管群成员名单中。",
  701. }
  702. )
  703. return {
  704. "user_id": int(user_id),
  705. "display_name": profile.get("display_name") or f"技师 {user_id}",
  706. "username": profile.get("username"),
  707. "eligibility": {"ready": not issues, "issues": issues},
  708. "service_profile": profile.get("service_profile") or None,
  709. "review_topic": await get_technician_review_topic(int(user_id)),
  710. }
  711. async def set_technician_accepting_requests(
  712. user_id: int,
  713. accepting_requests: bool,
  714. ) -> dict[str, Any]:
  715. profile = await get_technician_service_profile(user_id)
  716. values = dict(profile.get("service_profile") or {})
  717. if not values.get("is_complete"):
  718. raise ServiceDataError("technician_profile_required", "请先发布至少一个服务套餐。")
  719. values["accepting_requests"] = bool(accepting_requests)
  720. return await publish_technician_profile(
  721. user_id,
  722. values,
  723. actor_id=user_id,
  724. )
  725. async def instantiate_package_template(template_id: str) -> dict[str, Any]:
  726. template = await templatesdb.find_one(
  727. {"template_id": str(template_id), "status": "enabled"}
  728. )
  729. if not template:
  730. raise ServiceDataError("template_not_found", "套餐模板不存在或已停用。")
  731. package = dict(template["package"])
  732. package.update(
  733. {
  734. "package_id": uuid4().hex,
  735. "source_type": "template",
  736. "source_template_id": template["template_id"],
  737. "source_template_version": template["version"],
  738. "customized_from_template": False,
  739. }
  740. )
  741. return package
  742. def _fernet() -> Fernet:
  743. raw = str(getattr(wbb, "SERVICE_ADDRESS_ENCRYPTION_KEY", "") or "").strip()
  744. if not raw:
  745. raise ServiceDataError("address_encryption_required", "服务地址加密密钥尚未配置。")
  746. try:
  747. return Fernet(raw.encode())
  748. except (ValueError, TypeError) as exc:
  749. raise ServiceDataError("invalid_address_encryption_key", "服务地址加密密钥无效。") from exc
  750. def _encrypt_address(payload: dict[str, Any]) -> str:
  751. return _fernet().encrypt(json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode()).decode()
  752. def _decrypt_address(token: str) -> dict[str, Any]:
  753. try:
  754. return json.loads(_fernet().decrypt(str(token).encode()).decode())
  755. except (InvalidToken, ValueError, json.JSONDecodeError) as exc:
  756. raise ServiceDataError("address_unavailable", "精确地址无法解密。") from exc
  757. async def _schedule_address_redaction(order_id: str, closed_at: datetime) -> None:
  758. retention = max(
  759. 1,
  760. min(365, int(getattr(wbb, "SERVICE_ADDRESS_RETENTION_DAYS", 7) or 7)),
  761. )
  762. await addressesdb.update_one(
  763. {"order_id": str(order_id), "redacted": False},
  764. {
  765. "$set": {
  766. "redact_after": closed_at + timedelta(days=retention),
  767. "updated_at": closed_at,
  768. }
  769. },
  770. )
  771. def _distance_meters(lon1: float, lat1: float, lon2: float, lat2: float) -> float:
  772. radius = 6_371_000.0
  773. phi1, phi2 = math.radians(lat1), math.radians(lat2)
  774. delta_phi = math.radians(lat2 - lat1)
  775. delta_lambda = math.radians(lon2 - lon1)
  776. a = math.sin(delta_phi / 2) ** 2 + math.cos(phi1) * math.cos(phi2) * math.sin(delta_lambda / 2) ** 2
  777. return radius * 2 * math.atan2(math.sqrt(a), math.sqrt(1 - a))
  778. def distance_band(distance_meters: float) -> str:
  779. km = max(0.0, distance_meters / 1000)
  780. for threshold in (1, 3, 5, 10, 20, 50):
  781. if km <= threshold:
  782. return f"{threshold} 公里内"
  783. return "50 公里以上"
  784. async def _get_package(technician_id: int, package_id: str) -> tuple[dict[str, Any], dict[str, Any]]:
  785. profile = await profilesdb.find_one({"user_id": int(technician_id)})
  786. if not profile or profile.get("application_status") != "approved":
  787. raise ServiceDataError("technician_not_found", "未找到已认证技师。")
  788. if not profile.get("listed") or not profile.get("username"):
  789. raise ServiceDataError("technician_unavailable", "该技师暂未公开接单。")
  790. if not await membershipsdb.find_one({"user_id": int(technician_id), "active": True}):
  791. raise ServiceDataError("technician_unavailable", "该技师当前不满足接单资格。")
  792. service_profile = profile.get("service_profile") or {}
  793. if not service_profile.get("is_complete") or not service_profile.get("accepting_requests", True):
  794. raise ServiceDataError("technician_unavailable", "该技师暂未接单。")
  795. package = next((item for item in service_profile.get("packages", []) if str(item.get("package_id")) == str(package_id)), None)
  796. if not package:
  797. raise ServiceDataError("package_not_found", "未找到该服务套餐。")
  798. return profile, package
  799. async def list_available_technicians(
  800. *,
  801. longitude: float | None = None,
  802. latitude: float | None = None,
  803. max_distance_meters: float | None = None,
  804. query: str = "",
  805. page: int = 1,
  806. page_size: int = 10,
  807. ) -> tuple[list[dict[str, Any]], int]:
  808. use_distance = longitude is not None and latitude is not None
  809. lon = float(longitude) if longitude is not None else 0.0
  810. lat = float(latitude) if latitude is not None else 0.0
  811. if use_distance and (not -180 <= lon <= 180 or not -90 <= lat <= 90):
  812. raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
  813. filters: dict[str, Any] = {
  814. "application_status": "approved",
  815. "listed": True,
  816. "username": {"$nin": [None, ""]},
  817. "service_profile.is_complete": True,
  818. "service_profile.accepting_requests": True,
  819. }
  820. if query:
  821. pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
  822. filters["$or"] = [
  823. {"username": pattern},
  824. {"display_name": pattern},
  825. {"service_profile.tags": pattern},
  826. {"service_profile.packages.category": pattern},
  827. ]
  828. active_members = {
  829. int(item["user_id"])
  830. async for item in membershipsdb.find({"active": True}, {"user_id": 1})
  831. }
  832. profiles = {
  833. int(item["user_id"]): item
  834. async for item in profilesdb.find(filters)
  835. if int(item["user_id"]) in active_members
  836. }
  837. values: list[dict[str, Any]] = []
  838. locations = {}
  839. if profiles and use_distance:
  840. locations = {
  841. int(location["user_id"]): location
  842. async for location in locationsdb.find(
  843. {"user_id": {"$in": list(profiles)}}
  844. )
  845. }
  846. for technician_id, profile in profiles.items():
  847. service_profile = profile.get("service_profile") or {}
  848. value = {
  849. "user_id": int(profile["user_id"]),
  850. "username": profile.get("username"),
  851. "display_name": profile.get("display_name"),
  852. "headline": service_profile.get("headline"),
  853. "bio": service_profile.get("bio"),
  854. "tags": service_profile.get("tags", []),
  855. "public_area_text": service_profile.get("public_area_text"),
  856. "service_modes": service_profile.get("service_modes", []),
  857. "packages": service_profile.get("packages", []),
  858. "distance_meters": None,
  859. "distance_band": "",
  860. }
  861. if use_distance:
  862. location = locations.get(technician_id)
  863. if not location:
  864. continue
  865. point = location.get("point", {}).get("coordinates") or [
  866. location.get("longitude"),
  867. location.get("latitude"),
  868. ]
  869. if len(point) != 2 or point[0] is None or point[1] is None:
  870. continue
  871. distance = _distance_meters(lon, lat, float(point[0]), float(point[1]))
  872. if max_distance_meters is not None and distance > float(max_distance_meters):
  873. continue
  874. value.update(
  875. {
  876. "distance_meters": round(distance, 2),
  877. "distance_band": distance_band(distance),
  878. }
  879. )
  880. values.append(value)
  881. if use_distance:
  882. values.sort(key=lambda item: (item["distance_meters"], item["user_id"]))
  883. else:
  884. values.sort(
  885. key=lambda item: (
  886. str(item.get("public_area_text") or "").casefold(),
  887. str(item.get("display_name") or "").casefold(),
  888. item["user_id"],
  889. )
  890. )
  891. total = len(values)
  892. start = (max(1, page) - 1) * max(1, page_size)
  893. return values[start : start + max(1, page_size)], total
  894. async def create_service_order(
  895. *,
  896. customer: Any,
  897. technician_id: int,
  898. package_id: str,
  899. service_mode: str,
  900. requirements: str,
  901. longitude: float | None = None,
  902. latitude: float | None = None,
  903. address_text: str = "",
  904. ) -> dict[str, Any]:
  905. await ensure_service_indexes()
  906. await observe_customer(customer)
  907. await require_customer_allowed(int(customer.id))
  908. if int(customer.id) == int(technician_id):
  909. raise ServiceDataError("self_service_not_allowed", "不能向自己发起服务请求。")
  910. personal_block = await blocksdb.find_one({"scope": "technician", "technician_id": int(technician_id), "customer_id": int(customer.id), "active": True})
  911. if personal_block:
  912. raise ServiceDataError("customer_blocked", "该技师暂不接受你的服务请求。")
  913. service_settings = await get_service_settings()
  914. max_open = service_settings["customer_max_open_orders"]
  915. open_count = await ordersdb.count_documents({"customer_id": int(customer.id), "status": {"$in": list(ORDER_ACTIVE_STATUSES)}})
  916. if open_count >= max_open:
  917. raise ServiceDataError("too_many_open_orders", f"最多同时保留 {max_open} 个未关闭服务单。")
  918. day_start = utc_now().replace(hour=0, minute=0, second=0, microsecond=0)
  919. daily_limit = service_settings["customer_daily_request_limit"]
  920. if await ordersdb.count_documents({"customer_id": int(customer.id), "created_at": {"$gte": day_start}}) >= daily_limit:
  921. raise ServiceDataError("daily_request_limit", "今天发起的服务请求已达到上限。")
  922. profile, package = await _get_package(int(technician_id), str(package_id))
  923. if service_mode not in package["service_modes"]:
  924. raise ServiceDataError("service_mode_unavailable", "该套餐不支持所选服务方式。")
  925. technician_location = await locationsdb.find_one({"user_id": int(technician_id)})
  926. if not technician_location:
  927. raise ServiceDataError("technician_location_required", "技师尚未配置服务位置。")
  928. address_payload: dict[str, Any] | None = None
  929. distance = 0.0
  930. if service_mode == "onsite":
  931. if longitude is None or latitude is None:
  932. raise ServiceDataError("customer_location_required", "上门服务必须提供位置。")
  933. lon, lat = float(longitude), float(latitude)
  934. if not -180 <= lon <= 180 or not -90 <= lat <= 90:
  935. raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
  936. address = _clean_text(address_text, max_length=300, required=True)
  937. point = technician_location.get("point", {}).get("coordinates") or [technician_location.get("longitude"), technician_location.get("latitude")]
  938. distance = _distance_meters(float(point[0]), float(point[1]), lon, lat)
  939. radius = float(package.get("service_radius_km") or 0)
  940. if radius and distance > radius * 1000:
  941. raise ServiceDataError("outside_service_radius", "顾客位置超出该套餐的上门服务范围。")
  942. address_payload = {"kind": "customer", "longitude": lon, "latitude": lat, "address_text": address}
  943. else:
  944. if longitude is None or latitude is None:
  945. raise ServiceDataError(
  946. "customer_location_required", "到店服务必须先通过附近查找提供位置。"
  947. )
  948. lon, lat = float(longitude), float(latitude)
  949. if not -180 <= lon <= 180 or not -90 <= lat <= 90:
  950. raise ServiceDataError("invalid_coordinates", "位置坐标无效。")
  951. point = technician_location.get("point", {}).get("coordinates") or [
  952. technician_location.get("longitude"),
  953. technician_location.get("latitude"),
  954. ]
  955. distance = _distance_meters(float(point[0]), float(point[1]), lon, lat)
  956. venue = (profile.get("service_profile") or {}).get("venue") or {}
  957. address_payload = {
  958. "kind": "venue",
  959. "longitude": technician_location.get("longitude"),
  960. "latitude": technician_location.get("latitude"),
  961. "address_text": _clean_text(venue.get("address_hint"), max_length=300, required=True),
  962. }
  963. now = utc_now()
  964. order_id = uuid4().hex
  965. order = {
  966. "order_id": order_id,
  967. "customer_id": int(customer.id),
  968. "customer_name": _clean_text(getattr(customer, "first_name", ""), max_length=80) or f"用户 {customer.id}",
  969. "technician_id": int(technician_id),
  970. "technician_name": profile.get("display_name") or f"技师 {technician_id}",
  971. "package_id": str(package_id),
  972. "package_snapshot": package,
  973. "category": package["category"],
  974. "service_mode": service_mode,
  975. "requirements": _clean_text(requirements, max_length=800, required=True),
  976. "distance_meters": round(distance, 2),
  977. "distance_band": distance_band(distance),
  978. "status": "requested",
  979. "version": 1,
  980. "created_at": now,
  981. "updated_at": now,
  982. }
  983. encrypted_address = _encrypt_address(address_payload)
  984. await ordersdb.insert_one(order)
  985. try:
  986. await addressesdb.insert_one(
  987. {
  988. "order_id": order_id,
  989. "encrypted_payload": encrypted_address,
  990. "redacted": False,
  991. "redact_after": None,
  992. "created_at": now,
  993. "updated_at": now,
  994. }
  995. )
  996. except Exception:
  997. await ordersdb.delete_one({"order_id": order_id, "status": "requested"})
  998. raise
  999. await record_service_event(order_id, "order_requested", actor_id=int(customer.id))
  1000. return order
  1001. def normalize_quote(value: dict[str, Any], *, package: dict[str, Any]) -> dict[str, Any]:
  1002. currency = str(value.get("currency") or package.get("currency") or "CNY").upper()
  1003. if currency != str(package.get("currency") or currency):
  1004. raise ServiceDataError("currency_mismatch", "最终报价币种必须与套餐一致。")
  1005. base = _positive_decimal(value.get("base_amount"), "服务金额")
  1006. travel = _positive_decimal(value.get("travel_fee"), "交通费")
  1007. discount = _positive_decimal(value.get("discount"), "优惠金额")
  1008. addons = []
  1009. addon_total = Decimal("0")
  1010. for item in value.get("addons", []) or []:
  1011. quantity = max(1, min(999, int(item.get("quantity") or 1)))
  1012. unit_price = Decimal(_positive_decimal(item.get("unit_price"), "加项单价"))
  1013. addons.append({"name": _clean_text(item.get("name"), max_length=60, required=True), "quantity": quantity, "unit_price": format(unit_price, "f")})
  1014. addon_total += unit_price * quantity
  1015. total = Decimal(base) + Decimal(travel) + addon_total - Decimal(discount)
  1016. if total < 0:
  1017. raise ServiceDataError("invalid_total", "最终报价总额不能小于 0。")
  1018. provided = value.get("total_amount")
  1019. if provided not in (None, "") and Decimal(_positive_decimal(provided, "总价")) != total:
  1020. raise ServiceDataError("quote_total_mismatch", "报价明细与总价不一致。")
  1021. return {
  1022. "currency": currency,
  1023. "base_amount": base,
  1024. "travel_fee": travel,
  1025. "addons": addons,
  1026. "discount": discount,
  1027. "total_amount": format(total.quantize(Decimal("0.01")), "f"),
  1028. "note": _clean_text(value.get("note"), max_length=500),
  1029. }
  1030. async def submit_service_quote(order_id: str, technician_id: int, values: dict[str, Any]) -> tuple[dict[str, Any], dict[str, Any]]:
  1031. await ensure_service_indexes()
  1032. order = await ordersdb.find_one({"order_id": str(order_id)})
  1033. if not order or int(order["technician_id"]) != int(technician_id):
  1034. raise ServiceDataError("order_not_found", "未找到可报价服务单。")
  1035. if order["status"] not in {"requested", "quoted"}:
  1036. raise ServiceDataError("invalid_order_state", "当前服务单不能修改报价。")
  1037. quote = normalize_quote(values, package=order["package_snapshot"])
  1038. version = int(order.get("quote_version") or 0) + 1
  1039. quote.update({"quote_id": uuid4().hex, "order_id": str(order_id), "version": version, "created_by": int(technician_id), "created_at": utc_now()})
  1040. try:
  1041. await quotesdb.insert_one(quote)
  1042. except DuplicateKeyError as exc:
  1043. raise ServiceDataError(
  1044. "concurrent_order_update", "服务单状态已变化,请刷新后重试。"
  1045. ) from exc
  1046. updated = await ordersdb.find_one_and_update(
  1047. {
  1048. "order_id": str(order_id),
  1049. "status": {"$in": ["requested", "quoted"]},
  1050. "version": order["version"],
  1051. },
  1052. {
  1053. "$set": {
  1054. "status": "quoted",
  1055. "current_quote_id": quote["quote_id"],
  1056. "quote_snapshot": quote,
  1057. "quote_version": version,
  1058. "updated_at": utc_now(),
  1059. },
  1060. "$inc": {"version": 1},
  1061. },
  1062. return_document=ReturnDocument.AFTER,
  1063. )
  1064. if not updated:
  1065. await quotesdb.delete_one({"quote_id": quote["quote_id"]})
  1066. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  1067. await record_service_event(order_id, "quote_submitted", actor_id=int(technician_id), metadata={"quote_id": quote["quote_id"], "version": version})
  1068. return updated, quote
  1069. async def confirm_service_quote(order_id: str, customer_id: int) -> dict[str, Any]:
  1070. order = await ordersdb.find_one_and_update(
  1071. {"order_id": str(order_id), "customer_id": int(customer_id), "status": "quoted"},
  1072. {"$set": {"status": "confirmed", "confirmed_at": utc_now(), "updated_at": utc_now()}, "$inc": {"version": 1}},
  1073. return_document=ReturnDocument.AFTER,
  1074. )
  1075. if not order:
  1076. raise ServiceDataError("invalid_order_state", "当前报价无法确认。")
  1077. await record_service_event(order_id, "quote_confirmed", actor_id=int(customer_id))
  1078. return order
  1079. async def reject_service_order(order_id: str, technician_id: int, reason: str) -> dict[str, Any]:
  1080. now = utc_now()
  1081. order = await ordersdb.find_one_and_update(
  1082. {"order_id": str(order_id), "technician_id": int(technician_id), "status": {"$in": ["requested", "quoted"]}},
  1083. {"$set": {"status": "rejected", "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
  1084. return_document=ReturnDocument.AFTER,
  1085. )
  1086. if not order:
  1087. raise ServiceDataError("invalid_order_state", "当前服务单无法拒绝。")
  1088. await _schedule_address_redaction(order_id, now)
  1089. await record_service_event(order_id, "order_rejected", actor_id=int(technician_id), reason=reason)
  1090. return order
  1091. async def start_service_order(order_id: str, technician_id: int) -> dict[str, Any]:
  1092. order = await ordersdb.find_one_and_update(
  1093. {"order_id": str(order_id), "technician_id": int(technician_id), "status": "confirmed"},
  1094. {"$set": {"status": "in_progress", "started_at": utc_now(), "updated_at": utc_now()}, "$inc": {"version": 1}},
  1095. return_document=ReturnDocument.AFTER,
  1096. )
  1097. if not order:
  1098. raise ServiceDataError("invalid_order_state", "只有已确认服务单可以开始。")
  1099. await record_service_event(order_id, "service_started", actor_id=int(technician_id))
  1100. return order
  1101. async def cancel_service_order(order_id: str, actor_id: int, reason: str) -> dict[str, Any]:
  1102. order = await ordersdb.find_one({"order_id": str(order_id)})
  1103. if not order or int(actor_id) not in {int(order["customer_id"]), int(order["technician_id"])}:
  1104. raise ServiceDataError("order_not_found", "未找到该服务单。")
  1105. if order["status"] not in ORDER_ACTIVE_STATUSES - {"disputed"}:
  1106. raise ServiceDataError("invalid_order_state", "当前服务单不能取消。")
  1107. status = "canceled_customer" if int(actor_id) == int(order["customer_id"]) else "canceled_technician"
  1108. now = utc_now()
  1109. updated = await ordersdb.find_one_and_update(
  1110. {"order_id": str(order_id), "status": order["status"], "version": order["version"]},
  1111. {"$set": {"status": status, "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
  1112. return_document=ReturnDocument.AFTER,
  1113. )
  1114. if not updated:
  1115. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  1116. await _schedule_address_redaction(order_id, now)
  1117. await qrdb.update_many({"order_id": str(order_id), "status": "issued"}, {"$set": {"status": "revoked", "revoked_at": now, "revoke_reason": "服务单已取消"}})
  1118. await record_service_event(order_id, "order_canceled", actor_id=int(actor_id), reason=reason, metadata={"status": status})
  1119. return updated
  1120. async def dispute_service_order(order_id: str, actor_id: int, reason: str) -> dict[str, Any]:
  1121. current = await ordersdb.find_one(
  1122. {
  1123. "order_id": str(order_id),
  1124. "status": {"$in": list(ORDER_ACTIVE_STATUSES - {"disputed"})},
  1125. "$or": [{"customer_id": int(actor_id)}, {"technician_id": int(actor_id)}],
  1126. }
  1127. )
  1128. if not current:
  1129. raise ServiceDataError("invalid_order_state", "当前服务单不能发起争议。")
  1130. order = await ordersdb.find_one_and_update(
  1131. {"order_id": str(order_id), "status": current["status"], "version": current["version"]},
  1132. {"$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}},
  1133. return_document=ReturnDocument.AFTER,
  1134. )
  1135. if not order:
  1136. raise ServiceDataError("invalid_order_state", "当前服务单不能发起争议。")
  1137. await qrdb.update_many({"order_id": str(order_id), "status": "issued"}, {"$set": {"status": "revoked", "revoked_at": utc_now(), "revoke_reason": "服务单争议"}})
  1138. await record_service_event(order_id, "order_disputed", actor_id=int(actor_id), reason=reason)
  1139. return order
  1140. async def resolve_service_dispute(
  1141. order_id: str,
  1142. *,
  1143. actor_id: str,
  1144. action: str,
  1145. reason: str,
  1146. ) -> dict[str, Any]:
  1147. if action == "void":
  1148. return await admin_void_service_order(order_id, actor_id=actor_id, reason=reason)
  1149. if action != "resume":
  1150. raise ServiceDataError("invalid_dispute_action", "不支持该争议处理操作。")
  1151. current = await ordersdb.find_one({"order_id": str(order_id), "status": "disputed"})
  1152. if not current:
  1153. raise ServiceDataError("invalid_order_state", "该服务单当前不在争议中。")
  1154. resume_status = str(current.get("status_before_dispute") or "confirmed")
  1155. if resume_status == "completion_pending":
  1156. resume_status = "in_progress"
  1157. if resume_status not in {"requested", "quoted", "confirmed", "in_progress"}:
  1158. resume_status = "confirmed"
  1159. updated = await ordersdb.find_one_and_update(
  1160. {"order_id": str(order_id), "status": "disputed", "version": current["version"]},
  1161. {
  1162. "$set": {
  1163. "status": resume_status,
  1164. "dispute_resolved_at": utc_now(),
  1165. "dispute_resolution_reason": _clean_text(reason, max_length=500, required=True),
  1166. "updated_at": utc_now(),
  1167. },
  1168. "$unset": {"status_before_dispute": ""},
  1169. "$inc": {"version": 1},
  1170. },
  1171. return_document=ReturnDocument.AFTER,
  1172. )
  1173. if not updated:
  1174. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  1175. await record_service_event(order_id, "dispute_resolved", actor_id=actor_id, reason=reason, metadata={"status": resume_status})
  1176. return updated
  1177. async def admin_void_service_order(order_id: str, *, actor_id: str, reason: str) -> dict[str, Any]:
  1178. now = utc_now()
  1179. order = await ordersdb.find_one_and_update(
  1180. {"order_id": str(order_id), "status": {"$ne": "voided"}},
  1181. {"$set": {"status": "voided", "closed_at": now, "updated_at": now, "close_reason": _clean_text(reason, max_length=500, required=True)}, "$inc": {"version": 1}},
  1182. return_document=ReturnDocument.AFTER,
  1183. )
  1184. if not order:
  1185. raise ServiceDataError("order_not_found", "未找到可作废服务单。")
  1186. await _schedule_address_redaction(order_id, now)
  1187. await qrdb.update_many({"order_id": str(order_id), "status": {"$ne": "voided"}}, {"$set": {"status": "voided", "revoked_at": now}})
  1188. await reviewsdb.update_many({"order_id": str(order_id), "status": "approved"}, {"$set": {"status": "voided", "moderated_at": now, "moderation_reason": reason, "updated_at": now}})
  1189. await recompute_review_stats()
  1190. await record_service_event(order_id, "order_voided", actor_id=actor_id, reason=reason)
  1191. return order
  1192. async def issue_review_qr(order_id: str, technician_id: int) -> tuple[dict[str, Any], str]:
  1193. await ensure_service_indexes()
  1194. order = await ordersdb.find_one({"order_id": str(order_id), "technician_id": int(technician_id)})
  1195. if order and order.get("status") == "completion_pending":
  1196. previous = await qrdb.find_one(
  1197. {"order_id": str(order_id), "status": "issued"},
  1198. sort=[("created_at", DESCENDING)],
  1199. )
  1200. if previous and _as_utc(previous["expires_at"]) <= utc_now():
  1201. now = utc_now()
  1202. expired = await qrdb.find_one_and_update(
  1203. {"qr_id": previous["qr_id"], "status": "issued"},
  1204. {"$set": {"status": "expired", "expired_at": now}},
  1205. return_document=ReturnDocument.AFTER,
  1206. )
  1207. if expired:
  1208. await ordersdb.update_one(
  1209. {
  1210. "order_id": str(order_id),
  1211. "status": "completion_pending",
  1212. "current_review_qr_id": previous["qr_id"],
  1213. },
  1214. {
  1215. "$set": {"status": "in_progress", "updated_at": now},
  1216. "$unset": {"current_review_qr_id": ""},
  1217. "$inc": {"version": 1},
  1218. },
  1219. )
  1220. order = await ordersdb.find_one(
  1221. {"order_id": str(order_id), "technician_id": int(technician_id)}
  1222. )
  1223. if not order or order.get("status") != "in_progress":
  1224. raise ServiceDataError("invalid_order_state", "只有服务中的订单可以生成评价二维码。")
  1225. if await qrdb.find_one({"order_id": str(order_id), "status": {"$in": ["issued", "claimed"]}}):
  1226. raise ServiceDataError("review_qr_exists", "该服务单已经生成评价二维码。")
  1227. token = token_urlsafe(18)
  1228. now = utc_now()
  1229. hours = (await get_service_settings())["review_qr_expiry_hours"]
  1230. document = {
  1231. "qr_id": uuid4().hex,
  1232. "order_id": str(order_id),
  1233. "technician_id": int(technician_id),
  1234. "customer_id": int(order["customer_id"]),
  1235. "package_snapshot": order["package_snapshot"],
  1236. "quote_snapshot": order.get("quote_snapshot"),
  1237. "token_hash": hashlib.sha256(token.encode()).hexdigest(),
  1238. "status": "issued",
  1239. "expires_at": now + timedelta(hours=hours),
  1240. "created_at": now,
  1241. }
  1242. try:
  1243. await qrdb.insert_one(document)
  1244. except DuplicateKeyError as exc:
  1245. raise ServiceDataError("review_qr_exists", "该服务单已经生成评价二维码。") from exc
  1246. updated = await ordersdb.find_one_and_update(
  1247. {
  1248. "order_id": str(order_id),
  1249. "status": "in_progress",
  1250. "version": order["version"],
  1251. },
  1252. {
  1253. "$set": {
  1254. "status": "completion_pending",
  1255. "current_review_qr_id": document["qr_id"],
  1256. "completion_requested_at": now,
  1257. "updated_at": now,
  1258. },
  1259. "$inc": {"version": 1},
  1260. },
  1261. return_document=ReturnDocument.AFTER,
  1262. )
  1263. if not updated:
  1264. await qrdb.update_one(
  1265. {"qr_id": document["qr_id"], "status": "issued"},
  1266. {
  1267. "$set": {
  1268. "status": "voided",
  1269. "revoked_at": utc_now(),
  1270. "revoke_reason": "服务单状态已变化",
  1271. }
  1272. },
  1273. )
  1274. raise ServiceDataError("concurrent_order_update", "服务单状态已变化,请刷新后重试。")
  1275. await record_service_event(order_id, "review_qr_issued", actor_id=int(technician_id), metadata={"qr_id": document["qr_id"]})
  1276. return document, token
  1277. async def preview_review_qr(token: str, customer_id: int) -> dict[str, Any]:
  1278. document = await qrdb.find_one({"token_hash": hashlib.sha256(str(token).encode()).hexdigest()})
  1279. if not document:
  1280. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1281. if int(document["customer_id"]) != int(customer_id):
  1282. raise ServiceDataError("review_qr_customer_mismatch", "该二维码不属于当前顾客。")
  1283. if document.get("status") == "claimed":
  1284. if await reviewsdb.find_one({"qr_id": document["qr_id"]}):
  1285. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1286. return document
  1287. if document.get("status") != "issued":
  1288. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1289. if _as_utc(document["expires_at"]) <= utc_now():
  1290. now = utc_now()
  1291. expired = await qrdb.find_one_and_update(
  1292. {"qr_id": document["qr_id"], "status": "issued"},
  1293. {"$set": {"status": "expired", "expired_at": now}},
  1294. return_document=ReturnDocument.AFTER,
  1295. )
  1296. if expired:
  1297. await ordersdb.update_one(
  1298. {
  1299. "order_id": document["order_id"],
  1300. "status": "completion_pending",
  1301. "current_review_qr_id": document["qr_id"],
  1302. },
  1303. {
  1304. "$set": {"status": "in_progress", "updated_at": now},
  1305. "$unset": {"current_review_qr_id": ""},
  1306. "$inc": {"version": 1},
  1307. },
  1308. )
  1309. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1310. return document
  1311. async def claim_review_qr(token: str, customer_id: int) -> tuple[dict[str, Any], dict[str, Any]]:
  1312. token_hash = hashlib.sha256(str(token).encode()).hexdigest()
  1313. now = utc_now()
  1314. await preview_review_qr(token, customer_id)
  1315. document = await qrdb.find_one_and_update(
  1316. {"token_hash": token_hash, "customer_id": int(customer_id), "status": "issued"},
  1317. {"$set": {"status": "claimed", "claimed_at": now}},
  1318. return_document=ReturnDocument.AFTER,
  1319. )
  1320. if not document:
  1321. raise ServiceDataError("review_qr_unavailable", "评价二维码无效、已使用或已过期。")
  1322. order = await ordersdb.find_one_and_update(
  1323. {"order_id": document["order_id"], "status": "completion_pending", "customer_id": int(customer_id)},
  1324. {"$set": {"status": "completed", "completed_at": now, "closed_at": now, "updated_at": now}, "$inc": {"version": 1}},
  1325. return_document=ReturnDocument.AFTER,
  1326. )
  1327. if not order:
  1328. await qrdb.update_one({"qr_id": document["qr_id"]}, {"$set": {"status": "voided"}})
  1329. raise ServiceDataError("invalid_order_state", "服务单状态已变化,无法完成评价。")
  1330. await _schedule_address_redaction(order["order_id"], now)
  1331. await record_service_event(order["order_id"], "service_completed", actor_id=int(customer_id), metadata={"qr_id": document["qr_id"]})
  1332. return document, order
  1333. def _legacy_review_questions(values: dict[str, Any]) -> list[dict[str, Any]]:
  1334. questions = []
  1335. for item in values.get("rating_questions", []):
  1336. questions.append(
  1337. {
  1338. **item,
  1339. "type": "single_choice",
  1340. "options": [
  1341. {"value": str(score), "label": f"{score} 分", "score": score}
  1342. for score in range(5, 0, -1)
  1343. ],
  1344. }
  1345. )
  1346. for item in values.get("text_questions", []):
  1347. questions.append({**item, "type": "text"})
  1348. return questions
  1349. def normalize_review_template(values: dict[str, Any]) -> dict[str, Any]:
  1350. raw_questions = values.get("questions")
  1351. if raw_questions is None:
  1352. raw_questions = _legacy_review_questions(values)
  1353. if not isinstance(raw_questions, list) or not 1 <= len(raw_questions) <= 8:
  1354. raise ServiceDataError(
  1355. "invalid_review_questions",
  1356. "评价模板必须包含 1 到 8 个问题。",
  1357. )
  1358. questions = []
  1359. question_ids: set[str] = set()
  1360. scoring_weight = 0
  1361. for raw in raw_questions:
  1362. if not isinstance(raw, dict):
  1363. raise ServiceDataError("invalid_review_question", "评价问题格式无效。")
  1364. question_type = str(raw.get("type") or "").strip()
  1365. if question_type == "input":
  1366. question_type = "text"
  1367. if question_type not in {"single_choice", "multiple_choice", "text"}:
  1368. raise ServiceDataError(
  1369. "invalid_review_question_type",
  1370. "评价问题仅支持单选、多选和输入框。",
  1371. )
  1372. question_id = _clean_text(
  1373. raw.get("question_id") or uuid4().hex[:12],
  1374. max_length=32,
  1375. required=True,
  1376. )
  1377. if question_id in question_ids:
  1378. raise ServiceDataError("duplicate_review_question", "评价问题 ID 不能重复。")
  1379. question_ids.add(question_id)
  1380. question = {
  1381. "question_id": question_id,
  1382. "type": question_type,
  1383. "label": _clean_text(raw.get("label"), max_length=50, required=True),
  1384. "description": _clean_text(raw.get("description"), max_length=160),
  1385. "required": bool(raw.get("required")),
  1386. }
  1387. if question_type == "text":
  1388. question["max_length"] = max(
  1389. 50,
  1390. min(1000, int(raw.get("max_length") or 500)),
  1391. )
  1392. else:
  1393. raw_options = raw.get("options")
  1394. if not isinstance(raw_options, list) or not 2 <= len(raw_options) <= 8:
  1395. raise ServiceDataError(
  1396. "invalid_review_options",
  1397. "单选或多选问题必须设置 2 到 8 个选项。",
  1398. )
  1399. options = []
  1400. option_values: set[str] = set()
  1401. for index, raw_option in enumerate(raw_options):
  1402. if isinstance(raw_option, str):
  1403. raw_option = {"value": str(index + 1), "label": raw_option}
  1404. if not isinstance(raw_option, dict):
  1405. raise ServiceDataError(
  1406. "invalid_review_options",
  1407. "评价选项格式无效。",
  1408. )
  1409. option_value = _clean_text(
  1410. raw_option.get("value") or str(index + 1),
  1411. max_length=24,
  1412. required=True,
  1413. )
  1414. if option_value in option_values:
  1415. raise ServiceDataError(
  1416. "duplicate_review_option",
  1417. "同一问题的选项值不能重复。",
  1418. )
  1419. option_values.add(option_value)
  1420. option = {
  1421. "value": option_value,
  1422. "label": _clean_text(
  1423. raw_option.get("label"),
  1424. max_length=24,
  1425. required=True,
  1426. ),
  1427. }
  1428. if question_type == "single_choice":
  1429. score = int(raw_option.get("score") or 0)
  1430. if not 1 <= score <= 5:
  1431. raise ServiceDataError(
  1432. "invalid_review_option_score",
  1433. "单选项评分必须在 1 到 5 之间。",
  1434. )
  1435. option["score"] = score
  1436. options.append(option)
  1437. question["options"] = options
  1438. if question_type == "single_choice":
  1439. weight = int(raw.get("weight") or 0)
  1440. if weight <= 0:
  1441. raise ServiceDataError(
  1442. "invalid_review_weight",
  1443. "单选评分问题的权重必须大于 0。",
  1444. )
  1445. question["weight"] = weight
  1446. scoring_weight += weight
  1447. questions.append(question)
  1448. if scoring_weight != 100:
  1449. raise ServiceDataError(
  1450. "invalid_review_weight",
  1451. "全部单选评分问题的权重合计必须为 100%。",
  1452. )
  1453. return {"questions": questions}
  1454. async def get_active_review_template() -> dict[str, Any]:
  1455. await ensure_service_indexes()
  1456. current = await reviewtemplatesdb.find_one({"status": "active"}, sort=[("version", DESCENDING)])
  1457. if current:
  1458. return {**current, **normalize_review_template(current)}
  1459. now = utc_now()
  1460. document = {"template_id": uuid4().hex, "version": 1, "status": "active", **normalize_review_template(DEFAULT_REVIEW_TEMPLATE), "created_at": now, "published_at": now}
  1461. try:
  1462. await reviewtemplatesdb.insert_one(document)
  1463. except DuplicateKeyError:
  1464. return await reviewtemplatesdb.find_one({"status": "active"}, sort=[("version", DESCENDING)]) or document
  1465. return document
  1466. async def publish_review_template(values: dict[str, Any], *, actor_id: str) -> dict[str, Any]:
  1467. normalized = normalize_review_template(values)
  1468. await ensure_service_indexes()
  1469. while True:
  1470. current = await reviewtemplatesdb.find_one(sort=[("version", DESCENDING)])
  1471. version = int((current or {}).get("version") or 0) + 1
  1472. now = utc_now()
  1473. document = {
  1474. "template_id": uuid4().hex,
  1475. "version": version,
  1476. "status": "active",
  1477. **normalized,
  1478. "published_by": str(actor_id),
  1479. "created_at": now,
  1480. "published_at": now,
  1481. }
  1482. try:
  1483. await reviewtemplatesdb.insert_one(document)
  1484. break
  1485. except DuplicateKeyError:
  1486. continue
  1487. await reviewtemplatesdb.update_many(
  1488. {"status": "active", "template_id": {"$ne": document["template_id"]}},
  1489. {"$set": {"status": "retired", "retired_at": now}},
  1490. )
  1491. return document
  1492. async def begin_review_draft(
  1493. *,
  1494. customer_id: int,
  1495. technician_id: int,
  1496. source: str,
  1497. template: dict[str, Any],
  1498. package_id: str | None = None,
  1499. qr_id: str | None = None,
  1500. ) -> dict[str, Any]:
  1501. await ensure_service_indexes()
  1502. now = utc_now()
  1503. current = await reviewdraftsdb.find_one({"customer_id": int(customer_id)})
  1504. identity = {
  1505. "technician_id": int(technician_id),
  1506. "source": str(source),
  1507. "package_id": str(package_id) if package_id else None,
  1508. "qr_id": str(qr_id) if qr_id else None,
  1509. }
  1510. if (
  1511. current
  1512. and all(current.get(key) == value for key, value in identity.items())
  1513. and _as_utc(current["expires_at"]) > now
  1514. ):
  1515. snapshot = current.get("template_snapshot") or {}
  1516. if "questions" not in snapshot:
  1517. normalized_snapshot = {
  1518. "template_id": snapshot.get("template_id") or template["template_id"],
  1519. "version": snapshot.get("version") or template["version"],
  1520. "questions": normalize_review_template(snapshot)["questions"],
  1521. }
  1522. question_index = int(current.get("rating_index") or 0) + int(
  1523. current.get("text_index") or 0
  1524. )
  1525. await reviewdraftsdb.update_one(
  1526. {"draft_id": current["draft_id"]},
  1527. {
  1528. "$set": {
  1529. "template_snapshot": normalized_snapshot,
  1530. "question_index": question_index,
  1531. "updated_at": now,
  1532. },
  1533. "$unset": {"rating_index": "", "text_index": ""},
  1534. },
  1535. )
  1536. current.update(
  1537. {
  1538. "template_snapshot": normalized_snapshot,
  1539. "question_index": question_index,
  1540. }
  1541. )
  1542. return current
  1543. normalized_template = normalize_review_template(template)
  1544. document = {
  1545. "draft_id": uuid4().hex,
  1546. "customer_id": int(customer_id),
  1547. **identity,
  1548. "template_snapshot": {
  1549. "template_id": template["template_id"],
  1550. "version": template["version"],
  1551. "questions": normalized_template["questions"],
  1552. },
  1553. "question_index": 0,
  1554. "answers": {},
  1555. "created_at": now,
  1556. "updated_at": now,
  1557. "expires_at": now + timedelta(days=30),
  1558. }
  1559. await reviewdraftsdb.update_one(
  1560. {"customer_id": int(customer_id)}, {"$set": document}, upsert=True
  1561. )
  1562. return document
  1563. async def save_review_draft_progress(
  1564. customer_id: int,
  1565. *,
  1566. answers: dict[str, Any],
  1567. question_index: int | None = None,
  1568. rating_index: int = 0,
  1569. text_index: int = 0,
  1570. ) -> None:
  1571. current_index = (
  1572. max(0, int(question_index))
  1573. if question_index is not None
  1574. else max(0, int(rating_index) + int(text_index))
  1575. )
  1576. await reviewdraftsdb.update_one(
  1577. {"customer_id": int(customer_id)},
  1578. {
  1579. "$set": {
  1580. "answers": dict(answers),
  1581. "question_index": current_index,
  1582. "updated_at": utc_now(),
  1583. }
  1584. },
  1585. )
  1586. async def delete_review_draft(customer_id: int) -> None:
  1587. await reviewdraftsdb.delete_one({"customer_id": int(customer_id)})
  1588. def _score_review(
  1589. template: dict[str, Any],
  1590. answers: dict[str, Any],
  1591. ) -> tuple[float, dict[str, int], dict[str, Any], dict[str, str]]:
  1592. questions = normalize_review_template(template)["questions"]
  1593. ratings: dict[str, int] = {}
  1594. choices: dict[str, Any] = {}
  1595. texts: dict[str, str] = {}
  1596. total = 0.0
  1597. for question in questions:
  1598. raw = answers.get(question["question_id"])
  1599. if question["type"] == "text":
  1600. text = _clean_text(
  1601. raw,
  1602. max_length=question["max_length"],
  1603. required=question["required"],
  1604. )
  1605. if text:
  1606. texts[question["question_id"]] = text
  1607. continue
  1608. options = {
  1609. str(option["value"]): option for option in question["options"]
  1610. }
  1611. if question["type"] == "single_choice":
  1612. if raw in (None, "") and not question["required"]:
  1613. continue
  1614. option = options.get(str(raw))
  1615. if not option:
  1616. raise ServiceDataError("invalid_review_answer", "请选择有效的单选项。")
  1617. score = int(option["score"])
  1618. ratings[question["question_id"]] = score
  1619. choices[question["question_id"]] = option["label"]
  1620. total += score * question["weight"] / 100
  1621. continue
  1622. selected = raw if isinstance(raw, list) else ([] if raw in (None, "") else [raw])
  1623. selected_values = list(dict.fromkeys(str(item) for item in selected))
  1624. if question["required"] and not selected_values:
  1625. raise ServiceDataError("required_field", f"{question['label']} 为必选项。")
  1626. if set(selected_values) - set(options):
  1627. raise ServiceDataError("invalid_review_answer", "请选择有效的多选项。")
  1628. if selected_values:
  1629. choices[question["question_id"]] = [
  1630. options[value]["label"] for value in selected_values
  1631. ]
  1632. return round(total, 4), ratings, choices, texts
  1633. async def submit_review(
  1634. *,
  1635. customer: Any,
  1636. technician_id: int,
  1637. source: str,
  1638. answers: dict[str, Any],
  1639. anonymous: bool,
  1640. package_id: str | None = None,
  1641. qr_id: str | None = None,
  1642. template_snapshot: dict[str, Any] | None = None,
  1643. ) -> dict[str, Any]:
  1644. await observe_customer(customer)
  1645. await require_customer_allowed(int(customer.id))
  1646. if source not in REVIEW_SOURCES:
  1647. raise ServiceDataError("invalid_review_source", "评价来源无效。")
  1648. profile = await profilesdb.find_one({"user_id": int(technician_id), "application_status": "approved"})
  1649. if not profile:
  1650. raise ServiceDataError("technician_not_found", "未找到已认证技师。")
  1651. if int(customer.id) == int(technician_id):
  1652. raise ServiceDataError("self_review_not_allowed", "不能评价自己。")
  1653. order_id = None
  1654. package_snapshot = None
  1655. if source == "qr_verified":
  1656. qr = await qrdb.find_one({"qr_id": str(qr_id), "status": "claimed", "customer_id": int(customer.id), "technician_id": int(technician_id)})
  1657. if not qr:
  1658. raise ServiceDataError("review_qr_required", "缺少已完成服务单的评价资格。")
  1659. if await reviewsdb.find_one({"qr_id": str(qr_id)}):
  1660. raise ServiceDataError("review_already_submitted", "该服务单已经提交评价。")
  1661. order_id = qr["order_id"]
  1662. package_snapshot = qr["package_snapshot"]
  1663. else:
  1664. service_profile = profile.get("service_profile") or {}
  1665. package_snapshot = next((item for item in service_profile.get("packages", []) if str(item.get("package_id")) == str(package_id)), None)
  1666. if package_id and not package_snapshot:
  1667. raise ServiceDataError("package_not_found", "未找到所选套餐。")
  1668. if not package_snapshot:
  1669. packages = service_profile.get("packages") or []
  1670. primary_package = packages[0] if packages else {}
  1671. package_snapshot = {
  1672. "package_id": None,
  1673. "name": "技师服务",
  1674. "category": primary_package.get("category") or "其他",
  1675. }
  1676. template = template_snapshot or await get_active_review_template()
  1677. normalized_template = normalize_review_template(template)
  1678. score, rating_answers, choice_answers, text_answers = _score_review(
  1679. normalized_template,
  1680. answers,
  1681. )
  1682. normalized_content = "\n".join(
  1683. str(value).casefold()
  1684. for _, value in sorted(
  1685. {**choice_answers, **text_answers}.items()
  1686. )
  1687. if value
  1688. )
  1689. now = utc_now()
  1690. document = {
  1691. "review_id": uuid4().hex,
  1692. "technician_id": int(technician_id),
  1693. "technician_name": profile.get("display_name") or f"技师 {technician_id}",
  1694. "customer_id": int(customer.id),
  1695. "customer_name": _clean_text(getattr(customer, "first_name", ""), max_length=80) or f"用户 {customer.id}",
  1696. "source": source,
  1697. "package_snapshot": package_snapshot,
  1698. "category": package_snapshot.get("category") or "其他",
  1699. "template_snapshot": {
  1700. "template_id": template["template_id"],
  1701. "version": template["version"],
  1702. "questions": normalized_template["questions"],
  1703. },
  1704. "rating_answers": rating_answers,
  1705. "choice_answers": choice_answers,
  1706. "text_answers": text_answers,
  1707. "content_fingerprint": (
  1708. hashlib.sha256(normalized_content.encode()).hexdigest()
  1709. if normalized_content
  1710. else None
  1711. ),
  1712. "score": score,
  1713. "anonymous": bool(anonymous),
  1714. "status": "pending",
  1715. "created_at": now,
  1716. "updated_at": now,
  1717. }
  1718. if qr_id:
  1719. document["qr_id"] = str(qr_id)
  1720. document["order_id"] = order_id
  1721. try:
  1722. await reviewsdb.insert_one(document)
  1723. except DuplicateKeyError as exc:
  1724. if qr_id:
  1725. raise ServiceDataError(
  1726. "review_already_submitted", "该服务单已经提交评价。"
  1727. ) from exc
  1728. raise
  1729. return document
  1730. async def moderate_review(review_id: str, action: str, *, actor_id: str, reason: str = "") -> dict[str, Any]:
  1731. status_map = {"approve": "approved", "reject": "rejected", "void": "voided"}
  1732. source_status = {"approve": "pending", "reject": "pending", "void": "approved"}
  1733. if action not in status_map:
  1734. raise ServiceDataError("invalid_review_action", "不支持该评价操作。")
  1735. current = await reviewsdb.find_one({"review_id": str(review_id)})
  1736. if not current:
  1737. raise ServiceDataError("review_not_found", "未找到评价。")
  1738. if current["status"] != source_status[action]:
  1739. raise ServiceDataError("invalid_review_state", "当前评价状态不能执行该审核操作。")
  1740. if action in {"reject", "void"} and not _clean_text(reason, max_length=500):
  1741. raise ServiceDataError("reason_required", "拒绝或作废必须填写原因。")
  1742. updated = await reviewsdb.find_one_and_update(
  1743. {"review_id": str(review_id), "status": source_status[action]},
  1744. {"$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()}},
  1745. return_document=ReturnDocument.AFTER,
  1746. )
  1747. if not updated:
  1748. raise ServiceDataError("concurrent_review_update", "评价状态已变化,请刷新后重试。")
  1749. await recompute_review_stats()
  1750. return updated
  1751. async def recompute_review_stats() -> None:
  1752. await ensure_service_indexes()
  1753. approved = await reviewsdb.find({"status": "approved"}).to_list(length=100000)
  1754. grouped: dict[tuple[int, str], list[dict[str, Any]]] = defaultdict(list)
  1755. category_scores: dict[str, list[float]] = defaultdict(list)
  1756. for review in approved:
  1757. key = (int(review["technician_id"]), str(review.get("category") or "其他"))
  1758. grouped[key].append(review)
  1759. category_scores[key[1]].append(float(review["score"]))
  1760. now = utc_now()
  1761. await statsdb.delete_many({})
  1762. documents = []
  1763. for (technician_id, category), values in grouped.items():
  1764. count = len(values)
  1765. average = sum(float(item["score"]) for item in values) / count
  1766. category_average = sum(category_scores[category]) / len(category_scores[category])
  1767. rank_score = count / (count + 5) * average + 5 / (count + 5) * category_average
  1768. dimensions: dict[str, list[int]] = defaultdict(list)
  1769. for review in values:
  1770. for key, score in review.get("rating_answers", {}).items():
  1771. dimensions[key].append(int(score))
  1772. documents.append(
  1773. {
  1774. "technician_id": technician_id,
  1775. "category": category,
  1776. "review_count": count,
  1777. "average_score": round(average, 4),
  1778. "category_average": round(category_average, 4),
  1779. "rank_score": round(rank_score, 6),
  1780. "eligible": count >= 1,
  1781. "dimension_averages": {key: round(sum(items) / len(items), 4) for key, items in dimensions.items()},
  1782. "last_review_at": max(item.get("moderated_at") or item["created_at"] for item in values),
  1783. "updated_at": now,
  1784. }
  1785. )
  1786. if documents:
  1787. await statsdb.insert_many(documents)
  1788. async def list_leaderboard(*, category: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
  1789. await ensure_service_indexes()
  1790. filters: dict[str, Any] = {"review_count": {"$gte": 1}}
  1791. if category:
  1792. filters["category"] = str(category)
  1793. visible_ids = [
  1794. int(item["user_id"])
  1795. async for item in profilesdb.find(
  1796. {
  1797. "application_status": "approved",
  1798. "listed": True,
  1799. "username": {"$nin": [None, ""]},
  1800. },
  1801. {"user_id": 1},
  1802. )
  1803. ]
  1804. filters["technician_id"] = {"$in": visible_ids}
  1805. total = await statsdb.count_documents(filters)
  1806. 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)
  1807. 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})}
  1808. 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]
  1809. return visible, total
  1810. async def list_public_reviews(
  1811. technician_id: int,
  1812. *,
  1813. page: int = 1,
  1814. page_size: int = 10,
  1815. ) -> tuple[list[dict[str, Any]], int]:
  1816. filters = {"technician_id": int(technician_id), "status": "approved"}
  1817. total = await reviewsdb.count_documents(filters)
  1818. items = await reviewsdb.find(filters).sort("moderated_at", DESCENDING).skip(
  1819. (max(1, page) - 1) * page_size
  1820. ).limit(page_size).to_list(length=page_size)
  1821. public_items = []
  1822. for item in items:
  1823. public_items.append(
  1824. {
  1825. "review_id": item["review_id"],
  1826. "source": item["source"],
  1827. "source_label": (
  1828. "完成服务单评价"
  1829. if item["source"] == "qr_verified"
  1830. else "用户主动评价"
  1831. ),
  1832. "customer_name": "匿名顾客" if item.get("anonymous", True) else item.get("customer_name"),
  1833. "package_name": (item.get("package_snapshot") or {}).get("name"),
  1834. "category": item.get("category"),
  1835. "score": item.get("score"),
  1836. "rating_answers": item.get("rating_answers", {}),
  1837. "choice_answers": item.get("choice_answers", {}),
  1838. "text_answers": item.get("text_answers", {}),
  1839. "template_snapshot": item.get("template_snapshot", {}),
  1840. "approved_at": item.get("moderated_at"),
  1841. "topic_url": item.get("topic_url"),
  1842. }
  1843. )
  1844. return public_items, total
  1845. async def get_technician_review_topic(
  1846. technician_id: int,
  1847. ) -> dict[str, Any] | None:
  1848. await ensure_service_indexes()
  1849. return await reviewtopicsdb.find_one(
  1850. {"technician_id": int(technician_id), "status": "active"}
  1851. )
  1852. async def save_technician_review_topic(
  1853. *,
  1854. technician_id: int,
  1855. forum_chat_id: str,
  1856. forum_username: str,
  1857. message_thread_id: int,
  1858. topic_name: str,
  1859. topic_url: str,
  1860. ) -> dict[str, Any]:
  1861. await ensure_service_indexes()
  1862. now = utc_now()
  1863. await reviewtopicsdb.update_one(
  1864. {"technician_id": int(technician_id)},
  1865. {
  1866. "$set": {
  1867. "forum_chat_id": str(forum_chat_id),
  1868. "forum_username": str(forum_username),
  1869. "message_thread_id": int(message_thread_id),
  1870. "topic_name": _clean_text(
  1871. topic_name,
  1872. max_length=128,
  1873. required=True,
  1874. ),
  1875. "topic_url": _clean_text(
  1876. topic_url,
  1877. max_length=300,
  1878. required=True,
  1879. ),
  1880. "status": "active",
  1881. "updated_at": now,
  1882. },
  1883. "$setOnInsert": {
  1884. "technician_id": int(technician_id),
  1885. "created_at": now,
  1886. },
  1887. },
  1888. upsert=True,
  1889. )
  1890. return await get_technician_review_topic(int(technician_id)) or {}
  1891. async def mark_review_topic_published(
  1892. review_id: str,
  1893. *,
  1894. topic: dict[str, Any],
  1895. message_id: int,
  1896. ) -> dict[str, Any]:
  1897. updated = await reviewsdb.find_one_and_update(
  1898. {"review_id": str(review_id), "status": "approved"},
  1899. {
  1900. "$set": {
  1901. "topic_url": topic["topic_url"],
  1902. "topic_message_id": int(message_id),
  1903. "topic_published_at": utc_now(),
  1904. "updated_at": utc_now(),
  1905. },
  1906. "$unset": {"topic_publish_error": ""},
  1907. },
  1908. return_document=ReturnDocument.AFTER,
  1909. )
  1910. return updated or {}
  1911. async def get_review(review_id: str) -> dict[str, Any] | None:
  1912. return await reviewsdb.find_one({"review_id": str(review_id)})
  1913. async def mark_review_topic_publish_error(
  1914. review_id: str,
  1915. error: str,
  1916. ) -> None:
  1917. await reviewsdb.update_one(
  1918. {"review_id": str(review_id), "status": "approved"},
  1919. {
  1920. "$set": {
  1921. "topic_publish_error": _clean_text(error, max_length=300),
  1922. "updated_at": utc_now(),
  1923. }
  1924. },
  1925. )
  1926. async def get_technician_review_summary(technician_id: int) -> dict[str, Any]:
  1927. items = await reviewsdb.find(
  1928. {"technician_id": int(technician_id), "status": "approved"},
  1929. {"score": 1, "source": 1},
  1930. ).to_list(length=100000)
  1931. count = len(items)
  1932. return {
  1933. "review_count": count,
  1934. "average_score": round(
  1935. sum(float(item["score"]) for item in items) / count,
  1936. 4,
  1937. )
  1938. if count
  1939. else 0,
  1940. "source_counts": {
  1941. source: sum(1 for item in items if item.get("source") == source)
  1942. for source in REVIEW_SOURCES
  1943. },
  1944. "ranking_title": "审核评价排行",
  1945. }
  1946. async def list_service_orders(*, query: str = "", status: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
  1947. await ensure_service_indexes()
  1948. filters: dict[str, Any] = {}
  1949. if status:
  1950. filters["status"] = status
  1951. if query:
  1952. if query.isdigit():
  1953. filters["$or"] = [{"customer_id": int(query)}, {"technician_id": int(query)}]
  1954. else:
  1955. pattern = re.compile(re.escape(query), re.IGNORECASE)
  1956. filters["$or"] = [{"order_id": pattern}, {"customer_name": pattern}, {"technician_name": pattern}]
  1957. total = await ordersdb.count_documents(filters)
  1958. items = await ordersdb.find(filters).sort("created_at", DESCENDING).skip((max(1, page) - 1) * page_size).limit(page_size).to_list(length=page_size)
  1959. return items, total
  1960. async def list_actor_service_orders(
  1961. actor_id: int,
  1962. *,
  1963. role: str = "all",
  1964. page: int = 1,
  1965. page_size: int = 20,
  1966. ) -> tuple[list[dict[str, Any]], int]:
  1967. if role == "customer":
  1968. filters: dict[str, Any] = {"customer_id": int(actor_id)}
  1969. elif role == "technician":
  1970. filters = {"technician_id": int(actor_id)}
  1971. else:
  1972. filters = {
  1973. "$or": [
  1974. {"customer_id": int(actor_id)},
  1975. {"technician_id": int(actor_id)},
  1976. ]
  1977. }
  1978. total = await ordersdb.count_documents(filters)
  1979. items = await ordersdb.find(filters).sort("created_at", DESCENDING).skip(
  1980. (max(1, page) - 1) * page_size
  1981. ).limit(page_size).to_list(length=page_size)
  1982. return items, total
  1983. async def get_admin_service_order(order_id: str) -> dict[str, Any]:
  1984. order = await ordersdb.find_one({"order_id": str(order_id)})
  1985. if not order:
  1986. raise ServiceDataError("order_not_found", "未找到服务单。")
  1987. result = dict(order)
  1988. result["quotes"] = await quotesdb.find({"order_id": str(order_id)}).sort(
  1989. "version", ASCENDING
  1990. ).to_list(length=100)
  1991. result["events"] = await eventsdb.find({"order_id": str(order_id)}).sort(
  1992. "created_at", ASCENDING
  1993. ).to_list(length=500)
  1994. result["qr_records"] = await qrdb.find({"order_id": str(order_id)}).sort(
  1995. "created_at", DESCENDING
  1996. ).to_list(length=100)
  1997. return result
  1998. async def list_review_qr_records(
  1999. *,
  2000. status: str = "",
  2001. query: str = "",
  2002. page: int = 1,
  2003. page_size: int = 20,
  2004. ) -> tuple[list[dict[str, Any]], int]:
  2005. filters: dict[str, Any] = {}
  2006. if status:
  2007. filters["status"] = status
  2008. if query:
  2009. filters["$or"] = [
  2010. {"order_id": re.compile(re.escape(query), re.IGNORECASE)},
  2011. {"qr_id": re.compile(re.escape(query), re.IGNORECASE)},
  2012. {"customer_id": int(query) if query.isdigit() else -1},
  2013. {"technician_id": int(query) if query.isdigit() else -1},
  2014. ]
  2015. total = await qrdb.count_documents(filters)
  2016. items = await qrdb.find(filters).sort("created_at", DESCENDING).skip(
  2017. (max(1, page) - 1) * page_size
  2018. ).limit(page_size).to_list(length=page_size)
  2019. for item in items:
  2020. item.pop("token_hash", None)
  2021. return items, total
  2022. async def list_reviews(*, status: str = "", source: str = "", query: str = "", page: int = 1, page_size: int = 20) -> tuple[list[dict[str, Any]], int]:
  2023. await ensure_service_indexes()
  2024. filters: dict[str, Any] = {}
  2025. if status:
  2026. filters["status"] = status
  2027. if source:
  2028. filters["source"] = source
  2029. if query:
  2030. pattern = re.compile(re.escape(query), re.IGNORECASE)
  2031. filters["$or"] = [{"technician_name": pattern}, {"customer_name": pattern}, {"review_id": pattern}]
  2032. total = await reviewsdb.count_documents(filters)
  2033. items = await reviewsdb.find(filters).sort("created_at", DESCENDING).skip((max(1, page) - 1) * page_size).limit(page_size).to_list(length=page_size)
  2034. for item in items:
  2035. item["customer_review_count"] = await reviewsdb.count_documents({"customer_id": item["customer_id"], "technician_id": item["technician_id"]})
  2036. 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)}})
  2037. 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)}})
  2038. item["similar_content_count"] = (
  2039. await reviewsdb.count_documents(
  2040. {
  2041. "customer_id": item["customer_id"],
  2042. "technician_id": item["technician_id"],
  2043. "content_fingerprint": item["content_fingerprint"],
  2044. }
  2045. )
  2046. if item.get("content_fingerprint")
  2047. else 0
  2048. )
  2049. return items, total
  2050. async def get_order_for_actor(order_id: str, actor_id: int, *, reveal_address: bool = False) -> dict[str, Any]:
  2051. order = await ordersdb.find_one({"order_id": str(order_id)})
  2052. if not order or int(actor_id) not in {int(order["customer_id"]), int(order["technician_id"])}:
  2053. raise ServiceDataError("order_not_found", "未找到服务单。")
  2054. result = dict(order)
  2055. can_reveal = order["status"] in {"confirmed", "in_progress", "completion_pending", "completed", "disputed"}
  2056. if reveal_address and can_reveal:
  2057. address = await addressesdb.find_one({"order_id": str(order_id), "redacted": False})
  2058. if address:
  2059. result["exact_address"] = _decrypt_address(address["encrypted_payload"])
  2060. return result
  2061. async def admin_reveal_order_address(order_id: str, *, actor_id: str, reason: str) -> dict[str, Any]:
  2062. reason = _clean_text(reason, max_length=500, required=True)
  2063. order = await ordersdb.find_one({"order_id": str(order_id)})
  2064. address = await addressesdb.find_one({"order_id": str(order_id), "redacted": False})
  2065. if not order or not address:
  2066. raise ServiceDataError("address_unavailable", "精确地址已脱敏或不存在。")
  2067. await record_service_event(order_id, "admin_address_revealed", actor_id=actor_id, reason=reason)
  2068. return _decrypt_address(address["encrypted_payload"])
  2069. async def redact_expired_addresses() -> int:
  2070. now = utc_now()
  2071. candidates = await addressesdb.find({"redacted": False, "redact_after": {"$lte": now}}).to_list(length=1000)
  2072. count = 0
  2073. for item in candidates:
  2074. order = await ordersdb.find_one({"order_id": item["order_id"]})
  2075. if not order or order.get("status") not in ORDER_TERMINAL_STATUSES:
  2076. continue
  2077. result = await addressesdb.update_one({"_id": item["_id"], "redacted": False}, {"$set": {"redacted": True, "redacted_at": now, "updated_at": now}, "$unset": {"encrypted_payload": ""}})
  2078. if result.modified_count:
  2079. await ordersdb.update_one(
  2080. {"order_id": item["order_id"]},
  2081. {
  2082. "$set": {"address_redacted": True, "updated_at": now},
  2083. "$unset": {"distance_meters": ""},
  2084. },
  2085. )
  2086. count += result.modified_count
  2087. return count
  2088. 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]:
  2089. if scope not in {"global", "technician"}:
  2090. raise ServiceDataError("invalid_block_scope", "拉黑范围无效。")
  2091. if scope == "technician" and technician_id is None:
  2092. raise ServiceDataError("technician_required", "技师拉黑必须指定技师。")
  2093. now = utc_now()
  2094. await blocksdb.update_one(
  2095. {"scope": scope, "technician_id": int(technician_id or 0), "customer_id": int(customer_id)},
  2096. {"$set": {"active": bool(active), "reason": _clean_text(reason, max_length=500, required=active), "actor_id": actor_id, "updated_at": now}, "$setOnInsert": {"created_at": now}},
  2097. upsert=True,
  2098. )
  2099. return await blocksdb.find_one({"scope": scope, "technician_id": int(technician_id or 0), "customer_id": int(customer_id)}) or {}
  2100. async def set_technician_customer_block(
  2101. *,
  2102. technician_id: int,
  2103. customer_id: int,
  2104. active: bool,
  2105. reason: str,
  2106. ) -> dict[str, Any]:
  2107. profile = await profilesdb.find_one(
  2108. {"user_id": int(technician_id), "application_status": "approved"}
  2109. )
  2110. if not profile:
  2111. raise ServiceDataError("technician_required", "只有已认证技师可以管理个人拉黑。")
  2112. if not await ordersdb.find_one(
  2113. {"technician_id": int(technician_id), "customer_id": int(customer_id)}
  2114. ):
  2115. raise ServiceDataError("customer_relationship_required", "只能拉黑曾向你发起服务请求的顾客。")
  2116. return await set_customer_block(
  2117. customer_id=int(customer_id),
  2118. scope="technician",
  2119. technician_id=int(technician_id),
  2120. active=active,
  2121. actor_id=int(technician_id),
  2122. reason=reason,
  2123. )
  2124. async def list_customer_blocks(
  2125. *,
  2126. scope: str = "global",
  2127. active: bool | None = None,
  2128. page: int = 1,
  2129. page_size: int = 20,
  2130. ) -> tuple[list[dict[str, Any]], int]:
  2131. filters: dict[str, Any] = {"scope": scope}
  2132. if active is not None:
  2133. filters["active"] = active
  2134. total = await blocksdb.count_documents(filters)
  2135. items = await blocksdb.find(filters).sort("updated_at", DESCENDING).skip(
  2136. (max(1, page) - 1) * page_size
  2137. ).limit(page_size).to_list(length=page_size)
  2138. return items, total
  2139. async def create_service_report(*, reporter_id: int, target_type: str, target_id: str, reason: str) -> dict[str, Any]:
  2140. if target_type not in {"technician", "customer", "order", "review"}:
  2141. raise ServiceDataError("invalid_report_target", "举报对象无效。")
  2142. 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()}
  2143. await reportsdb.insert_one(document)
  2144. return document
  2145. async def list_service_reports(
  2146. *,
  2147. status: str = "",
  2148. page: int = 1,
  2149. page_size: int = 20,
  2150. ) -> tuple[list[dict[str, Any]], int]:
  2151. filters = {"status": status} if status else {}
  2152. total = await reportsdb.count_documents(filters)
  2153. items = await reportsdb.find(filters).sort("created_at", DESCENDING).skip(
  2154. (max(1, page) - 1) * page_size
  2155. ).limit(page_size).to_list(length=page_size)
  2156. return items, total
  2157. async def moderate_service_report(
  2158. report_id: str,
  2159. *,
  2160. action: str,
  2161. actor_id: str,
  2162. reason: str,
  2163. ) -> dict[str, Any]:
  2164. status_map = {"resolve": "resolved", "dismiss": "dismissed"}
  2165. if action not in status_map:
  2166. raise ServiceDataError("invalid_report_action", "不支持该举报操作。")
  2167. updated = await reportsdb.find_one_and_update(
  2168. {"report_id": str(report_id), "status": "open"},
  2169. {
  2170. "$set": {
  2171. "status": status_map[action],
  2172. "handled_by": str(actor_id),
  2173. "handled_at": utc_now(),
  2174. "handling_reason": _clean_text(reason, max_length=500, required=True),
  2175. }
  2176. },
  2177. return_document=ReturnDocument.AFTER,
  2178. )
  2179. if not updated:
  2180. raise ServiceDataError("report_not_found", "举报不存在或已处理。")
  2181. return updated
  2182. async def fulfillment_metrics() -> dict[str, Any]:
  2183. await ensure_service_indexes()
  2184. counts = {status: await ordersdb.count_documents({"status": status}) for status in ORDER_STATUSES}
  2185. valid_requests = sum(counts.values()) - counts["voided"]
  2186. quoted = await ordersdb.count_documents(
  2187. {"status": {"$ne": "voided"}, "current_quote_id": {"$exists": True}}
  2188. )
  2189. confirmed = await ordersdb.count_documents(
  2190. {"status": {"$ne": "voided"}, "confirmed_at": {"$exists": True}}
  2191. )
  2192. completed = counts["completed"]
  2193. canceled = counts["canceled_customer"] + counts["canceled_technician"]
  2194. directory_views = await eventsdb.count_documents({"event_type": "directory_viewed"})
  2195. technician_views = await eventsdb.count_documents({"event_type": "technician_viewed"})
  2196. request_starts = await eventsdb.count_documents({"event_type": "request_started"})
  2197. return {
  2198. "counts": counts,
  2199. "valid_requests": valid_requests,
  2200. "quoted_requests": quoted,
  2201. "confirmed_requests": confirmed,
  2202. "directory_views": directory_views,
  2203. "technician_views": technician_views,
  2204. "request_starts": request_starts,
  2205. "request_submission_rate": min(1.0, round(valid_requests / request_starts, 4))
  2206. if request_starts
  2207. else 0,
  2208. "quote_rate": round(quoted / valid_requests, 4) if valid_requests else 0,
  2209. "customer_confirmation_rate": round(confirmed / quoted, 4) if quoted else 0,
  2210. "match_success_rate": round(confirmed / valid_requests, 4) if valid_requests else 0,
  2211. "fulfillment_rate": round(completed / confirmed, 4) if confirmed else 0,
  2212. "cancellation_rate": round(canceled / valid_requests, 4) if valid_requests else 0,
  2213. "cancellation_distribution": {
  2214. "customer": counts["canceled_customer"],
  2215. "technician": counts["canceled_technician"],
  2216. },
  2217. "scope_note": "仅统计平台服务单;直接 Telegram 私聊成交不在统计范围。",
  2218. }