from __future__ import annotations import asyncio from pyrogram import filters from pyrogram.enums import ChatMemberStatus, ChatType from wbb import BOT_ID, BOT_PERMISSIONS, app, log from wbb.services.bot_permissions import has_permission from wbb.services.channel_management import ( ensure_channel, observe_channel_deletion, observe_channel_post, sweep_channel_posts, ) from wbb.utils.dbadmin import mark_managed_chat_unavailable _sweeper_started = False @app.on_message(filters.channel, group=-29) async def collect_channel_post(_, message): if not has_permission("channel.posts", BOT_PERMISSIONS): return try: await observe_channel_post(message) except Exception as exc: log.error(f"频道帖子采集失败:{exc}") @app.on_edited_message(filters.channel, group=-29) async def collect_edited_channel_post(_, message): if not has_permission("channel.posts", BOT_PERMISSIONS): return try: await observe_channel_post(message) except Exception as exc: log.error(f"频道帖子编辑采集失败:{exc}") @app.on_deleted_messages(filters.channel, group=-29) async def collect_deleted_channel_posts(_, messages): if not has_permission("channel.posts", BOT_PERMISSIONS): return for message in messages: chat_id = int(getattr(getattr(message, "chat", None), "id", 0) or 0) message_id = int(getattr(message, "id", 0) or 0) if chat_id and message_id: await observe_channel_deletion(chat_id, message_id) @app.on_chat_member_updated(filters.channel, group=-29) async def observe_channel_bot_membership(_, update): member = update.new_chat_member or update.old_chat_member if ( getattr(update.chat, "type", None) != ChatType.CHANNEL or not member or getattr(member.user, "id", None) != BOT_ID ): return if update.new_chat_member and update.new_chat_member.status in { ChatMemberStatus.OWNER, ChatMemberStatus.ADMINISTRATOR, }: try: await ensure_channel(int(update.chat.id)) except Exception as exc: log.error(f"频道自动登记失败:{exc}") else: await mark_managed_chat_unavailable(int(update.chat.id), "Bot is no longer a channel admin") async def _sweep_loop() -> None: await asyncio.sleep(5) while True: try: await sweep_channel_posts() except Exception as exc: log.error(f"频道定时发布失败:{exc}") await asyncio.sleep(15) def _start_sweeper() -> None: global _sweeper_started if _sweeper_started: return try: loop = asyncio.get_running_loop() except RuntimeError: return _sweeper_started = True loop.create_task(_sweep_loop(), name="channel-post-sweeper") _start_sweeper()