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