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