import pyrogram import aiohttp import json from async_worker.async_worker import AsyncTask from pyrogram.api.types import ChannelMessagesFilterEmpty from pyrogram.api.types.updates import ChannelDifference, ChannelDifferenceTooLong, ChannelDifferenceEmpty from pyrogram.api.functions.updates.get_channel_difference import GetChannelDifference from pyrogram.client.types.messages_and_media import Message as MessagePyrogram class JsonSerializerFixes: @staticmethod def user(obj): obj.type = "private" return obj @staticmethod def user_type(obj): obj.type = "channel" if obj.id < 0 else "private" return obj class JsonSerializer: fixes = { "from_user": { "new_name": "from", "patch": JsonSerializerFixes.user }, "user": { "new_name": "user", "patch": JsonSerializerFixes.user_type } } @staticmethod def default(obj): if isinstance(obj, bytes): return repr(obj) cls = JsonSerializer result = {} for name in filter(lambda x: not x.startswith("_"), obj.__dict__): value = getattr(obj, name) if value is None: continue if name in cls.fixes: value = cls.fixes[name]["patch"](value) name = cls.fixes[name]["new_name"] result[name] = value return result class ChannelHistoryReadTask(AsyncTask): channel: pyrogram.Chat client: pyrogram.Client pts: int webhook: str http: aiohttp.ClientSession def setup(self, client: pyrogram.Client, channel: pyrogram.Chat, webhook: str): self.client = client self.channel = channel self.pts = False self.webhook = webhook self.http = aiohttp.ClientSession() async def process(self): response = await self.client.send( GetChannelDifference( channel=self.channel, filter=ChannelMessagesFilterEmpty(), pts=self.pts if self.pts else 0xFFFFFFF, limit=0xFFFFFFF, force=True ) ) if isinstance(response, ChannelDifference): self.pts = response.pts users = {i.id: i for i in response.users} chats = {i.id: i for i in response.chats} for message in response.new_messages: message = await MessagePyrogram._parse(self.client, message, users, chats) message = {"update_id": 1, "message": message} data = json.dumps( message, default=JsonSerializer.default, ensure_ascii=True, allow_nan=False, check_circular=True, sort_keys=False ) result = await self.http.post( self.webhook, data=data, headers=[ ("Content-Type", "application/json") ] ) await result.read() result.close() if not response.final: return 1 return response.timeout if isinstance(response, ChannelDifferenceEmpty): self.pts = response.pts return response.timeout if isinstance(response, ChannelDifferenceTooLong): return False