from __future__ import annotations from datetime import UTC, datetime, timedelta from typing import Any from pyrogram.enums import ChatMembersFilter, ChatMemberStatus from pyrogram.errors import ChatNotModified from pyrogram.types import ChatPermissions, ChatPrivileges from wbb import BOT_ID, BOT_PROFILE_ID, SUDOERS, TELEGRAM_CONNECTED, app from wbb.services.blacklist_enforcement import ( RiskControlValidationError, normalize_risk_rules, replace_legacy_rule_keywords, risk_rule_keywords, sync_legacy_blacklist, ) from wbb.services.member_identity import observe_member_identity from wbb.utils.dbadmin import ( delete_unavailable_managed_chat, get_managed_chat_settings, list_managed_chats, list_stored_invite_links, managed_chat_exists, mark_managed_chat_unavailable, revoke_stored_invite_link, store_invite_link, update_managed_chat_settings, upsert_managed_chat, ) from wbb.utils.dbadmin import ( list_recent_chat_members as list_cached_recent_chat_members, ) from wbb.utils.dbfunctions import ( add_chatbot, add_warn, captcha_off, captcha_on, check_chatbot, del_welcome, delete_filter, flood_off, flood_on, get_blacklisted_words, get_served_chats, get_warn, int_to_alpha, remove_warns, rm_chatbot, save_filter, set_welcome, ) from wbb.utils.i18n import telegram_permission_label PERMISSION_NAMES = ( "can_manage_chat", "can_delete_messages", "can_manage_video_chats", "can_restrict_members", "can_promote_members", "can_change_info", "can_post_messages", "can_edit_messages", "can_invite_users", "can_pin_messages", ) DEFAULT_MEMBER_PERMISSIONS = { "can_send_messages": True, "can_send_media_messages": True, "can_send_other_messages": True, "can_send_polls": True, "can_add_web_page_previews": True, "can_change_info": False, "can_invite_users": True, "can_pin_messages": False, } class ChatManagementError(RuntimeError): def __init__(self, code: str, message: str, *, status: int = 400): super().__init__(message) self.code = code self.status = status def _ensure_telegram_connected() -> None: if not TELEGRAM_CONNECTED: raise ChatManagementError( "telegram_not_connected", "所选机器人尚未连接 Telegram。", status=409, ) def _status_value(status: Any) -> str: return str(getattr(status, "value", status)) def _privilege_list(member: Any) -> list[str]: if member.status == ChatMemberStatus.OWNER: return list(PERMISSION_NAMES) privileges = member.privileges if not privileges: return [] return [name for name in PERMISSION_NAMES if bool(getattr(privileges, name, False))] def serialize_user(user: Any) -> dict[str, Any]: return { "id": str(user.id), "username": user.username, "first_name": user.first_name, "last_name": user.last_name, "is_bot": bool(user.is_bot), "is_deleted": bool(user.is_deleted), } def serialize_member(member: Any) -> dict[str, Any]: return { "user": serialize_user(member.user), "status": _status_value(member.status), "custom_title": member.custom_title, "privileges": _privilege_list(member), "until_date": member.until_date, } async def sync_managed_chat(chat_id: int) -> dict[str, Any]: _ensure_telegram_connected() try: chat = await app.get_chat(int(chat_id)) bot_member = await app.get_chat_member(int(chat_id), BOT_ID) try: member_count = await app.get_chat_members_count(int(chat_id)) except Exception: member_count = None await upsert_managed_chat( chat_id=int(chat.id), title=chat.title, username=chat.username, chat_type=_status_value(chat.type), member_count=member_count, accessible=bot_member.status not in {ChatMemberStatus.LEFT, ChatMemberStatus.BANNED}, bot_status=_status_value(bot_member.status), bot_privileges=_privilege_list(bot_member), ) return { "chat_id": int(chat.id), "title": chat.title, "username": chat.username, "type": _status_value(chat.type), "description": chat.description, "member_count": member_count, "accessible": True, "bot_status": _status_value(bot_member.status), "bot_privileges": _privilege_list(bot_member), "permissions": { name: bool(getattr(chat.permissions, name, False)) for name in DEFAULT_MEMBER_PERMISSIONS } if chat.permissions else dict(DEFAULT_MEMBER_PERMISSIONS), } except Exception as exc: await mark_managed_chat_unavailable(int(chat_id), str(exc)) raise ChatManagementError( "chat_unavailable", "机器人当前无法访问该群。", status=404 ) from exc async def hydrate_managed_chats() -> None: if not TELEGRAM_CONNECTED: return _, total = await list_managed_chats(page=1, page_size=1, chat_type="group") if total: return if BOT_PROFILE_ID != "primary": return for item in await get_served_chats(): chat_id = int(item["chat_id"]) if chat_id >= 0: continue try: await sync_managed_chat(chat_id) except ChatManagementError: continue async def list_accessible_chats( *, query: str = "", page: int = 1, page_size: int = 20, actor_id: int | None = None, ) -> tuple[list[dict[str, Any]], int]: await hydrate_managed_chats() if actor_id is None or actor_id in SUDOERS: return await list_managed_chats(query=query, page=page, page_size=page_size, chat_type="group") all_items: list[dict[str, Any]] = [] source_page = 1 while True: batch, source_total = await list_managed_chats( query=query, page=source_page, page_size=100, chat_type="group" ) all_items.extend(batch) if not batch or len(all_items) >= source_total: break source_page += 1 visible: list[dict[str, Any]] = [] for item in all_items: try: member = await app.get_chat_member(int(item["chat_id"]), actor_id) except Exception: continue if member.status in {ChatMemberStatus.OWNER, ChatMemberStatus.ADMINISTRATOR}: visible.append(item) start = (max(1, page) - 1) * max(1, page_size) return visible[start : start + page_size], len(visible) async def remove_unavailable_chat(chat_id: int) -> dict[str, Any]: if await delete_unavailable_managed_chat(chat_id): return {"chat_id": str(chat_id), "removed": True} if not await managed_chat_exists(chat_id): raise ChatManagementError( "chat_not_found", "群记录不存在或已被移除。", status=404, ) raise ChatManagementError( "chat_still_accessible", "该群当前仍可用,不能从管理列表移除。", status=409, ) async def ensure_permission( chat_id: int, permission: str, *, actor_id: int | None = None, ) -> None: _ensure_telegram_connected() try: bot_member = await app.get_chat_member(int(chat_id), BOT_ID) except Exception as exc: raise ChatManagementError( "bot_not_member", "机器人不在该群或无法读取群状态。", status=403 ) from exc if bot_member.status != ChatMemberStatus.OWNER and permission not in _privilege_list( bot_member ): raise ChatManagementError( "bot_permission_missing", f"机器人缺少 Telegram 权限:{telegram_permission_label(permission)}", status=403, ) if actor_id is None or actor_id in SUDOERS: return try: actor = await app.get_chat_member(int(chat_id), int(actor_id)) except Exception as exc: raise ChatManagementError( "actor_not_admin", "你已不是该群管理员。", status=403 ) from exc if actor.status != ChatMemberStatus.OWNER and permission not in _privilege_list(actor): raise ChatManagementError( "actor_permission_missing", f"你缺少 Telegram 权限:{telegram_permission_label(permission)}", status=403, ) async def get_chat_overview(chat_id: int, actor_id: int | None = None) -> dict[str, Any]: _ensure_telegram_connected() if actor_id is not None and actor_id not in SUDOERS: try: member = await app.get_chat_member(chat_id, actor_id) except Exception as exc: raise ChatManagementError( "actor_not_admin", "你无权管理该群。", status=403 ) from exc if member.status not in { ChatMemberStatus.OWNER, ChatMemberStatus.ADMINISTRATOR, }: raise ChatManagementError( "actor_not_admin", "你无权管理该群。", status=403 ) overview = await sync_managed_chat(chat_id) overview["automation"] = await get_automation_settings(chat_id) return overview async def get_automation_settings(chat_id: int) -> dict[str, Any]: settings = dict(await get_managed_chat_settings(chat_id)) legacy_words = await get_blacklisted_words(chat_id) rules = normalize_risk_rules( settings.get("risk_rules") if "risk_rules" in settings else None, legacy_policy=( None if "risk_rules" in settings else settings.get("risk_control") ), fallback_keywords=legacy_words or settings.get("blacklist_words", []), ) settings.pop("risk_control", None) settings["risk_rules"] = rules settings["blacklist_words"] = risk_rule_keywords(rules) identity_monitor = settings.get("identity_monitor") or {} settings["identity_monitor"] = { "enabled": bool(identity_monitor.get("enabled", False)), "notify_in_chat": bool(identity_monitor.get("notify_in_chat", False)), } return settings async def update_chat_profile( chat_id: int, *, title: str | None = None, description: str | None = None, actor_id: int | None = None, ) -> dict[str, Any]: if title is not None: if not isinstance(title, str): raise ChatManagementError("invalid_title", "群标题必须是文本。") title = title.strip() if not title or len(title) > 128: raise ChatManagementError("invalid_title", "群标题长度需要在 1 到 128 之间。") if description is not None: if not isinstance(description, str): raise ChatManagementError("invalid_description", "群描述必须是文本。") if len(description) > 255: raise ChatManagementError("invalid_description", "群描述不能超过 255 个字符。") await ensure_permission(chat_id, "can_change_info", actor_id=actor_id) overview = await sync_managed_chat(chat_id) title_changed = title is not None and title != overview.get("title") description_changed = ( description is not None and description != (overview.get("description") or "") ) if title_changed: try: await app.set_chat_title(chat_id, title) except ChatNotModified: pass except Exception as exc: raise ChatManagementError( "chat_profile_update_failed", "Telegram 未能更新群标题,请稍后重试。", status=502, ) from exc overview["title"] = title if description_changed: try: await app.set_chat_description(chat_id, description) except ChatNotModified: pass except Exception as exc: message = ( "群标题已更新,但 Telegram 未能更新群描述,请重试。" if title_changed else "Telegram 未能更新群描述,请稍后重试。" ) raise ChatManagementError( "chat_profile_update_failed", message, status=502 ) from exc overview["description"] = description if title is not None: overview["title"] = title if description is not None: overview["description"] = description return overview async def update_chat_permissions( chat_id: int, values: dict[str, Any], *, actor_id: int | None = None ) -> dict[str, Any]: await ensure_permission(chat_id, "can_restrict_members", actor_id=actor_id) permissions = { key: bool(values.get(key, default)) for key, default in DEFAULT_MEMBER_PERMISSIONS.items() } await app.set_chat_permissions(chat_id, ChatPermissions(**permissions)) return await sync_managed_chat(chat_id) async def send_announcement( chat_id: int, *, text: str, media_type: str | None = None, file_id: str | None = None, pin: bool = False, actor_id: int | None = None, ) -> dict[str, Any]: await ensure_permission(chat_id, "can_change_info", actor_id=actor_id) if not text.strip() and not file_id: raise ChatManagementError("empty_announcement", "公告文字或媒体至少填写一项。") method_map = { "photo": app.send_photo, "animation": app.send_animation, "video": app.send_video, "document": app.send_document, } if file_id: method = method_map.get(media_type or "") if not method: raise ChatManagementError("invalid_media", "不支持的公告媒体类型。") keyword = "animation" if media_type == "animation" else media_type sent = await method(chat_id, **{keyword: file_id}, caption=text or None) else: sent = await app.send_message(chat_id, text) if pin: await ensure_permission(chat_id, "can_pin_messages", actor_id=actor_id) await app.pin_chat_message(chat_id, sent.id) return {"message_id": str(sent.id)} async def list_chat_admins(chat_id: int) -> list[dict[str, Any]]: await sync_managed_chat(chat_id) return [ serialize_member(member) async for member in app.get_chat_members( chat_id, filter=ChatMembersFilter.ADMINISTRATORS ) ] async def find_chat_member(chat_id: int, query: str) -> dict[str, Any]: _ensure_telegram_connected() query = query.strip() if not query: raise ChatManagementError("member_query_required", "请输入用户 ID 或 @username。") value: int | str = int(query) if query.lstrip("-").isdigit() else query try: member = await app.get_chat_member(chat_id, value) except Exception as exc: raise ChatManagementError("member_not_found", "未找到该群成员。", status=404) from exc await observe_member_identity(chat_id=chat_id, user=member.user) return serialize_member(member) async def search_chat_members( chat_id: int, query: str, *, limit: int = 20, ) -> list[dict[str, Any]]: _ensure_telegram_connected() normalized_query = query.strip() if not normalized_query: raise ChatManagementError( "member_query_required", "请输入成员昵称、用户名或用户 ID。" ) limit = max(1, min(int(limit), 30)) results: dict[int, dict[str, Any]] = {} if normalized_query.lstrip("-").isdigit() or normalized_query.startswith("@"): try: exact = await find_chat_member(chat_id, normalized_query) except ChatManagementError: exact = None if exact: results[int(exact["user"]["id"])] = exact else: try: members = app.get_chat_members( chat_id, query=normalized_query, limit=limit, filter=ChatMembersFilter.SEARCH, ) if members is not None: async for member in members: await observe_member_identity(chat_id=chat_id, user=member.user) results[int(member.user.id)] = serialize_member(member) except Exception: pass cached = await list_cached_recent_chat_members( chat_id, query=normalized_query, limit=limit, ) for item in cached: user_id = int(item["user_id"]) if user_id in results or len(results) >= limit: continue try: member = await app.get_chat_member(chat_id, user_id) except Exception: continue if member.status in {ChatMemberStatus.LEFT, ChatMemberStatus.BANNED}: continue await observe_member_identity(chat_id=chat_id, user=member.user) results[user_id] = serialize_member(member) if not results: raise ChatManagementError( "member_not_found", "未找到匹配的群成员。", status=404 ) return list(results.values())[:limit] def _serialize_cached_member(item: dict[str, Any]) -> dict[str, Any]: return { "user": { "id": str(item["user_id"]), "username": item.get("username"), "first_name": item.get("first_name"), "last_name": item.get("last_name"), "is_bot": bool(item.get("is_bot", False)), "is_deleted": False, }, "status": "recent", "custom_title": None, "privileges": [], "until_date": None, "last_seen_at": item.get("last_seen_at"), } async def list_recent_members( chat_id: int, *, limit: int = 30, ) -> list[dict[str, Any]]: _ensure_telegram_connected() limit = max(1, min(int(limit), 50)) results: dict[int, dict[str, Any]] = {} try: members = app.get_chat_members( chat_id, limit=limit, filter=ChatMembersFilter.RECENT, ) if members is not None: async for member in members: if member.user.is_bot: continue await observe_member_identity(chat_id=chat_id, user=member.user) results[int(member.user.id)] = serialize_member(member) except Exception: pass cached = await list_cached_recent_chat_members(chat_id, limit=limit) for item in cached: user_id = int(item["user_id"]) if user_id not in results and len(results) < limit: results[user_id] = _serialize_cached_member(item) return list(results.values())[:limit] async def execute_member_action( chat_id: int, *, user_id: int, action: str, reason: str = "", duration_seconds: int | None = None, privileges: dict[str, Any] | None = None, actor_id: int | None = None, ) -> dict[str, Any]: permission = ( "can_promote_members" if action in {"promote", "demote"} else "can_restrict_members" ) await ensure_permission(chat_id, permission, actor_id=actor_id) until_date = datetime.now(UTC) + timedelta(seconds=duration_seconds) if duration_seconds else None if action == "warn": key = await int_to_alpha(user_id) current = await get_warn(chat_id, key) count = int(current.get("warns", 0)) + 1 if current else 1 if count >= 3: await app.ban_chat_member(chat_id, user_id) await remove_warns(chat_id, key) return {"action": "ban", "warns": 3, "reason": reason} await add_warn(chat_id, key, {"warns": count}) return {"action": "warn", "warns": count, "reason": reason} if action == "ban": kwargs = {"until_date": until_date} if until_date else {} await app.ban_chat_member(chat_id, user_id, **kwargs) elif action == "unban": await app.unban_chat_member(chat_id, user_id) elif action == "kick": await app.ban_chat_member(chat_id, user_id) await app.unban_chat_member(chat_id, user_id) elif action == "mute": kwargs = {"until_date": until_date} if until_date else {} await app.restrict_chat_member( chat_id, user_id, ChatPermissions(), **kwargs ) elif action == "unmute": await app.restrict_chat_member( chat_id, user_id, ChatPermissions(**DEFAULT_MEMBER_PERMISSIONS) ) elif action == "promote": values = { key: bool((privileges or {}).get(key, False)) for key in PERMISSION_NAMES if key != "can_manage_chat" } await app.promote_chat_member(chat_id, user_id, ChatPrivileges(**values)) elif action == "demote": values = {key: False for key in PERMISSION_NAMES if key != "can_manage_chat"} await app.promote_chat_member(chat_id, user_id, ChatPrivileges(**values)) else: raise ChatManagementError("invalid_member_action", "不支持的成员操作。") return {"action": action, "user_id": str(user_id), "reason": reason} async def create_invite_link( chat_id: int, *, name: str, expires_at: datetime | None, member_limit: int | None, actor_id: int | None = None, ) -> dict[str, Any]: await ensure_permission(chat_id, "can_invite_users", actor_id=actor_id) link = await app.create_chat_invite_link( chat_id, name=name[:32] or None, expire_date=expires_at, member_limit=member_limit, ) await store_invite_link( chat_id=chat_id, invite_link=link.invite_link, name=link.name or name, expires_at=link.expire_date, ) return { "invite_link": link.invite_link, "name": link.name, "expires_at": link.expire_date, "member_limit": link.member_limit, "is_revoked": link.is_revoked, } async def revoke_invite_link( chat_id: int, invite_link: str, *, actor_id: int | None = None ) -> None: await ensure_permission(chat_id, "can_invite_users", actor_id=actor_id) await app.revoke_chat_invite_link(chat_id, invite_link) await revoke_stored_invite_link(chat_id, invite_link) async def get_invite_links(chat_id: int) -> list[dict[str, Any]]: return await list_stored_invite_links(chat_id) async def apply_automation_settings( chat_id: int, values: dict[str, Any], *, actor_id: int | None = None ) -> dict[str, Any]: await ensure_permission(chat_id, "can_change_info", actor_id=actor_id) previous = await get_automation_settings(chat_id) if "auto_replies" in values: for old in previous.get("auto_replies", []): await delete_filter(chat_id, str(old.get("keyword", ""))) normalized_replies = [] for item in values["auto_replies"][:100]: keyword = str(item.get("keyword") or "").strip().lower() if not keyword: continue reply = { "keyword": keyword[:100], "type": str(item.get("type") or "text"), "text": str(item.get("text") or "")[:4000], "file_id": item.get("file_id"), } normalized_replies.append(reply) await save_filter( chat_id, reply["keyword"], { "type": reply["type"], "data": reply["text"], "file_id": reply["file_id"], }, ) values["auto_replies"] = normalized_replies if "risk_rules" in values: try: rules = normalize_risk_rules(values["risk_rules"], strict=True) except RiskControlValidationError as exc: raise ChatManagementError("invalid_risk_rules", str(exc)) from exc values["risk_rules"] = rules values["blacklist_words"] = await sync_legacy_blacklist( chat_id, risk_rule_keywords(rules, enabled_only=True) ) values.pop("risk_control", None) elif "risk_control" in values: try: rules = normalize_risk_rules( None, legacy_policy=values.pop("risk_control"), fallback_keywords=values.get("blacklist_words", []), strict=True, ) except RiskControlValidationError as exc: raise ChatManagementError("invalid_risk_control", str(exc)) from exc values["risk_rules"] = rules values["blacklist_words"] = await sync_legacy_blacklist( chat_id, risk_rule_keywords(rules, enabled_only=True) ) elif "blacklist_words" in values: rules = replace_legacy_rule_keywords( previous.get("risk_rules", []), values["blacklist_words"] ) values["risk_rules"] = rules values["blacklist_words"] = await sync_legacy_blacklist( chat_id, risk_rule_keywords(rules, enabled_only=True) ) if "identity_monitor" in values: identity_monitor = values["identity_monitor"] if not isinstance(identity_monitor, dict): raise ChatManagementError( "invalid_identity_monitor", "成员资料监控配置格式不正确。" ) values["identity_monitor"] = { "enabled": bool(identity_monitor.get("enabled", False)), "notify_in_chat": bool(identity_monitor.get("notify_in_chat", False)), } if "welcome" in values: welcome = values["welcome"] or {} if welcome.get("enabled"): media = welcome.get("media") or {} media_type = str(media.get("type") or "text").capitalize() await set_welcome( chat_id, media_type, str(welcome.get("text") or "")[:4000], media.get("file_id"), ) else: await del_welcome(chat_id) if "captcha_enabled" in values: await (captcha_on(chat_id) if values["captcha_enabled"] else captcha_off(chat_id)) if "antiflood_enabled" in values: await (flood_on(chat_id) if values["antiflood_enabled"] else flood_off(chat_id)) if "chatbot_enabled" in values: state = await check_chatbot() enabled = chat_id in state.get("bot", []) if values["chatbot_enabled"] and not enabled: await add_chatbot(chat_id) elif not values["chatbot_enabled"] and enabled: await rm_chatbot(chat_id) await update_managed_chat_settings(chat_id, values) return await get_automation_settings(chat_id)