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)}