Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
telegram.py269 linesDownload Raw Back to document_loaders
1from __future__ import annotations2 3import asyncio4import json5from pathlib import Path6from typing import TYPE_CHECKING, Dict, List, Optional, Union7 8from langchain_core.documents import Document9 10from langchain_community.document_loaders.base import BaseLoader11 12if TYPE_CHECKING:13    import pandas as pd14    from telethon.hints import EntityLike15 16 17def concatenate_rows(row: dict) -> str:18    """Combine message information in a readable format ready to be used."""19    date = row["date"]20    sender = row["from"]21    text = row["text"]22    return f"{sender} on {date}: {text}\n\n"23 24 25class TelegramChatFileLoader(BaseLoader):26    """Load from `Telegram chat` dump."""27 28    def __init__(self, path: Union[str, Path]):29        """Initialize with a path."""30        self.file_path = path31 32    def load(self) -> List[Document]:33        """Load documents."""34        p = Path(self.file_path)35 36        with open(p, encoding="utf8") as f:37            d = json.load(f)38 39        text = "".join(40            concatenate_rows(message)41            for message in d["messages"]42            if message["type"] == "message" and isinstance(message["text"], str)43        )44        metadata = {"source": str(p)}45 46        return [Document(page_content=text, metadata=metadata)]47 48 49def text_to_docs(text: Union[str, List[str]]) -> List[Document]:50    """Convert a string or list of strings to a list of Documents with metadata."""51    from langchain_text_splitters import RecursiveCharacterTextSplitter52 53    text_splitter = RecursiveCharacterTextSplitter(54        chunk_size=800,55        separators=["\n\n", "\n", ".", "!", "?", ",", " ", ""],56        chunk_overlap=20,57    )58 59    if isinstance(text, str):60        # Take a single string as one page61        text = [text]62    page_docs = [Document(page_content=page) for page in text]63 64    # Add page numbers as metadata65    for i, doc in enumerate(page_docs):66        doc.metadata["page"] = i + 167 68    # Split pages into chunks69    doc_chunks = []70 71    for doc in page_docs:72        chunks = text_splitter.split_text(doc.page_content)73        for i, chunk in enumerate(chunks):74            doc = Document(75                page_content=chunk, metadata={"page": doc.metadata["page"], "chunk": i}76            )77            # Add sources a metadata78            doc.metadata["source"] = f"{doc.metadata['page']}-{doc.metadata['chunk']}"79            doc_chunks.append(doc)80    return doc_chunks81 82 83class TelegramChatApiLoader(BaseLoader):84    """Load `Telegram` chat json directory dump."""85 86    def __init__(87        self,88        chat_entity: Optional[EntityLike] = None,89        api_id: Optional[int] = None,90        api_hash: Optional[str] = None,91        username: Optional[str] = None,92        file_path: str = "telegram_data.json",93    ):94        """Initialize with API parameters.95 96        Args:97            chat_entity: The chat entity to fetch data from.98            api_id: The API ID.99            api_hash: The API hash.100            username: The username.101            file_path: The file path to save the data to. Defaults to102                 "telegram_data.json".103        """104        self.chat_entity = chat_entity105        self.api_id = api_id106        self.api_hash = api_hash107        self.username = username108        self.file_path = file_path109 110    async def fetch_data_from_telegram(self) -> None:111        """Fetch data from Telegram API and save it as a JSON file."""112        from telethon.sync import TelegramClient113 114        data = []115        async with TelegramClient(self.username, self.api_id, self.api_hash) as client:116            async for message in client.iter_messages(self.chat_entity):117                is_reply = message.reply_to is not None118                reply_to_id = message.reply_to.reply_to_msg_id if is_reply else None119                data.append(120                    {121                        "sender_id": message.sender_id,122                        "text": message.text,123                        "date": message.date.isoformat(),124                        "message.id": message.id,125                        "is_reply": is_reply,126                        "reply_to_id": reply_to_id,127                    }128                )129 130        with open(self.file_path, "w", encoding="utf-8") as f:131            json.dump(data, f, ensure_ascii=False, indent=4)132 133    def _get_message_threads(self, data: pd.DataFrame) -> dict:134        """Create a dictionary of message threads from the given data.135 136        Args:137            data (pd.DataFrame): A DataFrame containing the conversation \138                data with columns:139                - message.sender_id140                - text141                - date142                - message.id143                - is_reply144                - reply_to_id145 146        Returns:147            dict: A dictionary where the key is the parent message ID and \148                the value is a list of message IDs in ascending order.149        """150 151        def find_replies(parent_id: int, reply_data: pd.DataFrame) -> List[int]:152            """153            Recursively find all replies to a given parent message ID.154 155            Args:156                parent_id (int): The parent message ID.157                reply_data (pd.DataFrame): A DataFrame containing reply messages.158 159            Returns:160                list: A list of message IDs that are replies to the parent message ID.161            """162            # Find direct replies to the parent message ID163            direct_replies = reply_data[reply_data["reply_to_id"] == parent_id][164                "message.id"165            ].tolist()166 167            # Recursively find replies to the direct replies168            all_replies = []169            for reply_id in direct_replies:170                all_replies += [reply_id] + find_replies(reply_id, reply_data)171 172            return all_replies173 174        # Filter out parent messages175        parent_messages = data[~data["is_reply"]]176 177        # Filter out reply messages and drop rows with NaN in 'reply_to_id'178        reply_messages = data[data["is_reply"]].dropna(subset=["reply_to_id"])179 180        # Convert 'reply_to_id' to integer181        reply_messages["reply_to_id"] = reply_messages["reply_to_id"].astype(int)182 183        # Create a dictionary of message threads with parent message IDs as keys and \184        # lists of reply message IDs as values185        message_threads = {186            parent_id: [parent_id] + find_replies(parent_id, reply_messages)187            for parent_id in parent_messages["message.id"]188        }189 190        return message_threads191 192    def _combine_message_texts(193        self, message_threads: Dict[int, List[int]], data: pd.DataFrame194    ) -> str:195        """196        Combine the message texts for each parent message ID based \197            on the list of message threads.198 199        Args:200            message_threads (dict): A dictionary where the key is the parent message \201                ID and the value is a list of message IDs in ascending order.202            data (pd.DataFrame): A DataFrame containing the conversation data:203                - message.sender_id204                - text205                - date206                - message.id207                - is_reply208                - reply_to_id209 210        Returns:211            str: A combined string of message texts sorted by date.212        """213        combined_text = ""214 215        # Iterate through sorted parent message IDs216        for parent_id, message_ids in message_threads.items():217            # Get the message texts for the message IDs and sort them by date218            message_texts = (219                data[data["message.id"].isin(message_ids)]220                .sort_values(by="date")["text"]221                .tolist()222            )223            message_texts = [str(elem) for elem in message_texts]224 225            # Combine the message texts226            combined_text += " ".join(message_texts) + ".\n"227 228        return combined_text.strip()229 230    def load(self) -> List[Document]:231        """Load documents."""232 233        if self.chat_entity is not None:234            try:235                import nest_asyncio236 237                nest_asyncio.apply()238                asyncio.run(self.fetch_data_from_telegram())239            except ImportError:240                raise ImportError(241                    """`nest_asyncio` package not found.242                    please install with `pip install nest_asyncio`243                    """244                )245 246        p = Path(self.file_path)247 248        with open(p, encoding="utf8") as f:249            d = json.load(f)250        try:251            import pandas as pd252        except ImportError:253            raise ImportError(254                """`pandas` package not found. 255                please install with `pip install pandas`256                """257            )258        normalized_messages = pd.json_normalize(d)259        df = pd.DataFrame(normalized_messages)260 261        message_threads = self._get_message_threads(df)262        combined_texts = self._combine_message_texts(message_threads, df)263 264        return text_to_docs(combined_texts)265 266 267# For backwards compatibility268TelegramChatLoader = TelegramChatFileLoader269 
codekingpro/portable-devtools · Team Ai