| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528 |
- from __future__ import annotations
- import asyncio
- from contextlib import suppress
- from typing import Any
- from aiohttp import ClientError, ClientSession, ClientTimeout
- from pyrogram.enums import ChatMemberStatus
- import wbb
- from wbb import BOT_PROFILE_ID, SUDOERS, app, log
- from wbb.utils.dbdirectory import (
- DirectoryDataError,
- aware_utc,
- clear_directory_location,
- decide_teacher_application,
- expire_teacher_presence,
- get_directory_location,
- get_directory_profile,
- list_directory_teachers,
- list_membership_candidates,
- save_directory_location,
- set_teacher_state,
- submit_teacher_application,
- teacher_is_online,
- upsert_directory_identity,
- upsert_directory_membership,
- )
- ACTIVE_MEMBER_STATUSES = {
- ChatMemberStatus.OWNER,
- ChatMemberStatus.ADMINISTRATOR,
- ChatMemberStatus.MEMBER,
- ChatMemberStatus.RESTRICTED,
- }
- ADMIN_MEMBER_STATUSES = {
- ChatMemberStatus.OWNER,
- ChatMemberStatus.ADMINISTRATOR,
- }
- _presence_task: asyncio.Task[None] | None = None
- class DirectoryServiceError(RuntimeError):
- def __init__(self, code: str, message: str, *, status: int = 400):
- super().__init__(message)
- self.code = code
- self.status = status
- def status_value(value: Any) -> str:
- return str(getattr(value, "value", value or "")).lower()
- def is_sudoer(user_id: int) -> bool:
- try:
- return int(user_id) in SUDOERS
- except (TypeError, ValueError):
- return False
- def user_identity(user: Any) -> dict[str, Any]:
- return {
- "user_id": int(user.id),
- "username": getattr(user, "username", None),
- "first_name": getattr(user, "first_name", None),
- "last_name": getattr(user, "last_name", None),
- "display_name": " ".join(
- value.strip()
- for value in (
- getattr(user, "first_name", None),
- getattr(user, "last_name", None),
- )
- if value and value.strip()
- )
- or (
- f"@{user.username}"
- if getattr(user, "username", None)
- else f"用户 {user.id}"
- ),
- }
- async def observe_directory_user(user: Any, *, source: str = "telegram") -> dict[str, Any]:
- identity = user_identity(user)
- return await upsert_directory_identity(
- user_id=identity["user_id"],
- username=identity["username"],
- first_name=identity["first_name"],
- last_name=identity["last_name"],
- source=source,
- )
- async def observe_group_member(
- *,
- chat_id: int,
- chat_title: str,
- user: Any,
- status: str = "member",
- active: bool = True,
- verified: bool = False,
- ) -> dict[str, Any]:
- identity = user_identity(user)
- await observe_directory_user(user)
- return await upsert_directory_membership(
- bot_id=str(BOT_PROFILE_ID),
- chat_id=int(chat_id),
- user_id=identity["user_id"],
- status=status,
- active=active,
- chat_title=chat_title,
- username=identity["username"],
- display_name_value=identity["display_name"],
- verified=verified,
- )
- async def verify_local_membership(
- *,
- user_id: int,
- chat_ids: list[int],
- require_admin: bool = False,
- ) -> dict[str, Any]:
- accepted = ADMIN_MEMBER_STATUSES if require_admin else ACTIVE_MEMBER_STATUSES
- for chat_id in dict.fromkeys(int(value) for value in chat_ids):
- try:
- member = await app.get_chat_member(int(chat_id), int(user_id))
- except Exception:
- continue
- status = getattr(member, "status", None)
- active = status in ACTIVE_MEMBER_STATUSES
- user = getattr(member, "user", None)
- if user is not None:
- with suppress(Exception):
- await observe_group_member(
- chat_id=int(chat_id),
- chat_title="",
- user=user,
- status=status_value(status),
- active=active,
- verified=True,
- )
- if status in accepted:
- return {
- "allowed": True,
- "bot_id": str(BOT_PROFILE_ID),
- "chat_id": int(chat_id),
- "status": status_value(status),
- }
- return {"allowed": False}
- async def verify_platform_membership(
- *,
- user_id: int,
- require_admin: bool = False,
- ) -> dict[str, Any]:
- if require_admin and is_sudoer(int(user_id)):
- return {
- "allowed": True,
- "bot_id": str(BOT_PROFILE_ID),
- "chat_id": None,
- "status": "sudoer",
- }
- candidates = await list_membership_candidates(int(user_id))
- if not candidates:
- return {"allowed": False, "reason": "membership_evidence_required"}
- internal_url = str(
- getattr(wbb, "SUPERVISOR_INTERNAL_URL", "") or ""
- ).rstrip("/")
- internal_token = str(getattr(wbb, "INTERNAL_TOKEN", "") or "")
- if internal_url and internal_token:
- try:
- async with ClientSession(timeout=ClientTimeout(total=20)) as session:
- async with session.post(
- f"{internal_url}/api/internal/v1/directory/verify-platform",
- headers={"Authorization": f"Bearer {internal_token}"},
- json={
- "user_id": str(user_id),
- "require_admin": bool(require_admin),
- },
- ) as response:
- if response.status != 200:
- return {"allowed": False, "reason": "verification_unavailable"}
- payload = await response.json()
- return payload.get("data") or {"allowed": False}
- except (ClientError, TimeoutError, ValueError) as exc:
- log.error(f"跨 Bot 成员身份校验失败:{exc}")
- return {"allowed": False, "reason": "verification_unavailable"}
- local_chat_ids = [
- int(item["chat_id"])
- for item in candidates
- if str(item.get("bot_id")) == str(BOT_PROFILE_ID)
- ]
- if not local_chat_ids:
- return {"allowed": False, "reason": "verification_unavailable"}
- return await verify_local_membership(
- user_id=int(user_id),
- chat_ids=local_chat_ids,
- require_admin=require_admin,
- )
- async def require_platform_membership(
- user_id: int,
- *,
- require_admin: bool = False,
- ) -> dict[str, Any]:
- result = await verify_platform_membership(
- user_id=int(user_id),
- require_admin=require_admin,
- )
- if result.get("allowed"):
- return result
- reason = result.get("reason")
- if reason == "membership_evidence_required":
- raise DirectoryServiceError(
- "membership_verification_required",
- "请先在任一受管群发送 /directory 完成群成员身份验证。",
- status=403,
- )
- if reason == "verification_unavailable":
- raise DirectoryServiceError(
- "membership_verification_unavailable",
- "暂时无法通过对应机器人核验群身份,请稍后重试。",
- status=503,
- )
- raise DirectoryServiceError(
- "managed_group_member_required",
- "该功能仅对受管群的当前成员开放。",
- status=403,
- )
- async def update_user_location(
- *,
- user: Any,
- longitude: float,
- latitude: float,
- source: str = "telegram_location",
- ) -> dict[str, Any]:
- await require_platform_membership(int(user.id))
- identity = user_identity(user)
- await observe_directory_user(user)
- try:
- return await save_directory_location(
- user_id=int(user.id),
- longitude=longitude,
- latitude=latitude,
- source=source,
- actor_id=int(user.id),
- actor_name=identity["display_name"],
- bot_id=str(BOT_PROFILE_ID),
- )
- except DirectoryDataError as exc:
- raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
- async def clear_user_location(*, user: Any) -> bool:
- identity = user_identity(user)
- return await clear_directory_location(
- user_id=int(user.id),
- actor_id=int(user.id),
- actor_name=identity["display_name"],
- source="telegram_private",
- reason="用户主动清除位置",
- )
- async def apply_as_teacher(
- *,
- user: Any,
- chat_id: int | None = None,
- chat_title: str = "",
- source: str = "telegram_private",
- ) -> tuple[dict[str, Any], bool]:
- authorization = await require_platform_membership(int(user.id))
- await observe_directory_user(user)
- application_chat_id = int(chat_id or authorization.get("chat_id") or 0)
- if not application_chat_id:
- raise DirectoryServiceError(
- "application_chat_required",
- "请先在要申请技师的群里打开附近技师菜单。",
- status=409,
- )
- application_chat_title = str(chat_title or "").strip()
- if not application_chat_title:
- with suppress(Exception):
- chat = await app.get_chat(application_chat_id)
- application_chat_title = str(getattr(chat, "title", "") or "").strip()
- try:
- return await submit_teacher_application(
- user_id=int(user.id),
- source=source,
- bot_id=str(BOT_PROFILE_ID),
- chat_id=application_chat_id,
- chat_title=application_chat_title,
- )
- except DirectoryDataError as exc:
- raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
- async def change_own_teacher_state(
- *,
- user: Any,
- action: str,
- source: str = "telegram_private",
- idempotency_key: str | None = None,
- chat_id: int | None = None,
- ) -> tuple[dict[str, Any], bool]:
- await require_platform_membership(int(user.id))
- identity = user_identity(user)
- await observe_directory_user(user)
- try:
- return await set_teacher_state(
- user_id=int(user.id),
- action=action,
- actor_id=int(user.id),
- actor_name=identity["display_name"],
- source=source,
- reason="技师自助操作",
- idempotency_key=idempotency_key,
- chat_id=chat_id,
- )
- except DirectoryDataError as exc:
- raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
- async def admin_source_chat_ids(actor: Any) -> list[int] | None:
- if is_sudoer(int(actor.id)):
- return None
- candidates = await list_membership_candidates(int(actor.id))
- chat_ids: list[int] = []
- for item in candidates:
- if str(item.get("bot_id")) != str(BOT_PROFILE_ID):
- continue
- chat_id = int(item["chat_id"])
- try:
- member = await app.get_chat_member(chat_id, int(actor.id))
- except Exception:
- continue
- if getattr(member, "status", None) in ADMIN_MEMBER_STATUSES:
- chat_ids.append(chat_id)
- values = list(dict.fromkeys(chat_ids))
- if not values:
- raise DirectoryServiceError(
- "managed_group_admin_required",
- "只有申请来源群的管理员可以审批技师。",
- status=403,
- )
- return values
- async def require_teacher_admin(actor: Any, user_id: int) -> dict[str, Any]:
- profile = await get_directory_profile(int(user_id))
- if not profile:
- raise DirectoryServiceError("profile_not_found", "未找到该成员资料。", status=404)
- if is_sudoer(int(actor.id)):
- return profile
- source_chat_id = profile.get("application_source_chat_id")
- source_bot_id = profile.get("application_source_bot_id")
- if not source_chat_id:
- await require_platform_membership(int(actor.id), require_admin=True)
- return profile
- if source_bot_id and str(source_bot_id) != str(BOT_PROFILE_ID):
- raise DirectoryServiceError(
- "source_bot_required",
- "请使用接收技师申请提醒的机器人完成审批。",
- status=403,
- )
- try:
- member = await app.get_chat_member(int(source_chat_id), int(actor.id))
- except Exception as exc:
- raise DirectoryServiceError(
- "admin_verification_unavailable",
- "暂时无法核验申请来源群的管理员身份。",
- status=503,
- ) from exc
- if getattr(member, "status", None) not in ADMIN_MEMBER_STATUSES:
- raise DirectoryServiceError(
- "source_chat_admin_required",
- "只有申请来源群的管理员可以执行此操作。",
- status=403,
- )
- return profile
- async def admin_decide_teacher(
- *,
- actor: Any,
- user_id: int,
- action: str,
- reason: str = "",
- ) -> dict[str, Any]:
- profile = await require_teacher_admin(actor, int(user_id))
- identity = user_identity(actor)
- try:
- return await decide_teacher_application(
- user_id=int(user_id),
- action=action,
- actor_id=int(actor.id),
- actor_name=identity["display_name"],
- source="telegram_private",
- reason=reason,
- authorization_chat_id=profile.get("application_source_chat_id"),
- )
- except DirectoryDataError as exc:
- raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
- async def admin_change_teacher_state(
- *,
- actor: Any,
- user_id: int,
- action: str,
- reason: str,
- ) -> dict[str, Any]:
- await require_teacher_admin(actor, int(user_id))
- if not reason.strip():
- raise DirectoryServiceError("reason_required", "必须填写操作原因。")
- identity = user_identity(actor)
- try:
- profile, _ = await set_teacher_state(
- user_id=int(user_id),
- action=action,
- actor_id=int(actor.id),
- actor_name=identity["display_name"],
- source="telegram_private",
- reason=reason,
- )
- return profile
- except DirectoryDataError as exc:
- raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
- async def directory_for_user(
- *,
- user: Any,
- radius_km: int | None,
- page: int = 1,
- page_size: int = 10,
- ) -> tuple[list[dict[str, Any]], int]:
- await require_platform_membership(int(user.id))
- await observe_directory_user(user)
- location = await get_directory_location(int(user.id))
- if not location:
- raise DirectoryServiceError(
- "location_required",
- "请先使用 Telegram 定位或手动选择位置,再查看技师。",
- status=409,
- )
- try:
- return await list_directory_teachers(
- longitude=float(location["longitude"]),
- latitude=float(location["latitude"]),
- max_distance_meters=(
- float(radius_km) * 1000 if radius_km is not None else None
- ),
- page=max(1, int(page)),
- page_size=max(1, min(int(page_size), 20)),
- )
- except DirectoryDataError as exc:
- raise DirectoryServiceError(exc.code, str(exc)) from exc
- def public_teacher(profile: dict[str, Any]) -> dict[str, Any]:
- return {
- "user_id": int(profile["user_id"]),
- "username": profile.get("username"),
- "display_name": profile.get("display_name") or f"用户 {profile['user_id']}",
- "online": teacher_is_online(profile),
- "online_until": aware_utc(profile.get("online_until")),
- "distance_meters": float(profile.get("distance_meters") or 0),
- "location_updated_at": aware_utc(profile.get("location_updated_at")),
- }
- def profile_summary(profile: dict[str, Any], location: dict[str, Any] | None) -> dict[str, Any]:
- return {
- "user_id": int(profile["user_id"]),
- "username": profile.get("username"),
- "display_name": profile.get("display_name") or f"用户 {profile['user_id']}",
- "application_status": profile.get("application_status", "none"),
- "listed": bool(profile.get("listed")),
- "online": teacher_is_online(profile),
- "online_until": aware_utc(profile.get("online_until")),
- "has_location": bool(location),
- "location_source": location.get("source") if location else None,
- "location_updated_at": aware_utc(location.get("updated_at")) if location else None,
- }
- def format_distance(distance_meters: float) -> str:
- if distance_meters < 1000:
- return f"{max(0, round(distance_meters))} 米"
- return f"{distance_meters / 1000:.1f} 公里"
- async def _presence_sweeper() -> None:
- while True:
- try:
- await expire_teacher_presence()
- except Exception as exc:
- log.error(f"技师在线状态过期任务失败:{exc}")
- await asyncio.sleep(60)
- def start_presence_sweeper() -> None:
- global _presence_task
- if _presence_task and not _presence_task.done():
- return
- try:
- _presence_task = asyncio.create_task(
- _presence_sweeper(),
- name="teacher-presence-sweeper",
- )
- except RuntimeError:
- _presence_task = None
- async def stop_presence_sweeper() -> None:
- global _presence_task
- if _presence_task:
- _presence_task.cancel()
- with suppress(asyncio.CancelledError):
- await _presence_task
- _presence_task = None
|