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, list_directory_teachers, list_membership_candidates, save_directory_location, set_geocode_consent, 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_private", ) -> tuple[dict[str, Any], bool | None]: await require_platform_membership(int(user.id)) identity = user_identity(user) await observe_directory_user(user) previous = await get_directory_location(int(user.id)) consent = previous.get("geocode_consent") if previous else None location = 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"], geocode_consent=consent, bot_id=str(BOT_PROFILE_ID), ) return location, consent async def update_user_geocode_consent( *, user: Any, consent: bool, ) -> dict[str, Any]: identity = user_identity(user) return await set_geocode_consent( user_id=int(user.id), consent=consent, source="telegram_private", actor_id=int(user.id), actor_name=identity["display_name"], ) 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) -> dict[str, Any]: await require_platform_membership(int(user.id)) await observe_directory_user(user) try: return await submit_teacher_application( user_id=int(user.id), source="telegram_private", bot_id=str(BOT_PROFILE_ID), ) 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_decide_teacher( *, actor: Any, user_id: int, action: str, reason: str = "", ) -> dict[str, Any]: authorization = await require_platform_membership( int(actor.id), require_admin=True, ) 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=authorization.get("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_platform_membership(int(actor.id), require_admin=True) 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", "请先点击“更新位置”发送定位,再查看老师。", 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]: region = profile.get("region") or {} 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")), "region": str(region.get("label") or "区域暂不可用"), "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), "region": ( str((location.get("region") or {}).get("label") or "区域暂不可用") if location else "" ), "geocode_consent": location.get("geocode_consent") 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