from __future__ import annotations import asyncio from html import escape from typing import Any from aiohttp import ClientError, ClientSession, ClientTimeout from wbb.utils.dbservice import ( get_review, get_service_settings, get_technician_review_topic, get_technician_service_profile, mark_review_topic_publish_error, mark_review_topic_published, save_technician_review_topic, ) class ReviewTopicError(RuntimeError): pass _topic_locks: dict[int, asyncio.Lock] = {} async def _bot_api( bot_token: str, method: str, payload: dict[str, Any], ) -> dict[str, Any]: if not bot_token: raise ReviewTopicError("当前 Bot 凭据不可用。") url = f"https://api.telegram.org/bot{bot_token}/{method}" try: async with ClientSession(timeout=ClientTimeout(total=20)) as session: async with session.post(url, json=payload) as response: data = await response.json(content_type=None) except (ClientError, TimeoutError, ValueError) as exc: raise ReviewTopicError("Telegram 评价话题请求失败。") from exc if response.status >= 400 or not data.get("ok"): description = str(data.get("description") or "Telegram 返回未知错误") raise ReviewTopicError(description[:240]) result = data.get("result") return result if isinstance(result, dict) else {} def _topic_url(forum_username: str, message_thread_id: int) -> str: username = str(forum_username or "").strip().lstrip("@") if not username: raise ReviewTopicError("请先配置公开评价群用户名。") return f"https://t.me/{username}/{int(message_thread_id)}" async def ensure_technician_review_topic( bot_token: str, technician_id: int, ) -> dict[str, Any] | None: settings = await get_service_settings() forum_chat_id = str(settings.get("review_forum_chat_id") or "").strip() forum_username = str(settings.get("review_forum_username") or "").strip() if not forum_chat_id or not forum_username or not bot_token: return None technician_id = int(technician_id) current = await get_technician_review_topic(technician_id) if current and str(current.get("forum_chat_id")) == forum_chat_id: return current lock = _topic_locks.setdefault(technician_id, asyncio.Lock()) async with lock: current = await get_technician_review_topic(technician_id) if current and str(current.get("forum_chat_id")) == forum_chat_id: return current profile = await get_technician_service_profile(technician_id) display_name = str( profile.get("display_name") or f"技师 {technician_id}" ).strip() topic_name = f"{display_name} · 评价"[:128] created = await _bot_api( bot_token, "createForumTopic", { "chat_id": forum_chat_id, "name": topic_name, }, ) message_thread_id = int(created.get("message_thread_id") or 0) if message_thread_id <= 0: raise ReviewTopicError("Telegram 未返回有效的话题 ID。") topic = await save_technician_review_topic( technician_id=technician_id, forum_chat_id=forum_chat_id, forum_username=forum_username, message_thread_id=message_thread_id, topic_name=topic_name, topic_url=_topic_url(forum_username, message_thread_id), ) await _bot_api( bot_token, "sendMessage", { "chat_id": forum_chat_id, "message_thread_id": message_thread_id, "text": ( f"{escape(display_name)}的公开评价\n" "管理员审核通过的评价会自动发布在这里,任何群成员均可查看。" ), "parse_mode": "HTML", "disable_web_page_preview": True, }, ) return topic def _review_message(review: dict[str, Any]) -> str: template = review.get("template_snapshot") or {} labels = { str(question.get("question_id")): str(question.get("label") or "") for question in template.get("questions", []) } author = ( "匿名顾客" if review.get("anonymous", True) else str(review.get("customer_name") or "顾客") ) package_name = str( (review.get("package_snapshot") or {}).get("name") or "其他服务" ) lines = [ f"新评价 · {float(review.get('score') or 0):.1f} 分", f"技师:{escape(str(review.get('technician_name') or review['technician_id']))}", f"服务:{escape(package_name)}", f"评价人:{escape(author)}", ] choice_answers = review.get("choice_answers") or {} text_answers = review.get("text_answers") or {} for question_id, answer in choice_answers.items(): value = "、".join(str(item) for item in answer) if isinstance(answer, list) else str(answer) if value: lines.append( f"{escape(labels.get(str(question_id)) or '选择')}:{escape(value)}" ) for question_id, answer in text_answers.items(): if answer: lines.append( f"{escape(labels.get(str(question_id)) or '评价内容')}:" f"{escape(str(answer))}" ) return "\n".join(lines) async def publish_approved_review( bot_token: str, review: dict[str, Any], ) -> dict[str, Any]: stored = await get_review(str(review["review_id"])) if stored: review = stored if review.get("topic_message_id"): return { "status": "published", "topic_url": review.get("topic_url"), "message_id": review.get("topic_message_id"), } try: topic = await ensure_technician_review_topic( bot_token, int(review["technician_id"]), ) if not topic: return { "status": "unconfigured", "message": "尚未配置公开评价群。", } result = await _bot_api( bot_token, "sendMessage", { "chat_id": topic["forum_chat_id"], "message_thread_id": int(topic["message_thread_id"]), "text": _review_message(review), "parse_mode": "HTML", "disable_web_page_preview": True, }, ) message_id = int(result.get("message_id") or 0) if message_id <= 0: raise ReviewTopicError("Telegram 未返回有效的评价消息 ID。") await mark_review_topic_published( review["review_id"], topic=topic, message_id=message_id, ) return { "status": "published", "topic_url": topic["topic_url"], "message_id": message_id, } except ReviewTopicError as exc: await mark_review_topic_publish_error(review["review_id"], str(exc)) return {"status": "failed", "message": str(exc)}