import os import sys import re import asyncio import argparse from pathlib import Path from urllib.parse import urlparse from tqdm import tqdm # Add parent workspace directory to path to import scraper_core sys.path.append(str(Path(__file__).resolve().parents[1])) import scraper_core BASE_DIR = Path(__file__).resolve().parent DOWNLOAD_DIR = BASE_DIR / "videos" async def extract_telegram_media(page): """Scrape the current page for video and image links, returning (media_list, min_msg_id).""" messages = await page.query_selector_all(".tgme_widget_message") media_list = [] msg_ids = [] for msg in messages: # Extract message link to get the ID for pagination link_el = await msg.query_selector("a.tgme_widget_message_date") if not link_el: continue href = await link_el.get_attribute("href") if not href: continue parsed_path = urlparse(href).path.rstrip("/").split("/") if not parsed_path or not parsed_path[-1].isdigit(): continue msg_id = int(parsed_path[-1]) msg_ids.append(msg_id) # 1. Look for Video video_el = await msg.query_selector(".tgme_widget_message_video_player video") if video_el: video_src = await video_el.get_attribute("src") if video_src: media_list.append({ "id": msg_id, "url": video_src, "type": "video", "ext": ".mp4" }) continue # 2. Look for Image (Photo) photo_el = await msg.query_selector(".tgme_widget_message_photo_wrap") if photo_el: style = await photo_el.get_attribute("style") if style: # Extract URL from background-image: url('...') m = re.search(r"background-image:\s*url\(['\"]?(https://[^'\"]+)['\"]?\)", style) if m: media_list.append({ "id": msg_id, "url": m.group(1), "type": "photo", "ext": ".jpg" }) continue min_id = min(msg_ids) if msg_ids else None return media_list, min_id async def scrape_channel(page, channel_name: str, limit: int): """Crawl a public Telegram channel backwards in time to gather media URLs.""" tqdm.write(f"Scraping channel '{channel_name}' ...") base_url = f"https://t.me/s/{channel_name}" all_media = [] seen_ids = set() current_url = base_url while len(all_media) < limit: tqdm.write(f" Fetching page: {current_url}") try: await page.goto(current_url, wait_until="domcontentloaded", timeout=30000) await page.wait_for_timeout(3000) except Exception as e: tqdm.write(f" Error loading Telegram web page: {e}") break page_media, min_id = await extract_telegram_media(page) # Filter new media new_items = [] for item in page_media: if item["id"] not in seen_ids: seen_ids.add(item["id"]) new_items.append(item) if not new_items: tqdm.write(" No new media found on this page.") break all_media.extend(new_items) tqdm.write(f" Found {len(new_items)} new media items (Total collected: {len(all_media)})") if not min_id: break # Paginate to messages before the minimum ID we've seen current_url = f"{base_url}?before={min_id}" await page.wait_for_timeout(1000) return all_media[:limit] async def worker(queue, scraper, skip_existing, bar_pool, overall_bar, channel_name): """Worker task that downloads media concurrently.""" while True: item = await queue.get() if item is None: queue.task_done() break msg_id, url, mtype, ext = item filename = f"msg_{msg_id}{ext}" dest_path = DOWNLOAD_DIR / channel_name / filename if skip_existing and scraper_core.is_already_downloaded(dest_path): overall_bar.update(1) queue.task_done() continue pos = bar_pool.acquire() or 1 success = await asyncio.to_thread( scraper_core.download_file, url, dest_path, None, pos, f"https://t.me/s/{channel_name}" ) bar_pool.release(pos) overall_bar.update(1) queue.task_done() async def process_channel(channel_name, limit, concurrency, skip_existing): scraper = scraper_core.PlaywrightScraper() await scraper.start() page = await scraper.new_page() media_items = await scrape_channel(page, channel_name, limit) await page.close() if not media_items: tqdm.write("No media files found to download.") await scraper.close() return tqdm.write(f"Downloading {len(media_items)} files with concurrency {concurrency} ...") queue = asyncio.Queue() for item in media_items: await queue.put((item["id"], item["url"], item["type"], item["ext"])) for _ in range(concurrency): await queue.put(None) bar_pool = scraper_core.BarPositionPool(concurrency) overall_bar = tqdm( total=len(media_items), desc=f"Channel: {channel_name}", position=0, leave=True, ncols=80, ) workers = [ asyncio.create_task(worker(queue, scraper, skip_existing, bar_pool, overall_bar, channel_name)) for _ in range(concurrency) ] await asyncio.gather(*workers) overall_bar.close() sys.stdout.write("\n" * (concurrency + 1)) sys.stdout.flush() await scraper.close() def main(): parser = argparse.ArgumentParser( description="Download media from public Telegram channels (headless, concurrent, multithreaded)." ) parser.add_argument( "channel", help="Telegram channel username (e.g. 'durov')", ) parser.add_argument( "--limit", type=int, default=50, help="Maximum number of media items to download (default: 50).", ) parser.add_argument( "--concurrency", type=int, default=3, help="Number of concurrent downloads (default: 3).", ) parser.add_argument( "--skip-existing", action="store_true", default=True, help="Skip already-downloaded files (default: true).", ) parser.add_argument( "--no-skip-existing", action="store_false", dest="skip_existing", help="Re-download existing files.", ) args = parser.parse_args() asyncio.run(process_channel(args.channel, args.limit, args.concurrency, args.skip_existing)) if __name__ == "__main__": main()