From 7af597b8009fe45fd089d6fb62a1542e82ca1c12 Mon Sep 17 00:00:00 2001 From: mr-forust Date: Thu, 11 Jun 2026 18:19:09 +0200 Subject: [PATCH 1/3] sidebar toggle + css --- webui/index.html | 2 +- webui/style.css | 120 ++++++++++++++++++++++++++++++++++++++++++++-- webui/viewer.html | 3 +- webui/viewer.js | 14 ++++++ 4 files changed, 132 insertions(+), 7 deletions(-) diff --git a/webui/index.html b/webui/index.html index f5a2f49..faea273 100644 --- a/webui/index.html +++ b/webui/index.html @@ -4,7 +4,7 @@ Telegram Scraper Control Panel - +
diff --git a/webui/style.css b/webui/style.css index 492b881..3898a00 100644 --- a/webui/style.css +++ b/webui/style.css @@ -47,7 +47,8 @@ button { } .app-shell { - min-height: 100vh; + height: 100vh; + overflow: hidden; display: grid; grid-template-columns: 300px minmax(0, 1fr); } @@ -59,6 +60,7 @@ button { } .content { + overflow-y: auto; padding: 28px; max-width: 1320px; width: 100%; @@ -494,6 +496,10 @@ tbody tr:hover { overflow-y: auto; } +.sidebar-toggle { + display: none; +} + .viewer-sidebar-head { display: flex; align-items: flex-start; @@ -814,6 +820,10 @@ tbody tr:hover { grid-template-columns: 1fr; } + .app-shell { + grid-template-rows: auto 1fr; + } + .sidebar, .viewer-sidebar { border-right: 0; @@ -825,14 +835,114 @@ tbody tr:hover { } } -@media (max-width: 720px) { - .content, - .viewer-main, - .sidebar, +/* ---------- Mobile: tablet & phone ---------- */ + +@media (max-width: 768px) { + .sidebar-toggle { + display: inline-flex; + } + .viewer-sidebar { + position: fixed; + left: -300px; + top: 0; + bottom: 0; + width: 280px; + z-index: 100; + transition: left 0.2s ease; + border-right: 1px dashed var(--line); + border-bottom: 0; + } + + .viewer-sidebar.open { + left: 0; + } + + .viewer-sidebar.open::before { + content: ""; + position: fixed; + inset: 0; + background: rgba(0, 0, 0, 0.5); + z-index: -1; + } + + .viewer-header { + padding: 12px 14px; + } + + .messages-list { + padding: 8px 12px 16px; + } + + .chat-bubble { + max-width: 92%; + } + + .chat-bubble::before { + display: none; + } + + .content { padding: 18px; } + .viewer-sidebar-head h1 { + font-size: 1.2rem; + } +} + +@media (max-width: 480px) { + .content { + padding: 12px; + } + + .viewer-main { + padding: 0; + } + + .viewer-header { + padding: 10px 10px; + } + + .viewer-header h2 { + font-size: 1.1rem; + } + + .messages-list { + padding: 6px 8px 12px; + } + + .chat-bubble { + max-width: 96%; + padding: 6px 9px; + font-size: 0.92rem; + } + + .chat-footer { + font-size: 0.65rem; + } + + .chat-sender { + font-size: 0.75rem; + } + + .chat-reply-author { + font-size: 0.72rem; + } + + .chat-reply-text { + font-size: 0.75rem; + } + + .sidebar { + padding: 16px; + } + + .viewer-sidebar { + width: 260px; + padding: 16px; + } + .panel-header, .viewer-header, .dialog-header { diff --git a/webui/viewer.html b/webui/viewer.html index 718e9cf..1dfe75e 100644 --- a/webui/viewer.html +++ b/webui/viewer.html @@ -4,7 +4,7 @@ Telegram Scraper Viewer - +
@@ -14,6 +14,7 @@
Export Viewer

Messages

+ Dashboard
diff --git a/webui/viewer.js b/webui/viewer.js index deb7f83..4e211b0 100644 --- a/webui/viewer.js +++ b/webui/viewer.js @@ -81,6 +81,7 @@ function renderChannelList() { node.classList.add("active"); } node.addEventListener("click", () => { + document.querySelector(".viewer-sidebar")?.classList.remove("open"); viewerState.oldestMessageId = null; viewerState.newestMessageId = null; loadChannel(channel.channel_id); @@ -304,6 +305,19 @@ async function main() { } setupInfiniteScroll(); + + const toggle = document.getElementById("sidebar-toggle"); + const sidebar = document.querySelector(".viewer-sidebar"); + if (toggle && sidebar) { + toggle.addEventListener("click", () => { + sidebar.classList.toggle("open"); + }); + sidebar.addEventListener("click", (e) => { + if (e.target === sidebar) { + sidebar.classList.remove("open"); + } + }); + } } main().catch((error) => { From df644219dece740b5442b2bf20fb9c0bdbb2b8de Mon Sep 17 00:00:00 2001 From: mr-forust Date: Thu, 11 Jun 2026 18:19:17 +0200 Subject: [PATCH 2/3] k8s + image --- compose.yaml | 1 + k8s/telegram-scraper.yaml | 75 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 76 insertions(+) create mode 100644 k8s/telegram-scraper.yaml diff --git a/compose.yaml b/compose.yaml index 0e249e4..7cf7521 100644 --- a/compose.yaml +++ b/compose.yaml @@ -4,6 +4,7 @@ services: context: . dockerfile: Dockerfile container_name: telegram-scraper + image: gcr.forust.xyz/forust/telegram-scraper:latest user: 1000:1000 restart: unless-stopped tty: true diff --git a/k8s/telegram-scraper.yaml b/k8s/telegram-scraper.yaml new file mode 100644 index 0000000..04478d0 --- /dev/null +++ b/k8s/telegram-scraper.yaml @@ -0,0 +1,75 @@ +apiVersion: v1 +kind: Namespace +metadata: + name: telegram-scraper +--- +apiVersion: v1 +kind: Service +metadata: + name: telegram-scraper-service + namespace: telegram-scraper +spec: + selector: + app: telegram-scraper + ports: + - port: 8080 + targetPort: 8080 +--- +apiVersion: apps/v1 +kind: StatefulSet +metadata: + name: telegram-scraper-statefulset + namespace: telegram-scraper +spec: + selector: + matchLabels: + app: telegram-scraper + serviceName: telegram-scraper-service + replicas: 1 + template: + metadata: + labels: + app: telegram-scraper + spec: + containers: + - name: telegram-scraper + image: gcr.forust.xyz/forust/telegram-scraper:latest + tty: true + stdin: true + ports: + - containerPort: 8080 + volumeMounts: + - name: telegram-scraper-data + mountPath: /app/data + - name: telegram-scraper-session + mountPath: /app/session + volumeClaimTemplates: + - metadata: + name: telegram-scraper-data + spec: + accessModes: ["ReadWriteOnce"] + resources: + requests: + storage: 20Gi + - metadata: + name: telegram-scraper-session + spec: + accessModes: ["ReadWriteOnce"] + resources: + requests: + storage: 5Mi +--- +apiVersion: traefik.io/v1alpha1 +kind: IngressRoute +metadata: + name: telegram-scraper-route + namespace: telegram-scraper +spec: + entryPoints: + - websecure + routes: + - match: Host(`tg.workstation.internal`) + kind: Rule + services: + - name: telegram-scraper-service + port: 8080 From ac5beb831b8bcabc9495bc6771806873c2f09b3b Mon Sep 17 00:00:00 2001 From: mr-forust Date: Fri, 19 Jun 2026 11:52:30 +0200 Subject: [PATCH 3/3] lint: fix bare except and unused variables --- .editorconfig | 24 + .hadolint.yaml | 10 + .markdownlint.json | 8 + .prettierrc.yaml | 8 + .vscode/settings.json | 34 + .yamllint | 17 + .zed/settings.json | 39 + app_state.py | 8 +- compose.yaml | 6 + k8s/telegram-scraper.yaml | 12 + scraper_jobs.py | 3 +- scripts/smoke_test.py | 6 +- telegram-scraper-OLD.py | 676 +++++++++------ telegram_scraper_with_forwarding.py | 1188 ++++++++++++++++----------- webui_server.py | 48 +- 15 files changed, 1359 insertions(+), 728 deletions(-) create mode 100644 .editorconfig create mode 100644 .hadolint.yaml create mode 100644 .markdownlint.json create mode 100644 .prettierrc.yaml create mode 100644 .vscode/settings.json create mode 100644 .yamllint create mode 100644 .zed/settings.json diff --git a/.editorconfig b/.editorconfig new file mode 100644 index 0000000..1d95791 --- /dev/null +++ b/.editorconfig @@ -0,0 +1,24 @@ +root = true + +[*] +indent_style = space +indent_size = 2 +end_of_line = lf +charset = utf-8 +trim_trailing_whitespace = true +insert_final_newline = true + +[*.{yml,yaml}] +indent_size = 2 + +[*.{json,jsonc}] +indent_size = 2 + +[*.md] +trim_trailing_whitespace = false + +[*.py] +indent_size = 4 + +[{Makefile,makefile}] +indent_style = tab diff --git a/.hadolint.yaml b/.hadolint.yaml new file mode 100644 index 0000000..de4bacb --- /dev/null +++ b/.hadolint.yaml @@ -0,0 +1,10 @@ +ignored: + - DL3008 + - DL3042 + - DL3018 + - DL3059 +trustedRegistries: + - docker.io + - ghcr.io + - quay.io + - gcr.forust.xyz diff --git a/.markdownlint.json b/.markdownlint.json new file mode 100644 index 0000000..9dfeb53 --- /dev/null +++ b/.markdownlint.json @@ -0,0 +1,8 @@ +{ + "default": true, + "MD013": false, + "MD024": false, + "MD033": false, + "MD041": false, + "MD046": false +} diff --git a/.prettierrc.yaml b/.prettierrc.yaml new file mode 100644 index 0000000..e121c46 --- /dev/null +++ b/.prettierrc.yaml @@ -0,0 +1,8 @@ +bracketSameLine: true +htmlWhitespaceSensitivity: css +printWidth: 120 +tabWidth: 2 +singleQuote: true +trailingComma: all +proseWrap: preserve +endOfLine: lf diff --git a/.vscode/settings.json b/.vscode/settings.json new file mode 100644 index 0000000..4c8d6c5 --- /dev/null +++ b/.vscode/settings.json @@ -0,0 +1,34 @@ +{ + "[yaml]": { + "editor.tabSize": 2, + "editor.insertSpaces": true, + "editor.formatOnSave": true, + "editor.defaultFormatter": "redhat.vscode-yaml" + }, + "[python]": { + "editor.tabSize": 4, + "editor.formatOnSave": true, + "editor.defaultFormatter": "charliermarsh.ruff" + }, + "[dockerfile]": { + "editor.formatOnSave": true, + "editor.defaultFormatter": "exiasr.hadolint" + }, + "[markdown]": { + "editor.wordWrap": "wordWrapColumn", + "editor.wordWrapColumn": 120 + }, + "yaml.schemas": { + "kubernetes": ["**/k8s/*.yaml", "**/k8s/*.yml"] + }, + "yaml.format.enable": true, + "yaml.validate": true, + "yaml.completion": true, + "prettier.bracketSameLine": true, + "prettier.htmlWhitespaceSensitivity": "css", + "prettier.printWidth": 120, + "editor.formatOnSave": true, + "files.autoSave": "onFocusChange", + "editor.wordWrapColumn": 120, + "editor.wordWrap": "wordWrapColumn" +} diff --git a/.yamllint b/.yamllint new file mode 100644 index 0000000..0f84b6d --- /dev/null +++ b/.yamllint @@ -0,0 +1,17 @@ +extends: default + +rules: + comments: + min-spaces-from-content: 1 + comments-indentation: false + document-start: disable + line-length: disable + braces: + min-spaces-inside: 0 + max-spaces-inside: 1 + brackets: + min-spaces-inside: 0 + max-spaces-inside: 1 + indentation: + spaces: 2 + indent-sequences: consistent diff --git a/.zed/settings.json b/.zed/settings.json new file mode 100644 index 0000000..042dac8 --- /dev/null +++ b/.zed/settings.json @@ -0,0 +1,39 @@ +{ + "tab_size": 2, + "soft_wrap": "prefer_line", + "preferred_line_length": 120, + "format_on_save": "on", + "languages": { + "YAML": { + "tab_size": 2, + "hard_tabs": false, + "format_on_save": "on", + "formatter": { + "language_server": { "name": "yaml-language-server" } + } + }, + "Python": { + "tab_size": 4, + "format_on_save": "on", + "formatter": { + "language_server": { "name": "ruff" } + } + } + }, + "lsp": { + "yaml-language-server": { + "settings": { + "yaml": { + "schemas": { + "kubernetes": ["**/k8s/*.yaml", "**/k8s/*.yml"] + }, + "validate": true, + "completion": true, + "format": { + "enable": true + } + } + } + } + } +} diff --git a/app_state.py b/app_state.py index 8f02720..6d27411 100644 --- a/app_state.py +++ b/app_state.py @@ -57,7 +57,9 @@ class StateStore: self.path.parent.mkdir(parents=True, exist_ok=True) tmp_path = self.path.with_suffix(self.path.suffix + ".tmp") with tmp_path.open("w", encoding="utf-8") as handle: - json.dump(self._merge_defaults(state), handle, ensure_ascii=False, indent=2) + json.dump( + self._merge_defaults(state), handle, ensure_ascii=False, indent=2 + ) handle.write("\n") tmp_path.replace(self.path) self._cache = None @@ -71,7 +73,9 @@ class StateStore: def continuous_config(self) -> Dict[str, Any]: state = self.load() - return dict(state.get("continuous_scraping") or DEFAULT_STATE["continuous_scraping"]) + return dict( + state.get("continuous_scraping") or DEFAULT_STATE["continuous_scraping"] + ) def save_continuous_config(self, config: Dict[str, Any]) -> Dict[str, Any]: def mutate(state: Dict[str, Any]) -> None: diff --git a/compose.yaml b/compose.yaml index 7cf7521..e36b7ba 100644 --- a/compose.yaml +++ b/compose.yaml @@ -10,6 +10,12 @@ services: tty: true ports: - "7887:8080" + healthcheck: + test: ["CMD", "python", "-c", "import http.client,json;c=http.client.HTTPConnection('localhost',8080);c.request('GET','/health');r=c.getresponse();exit(0)if json.loads(r.read()).get('ok')else exit(1)"] + interval: 30s + timeout: 10s + retries: 3 + start_period: 15s volumes: - ./data:/app/data - ./session:/app/session diff --git a/k8s/telegram-scraper.yaml b/k8s/telegram-scraper.yaml index 04478d0..3bfa56d 100644 --- a/k8s/telegram-scraper.yaml +++ b/k8s/telegram-scraper.yaml @@ -38,6 +38,18 @@ spec: stdin: true ports: - containerPort: 8080 + livenessProbe: + httpGet: + path: /health + port: 8080 + initialDelaySeconds: 15 + periodSeconds: 30 + readinessProbe: + httpGet: + path: /health + port: 8080 + initialDelaySeconds: 5 + periodSeconds: 10 volumeMounts: - name: telegram-scraper-data mountPath: /app/data diff --git a/scraper_jobs.py b/scraper_jobs.py index 44ac6b2..2b50bde 100644 --- a/scraper_jobs.py +++ b/scraper_jobs.py @@ -41,7 +41,8 @@ class ScraperJobService: ) elif job_type == "scrape_selected": await self._scrape_channels( - scraper, [str(channel_id) for channel_id in payload.get("channels", [])] + scraper, + [str(channel_id) for channel_id in payload.get("channels", [])], ) elif job_type == "export_all": await scraper.export_data() diff --git a/scripts/smoke_test.py b/scripts/smoke_test.py index 33e2333..5160b9f 100644 --- a/scripts/smoke_test.py +++ b/scripts/smoke_test.py @@ -60,7 +60,11 @@ def main() -> int: status, body = fetch(endpoint) if status != 200: raise RuntimeError(f"{endpoint} returned HTTP {status}") - if endpoint.endswith(".json") or endpoint.startswith("/api") or endpoint.startswith("/health"): + if ( + endpoint.endswith(".json") + or endpoint.startswith("/api") + or endpoint.startswith("/health") + ): json.loads(body.decode("utf-8")) print(f"ok {endpoint}") return 0 diff --git a/telegram-scraper-OLD.py b/telegram-scraper-OLD.py index 5b5863d..f531db1 100644 --- a/telegram-scraper-OLD.py +++ b/telegram-scraper-OLD.py @@ -5,18 +5,28 @@ import csv import asyncio import time import sys -import uuid import warnings from dataclasses import dataclass from typing import Dict, List, Optional, Any from pathlib import Path from io import StringIO from telethon import TelegramClient -from telethon.tl.types import MessageMediaPhoto, MessageMediaDocument, MessageMediaWebPage, User, PeerChannel, Channel, Chat +from telethon.tl.types import ( + MessageMediaPhoto, + MessageMediaDocument, + MessageMediaWebPage, + User, + PeerChannel, + Channel, + Chat, +) from telethon.errors import FloodWaitError, SessionPasswordNeededError import qrcode -warnings.filterwarnings("ignore", message="Using async sessions support is an experimental feature") +warnings.filterwarnings( + "ignore", message="Using async sessions support is an experimental feature" +) + def display_ascii_art(): WHITE = "\033[97m" @@ -31,6 +41,7 @@ ___________________ _________ """ print(WHITE + art + RESET) + @dataclass class MessageData: message_id: int @@ -48,9 +59,10 @@ class MessageData: forwards: Optional[int] reactions: Optional[str] + class OptimizedTelegramScraper: def __init__(self): - self.STATE_FILE = 'state.json' + self.STATE_FILE = "state.json" self.state = self.load_state() self.client = None self.continuous_scraping_active = False @@ -58,25 +70,25 @@ class OptimizedTelegramScraper: self.batch_size = 100 self.state_save_interval = 50 self.db_connections = {} - + def load_state(self) -> Dict[str, Any]: if os.path.exists(self.STATE_FILE): try: - with open(self.STATE_FILE, 'r') as f: + with open(self.STATE_FILE, "r") as f: return json.load(f) - except: + except Exception: pass return { - 'api_id': None, - 'api_hash': None, - 'channels': {}, - 'channel_names': {}, - 'scrape_media': True, + "api_id": None, + "api_hash": None, + "channels": {}, + "channel_names": {}, + "scrape_media": True, } def save_state(self): try: - with open(self.STATE_FILE, 'w') as f: + with open(self.STATE_FILE, "w") as f: json.dump(self.state, f, indent=2) except Exception as e: print(f"Failed to save state: {e}") @@ -86,17 +98,19 @@ class OptimizedTelegramScraper: channel_dir = Path(channel) channel_dir.mkdir(exist_ok=True) - db_file = channel_dir / f'{channel}.db' + db_file = channel_dir / f"{channel}.db" conn = sqlite3.connect(str(db_file), check_same_thread=False) - conn.execute('''CREATE TABLE IF NOT EXISTS messages + conn.execute("""CREATE TABLE IF NOT EXISTS messages (id INTEGER PRIMARY KEY, message_id INTEGER UNIQUE, date TEXT, sender_id INTEGER, first_name TEXT, last_name TEXT, username TEXT, message TEXT, media_type TEXT, media_path TEXT, reply_to INTEGER, - post_author TEXT, views INTEGER, forwards INTEGER, reactions TEXT)''') - conn.execute('CREATE INDEX IF NOT EXISTS idx_message_id ON messages(message_id)') - conn.execute('CREATE INDEX IF NOT EXISTS idx_date ON messages(date)') - conn.execute('PRAGMA journal_mode=WAL') - conn.execute('PRAGMA synchronous=NORMAL') + post_author TEXT, views INTEGER, forwards INTEGER, reactions TEXT)""") + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_message_id ON messages(message_id)" + ) + conn.execute("CREATE INDEX IF NOT EXISTS idx_date ON messages(date)") + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA synchronous=NORMAL") conn.commit() self.migrate_database(conn) @@ -111,19 +125,19 @@ class OptimizedTelegramScraper: columns = {row[1] for row in cursor.fetchall()} migrations = [] - if 'post_author' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN post_author TEXT') - if 'views' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN views INTEGER') - if 'forwards' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN forwards INTEGER') - if 'reactions' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN reactions TEXT') + if "post_author" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN post_author TEXT") + if "views" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN views INTEGER") + if "forwards" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN forwards INTEGER") + if "reactions" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN reactions TEXT") for migration in migrations: try: conn.execute(migration) - except: + except Exception: pass if migrations: @@ -139,20 +153,38 @@ class OptimizedTelegramScraper: return conn = self.get_db_connection(channel) - data = [(msg.message_id, msg.date, msg.sender_id, msg.first_name, - msg.last_name, msg.username, msg.message, msg.media_type, - msg.media_path, msg.reply_to, msg.post_author, msg.views, - msg.forwards, msg.reactions) for msg in messages] + data = [ + ( + msg.message_id, + msg.date, + msg.sender_id, + msg.first_name, + msg.last_name, + msg.username, + msg.message, + msg.media_type, + msg.media_path, + msg.reply_to, + msg.post_author, + msg.views, + msg.forwards, + msg.reactions, + ) + for msg in messages + ] - conn.executemany('''INSERT OR IGNORE INTO messages + conn.executemany( + """INSERT OR IGNORE INTO messages (message_id, date, sender_id, first_name, last_name, username, message, media_type, media_path, reply_to, post_author, views, forwards, reactions) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''', data) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", + data, + ) conn.commit() async def download_media(self, channel: str, message) -> Optional[str]: - if not message.media or not self.state['scrape_media']: + if not message.media or not self.state["scrape_media"]: return None if isinstance(message.media, MessageMediaWebPage): @@ -160,23 +192,23 @@ class OptimizedTelegramScraper: try: channel_dir = Path(channel) - media_folder = channel_dir / 'media' + media_folder = channel_dir / "media" media_folder.mkdir(exist_ok=True) - + if isinstance(message.media, MessageMediaPhoto): - original_name = getattr(message.file, 'name', None) or "photo.jpg" + original_name = getattr(message.file, "name", None) or "photo.jpg" ext = "jpg" elif isinstance(message.media, MessageMediaDocument): - ext = getattr(message.file, 'ext', 'bin') if message.file else 'bin' - original_name = getattr(message.file, 'name', None) or f"document.{ext}" + ext = getattr(message.file, "ext", "bin") if message.file else "bin" + original_name = getattr(message.file, "name", None) or f"document.{ext}" else: return None - + base_name = Path(original_name).stem extension = Path(original_name).suffix or f".{ext}" unique_filename = f"{message.id}-{base_name}{extension}" media_path = media_folder / unique_filename - + existing_files = list(media_folder.glob(f"{message.id}-*")) if existing_files: return str(existing_files[0]) @@ -195,24 +227,30 @@ class OptimizedTelegramScraper: return None except Exception: if attempt < 2: - await asyncio.sleep(2 ** attempt) + await asyncio.sleep(2**attempt) else: return None - + return None except Exception: return None async def update_media_path(self, channel: str, message_id: int, media_path: str): conn = self.get_db_connection(channel) - conn.execute('UPDATE messages SET media_path = ? WHERE message_id = ?', - (media_path, message_id)) + conn.execute( + "UPDATE messages SET media_path = ? WHERE message_id = ?", + (media_path, message_id), + ) conn.commit() async def scrape_channel(self, channel: str, offset_id: int): try: - entity = await self.client.get_entity(PeerChannel(int(channel)) if channel.startswith('-') else channel) - result = await self.client.get_messages(entity, offset_id=offset_id, reverse=True, limit=0) + entity = await self.client.get_entity( + PeerChannel(int(channel)) if channel.startswith("-") else channel + ) + result = await self.client.get_messages( + entity, offset_id=offset_id, reverse=True, limit=0 + ) total_messages = result.total if total_messages == 0: @@ -227,7 +265,9 @@ class OptimizedTelegramScraper: last_message_id = offset_id semaphore = asyncio.Semaphore(self.max_concurrent_downloads) - async for message in self.client.iter_messages(entity, offset_id=offset_id, reverse=True): + async for message in self.client.iter_messages( + entity, offset_id=offset_id, reverse=True + ): try: sender = await message.get_sender() @@ -235,33 +275,45 @@ class OptimizedTelegramScraper: if message.reactions and message.reactions.results: reactions_parts = [] for reaction in message.reactions.results: - emoji = getattr(reaction.reaction, 'emoticon', '') + emoji = getattr(reaction.reaction, "emoticon", "") count = reaction.count if emoji: reactions_parts.append(f"{emoji} {count}") if reactions_parts: - reactions_str = ' '.join(reactions_parts) + reactions_str = " ".join(reactions_parts) msg_data = MessageData( message_id=message.id, - date=message.date.strftime('%Y-%m-%d %H:%M:%S'), + date=message.date.strftime("%Y-%m-%d %H:%M:%S"), sender_id=message.sender_id, - first_name=getattr(sender, 'first_name', None) if isinstance(sender, User) else None, - last_name=getattr(sender, 'last_name', None) if isinstance(sender, User) else None, - username=getattr(sender, 'username', None) if isinstance(sender, User) else None, - message=message.message or '', - media_type=message.media.__class__.__name__ if message.media else None, + first_name=getattr(sender, "first_name", None) + if isinstance(sender, User) + else None, + last_name=getattr(sender, "last_name", None) + if isinstance(sender, User) + else None, + username=getattr(sender, "username", None) + if isinstance(sender, User) + else None, + message=message.message or "", + media_type=message.media.__class__.__name__ + if message.media + else None, media_path=None, reply_to=message.reply_to_msg_id if message.reply_to else None, post_author=message.post_author, views=message.views, forwards=message.forwards, - reactions=reactions_str + reactions=reactions_str, ) message_batch.append(msg_data) - if self.state['scrape_media'] and message.media and not isinstance(message.media, MessageMediaWebPage): + if ( + self.state["scrape_media"] + and message.media + and not isinstance(message.media, MessageMediaWebPage) + ): media_tasks.append(message) last_message_id = message.id @@ -272,15 +324,19 @@ class OptimizedTelegramScraper: message_batch.clear() if processed_messages % self.state_save_interval == 0: - self.state['channels'][channel] = last_message_id + self.state["channels"][channel] = last_message_id self.save_state() progress = (processed_messages / total_messages) * 100 bar_length = 30 - filled_length = int(bar_length * processed_messages // total_messages) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - sys.stdout.write(f"\ršŸ“„ Messages: [{bar}] {progress:.1f}% ({processed_messages}/{total_messages})") + filled_length = int( + bar_length * processed_messages // total_messages + ) + bar = "ā–ˆ" * filled_length + "ā–‘" * (bar_length - filled_length) + + sys.stdout.write( + f"\ršŸ“„ Messages: [{bar}] {progress:.1f}% ({processed_messages}/{total_messages})" + ) sys.stdout.flush() except Exception as e: @@ -294,39 +350,47 @@ class OptimizedTelegramScraper: completed_media = 0 successful_downloads = 0 print(f"\nšŸ“„ Downloading {total_media} media files...") - + semaphore = asyncio.Semaphore(self.max_concurrent_downloads) - + async def download_single_media(message): async with semaphore: return await self.download_media(channel, message) - + batch_size = 10 for i in range(0, len(media_tasks), batch_size): - batch = media_tasks[i:i + batch_size] - tasks = [asyncio.create_task(download_single_media(msg)) for msg in batch] - + batch = media_tasks[i : i + batch_size] + tasks = [ + asyncio.create_task(download_single_media(msg)) for msg in batch + ] + for j, task in enumerate(tasks): try: media_path = await task if media_path: - await self.update_media_path(channel, batch[j].id, media_path) + await self.update_media_path( + channel, batch[j].id, media_path + ) successful_downloads += 1 except Exception: pass - + completed_media += 1 progress = (completed_media / total_media) * 100 bar_length = 30 filled_length = int(bar_length * completed_media // total_media) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - sys.stdout.write(f"\ršŸ“„ Media: [{bar}] {progress:.1f}% ({completed_media}/{total_media})") - sys.stdout.flush() - - print(f"\nāœ… Media download complete! ({successful_downloads}/{total_media} successful)") + bar = "ā–ˆ" * filled_length + "ā–‘" * (bar_length - filled_length) - self.state['channels'][channel] = last_message_id + sys.stdout.write( + f"\ršŸ“„ Media: [{bar}] {progress:.1f}% ({completed_media}/{total_media})" + ) + sys.stdout.flush() + + print( + f"\nāœ… Media download complete! ({successful_downloads}/{total_media} successful)" + ) + + self.state["channels"][channel] = last_message_id self.save_state() print(f"\nCompleted scraping channel {channel}") @@ -336,57 +400,78 @@ class OptimizedTelegramScraper: async def rescrape_media(self, channel: str): conn = self.get_db_connection(channel) cursor = conn.cursor() - cursor.execute('SELECT message_id FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NULL') + cursor.execute( + 'SELECT message_id FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NULL' + ) message_ids = [row[0] for row in cursor.fetchall()] - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') + channel_name = self.state.get("channel_names", {}).get(channel, "Unknown") if not message_ids: print(f"No media files to reprocess for {channel_name} (ID: {channel})") return - print(f"šŸ“„ Reprocessing {len(message_ids)} media files for {channel_name} (ID: {channel})") + print( + f"šŸ“„ Reprocessing {len(message_ids)} media files for {channel_name} (ID: {channel})" + ) try: - if channel.lstrip('-').isdigit(): + if channel.lstrip("-").isdigit(): entity = await self.client.get_entity(PeerChannel(int(channel))) else: entity = await self.client.get_entity(channel) semaphore = asyncio.Semaphore(self.max_concurrent_downloads) completed_media = 0 successful_downloads = 0 - + async def download_single_media(message): async with semaphore: return await self.download_media(channel, message) batch_size = 10 for i in range(0, len(message_ids), batch_size): - batch_ids = message_ids[i:i + batch_size] + batch_ids = message_ids[i : i + batch_size] messages = await self.client.get_messages(entity, ids=batch_ids) - - valid_messages = [msg for msg in messages if msg and msg.media and not isinstance(msg.media, MessageMediaWebPage)] - tasks = [asyncio.create_task(download_single_media(msg)) for msg in valid_messages] + + valid_messages = [ + msg + for msg in messages + if msg + and msg.media + and not isinstance(msg.media, MessageMediaWebPage) + ] + tasks = [ + asyncio.create_task(download_single_media(msg)) + for msg in valid_messages + ] for j, task in enumerate(tasks): try: media_path = await task if media_path: - await self.update_media_path(channel, valid_messages[j].id, media_path) + await self.update_media_path( + channel, valid_messages[j].id, media_path + ) successful_downloads += 1 except Exception: pass - + completed_media += 1 progress = (completed_media / len(message_ids)) * 100 bar_length = 30 - filled_length = int(bar_length * completed_media // len(message_ids)) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - sys.stdout.write(f"\ršŸ”„ Rescrape: [{bar}] {progress:.1f}% ({completed_media}/{len(message_ids)})") + filled_length = int( + bar_length * completed_media // len(message_ids) + ) + bar = "ā–ˆ" * filled_length + "ā–‘" * (bar_length - filled_length) + + sys.stdout.write( + f"\ršŸ”„ Rescrape: [{bar}] {progress:.1f}% ({completed_media}/{len(message_ids)})" + ) sys.stdout.flush() - print(f"\nāœ… Media reprocessing complete! ({successful_downloads}/{len(message_ids)} successful)") + print( + f"\nāœ… Media reprocessing complete! ({successful_downloads}/{len(message_ids)} successful)" + ) except Exception as e: print(f"Error reprocessing media: {e}") @@ -395,15 +480,19 @@ class OptimizedTelegramScraper: conn = self.get_db_connection(channel) cursor = conn.cursor() - cursor.execute('SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage"') + cursor.execute( + 'SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage"' + ) total_with_media = cursor.fetchone()[0] - cursor.execute('SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NOT NULL') + cursor.execute( + 'SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NOT NULL' + ) total_with_files = cursor.fetchone()[0] missing_count = total_with_media - total_with_files - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') + channel_name = self.state.get("channel_names", {}).get(channel, "Unknown") print(f"\nšŸ“Š Media Analysis for {channel_name} (ID: {channel}):") print(f"Messages with media: {total_with_media}") print(f"Media files downloaded: {total_with_files}") @@ -413,98 +502,119 @@ class OptimizedTelegramScraper: print("āœ… All media files are already downloaded!") return - cursor.execute('SELECT message_id, media_type FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND (media_path IS NULL OR media_path = "")') + cursor.execute( + 'SELECT message_id, media_type FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND (media_path IS NULL OR media_path = "")' + ) missing_media = cursor.fetchall() if not missing_media: print("āœ… No missing media found!") return - print(f"\nšŸ”§ Attempting to download {len(missing_media)} missing media files...") + print( + f"\nšŸ”§ Attempting to download {len(missing_media)} missing media files..." + ) try: - if channel.lstrip('-').isdigit(): + if channel.lstrip("-").isdigit(): entity = await self.client.get_entity(PeerChannel(int(channel))) else: entity = await self.client.get_entity(channel) semaphore = asyncio.Semaphore(self.max_concurrent_downloads) completed_media = 0 successful_downloads = 0 - + async def download_single_media(message): async with semaphore: return await self.download_media(channel, message) - + batch_size = 10 for i in range(0, len(missing_media), batch_size): - batch = missing_media[i:i + batch_size] + batch = missing_media[i : i + batch_size] message_ids = [msg[0] for msg in batch] - + messages = await self.client.get_messages(entity, ids=message_ids) - valid_messages = [msg for msg in messages if msg and msg.media and not isinstance(msg.media, MessageMediaWebPage)] - - tasks = [asyncio.create_task(download_single_media(msg)) for msg in valid_messages] + valid_messages = [ + msg + for msg in messages + if msg + and msg.media + and not isinstance(msg.media, MessageMediaWebPage) + ] + + tasks = [ + asyncio.create_task(download_single_media(msg)) + for msg in valid_messages + ] for j, task in enumerate(tasks): try: media_path = await task if media_path: - await self.update_media_path(channel, valid_messages[j].id, media_path) + await self.update_media_path( + channel, valid_messages[j].id, media_path + ) successful_downloads += 1 except Exception: pass - + completed_media += 1 progress = (completed_media / len(missing_media)) * 100 bar_length = 30 - filled_length = int(bar_length * completed_media // len(missing_media)) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - sys.stdout.write(f"\ršŸ”§ Fix Media: [{bar}] {progress:.1f}% ({completed_media}/{len(missing_media)})") + filled_length = int( + bar_length * completed_media // len(missing_media) + ) + bar = "ā–ˆ" * filled_length + "ā–‘" * (bar_length - filled_length) + + sys.stdout.write( + f"\ršŸ”§ Fix Media: [{bar}] {progress:.1f}% ({completed_media}/{len(missing_media)})" + ) sys.stdout.flush() - print(f"\nāœ… Media fix complete! ({successful_downloads}/{len(missing_media)} successful)") + print( + f"\nāœ… Media fix complete! ({successful_downloads}/{len(missing_media)} successful)" + ) except Exception as e: print(f"Error fixing missing media: {e}") async def continuous_scraping(self): self.continuous_scraping_active = True - + try: while self.continuous_scraping_active: start_time = time.time() - - for channel in self.state['channels']: + + for channel in self.state["channels"]: if not self.continuous_scraping_active: break print(f"\nChecking for new messages in channel: {channel}") - await self.scrape_channel(channel, self.state['channels'][channel]) - + await self.scrape_channel(channel, self.state["channels"][channel]) + elapsed = time.time() - start_time sleep_time = max(0, 60 - elapsed) if sleep_time > 0: await asyncio.sleep(sleep_time) - + except asyncio.CancelledError: print("Continuous scraping stopped") finally: self.continuous_scraping_active = False def get_export_filename(self, channel: str): - username = self.state.get('channel_names', {}).get(channel, 'no_username') + username = self.state.get("channel_names", {}).get(channel, "no_username") return f"{channel}_{username}" def export_to_csv(self, channel: str): conn = self.get_db_connection(channel) filename = self.get_export_filename(channel) - csv_file = Path(channel) / f'{filename}.csv' + csv_file = Path(channel) / f"{filename}.csv" cursor = conn.cursor() - cursor.execute('SELECT * FROM messages ORDER BY date') + cursor.execute("SELECT * FROM messages ORDER BY date") columns = [description[0] for description in cursor.description] - with open(csv_file, 'w', newline='', encoding='utf-8') as f: + with open(csv_file, "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(columns) @@ -517,14 +627,14 @@ class OptimizedTelegramScraper: def export_to_json(self, channel: str): conn = self.get_db_connection(channel) filename = self.get_export_filename(channel) - json_file = Path(channel) / f'{filename}.json' + json_file = Path(channel) / f"{filename}.json" cursor = conn.cursor() - cursor.execute('SELECT * FROM messages ORDER BY date') + cursor.execute("SELECT * FROM messages ORDER BY date") columns = [description[0] for description in cursor.description] - with open(json_file, 'w', encoding='utf-8') as f: - f.write('[\n') + with open(json_file, "w", encoding="utf-8") as f: + f.write("[\n") first_row = True while True: @@ -534,21 +644,21 @@ class OptimizedTelegramScraper: for row in rows: if not first_row: - f.write(',\n') + f.write(",\n") else: first_row = False data = dict(zip(columns, row)) json.dump(data, f, ensure_ascii=False, indent=2) - f.write('\n]') + f.write("\n]") async def export_data(self): - if not self.state['channels']: + if not self.state["channels"]: print("No channels to export") return - - for channel in self.state['channels']: + + for channel in self.state["channels"]: print(f"Exporting data for channel {channel}...") try: self.export_to_csv(channel) @@ -558,22 +668,30 @@ class OptimizedTelegramScraper: print(f"āŒ Export failed for channel {channel}: {e}") async def view_channels(self): - if not self.state['channels']: + if not self.state["channels"]: print("No channels saved") return print("\nCurrent channels:") - for i, (channel, last_id) in enumerate(self.state['channels'].items(), 1): + for i, (channel, last_id) in enumerate(self.state["channels"].items(), 1): try: conn = self.get_db_connection(channel) cursor = conn.cursor() - cursor.execute('SELECT COUNT(*) FROM messages') + cursor.execute("SELECT COUNT(*) FROM messages") count = cursor.fetchone()[0] - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') - print(f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}, Messages: {count}") - except: - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') - print(f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}") + channel_name = self.state.get("channel_names", {}).get( + channel, "Unknown" + ) + print( + f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}, Messages: {count}" + ) + except Exception: + channel_name = self.state.get("channel_names", {}).get( + channel, "Unknown" + ) + print( + f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}" + ) async def list_channels(self): try: @@ -582,23 +700,42 @@ class OptimizedTelegramScraper: channels_data = [] async for dialog in self.client.iter_dialogs(): entity = dialog.entity - if dialog.id != 777000 and (isinstance(entity, Channel) or isinstance(entity, Chat)): - channel_type = "Channel" if isinstance(entity, Channel) and entity.broadcast else "Group" - username = getattr(entity, 'username', None) or 'no_username' - print(f"[{count}] {dialog.title} (ID: {dialog.id}, Type: {channel_type}, Username: @{username})") - channels_data.append({ - 'number': count, - 'channel_name': dialog.title, - 'channel_id': str(dialog.id), - 'username': username, - 'type': channel_type - }) + if dialog.id != 777000 and ( + isinstance(entity, Channel) or isinstance(entity, Chat) + ): + channel_type = ( + "Channel" + if isinstance(entity, Channel) and entity.broadcast + else "Group" + ) + username = getattr(entity, "username", None) or "no_username" + print( + f"[{count}] {dialog.title} (ID: {dialog.id}, Type: {channel_type}, Username: @{username})" + ) + channels_data.append( + { + "number": count, + "channel_name": dialog.title, + "channel_id": str(dialog.id), + "username": username, + "type": channel_type, + } + ) count += 1 if channels_data: - csv_file = Path('channels_list.csv') - with open(csv_file, 'w', newline='', encoding='utf-8') as f: - writer = csv.DictWriter(f, fieldnames=['number', 'channel_name', 'channel_id', 'username', 'type']) + csv_file = Path("channels_list.csv") + with open(csv_file, "w", newline="", encoding="utf-8") as f: + writer = csv.DictWriter( + f, + fieldnames=[ + "number", + "channel_name", + "channel_id", + "username", + "type", + ], + ) writer.writeheader() writer.writerows(channels_data) print(f"\nāœ… Saved channels list to {csv_file}") @@ -613,7 +750,7 @@ class OptimizedTelegramScraper: qr = qrcode.QRCode(box_size=1, border=1) qr.add_data(qr_login.url) qr.make() - + f = StringIO() qr.print_ascii(out=f) f.seek(0) @@ -625,10 +762,10 @@ class OptimizedTelegramScraper: print("1. Open Telegram on your phone") print("2. Go to Settings > Devices > Scan QR") print("3. Scan the code below\n") - + qr_login = await self.client.qr_login() self.display_qr_code_ascii(qr_login) - + try: await qr_login.wait() print("\nāœ… Successfully logged in via QR code!") @@ -646,7 +783,7 @@ class OptimizedTelegramScraper: phone = input("Enter your phone number: ") await self.client.send_code_request(phone) code = input("Enter the code you received: ") - + try: await self.client.sign_in(phone, code) print("\nāœ… Successfully logged in via phone!") @@ -661,58 +798,62 @@ class OptimizedTelegramScraper: return False async def initialize_client(self): - if not all([self.state.get('api_id'), self.state.get('api_hash')]): + if not all([self.state.get("api_id"), self.state.get("api_hash")]): print("\n=== API Configuration Required ===") print("You need to provide API credentials from https://my.telegram.org") try: - self.state['api_id'] = int(input("Enter your API ID: ")) - self.state['api_hash'] = input("Enter your API Hash: ") + self.state["api_id"] = int(input("Enter your API ID: ")) + self.state["api_hash"] = input("Enter your API Hash: ") self.save_state() except ValueError: print("Invalid API ID. Must be a number.") return False - self.client = TelegramClient('session', self.state['api_id'], self.state['api_hash']) - + self.client = TelegramClient( + "session", self.state["api_id"], self.state["api_hash"] + ) + try: await self.client.connect() except Exception as e: print(f"Failed to connect: {e}") return False - + if not await self.client.is_user_authorized(): print("\n=== Choose Authentication Method ===") print("[1] QR Code (Recommended - No phone number needed)") print("[2] Phone Number (Traditional method)") - + while True: choice = input("Enter your choice (1 or 2): ").strip() - if choice in ['1', '2']: + if choice in ["1", "2"]: break print("Please enter 1 or 2") - - success = await self.qr_code_auth() if choice == '1' else await self.phone_auth() - + + success = ( + await self.qr_code_auth() if choice == "1" else await self.phone_auth() + ) + if not success: print("Authentication failed. Please try again.") await self.client.disconnect() return False else: print("āœ… Already authenticated!") - + return True def parse_channel_selection(self, choice): - channels_list = list(self.state['channels'].keys()) + channels_list = list(self.state["channels"].keys()) selected_channels = [] - - if choice.lower() == 'all': + + if choice.lower() == "all": return channels_list - - for selection in [x.strip() for x in choice.split(',')]: + + for selection in [x.strip() for x in choice.split(",")]: try: - if selection.startswith('-'): - if selection in self.state['channels']: + if selection.startswith("-"): + if selection in self.state["channels"]: selected_channels.append(selection) else: print(f"Channel ID {selection} not found in your channels") @@ -721,14 +862,18 @@ class OptimizedTelegramScraper: if 1 <= num <= len(channels_list): selected_channels.append(channels_list[num - 1]) else: - print(f"Invalid channel number: {num}. Valid range: 1-{len(channels_list)}") + print( + f"Invalid channel number: {num}. Valid range: 1-{len(channels_list)}" + ) except ValueError: - print(f"Invalid input: {selection}. Use numbers (1,2,3) or full IDs (-100123...)") - + print( + f"Invalid input: {selection}. Use numbers (1,2,3) or full IDs (-100123...)" + ) + return selected_channels async def scrape_specific_channels(self): - if not self.state['channels']: + if not self.state["channels"]: print("No channels available. Use [L] to add channels first") return @@ -737,60 +882,62 @@ class OptimizedTelegramScraper: print("• Single: 1 or -1001234567890") print("• Multiple: 1,3,5 or mix formats") print("• All channels: all") - + choice = input("\nEnter selection: ").strip() selected_channels = self.parse_channel_selection(choice) - + if selected_channels: print(f"\nšŸš€ Starting scrape of {len(selected_channels)} channel(s)...") for i, channel in enumerate(selected_channels, 1): print(f"\n[{i}/{len(selected_channels)}] Scraping: {channel}") - await self.scrape_channel(channel, self.state['channels'][channel]) + await self.scrape_channel(channel, self.state["channels"][channel]) print(f"\nāœ… Completed scraping {len(selected_channels)} channel(s)!") else: print("āŒ No valid channels selected") async def manage_channels(self): while True: - print("\n" + "="*40) + print("\n" + "=" * 40) print(" TELEGRAM SCRAPER") - print("="*40) + print("=" * 40) print("[S] Scrape channels") print("[C] Continuous scraping") - print(f"[M] Media scraping: {'ON' if self.state['scrape_media'] else 'OFF'}") + print( + f"[M] Media scraping: {'ON' if self.state['scrape_media'] else 'OFF'}" + ) print("[L] List & add channels") print("[R] Remove channels") print("[E] Export data") print("[T] Rescrape media") print("[F] Fix missing media") print("[Q] Quit") - print("="*40) + print("=" * 40) choice = input("Enter your choice: ").lower().strip() - + try: - if choice == 'r': - if not self.state['channels']: + if choice == "r": + if not self.state["channels"]: print("No channels to remove") continue - + await self.view_channels() print("\nTo remove channels:") print("• Single: 1 or -1001234567890") print("• Multiple: 1,2,3 or mix formats") selection = input("Enter selection: ").strip() selected_channels = self.parse_channel_selection(selection) - + if selected_channels: removed_count = 0 for channel in selected_channels: - if channel in self.state['channels']: - del self.state['channels'][channel] + if channel in self.state["channels"]: + del self.state["channels"][channel] print(f"āœ… Removed channel {channel}") removed_count += 1 else: print(f"āŒ Channel {channel} not found") - + if removed_count > 0: self.save_state() print(f"\nšŸŽ‰ Removed {removed_count} channel(s)!") @@ -799,20 +946,22 @@ class OptimizedTelegramScraper: print("No channels were removed") else: print("No valid channels selected") - - elif choice == 's': + + elif choice == "s": await self.scrape_specific_channels() - - elif choice == 'm': - self.state['scrape_media'] = not self.state['scrape_media'] + + elif choice == "m": + self.state["scrape_media"] = not self.state["scrape_media"] self.save_state() - print(f"\nāœ… Media scraping {'enabled' if self.state['scrape_media'] else 'disabled'}") - - elif choice == 'c': + print( + f"\nāœ… Media scraping {'enabled' if self.state['scrape_media'] else 'disabled'}" + ) + + elif choice == "c": task = asyncio.create_task(self.continuous_scraping()) print("Continuous scraping started. Press Ctrl+C to stop.") try: - await asyncio.sleep(float('inf')) + await asyncio.sleep(float("inf")) except KeyboardInterrupt: self.continuous_scraping_active = False task.cancel() @@ -821,11 +970,11 @@ class OptimizedTelegramScraper: await task except asyncio.CancelledError: pass - - elif choice == 'e': + + elif choice == "e": await self.export_data() - - elif choice == 'l': + + elif choice == "l": channels_data = await self.list_channels() if not channels_data: @@ -841,24 +990,37 @@ class OptimizedTelegramScraper: if selection: added_count = 0 - if selection.lower() == 'all': + if selection.lower() == "all": for channel_info in channels_data: - channel_id = channel_info['channel_id'] - if channel_id not in self.state['channels']: - self.state['channels'][channel_id] = 0 - if 'channel_names' not in self.state: - self.state['channel_names'] = {} - self.state['channel_names'][channel_id] = channel_info['username'] - print(f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})") + channel_id = channel_info["channel_id"] + if channel_id not in self.state["channels"]: + self.state["channels"][channel_id] = 0 + if "channel_names" not in self.state: + self.state["channel_names"] = {} + self.state["channel_names"][channel_id] = ( + channel_info["username"] + ) + print( + f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})" + ) added_count += 1 else: - print(f"Channel {channel_info['channel_name']} already added") + print( + f"Channel {channel_info['channel_name']} already added" + ) else: - for sel in [x.strip() for x in selection.split(',')]: + for sel in [x.strip() for x in selection.split(",")]: try: - if sel.startswith('-'): + if sel.startswith("-"): channel_id = sel - channel_info = next((c for c in channels_data if c['channel_id'] == channel_id), None) + channel_info = next( + ( + c + for c in channels_data + if c["channel_id"] == channel_id + ), + None, + ) if not channel_info: print(f"Channel ID {channel_id} not found") continue @@ -866,19 +1028,27 @@ class OptimizedTelegramScraper: num = int(sel) if 1 <= num <= len(channels_data): channel_info = channels_data[num - 1] - channel_id = channel_info['channel_id'] + channel_id = channel_info["channel_id"] else: - print(f"Invalid number: {num}. Choose 1-{len(channels_data)}") + print( + f"Invalid number: {num}. Choose 1-{len(channels_data)}" + ) continue - if channel_id in self.state['channels']: - print(f"Channel {channel_info['channel_name']} already added") + if channel_id in self.state["channels"]: + print( + f"Channel {channel_info['channel_name']} already added" + ) else: - self.state['channels'][channel_id] = 0 - if 'channel_names' not in self.state: - self.state['channel_names'] = {} - self.state['channel_names'][channel_id] = channel_info['username'] - print(f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})") + self.state["channels"][channel_id] = 0 + if "channel_names" not in self.state: + self.state["channel_names"] = {} + self.state["channel_names"][channel_id] = ( + channel_info["username"] + ) + print( + f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})" + ) added_count += 1 except ValueError: @@ -890,17 +1060,19 @@ class OptimizedTelegramScraper: await self.view_channels() else: print("No new channels were added") - - elif choice == 't': - if not self.state['channels']: + + elif choice == "t": + if not self.state["channels"]: print("No channels available. Add channels first") continue - + await self.view_channels() - print("\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)") + print( + "\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)" + ) selection = input("Enter your selection: ").strip() selected_channels = self.parse_channel_selection(selection) - + if len(selected_channels) == 1: channel = selected_channels[0] print(f"Rescaping media for channel: {channel}") @@ -909,17 +1081,19 @@ class OptimizedTelegramScraper: print("Please select only one channel for media rescaping") else: print("No valid channel selected") - - elif choice == 'f': - if not self.state['channels']: + + elif choice == "f": + if not self.state["channels"]: print("No channels available. Add channels first") continue - + await self.view_channels() - print("\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)") + print( + "\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)" + ) selection = input("Enter your selection: ").strip() selected_channels = self.parse_channel_selection(selection) - + if len(selected_channels) == 1: channel = selected_channels[0] await self.fix_missing_media(channel) @@ -927,17 +1101,17 @@ class OptimizedTelegramScraper: print("Please select only one channel for fixing missing media") else: print("No valid channel selected") - - elif choice == 'q': + + elif choice == "q": print("\nšŸ‘‹ Goodbye!") self.close_db_connections() if self.client: await self.client.disconnect() sys.exit() - + else: print("Invalid option") - + except Exception as e: print(f"Error: {e}") @@ -953,11 +1127,13 @@ class OptimizedTelegramScraper: else: print("Failed to initialize client. Exiting.") + async def main(): scraper = OptimizedTelegramScraper() await scraper.run() -if __name__ == '__main__': + +if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: diff --git a/telegram_scraper_with_forwarding.py b/telegram_scraper_with_forwarding.py index 9754594..b6bfe51 100644 --- a/telegram_scraper_with_forwarding.py +++ b/telegram_scraper_with_forwarding.py @@ -1,23 +1,32 @@ -import os import sqlite3 import json import csv import asyncio import time import sys -import uuid import warnings from dataclasses import dataclass from typing import Dict, List, Optional, Any from pathlib import Path from io import StringIO from telethon import TelegramClient, events -from telethon.tl.types import MessageMediaPhoto, MessageMediaDocument, MessageMediaWebPage, User, PeerChannel, Channel, Chat +from telethon.tl.types import ( + MessageMediaPhoto, + MessageMediaDocument, + MessageMediaWebPage, + User, + PeerChannel, + Channel, + Chat, +) from telethon.errors import FloodWaitError, SessionPasswordNeededError import qrcode from app_state import StateStore -warnings.filterwarnings("ignore", message="Using async sessions support is an experimental feature") +warnings.filterwarnings( + "ignore", message="Using async sessions support is an experimental feature" +) + def display_ascii_art(): WHITE = "\033[97m" @@ -32,6 +41,7 @@ ___________________ _________ """ print(WHITE + art + RESET) + @dataclass class MessageData: message_id: int @@ -49,6 +59,7 @@ class MessageData: forwards: Optional[int] reactions: Optional[str] + @dataclass class ForwardingRule: source_channel: str @@ -60,14 +71,15 @@ class ForwardingRule: forward_mode: str = "copy" enabled: bool = True + class OptimizedTelegramScraper: def __init__(self): - self.DATA_DIR = Path('data') - self.SESSION_DIR = Path('session') + self.DATA_DIR = Path("data") + self.SESSION_DIR = Path("session") self.DATA_DIR.mkdir(exist_ok=True) self.SESSION_DIR.mkdir(exist_ok=True) - self.STATE_FILE = str(self.DATA_DIR / 'state.json') - self.state_store = StateStore(self.DATA_DIR / 'state.json') + self.STATE_FILE = str(self.DATA_DIR / "state.json") + self.state_store = StateStore(self.DATA_DIR / "state.json") self.state = self.load_state() self.client = None self.continuous_scraping_active = False @@ -77,7 +89,7 @@ class OptimizedTelegramScraper: self.state_save_interval = 50 self.db_connections = {} self.forwarding_handler = None - + def load_state(self) -> Dict[str, Any]: return self.state_store.load() @@ -89,50 +101,54 @@ class OptimizedTelegramScraper: def get_forwarding_rules(self) -> List[ForwardingRule]: rules = [] - for rule_dict in self.state.get('forwarding_rules', []): - rules.append(ForwardingRule( - source_channel=rule_dict['source_channel'], - destination_channel=rule_dict['destination_channel'], - forward_text=rule_dict.get('forward_text', True), - forward_images=rule_dict.get('forward_images', True), - forward_videos=rule_dict.get('forward_videos', True), - forward_documents=rule_dict.get('forward_documents', True), - forward_mode=rule_dict.get('forward_mode', 'copy'), - enabled=rule_dict.get('enabled', True) - )) + for rule_dict in self.state.get("forwarding_rules", []): + rules.append( + ForwardingRule( + source_channel=rule_dict["source_channel"], + destination_channel=rule_dict["destination_channel"], + forward_text=rule_dict.get("forward_text", True), + forward_images=rule_dict.get("forward_images", True), + forward_videos=rule_dict.get("forward_videos", True), + forward_documents=rule_dict.get("forward_documents", True), + forward_mode=rule_dict.get("forward_mode", "copy"), + enabled=rule_dict.get("enabled", True), + ) + ) return rules def save_forwarding_rule(self, rule: ForwardingRule): rule_dict = { - 'source_channel': rule.source_channel, - 'destination_channel': rule.destination_channel, - 'forward_text': rule.forward_text, - 'forward_images': rule.forward_images, - 'forward_videos': rule.forward_videos, - 'forward_documents': rule.forward_documents, - 'forward_mode': rule.forward_mode, - 'enabled': rule.enabled + "source_channel": rule.source_channel, + "destination_channel": rule.destination_channel, + "forward_text": rule.forward_text, + "forward_images": rule.forward_images, + "forward_videos": rule.forward_videos, + "forward_documents": rule.forward_documents, + "forward_mode": rule.forward_mode, + "enabled": rule.enabled, } - + existing_idx = None - for i, existing in enumerate(self.state.get('forwarding_rules', [])): - if (existing['source_channel'] == rule.source_channel and - existing['destination_channel'] == rule.destination_channel): + for i, existing in enumerate(self.state.get("forwarding_rules", [])): + if ( + existing["source_channel"] == rule.source_channel + and existing["destination_channel"] == rule.destination_channel + ): existing_idx = i break - + if existing_idx is not None: - self.state['forwarding_rules'][existing_idx] = rule_dict + self.state["forwarding_rules"][existing_idx] = rule_dict else: - if 'forwarding_rules' not in self.state: - self.state['forwarding_rules'] = [] - self.state['forwarding_rules'].append(rule_dict) - + if "forwarding_rules" not in self.state: + self.state["forwarding_rules"] = [] + self.state["forwarding_rules"].append(rule_dict) + self.save_state() def remove_forwarding_rule(self, index: int) -> bool: - if 0 <= index < len(self.state.get('forwarding_rules', [])): - del self.state['forwarding_rules'][index] + if 0 <= index < len(self.state.get("forwarding_rules", [])): + del self.state["forwarding_rules"][index] self.save_state() return True return False @@ -142,17 +158,19 @@ class OptimizedTelegramScraper: channel_dir = self.DATA_DIR / channel channel_dir.mkdir(parents=True, exist_ok=True) - db_file = channel_dir / f'{channel}.db' + db_file = channel_dir / f"{channel}.db" conn = sqlite3.connect(str(db_file), check_same_thread=False, timeout=30) - conn.execute('''CREATE TABLE IF NOT EXISTS messages + conn.execute("""CREATE TABLE IF NOT EXISTS messages (id INTEGER PRIMARY KEY, message_id INTEGER UNIQUE, date TEXT, sender_id INTEGER, first_name TEXT, last_name TEXT, username TEXT, message TEXT, media_type TEXT, media_path TEXT, reply_to INTEGER, - post_author TEXT, views INTEGER, forwards INTEGER, reactions TEXT)''') - conn.execute('CREATE INDEX IF NOT EXISTS idx_message_id ON messages(message_id)') - conn.execute('CREATE INDEX IF NOT EXISTS idx_date ON messages(date)') - conn.execute('PRAGMA journal_mode=WAL') - conn.execute('PRAGMA synchronous=NORMAL') + post_author TEXT, views INTEGER, forwards INTEGER, reactions TEXT)""") + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_message_id ON messages(message_id)" + ) + conn.execute("CREATE INDEX IF NOT EXISTS idx_date ON messages(date)") + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA synchronous=NORMAL") conn.commit() self.migrate_database(conn) @@ -167,19 +185,19 @@ class OptimizedTelegramScraper: columns = {row[1] for row in cursor.fetchall()} migrations = [] - if 'post_author' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN post_author TEXT') - if 'views' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN views INTEGER') - if 'forwards' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN forwards INTEGER') - if 'reactions' not in columns: - migrations.append('ALTER TABLE messages ADD COLUMN reactions TEXT') + if "post_author" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN post_author TEXT") + if "views" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN views INTEGER") + if "forwards" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN forwards INTEGER") + if "reactions" not in columns: + migrations.append("ALTER TABLE messages ADD COLUMN reactions TEXT") for migration in migrations: try: conn.execute(migration) - except: + except Exception: pass if migrations: @@ -195,20 +213,38 @@ class OptimizedTelegramScraper: return conn = self.get_db_connection(channel) - data = [(msg.message_id, msg.date, msg.sender_id, msg.first_name, - msg.last_name, msg.username, msg.message, msg.media_type, - msg.media_path, msg.reply_to, msg.post_author, msg.views, - msg.forwards, msg.reactions) for msg in messages] + data = [ + ( + msg.message_id, + msg.date, + msg.sender_id, + msg.first_name, + msg.last_name, + msg.username, + msg.message, + msg.media_type, + msg.media_path, + msg.reply_to, + msg.post_author, + msg.views, + msg.forwards, + msg.reactions, + ) + for msg in messages + ] - conn.executemany('''INSERT OR IGNORE INTO messages + conn.executemany( + """INSERT OR IGNORE INTO messages (message_id, date, sender_id, first_name, last_name, username, message, media_type, media_path, reply_to, post_author, views, forwards, reactions) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''', data) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", + data, + ) conn.commit() async def download_media(self, channel: str, message) -> Optional[str]: - if not message.media or not self.state['scrape_media']: + if not message.media or not self.state["scrape_media"]: return None if isinstance(message.media, MessageMediaWebPage): @@ -216,23 +252,23 @@ class OptimizedTelegramScraper: try: channel_dir = self.DATA_DIR / channel - media_folder = channel_dir / 'media' + media_folder = channel_dir / "media" media_folder.mkdir(exist_ok=True) - + if isinstance(message.media, MessageMediaPhoto): - original_name = getattr(message.file, 'name', None) or "photo.jpg" + original_name = getattr(message.file, "name", None) or "photo.jpg" ext = "jpg" elif isinstance(message.media, MessageMediaDocument): - ext = getattr(message.file, 'ext', 'bin') if message.file else 'bin' - original_name = getattr(message.file, 'name', None) or f"document.{ext}" + ext = getattr(message.file, "ext", "bin") if message.file else "bin" + original_name = getattr(message.file, "name", None) or f"document.{ext}" else: return None - + base_name = Path(original_name).stem extension = Path(original_name).suffix or f".{ext}" unique_filename = f"{message.id}-{base_name}{extension}" media_path = media_folder / unique_filename - + existing_files = list(media_folder.glob(f"{message.id}-*")) if existing_files: return str(existing_files[0]) @@ -251,75 +287,85 @@ class OptimizedTelegramScraper: return None except Exception: if attempt < 2: - await asyncio.sleep(2 ** attempt) + await asyncio.sleep(2**attempt) else: return None - + return None except Exception: return None async def update_media_path(self, channel: str, message_id: int, media_path: str): conn = self.get_db_connection(channel) - conn.execute('UPDATE messages SET media_path = ? WHERE message_id = ?', - (media_path, message_id)) + conn.execute( + "UPDATE messages SET media_path = ? WHERE message_id = ?", + (media_path, message_id), + ) conn.commit() def should_forward_message(self, message, rule: ForwardingRule) -> bool: if not rule.enabled: return False - + if not message.media: return rule.forward_text and bool(message.message) - + if isinstance(message.media, MessageMediaWebPage): return rule.forward_text and bool(message.message) - + if isinstance(message.media, MessageMediaPhoto): return rule.forward_images - + if isinstance(message.media, MessageMediaDocument): if message.file: - mime_type = getattr(message.file, 'mime_type', '') or '' - if mime_type.startswith('video/'): + mime_type = getattr(message.file, "mime_type", "") or "" + if mime_type.startswith("video/"): return rule.forward_videos else: return rule.forward_documents return rule.forward_documents - + return False - async def forward_message(self, message, rule: ForwardingRule, source_channel_id: int = None): + async def forward_message( + self, message, rule: ForwardingRule, source_channel_id: int = None + ): try: dest_entity = await self._resolve_entity(rule.destination_channel) - + if rule.forward_mode == "forward": await self.client.forward_messages(dest_entity, message) print(f" Forwarded message {message.id}") else: - source_name = self.state.get('channel_names', {}).get(rule.source_channel, '') + source_name = self.state.get("channel_names", {}).get( + rule.source_channel, "" + ) if source_channel_id: - full_id = f"-100{abs(source_channel_id)}" if not str(source_channel_id).startswith('-100') else str(source_channel_id) + full_id = ( + f"-100{abs(source_channel_id)}" + if not str(source_channel_id).startswith("-100") + else str(source_channel_id) + ) else: full_id = rule.source_channel - - if source_name and source_name != 'no_username': - source_line = f"From: @{source_name} ({full_id})\n─────────────────\n" + + if source_name and source_name != "no_username": + source_line = ( + f"From: @{source_name} ({full_id})\n─────────────────\n" + ) else: source_line = f"From: {full_id}\n─────────────────\n" - - copy_text = source_line + (message.message or '') - + + copy_text = source_line + (message.message or "") + if message.media and not isinstance(message.media, MessageMediaWebPage): await self.client.send_message( - dest_entity, - copy_text, - file=message.media + dest_entity, copy_text, file=message.media ) else: await self.client.send_message(dest_entity, copy_text) print(f" Copied message {message.id}") - + return True except FloodWaitError as e: print(f" Rate limited, waiting {e.seconds}s...") @@ -334,17 +380,17 @@ class OptimizedTelegramScraper: if not rules: print("No forwarding rules configured") return False - + enabled_rules = [r for r in rules if r.enabled] if not enabled_rules: print("No enabled forwarding rules") return False - + source_channels = [] source_id_map = {} for rule in enabled_rules: try: - if rule.source_channel.lstrip('-').isdigit(): + if rule.source_channel.lstrip("-").isdigit(): channel_id = int(rule.source_channel) else: entity = await self.client.get_entity(rule.source_channel) @@ -353,63 +399,88 @@ class OptimizedTelegramScraper: source_id_map[channel_id] = rule.source_channel except Exception as e: print(f"Failed to get entity for {rule.source_channel}: {e}") - + if not source_channels: print("No valid source channels") return False - - @self.client.on(events.NewMessage(chats=source_channels, incoming=True, outgoing=True)) + + @self.client.on( + events.NewMessage(chats=source_channels, incoming=True, outgoing=True) + ) async def forwarding_handler(event): message = event.message chat_id = event.chat_id - + for rule in enabled_rules: rule_source_id = None - if rule.source_channel.lstrip('-').isdigit(): + if rule.source_channel.lstrip("-").isdigit(): rule_source_id = int(rule.source_channel) else: - rule_source_id = next((k for k, v in source_id_map.items() if v == rule.source_channel), None) - + rule_source_id = next( + ( + k + for k, v in source_id_map.items() + if v == rule.source_channel + ), + None, + ) + if rule_source_id == chat_id: if self.should_forward_message(message, rule): - source_name = self.state.get('channel_names', {}).get(rule.source_channel, rule.source_channel) - dest_name = self.state.get('channel_names', {}).get(rule.destination_channel, rule.destination_channel) - print(f"\n[{time.strftime('%H:%M:%S')}] New message in {source_name} -> forwarding to {dest_name}") + source_name = self.state.get("channel_names", {}).get( + rule.source_channel, rule.source_channel + ) + dest_name = self.state.get("channel_names", {}).get( + rule.destination_channel, rule.destination_channel + ) + print( + f"\n[{time.strftime('%H:%M:%S')}] New message in {source_name} -> forwarding to {dest_name}" + ) await self.forward_message(message, rule, chat_id) - + self.forwarding_handler = forwarding_handler return True async def start_forwarding(self): self.forwarding_active = True - + if not self.client.is_connected(): await self.client.connect() - + if not await self.setup_forwarding_handler(): self.forwarding_active = False return - + rules = self.get_forwarding_rules() enabled_rules = [r for r in rules if r.enabled] - - print(f"\nForwarding service started!") + + print("\nForwarding service started!") print(f" Monitoring {len(enabled_rules)} rule(s)") print(" Press Ctrl+C to stop\n") - + for rule in enabled_rules: - source_name = self.state.get('channel_names', {}).get(rule.source_channel, rule.source_channel) - dest_name = self.state.get('channel_names', {}).get(rule.destination_channel, rule.destination_channel) + source_name = self.state.get("channel_names", {}).get( + rule.source_channel, rule.source_channel + ) + dest_name = self.state.get("channel_names", {}).get( + rule.destination_channel, rule.destination_channel + ) content_types = [] - if rule.forward_text: content_types.append("text") - if rule.forward_images: content_types.append("images") - if rule.forward_videos: content_types.append("videos") - if rule.forward_documents: content_types.append("documents") + if rule.forward_text: + content_types.append("text") + if rule.forward_images: + content_types.append("images") + if rule.forward_videos: + content_types.append("videos") + if rule.forward_documents: + content_types.append("documents") print(f" - {source_name} -> {dest_name}") - print(f" Mode: {rule.forward_mode} | Content: {', '.join(content_types)}") - + print( + f" Mode: {rule.forward_mode} | Content: {', '.join(content_types)}" + ) + print("\nWatching for new messages...\n") - + try: await self.client.run_until_disconnected() except asyncio.CancelledError: @@ -424,46 +495,56 @@ class OptimizedTelegramScraper: async def manage_forwarding_rules(self): while True: - print("\n" + "="*45) + print("\n" + "=" * 45) print(" FORWARDING RULES MANAGER") - print("="*45) - + print("=" * 45) + rules = self.get_forwarding_rules() if rules: print("\nCurrent Rules:") for i, rule in enumerate(rules, 1): - source_name = self.state.get('channel_names', {}).get(rule.source_channel, rule.source_channel) - dest_name = self.state.get('channel_names', {}).get(rule.destination_channel, rule.destination_channel) + source_name = self.state.get("channel_names", {}).get( + rule.source_channel, rule.source_channel + ) + dest_name = self.state.get("channel_names", {}).get( + rule.destination_channel, rule.destination_channel + ) status = "āœ…" if rule.enabled else "āŒ" content_types = [] - if rule.forward_text: content_types.append("T") - if rule.forward_images: content_types.append("I") - if rule.forward_videos: content_types.append("V") - if rule.forward_documents: content_types.append("D") + if rule.forward_text: + content_types.append("T") + if rule.forward_images: + content_types.append("I") + if rule.forward_videos: + content_types.append("V") + if rule.forward_documents: + content_types.append("D") print(f" [{i}] {status} {source_name} → {dest_name}") - print(f" Mode: {rule.forward_mode} | Content: {'/'.join(content_types)}") + print( + f" Mode: {rule.forward_mode} | Content: {'/'.join(content_types)}" + ) else: print("\nNo forwarding rules configured") - + print("\n[A] Add new rule") print("[E] Edit rule") print("[T] Toggle rule on/off") print("[D] Delete rule") print("[S] Start forwarding") print("[B] Back to main menu") - print("="*45) - + print("=" * 45) + choice = input("Enter your choice: ").lower().strip() - - if choice == 'a': + + if choice == "a": await self.add_forwarding_rule_interactive() - elif choice == 'e': + elif choice == "e": await self.edit_forwarding_rule_interactive() - elif choice == 't': + elif choice == "t": await self.toggle_forwarding_rule() - elif choice == 'd': + elif choice == "d": await self.delete_forwarding_rule_interactive() - elif choice == 's': + elif choice == "s": if not rules or not any(r.enabled for r in rules): print("āŒ No enabled forwarding rules. Add or enable rules first.") continue @@ -472,7 +553,7 @@ class OptimizedTelegramScraper: except KeyboardInterrupt: self.forwarding_active = False print("\nForwarding stopped") - elif choice == 'b': + elif choice == "b": break else: print("Invalid option") @@ -480,29 +561,31 @@ class OptimizedTelegramScraper: async def add_forwarding_rule_interactive(self): print("\nšŸ“ ADD NEW FORWARDING RULE") print("-" * 40) - + await self.view_channels() - - if not self.state['channels']: - print("\nāŒ No channels available. Add channels first using [L] in the main menu.") + + if not self.state["channels"]: + print( + "\nāŒ No channels available. Add channels first using [L] in the main menu." + ) return - - channels_list = list(self.state['channels'].keys()) - + print("\n1ļøāƒ£ Select SOURCE channel (messages will be forwarded FROM here):") source_input = input("Enter channel number or ID: ").strip() source_channels = self.parse_channel_selection(source_input) - + if not source_channels: print("āŒ Invalid source channel selection") return source_channel = source_channels[0] - + print("\n2ļøāƒ£ Select DESTINATION channel (messages will be forwarded TO here):") - print(" You can enter a channel number from the list, or a channel username (@channelname)") + print( + " You can enter a channel number from the list, or a channel username (@channelname)" + ) dest_input = input("Enter channel number, ID, or @username: ").strip() - - if dest_input.startswith('@'): + + if dest_input.startswith("@"): dest_channel = dest_input else: # First try parsing as a channel from the tracked list @@ -512,16 +595,30 @@ class OptimizedTelegramScraper: else: # If not in tracked list, try as a raw channel ID (for destination-only channels) try: - if dest_input.lstrip('-').isdigit(): - test_id = int(dest_input) + if dest_input.lstrip("-").isdigit(): + int(dest_input) # Try to resolve the channel via Telethon to verify it exists try: entity = await self._resolve_entity(dest_input) dest_channel = dest_input - name = getattr(entity, 'title', None) or ' '.join(filter(None, [getattr(entity, 'first_name', None), getattr(entity, 'last_name', None)])) or dest_input + name = ( + getattr(entity, "title", None) + or " ".join( + filter( + None, + [ + getattr(entity, "first_name", None), + getattr(entity, "last_name", None), + ], + ) + ) + or dest_input + ) print(f"āœ… Found: {name}") except Exception: - print(f"āŒ Could not access channel ID {dest_input}. Make sure this account is a member.") + print( + f"āŒ Could not access channel ID {dest_input}. Make sure this account is a member." + ) return else: print("āŒ Invalid destination channel selection") @@ -529,7 +626,7 @@ class OptimizedTelegramScraper: except ValueError: print("āŒ Invalid destination channel selection") return - + print("\n3ļøāƒ£ Select content types to forward:") print(" [1] Text messages") print(" [2] Images/Photos") @@ -537,29 +634,29 @@ class OptimizedTelegramScraper: print(" [4] Documents/Files") print(" [A] All types") print(" Example: 1,2,3 or A for all") - + content_input = input("Enter selection: ").strip().lower() - - if content_input == 'a': + + if content_input == "a": forward_text = forward_images = forward_videos = forward_documents = True else: - selections = [x.strip() for x in content_input.split(',')] - forward_text = '1' in selections - forward_images = '2' in selections - forward_videos = '3' in selections - forward_documents = '4' in selections - + selections = [x.strip() for x in content_input.split(",")] + forward_text = "1" in selections + forward_images = "2" in selections + forward_videos = "3" in selections + forward_documents = "4" in selections + if not any([forward_text, forward_images, forward_videos, forward_documents]): print("āŒ You must select at least one content type") return - + print("\n4ļøāƒ£ Select forwarding mode:") print(" [1] Copy - Send as new message (no 'Forwarded from' header)") print(" [2] Forward - Keep 'Forwarded from' header") - + mode_input = input("Enter 1 or 2: ").strip() - forward_mode = "forward" if mode_input == '2' else "copy" - + forward_mode = "forward" if mode_input == "2" else "copy" + rule = ForwardingRule( source_channel=source_channel, destination_channel=dest_channel, @@ -568,21 +665,31 @@ class OptimizedTelegramScraper: forward_videos=forward_videos, forward_documents=forward_documents, forward_mode=forward_mode, - enabled=True + enabled=True, ) - + self.save_forwarding_rule(rule) - - source_name = self.state.get('channel_names', {}).get(source_channel, source_channel) - dest_name = dest_channel if dest_channel.startswith('@') else self.state.get('channel_names', {}).get(dest_channel, dest_channel) - - print(f"\nāœ… Forwarding rule created!") + + source_name = self.state.get("channel_names", {}).get( + source_channel, source_channel + ) + dest_name = ( + dest_channel + if dest_channel.startswith("@") + else self.state.get("channel_names", {}).get(dest_channel, dest_channel) + ) + + print("\nāœ… Forwarding rule created!") print(f" {source_name} → {dest_name}") content_types = [] - if forward_text: content_types.append("text") - if forward_images: content_types.append("images") - if forward_videos: content_types.append("videos") - if forward_documents: content_types.append("documents") + if forward_text: + content_types.append("text") + if forward_images: + content_types.append("images") + if forward_videos: + content_types.append("videos") + if forward_documents: + content_types.append("documents") print(f" Content: {', '.join(content_types)}") print(f" Mode: {forward_mode}") @@ -591,7 +698,7 @@ class OptimizedTelegramScraper: if not rules: print("No rules to edit") return - + try: idx = int(input("Enter rule number to edit: ")) - 1 if not (0 <= idx < len(rules)): @@ -600,41 +707,43 @@ class OptimizedTelegramScraper: except ValueError: print("Invalid input") return - + rule = rules[idx] - + print(f"\nEditing rule {idx + 1}:") print("Press Enter to keep current value\n") - + print("Current content types:") print(f" Text: {'Yes' if rule.forward_text else 'No'}") print(f" Images: {'Yes' if rule.forward_images else 'No'}") print(f" Videos: {'Yes' if rule.forward_videos else 'No'}") print(f" Documents: {'Yes' if rule.forward_documents else 'No'}") - + update_content = input("\nUpdate content types? (y/n): ").lower().strip() - if update_content == 'y': + if update_content == "y": print("Enter 1,2,3,4 or A for all:") print(" [1] Text [2] Images [3] Videos [4] Documents") content_input = input("Selection: ").strip().lower() - - if content_input == 'a': - rule.forward_text = rule.forward_images = rule.forward_videos = rule.forward_documents = True + + if content_input == "a": + rule.forward_text = rule.forward_images = rule.forward_videos = ( + rule.forward_documents + ) = True elif content_input: - selections = [x.strip() for x in content_input.split(',')] - rule.forward_text = '1' in selections - rule.forward_images = '2' in selections - rule.forward_videos = '3' in selections - rule.forward_documents = '4' in selections - + selections = [x.strip() for x in content_input.split(",")] + rule.forward_text = "1" in selections + rule.forward_images = "2" in selections + rule.forward_videos = "3" in selections + rule.forward_documents = "4" in selections + print(f"\nCurrent mode: {rule.forward_mode}") print(" [1] Copy [2] Forward") mode_input = input("New mode (or Enter to keep): ").strip() - if mode_input == '1': - rule.forward_mode = 'copy' - elif mode_input == '2': - rule.forward_mode = 'forward' - + if mode_input == "1": + rule.forward_mode = "copy" + elif mode_input == "2": + rule.forward_mode = "forward" + self.save_forwarding_rule(rule) print("āœ… Rule updated!") @@ -643,7 +752,7 @@ class OptimizedTelegramScraper: if not rules: print("No rules to toggle") return - + try: idx = int(input("Enter rule number to toggle: ")) - 1 if not (0 <= idx < len(rules)): @@ -652,11 +761,11 @@ class OptimizedTelegramScraper: except ValueError: print("Invalid input") return - + rule = rules[idx] rule.enabled = not rule.enabled self.save_forwarding_rule(rule) - + status = "enabled" if rule.enabled else "disabled" print(f"āœ… Rule {idx + 1} {status}") @@ -665,7 +774,7 @@ class OptimizedTelegramScraper: if not rules: print("No rules to delete") return - + try: idx = int(input("Enter rule number to delete: ")) - 1 if not (0 <= idx < len(rules)): @@ -674,9 +783,9 @@ class OptimizedTelegramScraper: except ValueError: print("Invalid input") return - + confirm = input(f"Delete rule {idx + 1}? (y/n): ").lower().strip() - if confirm == 'y': + if confirm == "y": if self.remove_forwarding_rule(idx): print("āœ… Rule deleted") else: @@ -684,7 +793,7 @@ class OptimizedTelegramScraper: async def _resolve_entity(self, channel: str): """Resolve a channel/chat/user entity from an ID string or username.""" - if channel.lstrip('-').isdigit(): + if channel.lstrip("-").isdigit(): cid = int(channel) # Channel/supergroup IDs are negative and large; user IDs are positive or small negative try: @@ -699,9 +808,11 @@ class OptimizedTelegramScraper: try: if not self.client.is_connected(): await self.client.connect() - + entity = await self._resolve_entity(channel) - result = await self.client.get_messages(entity, offset_id=offset_id, reverse=True, limit=0) + result = await self.client.get_messages( + entity, offset_id=offset_id, reverse=True, limit=0 + ) total_messages = result.total if total_messages == 0: @@ -716,7 +827,9 @@ class OptimizedTelegramScraper: last_message_id = offset_id semaphore = asyncio.Semaphore(self.max_concurrent_downloads) - async for message in self.client.iter_messages(entity, offset_id=offset_id, reverse=True): + async for message in self.client.iter_messages( + entity, offset_id=offset_id, reverse=True + ): try: sender = await message.get_sender() @@ -724,33 +837,45 @@ class OptimizedTelegramScraper: if message.reactions and message.reactions.results: reactions_parts = [] for reaction in message.reactions.results: - emoji = getattr(reaction.reaction, 'emoticon', '') + emoji = getattr(reaction.reaction, "emoticon", "") count = reaction.count if emoji: reactions_parts.append(f"{emoji} {count}") if reactions_parts: - reactions_str = ' '.join(reactions_parts) + reactions_str = " ".join(reactions_parts) msg_data = MessageData( message_id=message.id, - date=message.date.strftime('%Y-%m-%d %H:%M:%S'), + date=message.date.strftime("%Y-%m-%d %H:%M:%S"), sender_id=message.sender_id, - first_name=getattr(sender, 'first_name', None) if isinstance(sender, User) else None, - last_name=getattr(sender, 'last_name', None) if isinstance(sender, User) else None, - username=getattr(sender, 'username', None) if isinstance(sender, User) else None, - message=message.message or '', - media_type=message.media.__class__.__name__ if message.media else None, + first_name=getattr(sender, "first_name", None) + if isinstance(sender, User) + else None, + last_name=getattr(sender, "last_name", None) + if isinstance(sender, User) + else None, + username=getattr(sender, "username", None) + if isinstance(sender, User) + else None, + message=message.message or "", + media_type=message.media.__class__.__name__ + if message.media + else None, media_path=None, reply_to=message.reply_to_msg_id if message.reply_to else None, post_author=message.post_author, views=message.views, forwards=message.forwards, - reactions=reactions_str + reactions=reactions_str, ) message_batch.append(msg_data) - if self.state['scrape_media'] and message.media and not isinstance(message.media, MessageMediaWebPage): + if ( + self.state["scrape_media"] + and message.media + and not isinstance(message.media, MessageMediaWebPage) + ): media_tasks.append(message) last_message_id = message.id @@ -761,15 +886,19 @@ class OptimizedTelegramScraper: message_batch.clear() if processed_messages % self.state_save_interval == 0: - self.state['channels'][channel] = last_message_id + self.state["channels"][channel] = last_message_id self.save_state() progress = (processed_messages / total_messages) * 100 bar_length = 30 - filled_length = int(bar_length * processed_messages // total_messages) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - sys.stdout.write(f"\ršŸ“„ Messages: [{bar}] {progress:.1f}% ({processed_messages}/{total_messages})") + filled_length = int( + bar_length * processed_messages // total_messages + ) + bar = "ā–ˆ" * filled_length + "ā–‘" * (bar_length - filled_length) + + sys.stdout.write( + f"\ršŸ“„ Messages: [{bar}] {progress:.1f}% ({processed_messages}/{total_messages})" + ) sys.stdout.flush() except Exception as e: @@ -783,39 +912,47 @@ class OptimizedTelegramScraper: completed_media = 0 successful_downloads = 0 print(f"\nšŸ“„ Downloading {total_media} media files...") - + semaphore = asyncio.Semaphore(self.max_concurrent_downloads) - + async def download_single_media(message): async with semaphore: return await self.download_media(channel, message) - + batch_size = 10 for i in range(0, len(media_tasks), batch_size): - batch = media_tasks[i:i + batch_size] - tasks = [asyncio.create_task(download_single_media(msg)) for msg in batch] - + batch = media_tasks[i : i + batch_size] + tasks = [ + asyncio.create_task(download_single_media(msg)) for msg in batch + ] + for j, task in enumerate(tasks): try: media_path = await task if media_path: - await self.update_media_path(channel, batch[j].id, media_path) + await self.update_media_path( + channel, batch[j].id, media_path + ) successful_downloads += 1 except Exception: pass - + completed_media += 1 progress = (completed_media / total_media) * 100 bar_length = 30 filled_length = int(bar_length * completed_media // total_media) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - sys.stdout.write(f"\ršŸ“„ Media: [{bar}] {progress:.1f}% ({completed_media}/{total_media})") - sys.stdout.flush() - - print(f"\nāœ… Media download complete! ({successful_downloads}/{total_media} successful)") + bar = "ā–ˆ" * filled_length + "ā–‘" * (bar_length - filled_length) - self.state['channels'][channel] = last_message_id + sys.stdout.write( + f"\ršŸ“„ Media: [{bar}] {progress:.1f}% ({completed_media}/{total_media})" + ) + sys.stdout.flush() + + print( + f"\nāœ… Media download complete! ({successful_downloads}/{total_media} successful)" + ) + + self.state["channels"][channel] = last_message_id self.save_state() print(f"Completed scraping channel {channel}") @@ -825,11 +962,11 @@ class OptimizedTelegramScraper: async def rescrape_media(self, channel: str): conn = self.get_db_connection(channel) cursor = conn.cursor() - cursor.execute('SELECT message_id FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NULL') + cursor.execute( + 'SELECT message_id FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NULL' + ) message_ids = [row[0] for row in cursor.fetchall()] - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') - if not message_ids: # print(f"No media files to reprocess for {channel_name} (ID: {channel})") return @@ -848,28 +985,33 @@ class OptimizedTelegramScraper: batch_size = 10 for i in range(0, len(message_ids), batch_size): - batch_ids = message_ids[i:i + batch_size] + batch_ids = message_ids[i : i + batch_size] messages = await self.client.get_messages(entity, ids=batch_ids) - valid_messages = [msg for msg in messages if msg and msg.media and not isinstance(msg.media, MessageMediaWebPage)] - tasks = [asyncio.create_task(download_single_media(msg)) for msg in valid_messages] + valid_messages = [ + msg + for msg in messages + if msg + and msg.media + and not isinstance(msg.media, MessageMediaWebPage) + ] + tasks = [ + asyncio.create_task(download_single_media(msg)) + for msg in valid_messages + ] for j, task in enumerate(tasks): try: media_path = await task if media_path: - await self.update_media_path(channel, valid_messages[j].id, media_path) + await self.update_media_path( + channel, valid_messages[j].id, media_path + ) successful_downloads += 1 except Exception: pass completed_media += 1 - progress = (completed_media / len(message_ids)) * 100 - bar_length = 30 - filled_length = int(bar_length * completed_media // len(message_ids)) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - # sys.stdout.write(f"\ršŸ”„ Rescrape: [{bar}] {progress:.1f}% ({completed_media}/{len(message_ids)})") # sys.stdout.flush() # print(f"\nāœ… Media reprocessing complete! ({successful_downloads}/{len(message_ids)} successful)") @@ -881,15 +1023,19 @@ class OptimizedTelegramScraper: conn = self.get_db_connection(channel) cursor = conn.cursor() - cursor.execute('SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage"') + cursor.execute( + 'SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage"' + ) total_with_media = cursor.fetchone()[0] - cursor.execute('SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NOT NULL') + cursor.execute( + 'SELECT COUNT(*) FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND media_path IS NOT NULL' + ) total_with_files = cursor.fetchone()[0] missing_count = total_with_media - total_with_files - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') + channel_name = self.state.get("channel_names", {}).get(channel, "Unknown") print(f"\nšŸ“Š Media Analysis for {channel_name} (ID: {channel}):") print(f"Messages with media: {total_with_media}") print(f"Media files downloaded: {total_with_files}") @@ -899,14 +1045,18 @@ class OptimizedTelegramScraper: print("āœ… All media files are already downloaded!") return - cursor.execute('SELECT message_id, media_type FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND (media_path IS NULL OR media_path = "")') + cursor.execute( + 'SELECT message_id, media_type FROM messages WHERE media_type IS NOT NULL AND media_type != "MessageMediaWebPage" AND (media_path IS NULL OR media_path = "")' + ) missing_media = cursor.fetchall() if not missing_media: print("āœ… No missing media found!") return - print(f"\nšŸ”§ Attempting to download {len(missing_media)} missing media files...") + print( + f"\nšŸ”§ Attempting to download {len(missing_media)} missing media files..." + ) try: entity = await self._resolve_entity(channel) @@ -917,58 +1067,75 @@ class OptimizedTelegramScraper: async def download_single_media(message): async with semaphore: return await self.download_media(channel, message) - + batch_size = 10 for i in range(0, len(missing_media), batch_size): - batch = missing_media[i:i + batch_size] + batch = missing_media[i : i + batch_size] message_ids = [msg[0] for msg in batch] - + messages = await self.client.get_messages(entity, ids=message_ids) - valid_messages = [msg for msg in messages if msg and msg.media and not isinstance(msg.media, MessageMediaWebPage)] - - tasks = [asyncio.create_task(download_single_media(msg)) for msg in valid_messages] + valid_messages = [ + msg + for msg in messages + if msg + and msg.media + and not isinstance(msg.media, MessageMediaWebPage) + ] + + tasks = [ + asyncio.create_task(download_single_media(msg)) + for msg in valid_messages + ] for j, task in enumerate(tasks): try: media_path = await task if media_path: - await self.update_media_path(channel, valid_messages[j].id, media_path) + await self.update_media_path( + channel, valid_messages[j].id, media_path + ) successful_downloads += 1 except Exception: pass - + completed_media += 1 progress = (completed_media / len(missing_media)) * 100 bar_length = 30 - filled_length = int(bar_length * completed_media // len(missing_media)) - bar = 'ā–ˆ' * filled_length + 'ā–‘' * (bar_length - filled_length) - - sys.stdout.write(f"\ršŸ”§ Fix Media: [{bar}] {progress:.1f}% ({completed_media}/{len(missing_media)})") + filled_length = int( + bar_length * completed_media // len(missing_media) + ) + bar = "ā–ˆ" * filled_length + "ā–‘" * (bar_length - filled_length) + + sys.stdout.write( + f"\ršŸ”§ Fix Media: [{bar}] {progress:.1f}% ({completed_media}/{len(missing_media)})" + ) sys.stdout.flush() - print(f"\nāœ… Media fix complete! ({successful_downloads}/{len(missing_media)} successful)") + print( + f"\nāœ… Media fix complete! ({successful_downloads}/{len(missing_media)} successful)" + ) except Exception as e: print(f"Error fixing missing media: {e}") async def continuous_scraping(self): self.continuous_scraping_active = True - + try: while self.continuous_scraping_active: start_time = time.time() - - for channel in self.state['channels']: + + for channel in self.state["channels"]: if not self.continuous_scraping_active: break print(f"\nChecking for new messages in channel: {channel}") - await self.scrape_channel(channel, self.state['channels'][channel]) - + await self.scrape_channel(channel, self.state["channels"][channel]) + elapsed = time.time() - start_time sleep_time = max(0, 60 - elapsed) if sleep_time > 0: await asyncio.sleep(sleep_time) - + except asyncio.CancelledError: print("Continuous scraping stopped") finally: @@ -977,71 +1144,92 @@ class OptimizedTelegramScraper: async def continuous_scraping_with_forwarding(self): self.continuous_scraping_active = True self.forwarding_active = True - + if not self.client.is_connected(): await self.client.connect() - + rules = self.get_forwarding_rules() enabled_rules = [r for r in rules if r.enabled] forwarding_enabled = False - + if enabled_rules: forwarding_enabled = await self.setup_forwarding_handler() - - print("\n" + "="*50) + + print("\n" + "=" * 50) print(" CONTINUOUS SCRAPING + FORWARDING SERVICE") - print("="*50) - - if self.state['channels']: - print(f"\nScraping {len(self.state['channels'])} channel(s) every 60 seconds") - for channel in self.state['channels']: - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') + print("=" * 50) + + if self.state["channels"]: + print( + f"\nScraping {len(self.state['channels'])} channel(s) every 60 seconds" + ) + for channel in self.state["channels"]: + channel_name = self.state.get("channel_names", {}).get( + channel, "Unknown" + ) print(f" - {channel_name} ({channel})") else: print("\nNo channels configured for scraping") - + if forwarding_enabled: print(f"\nForwarding {len(enabled_rules)} rule(s):") for rule in enabled_rules: - source_name = self.state.get('channel_names', {}).get(rule.source_channel, rule.source_channel) - dest_name = self.state.get('channel_names', {}).get(rule.destination_channel, rule.destination_channel) + source_name = self.state.get("channel_names", {}).get( + rule.source_channel, rule.source_channel + ) + dest_name = self.state.get("channel_names", {}).get( + rule.destination_channel, rule.destination_channel + ) content_types = [] - if rule.forward_text: content_types.append("T") - if rule.forward_images: content_types.append("I") - if rule.forward_videos: content_types.append("V") - if rule.forward_documents: content_types.append("D") + if rule.forward_text: + content_types.append("T") + if rule.forward_images: + content_types.append("I") + if rule.forward_videos: + content_types.append("V") + if rule.forward_documents: + content_types.append("D") print(f" - {source_name} -> {dest_name} [{'/'.join(content_types)}]") else: print("\nNo forwarding rules enabled") - - print("\n" + "="*50) + + print("\n" + "=" * 50) print("Press Ctrl+C to stop") - print("="*50 + "\n") - + print("=" * 50 + "\n") + try: + async def scraping_loop(): while self.continuous_scraping_active: start_time = time.time() - - for channel in list(self.state['channels'].keys()): + + for channel in list(self.state["channels"].keys()): if not self.continuous_scraping_active: break - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') - print(f"\n[{time.strftime('%H:%M:%S')}] Checking: {channel_name}") - await self.scrape_channel(channel, self.state['channels'][channel]) - + channel_name = self.state.get("channel_names", {}).get( + channel, "Unknown" + ) + print( + f"\n[{time.strftime('%H:%M:%S')}] Checking: {channel_name}" + ) + await self.scrape_channel( + channel, self.state["channels"][channel] + ) + if self.continuous_scraping_active: elapsed = time.time() - start_time sleep_time = max(0, 60 - elapsed) if sleep_time > 0: - next_check = time.strftime('%H:%M:%S', time.localtime(time.time() + sleep_time)) + next_check = time.strftime( + "%H:%M:%S", time.localtime(time.time() + sleep_time) + ) print(f"\nNext scrape cycle at {next_check}") await asyncio.sleep(sleep_time) - - scraping_task = asyncio.create_task(scraping_loop()) - + + asyncio.create_task(scraping_loop()) + await self.client.run_until_disconnected() - + except asyncio.CancelledError: pass except ConnectionError: @@ -1054,19 +1242,19 @@ class OptimizedTelegramScraper: print("\nCombined service stopped") def get_export_filename(self, channel: str): - username = self.state.get('channel_names', {}).get(channel, 'no_username') + username = self.state.get("channel_names", {}).get(channel, "no_username") return f"{channel}_{username}" def export_to_csv(self, channel: str): conn = self.get_db_connection(channel) filename = self.get_export_filename(channel) - csv_file = self.DATA_DIR / channel / f'{filename}.csv' + csv_file = self.DATA_DIR / channel / f"{filename}.csv" cursor = conn.cursor() - cursor.execute('SELECT * FROM messages ORDER BY date') + cursor.execute("SELECT * FROM messages ORDER BY date") columns = [description[0] for description in cursor.description] - with open(csv_file, 'w', newline='', encoding='utf-8') as f: + with open(csv_file, "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(columns) @@ -1079,14 +1267,14 @@ class OptimizedTelegramScraper: def export_to_json(self, channel: str): conn = self.get_db_connection(channel) filename = self.get_export_filename(channel) - json_file = self.DATA_DIR / channel / f'{filename}.json' + json_file = self.DATA_DIR / channel / f"{filename}.json" cursor = conn.cursor() - cursor.execute('SELECT * FROM messages ORDER BY date') + cursor.execute("SELECT * FROM messages ORDER BY date") columns = [description[0] for description in cursor.description] - with open(json_file, 'w', encoding='utf-8') as f: - f.write('[\n') + with open(json_file, "w", encoding="utf-8") as f: + f.write("[\n") first_row = True while True: @@ -1096,21 +1284,21 @@ class OptimizedTelegramScraper: for row in rows: if not first_row: - f.write(',\n') + f.write(",\n") else: first_row = False data = dict(zip(columns, row)) json.dump(data, f, ensure_ascii=False, indent=2) - f.write('\n]') + f.write("\n]") async def export_data(self): - if not self.state['channels']: + if not self.state["channels"]: print("No channels to export") return - - for channel in self.state['channels']: + + for channel in self.state["channels"]: print(f"Exporting data for channel {channel}...") try: self.export_to_csv(channel) @@ -1120,7 +1308,7 @@ class OptimizedTelegramScraper: print(f"āŒ Export failed for channel {channel}: {e}") async def view_channels(self): - if not self.state['channels']: + if not self.state["channels"]: print("No channels saved") return @@ -1128,38 +1316,51 @@ class OptimizedTelegramScraper: await self.client.connect() print("\nCurrent channels:") - for i, (channel, last_id) in enumerate(self.state['channels'].items(), 1): + for i, (channel, last_id) in enumerate(self.state["channels"].items(), 1): try: - channel_name = self.state.get('channel_names', {}).get(channel) - - if not channel_name or channel_name in ['Unknown', 'no_username']: + channel_name = self.state.get("channel_names", {}).get(channel) + + if not channel_name or channel_name in ["Unknown", "no_username"]: try: entity = await self._resolve_entity(channel) - channel_name = getattr(entity, 'title', None) or getattr(entity, 'first_name', None) or getattr(entity, 'username', None) or 'Unknown' + channel_name = ( + getattr(entity, "title", None) + or getattr(entity, "first_name", None) + or getattr(entity, "username", None) + or "Unknown" + ) # For users, combine first+last name if isinstance(entity, User) and entity.first_name: - channel_name = ' '.join(filter(None, [entity.first_name, entity.last_name])) - if 'channel_names' not in self.state: - self.state['channel_names'] = {} - self.state['channel_names'][channel] = channel_name + channel_name = " ".join( + filter(None, [entity.first_name, entity.last_name]) + ) + if "channel_names" not in self.state: + self.state["channel_names"] = {} + self.state["channel_names"][channel] = channel_name self.save_state() - except: - channel_name = 'Unknown' - + except Exception: + channel_name = "Unknown" + conn = self.get_db_connection(channel) cursor = conn.cursor() - cursor.execute('SELECT COUNT(*) FROM messages') + cursor.execute("SELECT COUNT(*) FROM messages") count = cursor.fetchone()[0] - print(f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}, Messages: {count}") - except: - channel_name = self.state.get('channel_names', {}).get(channel, 'Unknown') - print(f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}") + print( + f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}, Messages: {count}" + ) + except Exception: + channel_name = self.state.get("channel_names", {}).get( + channel, "Unknown" + ) + print( + f"[{i}] {channel_name} (ID: {channel}), Last Message ID: {last_id}" + ) async def list_channels(self): try: if not self.client.is_connected(): await self.client.connect() - + print("\nList of channels, groups and private chats joined by account:") count = 1 channels_data = [] @@ -1172,22 +1373,37 @@ class OptimizedTelegramScraper: channel_type = "Group" else: channel_type = "Private" - username = getattr(entity, 'username', None) or 'no_username' - display_username = f"@{username}" if username != 'no_username' else '(no username)' - print(f"[{count}] {dialog.title} (ID: {dialog.id}, Type: {channel_type}, Username: {display_username})") - channels_data.append({ - 'number': count, - 'channel_name': dialog.title, - 'channel_id': str(dialog.id), - 'username': username, - 'type': channel_type - }) + username = getattr(entity, "username", None) or "no_username" + display_username = ( + f"@{username}" if username != "no_username" else "(no username)" + ) + print( + f"[{count}] {dialog.title} (ID: {dialog.id}, Type: {channel_type}, Username: {display_username})" + ) + channels_data.append( + { + "number": count, + "channel_name": dialog.title, + "channel_id": str(dialog.id), + "username": username, + "type": channel_type, + } + ) count += 1 if channels_data: - csv_file = self.DATA_DIR / 'channels_list.csv' - with open(csv_file, 'w', newline='', encoding='utf-8') as f: - writer = csv.DictWriter(f, fieldnames=['number', 'channel_name', 'channel_id', 'username', 'type']) + csv_file = self.DATA_DIR / "channels_list.csv" + with open(csv_file, "w", newline="", encoding="utf-8") as f: + writer = csv.DictWriter( + f, + fieldnames=[ + "number", + "channel_name", + "channel_id", + "username", + "type", + ], + ) writer.writeheader() writer.writerows(channels_data) print(f"\nāœ… Saved channels list to {csv_file}") @@ -1202,7 +1418,7 @@ class OptimizedTelegramScraper: qr = qrcode.QRCode(box_size=1, border=1) qr.add_data(qr_login.url) qr.make() - + f = StringIO() qr.print_ascii(out=f) f.seek(0) @@ -1214,10 +1430,10 @@ class OptimizedTelegramScraper: print("1. Open Telegram on your phone") print("2. Go to Settings > Devices > Scan QR") print("3. Scan the code below\n") - + qr_login = await self.client.qr_login() self.display_qr_code_ascii(qr_login) - + try: await qr_login.wait() print("\nāœ… Successfully logged in via QR code!") @@ -1235,7 +1451,7 @@ class OptimizedTelegramScraper: phone = input("Enter your phone number: ") await self.client.send_code_request(phone) code = input("Enter the code you received: ") - + try: await self.client.sign_in(phone, code) print("\nāœ… Successfully logged in via phone!") @@ -1250,45 +1466,53 @@ class OptimizedTelegramScraper: return False async def initialize_client(self, interactive: bool = True): - if not all([self.state.get('api_id'), self.state.get('api_hash')]): + if not all([self.state.get("api_id"), self.state.get("api_hash")]): if not interactive: print("API credentials are missing in state.json") return False print("\n=== API Configuration Required ===") print("You need to provide API credentials from https://my.telegram.org") try: - self.state['api_id'] = int(input("Enter your API ID: ")) - self.state['api_hash'] = input("Enter your API Hash: ") + self.state["api_id"] = int(input("Enter your API ID: ")) + self.state["api_hash"] = input("Enter your API Hash: ") self.save_state() except ValueError: print("Invalid API ID. Must be a number.") return False - self.client = TelegramClient(str(self.SESSION_DIR / 'session'), self.state['api_id'], self.state['api_hash']) - + self.client = TelegramClient( + str(self.SESSION_DIR / "session"), + self.state["api_id"], + self.state["api_hash"], + ) + try: await self.client.connect() except Exception as e: print(f"Failed to connect: {e}") return False - + if not await self.client.is_user_authorized(): if not interactive: - print("Telegram session is not authorized. Run the CLI once to sign in.") + print( + "Telegram session is not authorized. Run the CLI once to sign in." + ) await self.client.disconnect() return False print("\n=== Choose Authentication Method ===") print("[1] QR Code (Recommended - No phone number needed)") print("[2] Phone Number (Traditional method)") - + while True: choice = input("Enter your choice (1 or 2): ").strip() - if choice in ['1', '2']: + if choice in ["1", "2"]: break print("Please enter 1 or 2") - - success = await self.qr_code_auth() if choice == '1' else await self.phone_auth() - + + success = ( + await self.qr_code_auth() if choice == "1" else await self.phone_auth() + ) + if not success: print("Authentication failed. Please try again.") await self.client.disconnect() @@ -1296,20 +1520,20 @@ class OptimizedTelegramScraper: else: pass # print("āœ… Already authenticated!") - + return True def parse_channel_selection(self, choice): - channels_list = list(self.state['channels'].keys()) + channels_list = list(self.state["channels"].keys()) selected_channels = [] - - if choice.lower() == 'all': + + if choice.lower() == "all": return channels_list - - for selection in [x.strip() for x in choice.split(',')]: + + for selection in [x.strip() for x in choice.split(",")]: try: - if selection.startswith('-'): - if selection in self.state['channels']: + if selection.startswith("-"): + if selection in self.state["channels"]: selected_channels.append(selection) else: print(f"Channel ID {selection} not found in your channels") @@ -1318,14 +1542,18 @@ class OptimizedTelegramScraper: if 1 <= num <= len(channels_list): selected_channels.append(channels_list[num - 1]) else: - print(f"Invalid channel number: {num}. Valid range: 1-{len(channels_list)}") + print( + f"Invalid channel number: {num}. Valid range: 1-{len(channels_list)}" + ) except ValueError: - print(f"Invalid input: {selection}. Use numbers (1,2,3) or full IDs (-100123...)") - + print( + f"Invalid input: {selection}. Use numbers (1,2,3) or full IDs (-100123...)" + ) + return selected_channels async def scrape_specific_channels(self): - if not self.state['channels']: + if not self.state["channels"]: print("No channels available. Use [L] to add channels first") return @@ -1334,28 +1562,30 @@ class OptimizedTelegramScraper: print("• Single: 1 or -1001234567890") print("• Multiple: 1,3,5 or mix formats") print("• All channels: all") - + choice = input("\nEnter selection: ").strip() selected_channels = self.parse_channel_selection(choice) - + if selected_channels: print(f"\nšŸš€ Starting scrape of {len(selected_channels)} channel(s)...") for i, channel in enumerate(selected_channels, 1): print(f"\n[{i}/{len(selected_channels)}] Scraping: {channel}") - await self.scrape_channel(channel, self.state['channels'][channel]) + await self.scrape_channel(channel, self.state["channels"][channel]) print(f"\nāœ… Completed scraping {len(selected_channels)} channel(s)!") else: print("āŒ No valid channels selected") async def manage_channels(self): while True: - print("\n" + "="*40) + print("\n" + "=" * 40) print(" TELEGRAM SCRAPER") - print("="*40) + print("=" * 40) print("[S] Scrape channels") print("[C] Continuous scraping") print("[B] Combined scrape + forward") - print(f"[M] Media scraping: {'ON' if self.state['scrape_media'] else 'OFF'}") + print( + f"[M] Media scraping: {'ON' if self.state['scrape_media'] else 'OFF'}" + ) print("[L] List & add channels") print("[R] Remove channels") print("[E] Export data") @@ -1363,33 +1593,33 @@ class OptimizedTelegramScraper: print("[F] Fix missing media") print("[W] Forwarding rules") print("[Q] Quit") - print("="*40) + print("=" * 40) choice = input("Enter your choice: ").lower().strip() - + try: - if choice == 'r': - if not self.state['channels']: + if choice == "r": + if not self.state["channels"]: print("No channels to remove") continue - + await self.view_channels() print("\nTo remove channels:") print("• Single: 1 or -1001234567890") print("• Multiple: 1,2,3 or mix formats") selection = input("Enter selection: ").strip() selected_channels = self.parse_channel_selection(selection) - + if selected_channels: removed_count = 0 for channel in selected_channels: - if channel in self.state['channels']: - del self.state['channels'][channel] + if channel in self.state["channels"]: + del self.state["channels"][channel] print(f"āœ… Removed channel {channel}") removed_count += 1 else: print(f"āŒ Channel {channel} not found") - + if removed_count > 0: self.save_state() print(f"\nšŸŽ‰ Removed {removed_count} channel(s)!") @@ -1398,20 +1628,22 @@ class OptimizedTelegramScraper: print("No channels were removed") else: print("No valid channels selected") - - elif choice == 's': + + elif choice == "s": await self.scrape_specific_channels() - - elif choice == 'm': - self.state['scrape_media'] = not self.state['scrape_media'] + + elif choice == "m": + self.state["scrape_media"] = not self.state["scrape_media"] self.save_state() - print(f"\nāœ… Media scraping {'enabled' if self.state['scrape_media'] else 'disabled'}") - - elif choice == 'c': + print( + f"\nāœ… Media scraping {'enabled' if self.state['scrape_media'] else 'disabled'}" + ) + + elif choice == "c": task = asyncio.create_task(self.continuous_scraping()) print("Continuous scraping started. Press Ctrl+C to stop.") try: - await asyncio.sleep(float('inf')) + await asyncio.sleep(float("inf")) except KeyboardInterrupt: self.continuous_scraping_active = False task.cancel() @@ -1420,11 +1652,18 @@ class OptimizedTelegramScraper: await task except asyncio.CancelledError: pass - - elif choice == 'b': - if not self.state['channels'] and not any(r.get('enabled', True) for r in self.state.get('forwarding_rules', [])): - print("āŒ No channels to scrape and no forwarding rules enabled.") - print(" Add channels with [L] or configure forwarding with [W] first.") + + elif choice == "b": + if not self.state["channels"] and not any( + r.get("enabled", True) + for r in self.state.get("forwarding_rules", []) + ): + print( + "āŒ No channels to scrape and no forwarding rules enabled." + ) + print( + " Add channels with [L] or configure forwarding with [W] first." + ) continue try: await self.continuous_scraping_with_forwarding() @@ -1432,11 +1671,11 @@ class OptimizedTelegramScraper: self.continuous_scraping_active = False self.forwarding_active = False print("\nStopping combined service...") - - elif choice == 'e': + + elif choice == "e": await self.export_data() - - elif choice == 'l': + + elif choice == "l": channels_data = await self.list_channels() if not channels_data: @@ -1452,24 +1691,37 @@ class OptimizedTelegramScraper: if selection: added_count = 0 - if selection.lower() == 'all': + if selection.lower() == "all": for channel_info in channels_data: - channel_id = channel_info['channel_id'] - if channel_id not in self.state['channels']: - self.state['channels'][channel_id] = 0 - if 'channel_names' not in self.state: - self.state['channel_names'] = {} - self.state['channel_names'][channel_id] = channel_info['username'] - print(f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})") + channel_id = channel_info["channel_id"] + if channel_id not in self.state["channels"]: + self.state["channels"][channel_id] = 0 + if "channel_names" not in self.state: + self.state["channel_names"] = {} + self.state["channel_names"][channel_id] = ( + channel_info["username"] + ) + print( + f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})" + ) added_count += 1 else: - print(f"Channel {channel_info['channel_name']} already added") + print( + f"Channel {channel_info['channel_name']} already added" + ) else: - for sel in [x.strip() for x in selection.split(',')]: + for sel in [x.strip() for x in selection.split(",")]: try: - if sel.startswith('-'): + if sel.startswith("-"): channel_id = sel - channel_info = next((c for c in channels_data if c['channel_id'] == channel_id), None) + channel_info = next( + ( + c + for c in channels_data + if c["channel_id"] == channel_id + ), + None, + ) if not channel_info: print(f"Channel ID {channel_id} not found") continue @@ -1477,19 +1729,27 @@ class OptimizedTelegramScraper: num = int(sel) if 1 <= num <= len(channels_data): channel_info = channels_data[num - 1] - channel_id = channel_info['channel_id'] + channel_id = channel_info["channel_id"] else: - print(f"Invalid number: {num}. Choose 1-{len(channels_data)}") + print( + f"Invalid number: {num}. Choose 1-{len(channels_data)}" + ) continue - if channel_id in self.state['channels']: - print(f"Channel {channel_info['channel_name']} already added") + if channel_id in self.state["channels"]: + print( + f"Channel {channel_info['channel_name']} already added" + ) else: - self.state['channels'][channel_id] = 0 - if 'channel_names' not in self.state: - self.state['channel_names'] = {} - self.state['channel_names'][channel_id] = channel_info['username'] - print(f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})") + self.state["channels"][channel_id] = 0 + if "channel_names" not in self.state: + self.state["channel_names"] = {} + self.state["channel_names"][channel_id] = ( + channel_info["username"] + ) + print( + f"āœ… Added channel {channel_info['channel_name']} (ID: {channel_id})" + ) added_count += 1 except ValueError: @@ -1501,17 +1761,19 @@ class OptimizedTelegramScraper: await self.view_channels() else: print("No new channels were added") - - elif choice == 't': - if not self.state['channels']: + + elif choice == "t": + if not self.state["channels"]: print("No channels available. Add channels first") continue - + await self.view_channels() - print("\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)") + print( + "\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)" + ) selection = input("Enter your selection: ").strip() selected_channels = self.parse_channel_selection(selection) - + if len(selected_channels) == 1: channel = selected_channels[0] print(f"Rescaping media for channel: {channel}") @@ -1520,17 +1782,19 @@ class OptimizedTelegramScraper: print("Please select only one channel for media rescaping") else: print("No valid channel selected") - - elif choice == 'f': - if not self.state['channels']: + + elif choice == "f": + if not self.state["channels"]: print("No channels available. Add channels first") continue - + await self.view_channels() - print("\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)") + print( + "\nEnter channel NUMBER (1,2,3...) or full channel ID (-100123...)" + ) selection = input("Enter your selection: ").strip() selected_channels = self.parse_channel_selection(selection) - + if len(selected_channels) == 1: channel = selected_channels[0] await self.fix_missing_media(channel) @@ -1538,20 +1802,20 @@ class OptimizedTelegramScraper: print("Please select only one channel for fixing missing media") else: print("No valid channel selected") - - elif choice == 'w': + + elif choice == "w": await self.manage_forwarding_rules() - - elif choice == 'q': + + elif choice == "q": print("\nšŸ‘‹ Goodbye!") self.close_db_connections() if self.client: await self.client.disconnect() sys.exit() - + else: print("Invalid option") - + except Exception as e: print(f"Error: {e}") @@ -1567,11 +1831,13 @@ class OptimizedTelegramScraper: else: print("Failed to initialize client. Exiting.") + async def main(): scraper = OptimizedTelegramScraper() await scraper.run() -if __name__ == '__main__': + +if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: diff --git a/webui_server.py b/webui_server.py index c886d7a..bac2aaf 100644 --- a/webui_server.py +++ b/webui_server.py @@ -306,7 +306,9 @@ class JobRunner: if job.status == "running": running_jobs.append(job) if running_jobs: - logger.warning("Waiting for %d running job(s) to finish...", len(running_jobs)) + logger.warning( + "Waiting for %d running job(s) to finish...", len(running_jobs) + ) deadline = time.time() + timeout while time.time() < deadline: with self.lock: @@ -561,6 +563,7 @@ class TelegramAuthManager: async def _disconnect(): if self.client: await self.client.disconnect() + if self.client: try: future = asyncio.run_coroutine_threadsafe(_disconnect(), self.loop) @@ -831,7 +834,10 @@ def openapi_payload() -> Dict[str, Any]: "requestBody": json_body( { "api_id": {"type": "integer", "example": 123456}, - "api_hash": {"type": "string", "example": "0123456789abcdef"}, + "api_hash": { + "type": "string", + "example": "0123456789abcdef", + }, }, ["api_id", "api_hash"], ), @@ -929,7 +935,11 @@ def openapi_payload() -> Dict[str, Any]: { "name": "limit", "in": "query", - "schema": {"type": "integer", "default": 120, "maximum": 300}, + "schema": { + "type": "integer", + "default": 120, + "maximum": 300, + }, }, { "name": "before", @@ -945,7 +955,10 @@ def openapi_payload() -> Dict[str, Any]: "summary": "Track a channel", "requestBody": json_body( { - "channel_id": {"type": "string", "example": "-1001234567890"}, + "channel_id": { + "type": "string", + "example": "-1001234567890", + }, "name": {"type": "string", "example": "Research feed"}, }, ["channel_id"], @@ -990,7 +1003,10 @@ def openapi_payload() -> Dict[str, Any]: "schema": {"type": "string"}, } ], - "responses": {**json_response, "404": {"description": "Job not found"}}, + "responses": { + **json_response, + "404": {"description": "Job not found"}, + }, } }, "/api/jobs/scrape": { @@ -1416,15 +1432,18 @@ class TelegramScraperRequestHandler(BaseHTTPRequestHandler): try: start, end = self._parse_range(range_header, file_size) except ValueError: - self.send_error_json(HTTPStatus.REQUESTED_RANGE_NOT_SATISFIABLE, "Invalid range") + self.send_error_json( + HTTPStatus.REQUESTED_RANGE_NOT_SATISFIABLE, "Invalid range" + ) return content_length = end - start + 1 self.send_response(HTTPStatus.PARTIAL_CONTENT) self.send_header("Content-Type", content_type) self.send_header("Accept-Ranges", "bytes") - self.send_header("Last-Modified", time.strftime( - "%a, %d %b %Y %H:%M:%S GMT", time.gmtime(last_modified) - )) + self.send_header( + "Last-Modified", + time.strftime("%a, %d %b %Y %H:%M:%S GMT", time.gmtime(last_modified)), + ) self.send_header("Content-Range", f"bytes {start}-{end}/{file_size}") self.send_header("Content-Length", str(content_length)) self.end_headers() @@ -1434,9 +1453,10 @@ class TelegramScraperRequestHandler(BaseHTTPRequestHandler): self.send_response(HTTPStatus.OK) self.send_header("Content-Type", content_type) self.send_header("Accept-Ranges", "bytes") - self.send_header("Last-Modified", time.strftime( - "%a, %d %b %Y %H:%M:%S GMT", time.gmtime(last_modified) - )) + self.send_header( + "Last-Modified", + time.strftime("%a, %d %b %Y %H:%M:%S GMT", time.gmtime(last_modified)), + ) self.send_header("Content-Length", str(file_size)) self.end_headers() if not head_only: @@ -1495,7 +1515,9 @@ class TelegramScraperWebServer(ThreadingHTTPServer): self.job_runner = JobRunner() self.auth_manager = TelegramAuthManager() self.continuous_manager = ContinuousScrapeManager() - if START_CONTINUOUS and self.continuous_manager.snapshot()["config"].get("enabled", True): + if START_CONTINUOUS and self.continuous_manager.snapshot()["config"].get( + "enabled", True + ): self.continuous_manager.start()