codekingpro/portable-devtools
114k
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 