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