from __future__ import annotations import base64 from datetime import UTC, datetime, timedelta from pathlib import Path from typing import Any from uuid import uuid4 from aiohttp import ClientSession, ClientTimeout from pymongo.errors import DuplicateKeyError from pyrogram.enums import ChatMemberStatus, ChatType, ParseMode from pyrogram.errors import RPCError from pyrogram.types import ( InputMediaAnimation, InputMediaDocument, InputMediaPhoto, InputMediaVideo, ) import wbb from wbb import BOT_ID, BOT_PERMISSIONS, BOT_PROFILE_ID, app, log from wbb.services.bot_permissions import has_permission from wbb.services.chat_management import ( ChatManagementError, _privilege_list, ensure_permission, serialize_member, sync_managed_chat, update_chat_profile, ) from wbb.utils.dbadmin import list_managed_chats, record_audit, upsert_managed_chat from wbb.utils.dbchannel import ( claim_due_posts, create_post, ensure_channel_indexes, get_post, mark_stale_sends_uncertain, now_utc, postsdb, telegram_key, update_scheduled_post, ) MEDIA_METHODS = { "photo": ("send_photo", "photo", InputMediaPhoto), "animation": ("send_animation", "animation", InputMediaAnimation), "video": ("send_video", "video", InputMediaVideo), "document": ("send_document", "document", InputMediaDocument), } CHANNEL_ADMIN_PRIVILEGES = { "can_change_info", "can_post_messages", "can_edit_messages", "can_delete_messages", "can_invite_users", "can_promote_members", "can_manage_video_chats", } def telegram_error_detail(exc: Exception) -> str: return f"{type(exc).__name__}: {str(exc)[:180]}" if str(exc) else type(exc).__name__ async def channel_bot_api(method: str, **params: Any) -> Any: token = str(getattr(wbb, "BOT_TOKEN", "") or "") if not token: raise ChatManagementError("bot_token_missing", "机器人令牌不可用。", status=503) data = { key: ("true" if value else "false") if isinstance(value, bool) else str(value) for key, value in params.items() } try: async with ClientSession(timeout=ClientTimeout(total=15)) as session: async with session.post(f"https://api.telegram.org/bot{token}/{method}", data=data) as response: payload = await response.json(content_type=None) except Exception as exc: raise ChatManagementError( "channel_telegram_unavailable", "Telegram 管理员接口暂时不可用。", status=502 ) from exc if not isinstance(payload, dict) or not payload.get("ok"): detail = str(payload.get("description") or "请求失败")[:180] if isinstance(payload, dict) else "响应无效" raise ChatManagementError( "channel_telegram_failed", f"Telegram 拒绝管理员操作:{detail}", status=502 ) return payload.get("result") def serialize_bot_api_member(member: dict[str, Any]) -> dict[str, Any]: user = member.get("user") or {} return { "user": { "id": str(user.get("id") or ""), "username": user.get("username"), "first_name": user.get("first_name"), "last_name": user.get("last_name"), "is_bot": bool(user.get("is_bot")), "is_deleted": bool(user.get("is_deleted")), }, "status": "owner" if member.get("status") == "creator" else member.get("status"), "custom_title": member.get("custom_title"), "privileges": [name for name in sorted(CHANNEL_ADMIN_PRIVILEGES) if member.get(name)], "until_date": member.get("until_date"), } def public_post(post: dict[str, Any]) -> dict[str, Any]: return {key: value for key, value in post.items() if key not in {"_id", "telegram_key"}} async def warm_channel_peer(chat_id: int) -> None: try: resolver = getattr(app, "resolve_peer", None) if resolver: await resolver(chat_id) else: await app.get_chat(chat_id) return except ValueError as exc: if "Peer id invalid" not in str(exc): raise token = str(getattr(wbb, "BOT_TOKEN", "") or "") if not token: raise ChatManagementError("channel_peer_unavailable", "Bot 尚未识别该频道,请等待频道新帖后重试。", status=409) try: async with ClientSession(timeout=ClientTimeout(total=10)) as session: async with session.post( f"https://api.telegram.org/bot{token}/getChat", data={"chat_id": str(chat_id)} ) as response: data = await response.json() chat = data.get("result") if data.get("ok") else None if not chat or int(chat.get("id", 0)) != chat_id or chat.get("type") != "channel": raise ChatManagementError("channel_unavailable", "Telegram 未能确认该频道及 Bot 的访问权限。", status=404) username = chat.get("username") if not username: raise ChatManagementError( "channel_peer_unavailable", "Bot 尚未缓存该私有频道;请在频道发布新帖后重试。", status=409 ) await app.get_chat(f"@{username}") except ChatManagementError: raise except Exception as exc: raise ChatManagementError("channel_peer_unavailable", "Bot 暂时无法解析频道 ID,请稍后重试。", status=502) from exc async def ensure_channel( chat_id: int, *, admin: bool = True, include_photo: bool = False ) -> dict[str, Any]: await warm_channel_peer(chat_id) overview = await sync_managed_chat(chat_id) if overview["type"] != "channel": raise ChatManagementError("not_channel", "目标不是频道。", status=400) if admin and overview["bot_status"] not in {"owner", "administrator"}: raise ChatManagementError("bot_not_admin", "机器人不是该频道管理员。", status=403) chat = await app.get_chat(chat_id) photo = getattr(chat, "photo", None) overview["photo_file_id"] = getattr(photo, "small_file_id", None) if include_photo and overview["photo_file_id"]: try: content = await app.download_media(overview["photo_file_id"], in_memory=True) overview["photo_data_url"] = "data:image/jpeg;base64," + base64.b64encode( content.getvalue() ).decode("ascii") except Exception: overview["photo_data_url"] = None return overview async def list_channels(*, query: str, page: int, page_size: int) -> tuple[list[dict], int]: return await list_managed_chats(query=query, page=page, page_size=page_size, chat_type="channel") async def update_channel_profile(chat_id: int, *, title: str | None, description: str | None) -> dict: await ensure_channel(chat_id) try: return await update_chat_profile(chat_id, title=title, description=description) except ChatManagementError as exc: raise ChatManagementError(exc.code, str(exc).replace("群", "频道"), status=exc.status) from exc async def update_channel_photo(chat_id: int, photo_path: Path) -> dict: await ensure_channel(chat_id) await ensure_permission(chat_id, "can_change_info") try: await app.set_chat_photo(chat_id, photo=str(photo_path)) except RPCError as exc: raise ChatManagementError( "channel_photo_failed", f"Telegram 拒绝更新频道头像:{telegram_error_detail(exc)}", status=502 ) from exc except Exception as exc: raise ChatManagementError("channel_photo_failed", "Telegram 未能更新频道头像。", status=502) from exc return await ensure_channel(chat_id, include_photo=True) async def list_channel_admins(chat_id: int) -> list[dict]: await ensure_channel(chat_id) from pyrogram.enums import ChatMembersFilter try: return [ serialize_member(member) async for member in app.get_chat_members(chat_id, filter=ChatMembersFilter.ADMINISTRATORS) ] except Exception as exc: raise ChatManagementError("admin_list_failed", "Telegram 未能读取频道管理员。", status=502) from exc async def list_channel_admin_candidates(chat_id: int, *, query: str = "", limit: int = 30) -> list[dict]: await ensure_channel(chat_id) await ensure_permission(chat_id, "can_promote_members") from pyrogram.enums import ChatMembersFilter query = query.strip() if len(query) > 64: raise ChatManagementError("member_query_too_long", "搜索内容不能超过 64 个字符。") limit = max(1, min(limit, 30)) if query.isdigit(): member = await channel_bot_api("getChatMember", chat_id=chat_id, user_id=int(query)) if member.get("status") != "member" or (member.get("user") or {}).get("is_bot"): return [] return [serialize_bot_api_member(member)] results: list[dict] = [] try: members = app.get_chat_members( chat_id, query=query.lstrip("@"), limit=limit, filter=ChatMembersFilter.SEARCH if query else ChatMembersFilter.RECENT, ) async for member in members: if member.status != ChatMemberStatus.MEMBER or member.user.is_bot: continue results.append(serialize_member(member)) except RPCError as exc: raise ChatManagementError( "member_list_failed", f"Telegram 拒绝读取频道订阅者:{telegram_error_detail(exc)}", status=502 ) from exc except Exception as exc: raise ChatManagementError("member_list_failed", "Telegram 未能读取频道订阅者。", status=502) from exc return results async def set_channel_admin(chat_id: int, user_id: int, privileges: list[str]) -> dict: await ensure_channel(chat_id) await ensure_permission(chat_id, "can_promote_members") requested = set(privileges) if not requested or requested - CHANNEL_ADMIN_PRIVILEGES: raise ChatManagementError("invalid_privileges", "请选择有效的频道管理员权限。") try: bot = await app.get_chat_member(chat_id, BOT_ID) except Exception as exc: raise ChatManagementError("bot_permission_unavailable", "Telegram 未能读取机器人权限。", status=502) from exc allowed = set(_privilege_list(bot)) if bot.status != ChatMemberStatus.OWNER and requested - allowed: raise ChatManagementError("privilege_exceeds_bot", "不能授予超出机器人自身的权限。", status=403) target = await channel_bot_api("getChatMember", chat_id=chat_id, user_id=user_id) if target.get("status") in {"creator", "owner"}: raise ChatManagementError("owner_protected", "不能修改频道所有者。", status=403) values = {key: key in requested for key in CHANNEL_ADMIN_PRIVILEGES} await channel_bot_api("promoteChatMember", chat_id=chat_id, user_id=user_id, **values) return serialize_bot_api_member( await channel_bot_api("getChatMember", chat_id=chat_id, user_id=user_id) ) async def remove_channel_admin(chat_id: int, user_id: int) -> dict: await ensure_channel(chat_id) await ensure_permission(chat_id, "can_promote_members") target = await channel_bot_api("getChatMember", chat_id=chat_id, user_id=user_id) if target.get("status") in {"creator", "owner"}: raise ChatManagementError("owner_protected", "不能移除频道所有者。", status=403) if target.get("status") != "administrator": raise ChatManagementError("not_admin", "该用户不是频道管理员。", status=409) await channel_bot_api( "promoteChatMember", chat_id=chat_id, user_id=user_id, **{key: False for key in CHANNEL_ADMIN_PRIVILEGES}, ) return {"user_id": str(user_id), "removed": True} def validate_post_values( values: dict[str, Any], *, scheduled: bool = False, keep_publish_at: bool = False ) -> dict[str, Any]: text = str(values.get("text") or "").strip() media_type = values.get("media_type") or None file_id = values.get("file_id") or None if media_type not in (None, *MEDIA_METHODS) or bool(media_type) != bool(file_id): raise ChatManagementError("invalid_media", "媒体类型与文件 ID 不匹配。") if not text and not file_id: raise ChatManagementError("empty_post", "帖子文字或媒体至少填写一项。") if len(text) > (1024 if file_id else 4096): raise ChatManagementError("post_too_long", "文字超过 Telegram 长度限制。") media_filename = str(values.get("media_filename") or "").strip() if file_id else "" media_mime_type = str(values.get("media_mime_type") or "").strip() if file_id else "" if len(media_filename) > 255 or len(media_mime_type) > 100: raise ChatManagementError("invalid_media_metadata", "媒体文件信息过长。") try: media_size = int(values["media_size"]) if file_id and values.get("media_size") is not None else None except (TypeError, ValueError) as exc: raise ChatManagementError("invalid_media_metadata", "媒体文件大小无效。") from exc if media_size is not None and media_size < 0: raise ChatManagementError("invalid_media_metadata", "媒体文件大小无效。") publish_at = values.get("publish_at") if publish_at is not None: if isinstance(publish_at, str): try: publish_at = datetime.fromisoformat(publish_at.replace("Z", "+00:00")) except ValueError as exc: raise ChatManagementError("invalid_publish_at", "定时时间无效。") from exc if publish_at.tzinfo is None: raise ChatManagementError("invalid_publish_at", "定时时间须包含时区。") if not isinstance(publish_at, datetime): raise ChatManagementError("invalid_publish_at", "定时时间须包含时区。") if publish_at.tzinfo is None: publish_at = publish_at.replace(tzinfo=UTC) publish_at = publish_at.astimezone(UTC) if not keep_publish_at and publish_at <= now_utc() + timedelta(seconds=5): raise ChatManagementError("invalid_publish_at", "定时时间至少晚于当前 5 秒。") elif scheduled: raise ChatManagementError("invalid_publish_at", "请填写定时时间。") return { "text": text, "media_type": media_type, "file_id": file_id, "media_filename": media_filename or None, "media_size": media_size, "media_mime_type": media_mime_type or None, "pin": bool(values.get("pin", False)), "publish_at": publish_at, } def same_publish_at(value: Any, current: Any) -> bool: if isinstance(value, str): try: value = datetime.fromisoformat(value.replace("Z", "+00:00")) except ValueError: return False if not isinstance(value, datetime) or not isinstance(current, datetime): return False value = value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC) current = current.replace(tzinfo=UTC) if current.tzinfo is None else current.astimezone(UTC) return value == current async def publish_post(post: dict[str, Any]) -> dict: if not has_permission("channel.posts", BOT_PERMISSIONS): raise ChatManagementError("bot_role_permission_denied", "所选机器人的职责角色不允许发布频道帖子。", status=403) chat_id = int(post["chat_id"]) await ensure_channel(chat_id) await ensure_permission(chat_id, "can_post_messages") if post.get("pin"): await ensure_permission(chat_id, "can_edit_messages") try: if post.get("file_id"): method_name, keyword, _ = MEDIA_METHODS[post["media_type"]] sent = await getattr(app, method_name)( chat_id, **{keyword: post["file_id"]}, caption=post["text"] or None, parse_mode=ParseMode.DISABLED, ) else: sent = await app.send_message(chat_id, post["text"], parse_mode=ParseMode.DISABLED) except RPCError as exc: detail = telegram_error_detail(exc) await postsdb.update_one( {"post_id": post["post_id"], "status": "sending"}, {"$set": {"status": "failed", "last_error": detail, "updated_at": now_utc()}}, ) raise ChatManagementError("post_send_failed", f"Telegram 拒绝发布:{detail}", status=502) from exc except Exception as exc: await postsdb.update_one( {"post_id": post["post_id"], "status": "sending"}, {"$set": {"status": "uncertain", "last_error": "Telegram 发送返回异常,结果需要人工核实。", "updated_at": now_utc()}}, ) raise ChatManagementError( "post_delivery_uncertain", "Telegram 发送结果不明,已停止自动重试;请在频道核实。", status=502 ) from exc key = telegram_key(chat_id, int(sent.id)) try: for attempt in range(3): await postsdb.delete_one({"telegram_key": key, "source": "telegram"}) try: await postsdb.update_one( {"post_id": post["post_id"]}, {"$set": { "telegram_key": key, "message_id": int(sent.id), "status": "published", "published_at": now_utc(), "updated_at": now_utc(), }}, ) break except DuplicateKeyError: if attempt == 2: raise except Exception as exc: raise ChatManagementError( "post_delivery_uncertain", "消息已发出,但记录更新失败;请在频道核实,系统不会自动重发。", status=502 ) from exc if post.get("pin"): try: await app.pin_chat_message(chat_id, int(sent.id)) except Exception: await postsdb.update_one( {"post_id": post["post_id"]}, {"$set": {"pin_error": "Telegram 未能置顶该消息。"}}, ) return public_post(await get_post(chat_id, post["post_id"])) async def create_channel_post(chat_id: int, values: dict[str, Any]) -> dict: await ensure_channel(chat_id) await ensure_permission(chat_id, "can_post_messages") data = validate_post_values(values) if data["pin"]: await ensure_permission(chat_id, "can_edit_messages") post = await create_post(chat_id, **data) if data["publish_at"]: return public_post(post) try: return await publish_post(post) except ChatManagementError as exc: if exc.code != "post_delivery_uncertain": await postsdb.update_one( {"post_id": post["post_id"], "status": "sending"}, {"$set": {"status": "failed", "last_error": str(exc)[:300], "updated_at": now_utc()}}, ) raise async def edit_channel_post(chat_id: int, post_id: str, values: dict[str, Any]) -> dict: await ensure_channel(chat_id) current = await get_post(chat_id, post_id) if not current: raise ChatManagementError("post_not_found", "帖子不存在。", status=404) merged = {key: current.get(key) for key in ( "text", "media_type", "file_id", "media_filename", "media_size", "media_mime_type", "pin", "publish_at" )} if current["status"] != "scheduled": merged["publish_at"] = None merged.update(values) if values.get("file_id") and values["file_id"] != current.get("file_id"): for key in ("media_filename", "media_size", "media_mime_type"): if key not in values: merged[key] = None data = validate_post_values( merged, scheduled=current["status"] == "scheduled", keep_publish_at=current["status"] == "scheduled" and ( "publish_at" not in values or same_publish_at(values["publish_at"], current.get("publish_at")) ), ) if current["status"] == "scheduled": await ensure_permission(chat_id, "can_post_messages") updated = await update_scheduled_post(chat_id, post_id, data) if not updated: raise ChatManagementError("post_in_progress", "帖子已进入发送流程,不能编辑。", status=409) return public_post(updated) if current["status"] != "published": raise ChatManagementError("post_not_editable", "当前状态不能编辑帖子。", status=409) if "publish_at" in values or "pin" in values: raise ChatManagementError("invalid_edit", "已发布帖子不能修改定时或置顶设置。") await ensure_permission(chat_id, "can_edit_messages") message_id = int(current["message_id"]) try: if data["media_type"]: if not current.get("media_type"): raise ChatManagementError("media_change_unsupported", "文字帖不能改为媒体帖。") if data["file_id"] != current.get("file_id") or data["media_type"] != current.get("media_type"): _, _, media_class = MEDIA_METHODS[data["media_type"]] await app.edit_message_media( chat_id, message_id, media_class( data["file_id"], caption=data["text"], parse_mode=ParseMode.DISABLED ) ) else: await app.edit_message_caption( chat_id, message_id, data["text"], parse_mode=ParseMode.DISABLED ) else: if current.get("media_type"): raise ChatManagementError("media_change_unsupported", "媒体帖不能改为纯文字帖。") await app.edit_message_text( chat_id, message_id, data["text"], parse_mode=ParseMode.DISABLED ) except ChatManagementError: raise except RPCError as exc: raise ChatManagementError( "post_edit_failed", f"Telegram 拒绝编辑帖子:{telegram_error_detail(exc)}", status=502 ) from exc except Exception as exc: raise ChatManagementError("post_edit_failed", "Telegram 未能编辑帖子;请检查 Bot 权限与消息状态。", status=502) from exc await postsdb.update_one( {"post_id": post_id, "bot_id": BOT_PROFILE_ID}, {"$set": {key: data[key] for key in ( "text", "media_type", "file_id", "media_filename", "media_size", "media_mime_type" )} | {"updated_at": now_utc()}}, ) return public_post(await get_post(chat_id, post_id)) async def delete_channel_post(chat_id: int, post_id: str) -> dict: await ensure_channel(chat_id) current = await get_post(chat_id, post_id) if not current: raise ChatManagementError("post_not_found", "帖子不存在。", status=404) if current["status"] == "scheduled": updated = await update_scheduled_post(chat_id, post_id, {"status": "canceled"}) if not updated: raise ChatManagementError("post_in_progress", "帖子已进入发送流程,不能取消。", status=409) return public_post(updated) if current["status"] != "published": raise ChatManagementError("post_not_deletable", "当前状态不能删除帖子。", status=409) await ensure_permission(chat_id, "can_delete_messages") try: deleted = await app.delete_messages(chat_id, int(current["message_id"])) if deleted == 0: raise RuntimeError("Telegram did not delete the message") except RPCError as exc: raise ChatManagementError( "post_delete_failed", f"Telegram 拒绝删除帖子:{telegram_error_detail(exc)}", status=502 ) from exc except Exception as exc: raise ChatManagementError("post_delete_failed", "Telegram 未能删除帖子;请检查 Bot 权限与消息状态。", status=502) from exc await postsdb.update_one( {"post_id": post_id, "bot_id": BOT_PROFILE_ID}, {"$set": {"status": "deleted", "updated_at": now_utc()}}, ) return public_post(await get_post(chat_id, post_id)) async def observe_channel_post(message: Any) -> None: if getattr(message.chat, "type", None) != ChatType.CHANNEL: return chat_id, message_id = int(message.chat.id), int(message.id) bot = await app.get_chat_member(chat_id, BOT_ID) if bot.status not in {ChatMemberStatus.OWNER, ChatMemberStatus.ADMINISTRATOR}: return await upsert_managed_chat( chat_id=chat_id, title=message.chat.title, username=message.chat.username, chat_type="channel", bot_status=str(getattr(bot.status, "value", bot.status)), bot_privileges=_privilege_list(bot), ) media_type = next((name for name in MEDIA_METHODS if getattr(message, name, None)), None) media = getattr(message, media_type, None) if media_type else None key = telegram_key(chat_id, message_id) values = { "text": getattr(message, "text", None) or getattr(message, "caption", None) or "", "media_type": media_type or ("other" if getattr(message, "media", None) else None), "file_id": getattr(media, "file_id", None), "media_filename": getattr(media, "file_name", None), "media_size": getattr(media, "file_size", None), "media_mime_type": getattr(media, "mime_type", None), "updated_at": now_utc(), } await ensure_channel_indexes() update = { "$set": values, "$setOnInsert": { "post_id": uuid4().hex, "bot_id": BOT_PROFILE_ID, "chat_id": chat_id, "telegram_key": key, "message_id": message_id, "source": "telegram", "status": "published", "pin": False, "created_at": getattr(message, "date", None) or now_utc(), "published_at": getattr(message, "date", None) or now_utc(), }, } try: await postsdb.update_one({"telegram_key": key}, update, upsert=True) except DuplicateKeyError: await postsdb.update_one({"telegram_key": key}, {"$set": values}) async def observe_channel_deletion(chat_id: int, message_id: int) -> None: await ensure_channel_indexes() await postsdb.update_one( {"telegram_key": telegram_key(chat_id, message_id)}, {"$set": {"status": "deleted", "updated_at": now_utc()}}, ) async def sweep_channel_posts() -> None: await mark_stale_sends_uncertain() for post in await claim_due_posts(): try: await publish_post(post) except Exception as exc: if isinstance(exc, ChatManagementError) and exc.code != "post_delivery_uncertain": await postsdb.update_one( {"post_id": post["post_id"], "status": "sending"}, {"$set": {"status": "failed", "last_error": str(exc)[:300], "updated_at": now_utc()}}, ) try: await record_audit( source="system", actor_id=None, actor_name="频道定时发布", action="channel.post.publish", chat_id=int(post["chat_id"]), target_id=post["post_id"], summary="定时发布失败或结果不明", success=False, error=str(exc)[:300], ) except Exception as audit_exc: log.error(f"频道定时发布审计写入失败:{audit_exc}") else: try: await record_audit( source="system", actor_id=None, actor_name="频道定时发布", action="channel.post.publish", chat_id=int(post["chat_id"]), target_id=post["post_id"], summary="定时帖子已发布", ) except Exception as audit_exc: log.error(f"频道定时发布审计写入失败:{audit_exc}")