diff --git a/.claude/skills/telegram-scraper/SKILL.md b/.claude/skills/telegram-scraper/SKILL.md new file mode 100644 index 0000000..f7e786c --- /dev/null +++ b/.claude/skills/telegram-scraper/SKILL.md @@ -0,0 +1,75 @@ +--- +name: telegram-scraper +description: Read, search, export and analyze PUBLIC Telegram channels (t.me/) without an API key or login, using the tgscraper CLI / Python library. Use when the user mentions a Telegram channel, a t.me link, wants the latest posts, news or announcements from Telegram, wants to monitor a channel, export posts to CSV/JSON/Excel/SQLite, download channel photos/videos, or get channel statistics (views, top posts, posting times, hashtags, sentiment). +--- + +# Telegram Scraper + +`tgscraper` reads the public web preview of Telegram channels (`https://t.me/s/`). +No API key, phone number or login. Only **public channels with web preview** work — not private +groups, chats or bots. + +## 0. Setup (once) + +```bash +tgscraper --version || pip install "tgscraper[all] @ git+https://github.com/specialteam/TelegramScraper" +``` + +If the `telegram-scraper` MCP tools (`get_messages`, `search_messages`, `get_channel_info`, +`analyze_channel`, ...) are available, prefer them over the shell — same features. + +Channel arguments accept `durov`, `@durov`, `t.me/durov` or `https://t.me/s/durov`. +For a post link `https://t.me/durov/123` the channel is `durov` and the message id is `123`. + +## 1. Pick the right command + +| User wants | Command | +|---|---| +| Latest posts | `tgscraper durov -n 20 --json` | +| Posts in a date range | `tgscraper durov -n 0 --since 2026-01-01 --until 2026-01-31 --json` | +| Posts about a topic (whole history) | `tgscraper search durov "privacy" -n 30 --json` | +| Filter recent posts by words | `tgscraper durov -n 200 -k bitcoin -k btc --json` | +| Only posts with photos/videos | `tgscraper durov --media-only --media-type photo --json` | +| Popular posts | `tgscraper durov -n 300 --min-views 100000 --json` | +| Channel info (subscribers, description) | `tgscraper info durov --json` | +| Statistics / analysis | `tgscraper stats durov -n 300 --json` | +| Save to a file | `tgscraper durov -n 500 -o durov.csv` (`.json .jsonl .xlsx .db .md`) | +| Several channels | `tgscraper chan1 chan2 chan3 -n 50 -o all.db` | +| Only new posts since last run | `tgscraper durov --incremental -o archive.db` | +| Download photos/videos | `tgscraper media durov -n 30 -d ./media` | +| Monitor for new posts | `tgscraper watch durov -i 120 --webhook URL` (long-running; run in background) | + +`-n 0` means "no limit" (whole history — can be slow for big channels; combine with `--since`). + +## 2. Output (`--json`) + +A JSON array, newest first. Each message: + +```json +{"id": 123, "channel": "durov", "url": "https://t.me/durov/123", "date": "2026-01-10T09:30:00+00:00", + "text": "...", "views": 1250000, "author": null, "edited": false, "forwarded_from": null, "reply_to": null, + "media": [{"type": "photo", "url": "https://cdn.../x.jpg"}], "reactions": {"👍": 15000}, + "hashtags": ["news"], "mentions": ["telegram"], "links": ["https://..."], "html": "..."} +``` + +Pipe large results through `jq` instead of reading everything, e.g. +`tgscraper durov -n 200 --json | jq '[.[] | {url, views, text: .text[:120]}]'`. + +## 3. Python (for custom processing) + +```python +import tgscraper as tg +posts = tg.scrape("durov", limit=100, since="2026-01-01", keywords=["ton"]) +tg.export(posts, "out.xlsx") +stats = tg.summarize(posts) # dict: top_posts, posts_by_hour, top_hashtags, sentiment... +results = tg.scrape_many(["a", "b"]) # concurrent, {channel: [Message] | Exception} +``` + +## 4. Answering well + +- Always cite posts with their `url` and date. +- Report views as numbers (`1.2M`), and say how many posts you analyzed. +- `sentiment` is a rough lexicon score (-1..1), mention that it is approximate. +- Errors: `ChannelNotFound` → the channel is private, misspelled, or has web preview disabled. + Network / HTTP 429 errors → wait and retry, or pass a proxy with `-p socks5://host:port`. +- Respect privacy and Telegram's terms: only public data, reasonable request volume. diff --git a/.github/workflows/live-test.yml b/.github/workflows/live-test.yml new file mode 100644 index 0000000..b58c7f7 --- /dev/null +++ b/.github/workflows/live-test.yml @@ -0,0 +1,67 @@ +name: Live test + +# Scrapes the real t.me to catch Telegram HTML changes. +on: + pull_request: + workflow_dispatch: + inputs: + channel: + description: "Channel to test" + default: "durov" + schedule: + - cron: "17 6 * * 1" # weekly + +jobs: + live: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: "3.12" + - run: pip install -e ".[dev]" + - name: Live pytest + env: + TGSCRAPER_LIVE: "1" + TGSCRAPER_LIVE_CHANNEL: ${{ inputs.channel || 'durov' }} + run: pytest tests/test_live.py -v -s + - name: Show raw reaction markup + run: | + python - <<'EOF' + import httpx + from bs4 import BeautifulSoup + soup = BeautifulSoup(httpx.get("https://t.me/s/durov").text, "html.parser") + for r in soup.select(".tgme_reaction")[:6]: + print(r) + EOF + - name: CLI against real channels + run: | + tgscraper info durov telegram + tgscraper durov -n 5 + tgscraper stats durov -n 50 + tgscraper durov -n 30 -o out/durov.xlsx + tgscraper durov -n 30 -o out/durov.db + tgscraper durov -n 3 --json | python -c "import json,sys; d=json.load(sys.stdin); print(json.dumps(d[0], ensure_ascii=False, indent=2))" + - name: MCP server over stdio + run: | + python - <<'EOF' + import asyncio + from mcp import ClientSession, StdioServerParameters + from mcp.client.stdio import stdio_client + + async def main(): + async with stdio_client(StdioServerParameters(command="tgscraper-mcp")) as (r, w): + async with ClientSession(r, w) as s: + await s.initialize() + res = await s.call_tool("get_messages", {"channel": "durov", "limit": 2}) + text = res.content[0].text + print(text[:1500]) + assert '"error"' not in text, text + + asyncio.run(main()) + EOF + - uses: actions/upload-artifact@v4 + if: always() + with: + name: live-output + path: out/ diff --git a/.github/workflows/python-package.yml b/.github/workflows/python-package.yml index e56abb6..d3270e5 100644 --- a/.github/workflows/python-package.yml +++ b/.github/workflows/python-package.yml @@ -1,6 +1,3 @@ -# This workflow will install Python dependencies, run tests and lint with a variety of Python versions -# For more information see: https://docs.github.com/en/actions/automating-builds-and-tests/building-and-testing-python - name: Python package on: @@ -11,30 +8,27 @@ on: jobs: build: - runs-on: ubuntu-latest strategy: fail-fast: false matrix: - python-version: ["3.9", "3.10", "3.11"] + python-version: ["3.9", "3.10", "3.11", "3.12", "3.13"] steps: - uses: actions/checkout@v4 - name: Set up Python ${{ matrix.python-version }} - uses: actions/setup-python@v3 + uses: actions/setup-python@v5 with: python-version: ${{ matrix.python-version }} - - name: Install dependencies + - name: Install run: | python -m pip install --upgrade pip - python -m pip install flake8 pytest - if [ -f requirements.txt ]; then pip install -r requirements.txt; fi + pip install -e ".[dev]" - name: Lint with flake8 run: | - # stop the build if there are Python syntax errors or undefined names flake8 . --count --select=E9,F63,F7,F82 --show-source --statistics - # exit-zero treats all errors as warnings. The GitHub editor is 127 chars wide - flake8 . --count --exit-zero --max-complexity=10 --max-line-length=127 --statistics - - name: Test with pytest - run: | - pytest + flake8 . --count --exit-zero --max-complexity=12 --max-line-length=127 --statistics + - name: Test + run: pytest -q + - name: CLI smoke test + run: tgscraper --help diff --git a/.mcp.json b/.mcp.json new file mode 100644 index 0000000..47f8da4 --- /dev/null +++ b/.mcp.json @@ -0,0 +1,8 @@ +{ + "mcpServers": { + "telegram-scraper": { + "command": "uvx", + "args": ["--from", ".[mcp]", "tgscraper-mcp"] + } + } +} diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..cc45518 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,34 @@ +# AGENTS.md + +Guide for AI coding agents working on this repository. Users of the tool should read +[README.md](README.md) or the skill in `.claude/skills/telegram-scraper/SKILL.md`. + +## Setup & checks + +```bash +pip install -e ".[dev]" +pytest -q # offline; HTTP is mocked with respx +flake8 . --max-line-length=127 +``` + +The public site `t.me` is never called in tests. Parser tests use `tests/fixtures/page1.html`; +paging tests build pages with `tests/conftest.py::make_page`. + +## Architecture + +- `tgscraper/parser.py` — the only place that knows Telegram's HTML classes (`tgme_widget_message*`). + If Telegram changes its markup, fix it here and update the fixture. +- `tgscraper/client.py` — `Scraper` (sync) and `AsyncScraper` share paging logic in `_Walk.feed`. + Pages come oldest→newest; the public API yields newest→oldest. +- `tgscraper/__init__.py` — the simple functional API (`scrape`, `search`, ...). Keep it simple. +- `tgscraper/cli.py` — `tgscraper ` is shorthand for `tgscraper scrape `. + Every command should support `--json`. +- `tgscraper/mcp_server.py` — MCP tools; return JSON-friendly dicts, never raise to the client + (return `{"error": ...}`), strip `html`, cap limits. Works with `mcp` 1.x (FastMCP) and 2.x (MCPServer). + +## Conventions + +- Python 3.9+ (`from __future__ import annotations`, `typing.Optional`), core deps only `httpx` + `beautifulsoup4`; + everything else is an optional extra (`excel`, `mcp`, `dashboard`). +- New features need a test and a line in README (and the skill / MCP tool if agents should use them). +- `telegram_scraper.py` is a backward-compatibility shim; don't remove it. diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..91ca633 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,14 @@ +FROM python:3.12-slim + +WORKDIR /app +COPY pyproject.toml README.md LICENSE telegram_scraper.py ./ +COPY tgscraper ./tgscraper +RUN pip install --no-cache-dir ".[all]" + +# Data (exports, state file) goes here: docker run -v "$PWD/data:/data" ... +WORKDIR /data +ENV TGSCRAPER_STATE=/data/.tgscraper-state.json +EXPOSE 8501 + +ENTRYPOINT ["tgscraper"] +CMD ["--help"] diff --git a/README.md b/README.md index 73d9491..9227765 100644 --- a/README.md +++ b/README.md @@ -1,125 +1,461 @@ +
-# Telegram Scraper +# 📡 Telegram Scraper -Telegram Scraper is a Python tool for scraping messages from a Telegram channel using web scraping techniques. This tool extracts a specified number of messages from a given Telegram channel URL and supports the use of a proxy. +**Scrape, search, monitor and analyze any public Telegram channel — no API key, no login, no phone number.** -## Features +Python library · CLI · MCP server for AI agents · Claude Skill · Web dashboard · Docker -- Extract messages from a Telegram channel. -- Supports pagination to retrieve older messages. -- Optional proxy support for making requests. +[![Tests](https://github.com/specialteam/TelegramScraper/actions/workflows/python-package.yml/badge.svg)](https://github.com/specialteam/TelegramScraper/actions) +![Python](https://img.shields.io/badge/python-3.9%2B-blue) +![MCP](https://img.shields.io/badge/MCP-server-8A2BE2) +![License](https://img.shields.io/badge/license-MIT-green) -## Requirements +[Quick start](#-quick-start) · [CLI](#-command-line) · [Python](#-python-library) · [AI / MCP](#-use-it-from-ai-assistants-mcp--skill) · [Dashboard](#-web-dashboard) · [فارسی](#-راهنمای-فارسی) -- Python 3.x -- `requests` library -- `beautifulsoup4` library +
-## Installation +```bash +pip install "tgscraper[all] @ git+https://github.com/specialteam/TelegramScraper" +tgscraper durov +``` + +That's it — the latest posts of `t.me/durov`, in your terminal. + +--- + +## ✨ Features + +| | | +|---|---| +| 🔓 **Zero setup** | Uses the public web preview `t.me/s/`. No `api_id`, no session files, no account ban risk. | +| 🧾 **Rich data** | id, date, text, HTML, views, reactions, author, edited, forwards, replies, photos, videos, voice, documents, link previews, hashtags, mentions, links. | +| 🔎 **Search & filters** | Telegram's server-side search, date ranges, keywords, regex, hashtags, media type, minimum views. | +| 💾 **Export anywhere** | JSON, JSON Lines, CSV (Excel-friendly UTF-8), Excel `.xlsx`, SQLite (upsert archive), Markdown. | +| ⚡ **Fast & robust** | Async + concurrent multi-channel scraping, retries with backoff, `Retry-After` handling, rate limiting, rotating HTTP/SOCKS proxies. | +| 🔁 **Incremental** | Remembers the last post per channel — the next run fetches only new posts. | +| 👀 **Monitor** | Watch channels and push new posts to a webhook (Slack, Discord, n8n…) or a Telegram bot. | +| 🖼 **Media download** | Save photos, videos and voice notes of any post. | +| 📊 **Analytics** | Top posts, posting frequency by day/hour/weekday, hashtags, top words, reactions, EN/FA sentiment. | +| 🤖 **AI-native** | MCP server with 8 tools, a ready-made Claude Skill, and `--json` output for every command. | +| 🖥 **Dashboard** | Streamlit UI with charts and one-click export. | +| 🐳 **Docker** | Run the CLI, the dashboard or a 24/7 monitor in a container. | -First, clone the repository: +--- + +## 🚀 Quick start + +### Install ```bash -git clone https://github.com/yourusername/telegram-scraper.git -cd telegram-scraper +# everything (CLI + Excel + MCP server + dashboard) +pip install "tgscraper[all] @ git+https://github.com/specialteam/TelegramScraper" + +# or minimal (CLI + library only: httpx + beautifulsoup4) +pip install "git+https://github.com/specialteam/TelegramScraper" + +# or from a clone +git clone https://github.com/specialteam/TelegramScraper && cd TelegramScraper && pip install -e ".[all]" ``` -Install the required Python packages: +### Three ways to use it ```bash -pip install -r requirements.txt +tgscraper durov -n 100 -o durov.csv # 1. command line +``` + +```python +import tgscraper as tg # 2. Python +posts = tg.scrape("durov", limit=100) +``` + +```text +"What did @durov post this week?" # 3. ask your AI assistant (MCP / Skill) +``` + +--- + +## 💻 Command line + +Anything that looks like a channel works: `durov`, `@durov`, `t.me/durov`, `https://t.me/s/durov`. + +```bash +tgscraper durov # latest 20 posts, pretty output +tgscraper durov -n 500 -o durov.xlsx # save (.json .jsonl .csv .xlsx .db .md) +tgscraper durov -n 0 -o full_history.db # the whole channel history (-n 0 = no limit) +tgscraper durov telegram tginfo -n 50 -o all.db # several channels into one SQLite file + +tgscraper durov --since 2026-01-01 --until 2026-01-31 +tgscraper durov -n 300 -k ton -k bitcoin # keyword filter (any of them) +tgscraper durov --regex "v\d+\.\d+" # regex filter +tgscraper durov --hashtag news --min-views 50000 +tgscraper durov --media-only --media-type video + +tgscraper search durov "privacy" -n 30 # Telegram's own full-history search +tgscraper info durov # title, description, subscribers, counters +tgscraper stats durov -n 300 # analytics report (or: tgscraper stats durov.json) +tgscraper media durov -n 50 -d ./media # download photos & videos +tgscraper durov --incremental -o archive.db # only posts newer than the last run + +tgscraper watch durov telegram -i 120 # print new posts live +tgscraper watch durov --webhook https://hooks.slack.com/services/... # push to a webhook +tgscraper watch durov -k airdrop --bot-token 123:ABC --chat-id 42 # alert via your Telegram bot + +tgscraper durov --json | jq '.[] | {url, views}' # machine-readable output for scripts & agents +tgscraper durov -p socks5://127.0.0.1:1080 # proxy (repeat -p to rotate several) +``` + +
+Example: tgscraper stats (illustrative output) + +```text +📊 300 messages from durov + 2025-03-02T10:14:00+00:00 → 2026-09-20T16:40:00+00:00 (1.3 posts/day) +👁 total views 412,905,120 · average 1,376,350 +📎 with media 121 {'photo': 88, 'video': 33} · forwarded 4 +🙂 sentiment avg +0.21 (+97 / =180 / -23) + +🔥 Top posts: + 4,812,000 https://t.me/durov/301 'Telegram now has ...' +... +🕒 Posts by hour (UTC): + 14 ██████████████ 41 + 15 ██████████████████████████████ 87 ``` +
-## Usage +Run `tgscraper --help` or `tgscraper --help` for every option. -### Example Script +--- -Here's an example script (`example.py`) to demonstrate how to use the `TelegramScraper` class: +## 🐍 Python library ```python -from telegram_scraper import TelegramScraper - -def main(): - # URL of the Telegram channel - channel_url = 'https://t.me/s/mobydick_crypto' - - # Create an instance of TelegramScraper - scraper = TelegramScraper(channel_url, 100) - - # Set proxy if needed - # scraper.set_proxy('http://your_proxy_address:port') - - # Fetch messages - messages = scraper.fetch_messages() - - # Display messages - for i, message in enumerate(messages, 1): - print(f"{i}: {message}") - -if __name__ == "__main__": - main() -``` - -To run the example script: +import tgscraper as tg + +# Channel info +info = tg.channel_info("durov") +print(info.title, info.subscribers, info.description) + +# Latest posts — list of Message objects, newest first +posts = tg.scrape("durov", limit=100) +for p in posts: + print(p.date, p.views, p.url, p.text[:80], p.media_types, p.reactions) + +# Filters (all optional, combine freely) +posts = tg.scrape( + "durov", limit=None, # None = whole history + since="2026-01-01", until="2026-06-30", + keywords=["ton", "wallet"], regex=r"\bv\d+", hashtag="update", + media_only=True, media_types=["photo"], min_views=100_000, +) + +tg.search("durov", "privacy", limit=20) # Telegram server-side search +tg.get_message("durov", 123) # one post +tg.scrape("durov", incremental=True) # only new posts since last incremental call + +# Many channels concurrently +results = tg.scrape_many(["durov", "telegram", "tginfo"], limit=200) # {channel: [Message] | Exception} + +# Export / load +tg.export(posts, "posts.xlsx") # .json .jsonl .csv .xlsx .db .md +posts = tg.load("posts.json") # from .json .jsonl .db + +# Analytics +stats = tg.summarize(posts) # JSON-friendly dict +print(tg.format_summary(stats)) +tg.sentiment("Great news, bullish!") # -1 .. 1 (EN + FA lexicon) + +# Media +tg.download_media(posts, "media/", types=["photo", "video"]) + +# Monitor forever +tg.watch(["durov"], [tg.webhook_notifier("https://example.com/hook"), print], interval=60) +``` + +
+Advanced: reusable / async clients, proxies, streaming + +```python +from tgscraper import Scraper, AsyncScraper, MessageFilter + +with Scraper(proxies=["socks5://p1:1080", "http://p2:8080"], timeout=20, retries=3, delay=0.5) as s: + for msg in s.iter_messages("durov", limit=None, filter=MessageFilter(since="2026-01-01")): + print(msg.id) # streams page by page, low memory + +async with AsyncScraper(concurrency=10) as s: + posts = await s.get_messages("durov", 1000) + async for msg in s.iter_messages("telegram", 50, query="stories"): + ... +``` +
+ +### Message fields + +```json +{ + "id": 123, "channel": "durov", "url": "https://t.me/durov/123", + "date": "2026-01-10T09:30:00+00:00", "text": "…", "html": "…", + "views": 1250000, "author": null, "edited": false, + "forwarded_from": null, "reply_to": 120, + "media": [{"type": "photo", "url": "https://cdn4.telesco.pe/…jpg", "thumbnail": "…", "duration": null, "title": null}], + "reactions": {"👍": 15000, "🔥": 3200}, + "hashtags": ["news"], "mentions": ["telegram"], "links": ["https://telegram.org/blog"] +} +``` + +Media types: `photo`, `video`, `round_video`, `voice`, `audio`, `document`, `sticker`, `link_preview`. + +--- + +## 🤖 Use it from AI assistants (MCP + Skill) + +Telegram Scraper ships an **[MCP](https://modelcontextprotocol.io) server**, so Claude, Cursor, VS Code Copilot, +Windsurf, ChatGPT and any MCP client can read Telegram channels for you. + +```mermaid +flowchart LR + U["You: 'Summarize @durov this week'"] --> AI[AI assistant] + AI -- MCP tools --> S[tgscraper-mcp] + S -- HTTPS --> T["t.me/s/durov"] + S -- JSON --> AI --> A[Answer with links & stats] +``` + +### MCP tools + +| Tool | What it does | +|---|---| +| `get_channel_info` | Title, description, subscribers, photo, media counters | +| `get_messages` | Latest posts with filters (dates, keywords, hashtag, media, views, paging with `before_id`) | +| `search_messages` | Full-history search inside a channel | +| `get_message` | One post by id (for `t.me//` links) | +| `get_new_messages` | Only posts newer than an id — follow a channel over time | +| `analyze_channel` | Stats: top posts, activity, hours, hashtags, words, reactions, sentiment | +| `compare_channels` | Side-by-side: subscribers, posts/day, avg views, engagement rate | +| `export_messages` | Save posts to a local `.json/.csv/.xlsx/.db/.md` file | + +Plus prompts `summarize_channel` and `track_topic`, and the resource `telegram://channel/{channel}`. + +### Connect it + +The only requirement is [uv](https://docs.astral.sh/uv/) (`pip install uv`) — `uvx` downloads and runs the server +on demand. Or `pip install "tgscraper[mcp] @ git+…"` and use `"command": "tgscraper-mcp"` with no args. + +
+Claude Code + +```bash +claude mcp add telegram-scraper -- uvx --from "tgscraper[mcp] @ git+https://github.com/specialteam/TelegramScraper" tgscraper-mcp +``` +Inside this repository it is automatic: [`.mcp.json`](.mcp.json) registers the server and +[`.claude/skills/telegram-scraper`](.claude/skills/telegram-scraper/SKILL.md) loads the skill. +
+ +
+Claude Desktop · Cursor · Windsurf · any JSON-configured client + +Add to `claude_desktop_config.json` (Settings → Developer → Edit config), `~/.cursor/mcp.json`, +or `~/.codeium/windsurf/mcp_config.json`: + +```json +{ + "mcpServers": { + "telegram-scraper": { + "command": "uvx", + "args": ["--from", "tgscraper[mcp] @ git+https://github.com/specialteam/TelegramScraper", "tgscraper-mcp"], + "env": { "TGSCRAPER_PROXY": "" } + } + } +} +``` +
+ +
+VS Code (Copilot agent mode) + +`.vscode/mcp.json`: +```json +{ + "servers": { + "telegram-scraper": { + "type": "stdio", + "command": "uvx", + "args": ["--from", "tgscraper[mcp] @ git+https://github.com/specialteam/TelegramScraper", "tgscraper-mcp"] + } + } +} +``` +
+ +
+Remote / HTTP (ChatGPT connectors, n8n, other hosts) ```bash -python example.py +tgscraper mcp --transport streamable-http # serves MCP over HTTP ``` +
+ +Environment variables: `TGSCRAPER_PROXY` (proxy URL for all requests), `TGSCRAPER_MCP_MAX_LIMIT` (default 500). + +### Claude Skill -### Class Details +[`.claude/skills/telegram-scraper/SKILL.md`](.claude/skills/telegram-scraper/SKILL.md) teaches an agent when and how +to use the CLI (commands, JSON schema, how to cite results). Install it for all your projects: -#### TelegramScraper +```bash +mkdir -p ~/.claude/skills && cp -r .claude/skills/telegram-scraper ~/.claude/skills/ +``` -A class for scraping messages from a Telegram channel. +For claude.ai, zip the `telegram-scraper` folder and upload it under **Settings → Capabilities → Skills**. -##### `__init__(self, base_url, num_messages=100)` +**Try asking:** +- *"What are the 5 most viewed posts on @durov this year?"* +- *"Compare the engagement of these three crypto channels: …"* +- *"Search @xyz for 'airdrop' and give me the dates and links."* +- *"Export the last 1000 posts of t.me/abc to Excel."* -- `base_url` (str): The URL of the Telegram channel. -- `num_messages` (int): The number of messages to retrieve. +Other agents: [`AGENTS.md`](AGENTS.md) and [`llms.txt`](llms.txt) describe the project for LLMs. -##### `set_proxy(self, proxy)` +--- -- `proxy` (str): The proxy address (e.g., `'http://127.0.0.1:8080'`). +## 🖥 Web dashboard -Sets the proxy for making requests. +```bash +pip install "tgscraper[dashboard] @ git+https://github.com/specialteam/TelegramScraper" +tgscraper dashboard # → http://localhost:8501 +``` -##### `fetch_messages(self)` +Channel metrics, posts table with links, activity/views charts, top words & hashtags, CSV/JSON/Markdown download. -Fetches messages from the Telegram channel. +## 🐳 Docker -- Returns: `list` of messages. +```bash +docker build -t tgscraper . +docker run --rm -v "$PWD/data:/data" tgscraper durov -n 100 -o durov.csv +docker compose up dashboard # dashboard on :8501 +docker compose --profile watch up -d watch # 24/7 monitor archiving to data/archive.db +``` -## Tests +--- -Unit tests for the `TelegramScraper` class are provided in `test_telegram_scraper.py`. +## ❓ FAQ -### Running Tests +**Does it need a Telegram account or API key?** No. It reads the same public page you see at `https://t.me/s/durov`. -To run the tests, use the following command: +**Which channels work?** Public channels with web preview enabled. Private channels, groups, DMs and bots don't have a +public preview — use [Telethon](https://github.com/LonamiWebs/Telethon) for those. + +**`ChannelNotFound`?** The name is wrong, the channel is private, or its owner disabled the web preview. + +**Getting HTTP 429 / blocked?** The scraper already retries with backoff. Increase `delay`, lower concurrency, or +rotate proxies (`-p` multiple times). In regions where Telegram is filtered, use `-p socks5://…`. + +**Can I get comments / member lists?** No — they are not part of the public preview. + +**How accurate is sentiment?** It's a small English/Persian word list: good for trends, not for single posts. + +## 🧑‍💻 Development ```bash -python -m unittest test_telegram_scraper.py +pip install -e ".[dev]" +pytest -q # offline tests with HTML fixtures, no network needed ``` -## Project Structure +Project layout: ``` -. -├── telegram_scraper.py -├── test_telegram_scraper.py -├── example.py -├── requirements.txt -└── README.md +tgscraper/ + client.py Scraper / AsyncScraper: paging, retries, proxies + parser.py t.me/s HTML → Message / Channel + models.py Message, Media, Channel dataclasses + filters.py MessageFilter + exporters.py json, jsonl, csv, xlsx, sqlite, md + analytics.py summarize(), sentiment() + monitor.py watch() + webhook / Telegram bot notifiers + media.py download_media() + state.py incremental state + cli.py `tgscraper` command + mcp_server.py `tgscraper-mcp` MCP server + dashboard.py Streamlit app ``` -## License +The old `from telegram_scraper import TelegramScraper` API still works. + +## ⚖️ Responsible use + +Only public data is accessed. Respect Telegram's Terms of Service, local laws and people's privacy; keep request +rates reasonable. This project is not affiliated with Telegram. + +--- + +## 🇮🇷 راهنمای فارسی + +
+ +**Telegram Scraper** ابزاری برای خواندن، جست‌وجو، مانیتور و تحلیل **کانال‌های عمومی تلگرام** است؛ بدون API، +بدون لاگین و بدون شماره تلفن. + +### نصب + +
+ +```bash +pip install "tgscraper[all] @ git+https://github.com/specialteam/TelegramScraper" +``` + +
+ +### مهم‌ترین دستورها + +
+ +```bash +tgscraper durov # ۲۰ پست آخر +tgscraper durov -n 500 -o durov.xlsx # ذخیره در اکسل (یا csv / json / db / md) +tgscraper durov --since 2026-01-01 -k بیت‌کوین +tgscraper search durov "privacy" # جست‌وجو در کل تاریخچه +tgscraper info durov # اطلاعات و تعداد اعضای کانال +tgscraper stats durov -n 300 # آمار: پربازدیدها، ساعت‌های فعالیت، هشتگ‌ها، احساسات +tgscraper media durov -d ./media # دانلود عکس و ویدیو +tgscraper watch durov --bot-token TOKEN --chat-id ID # اعلان پست جدید با ربات تلگرام +tgscraper durov -p socks5://127.0.0.1:1080 # استفاده از پراکسی +tgscraper dashboard # داشبورد وب +``` + +
+ +### پایتون + +
+ +```python +import tgscraper as tg +posts = tg.scrape("durov", limit=100, since="2026-01-01") +tg.export(posts, "posts.csv") +print(tg.format_summary(tg.summarize(posts))) +``` + +
+ +### اتصال به هوش مصنوعی + +- **MCP:** با تنظیمات بخش [AI / MCP](#-use-it-from-ai-assistants-mcp--skill) به Claude، Cursor، VS Code و … وصل + کنید. بعد کافی است بپرسید: «پربازدیدترین پست‌های این هفته‌ی @durov چی بوده؟» +- **Skill:** پوشه‌ی `.claude/skills/telegram-scraper` را در `~/.claude/skills/` کپی کنید. +- خروجی همه‌ی دستورها با `--json` برای ایجنت‌ها و اسکریپت‌ها قابل خواندن است. + +فقط کانال‌هایی که پیش‌نمایش وب (`t.me/s/...`) دارند پشتیبانی می‌شوند. در ایران برای دسترسی از پراکسی استفاده کنید. + +
-This project is licensed under the MIT License - see the [LICENSE](LICENSE) file for details. +--- -## Contributing +
-Contributions are welcome! Please feel free to submit a Pull Request. +MIT License · If this project helps you, give it a ⭐ -## Contact +Keywords: telegram scraper, telegram channel scraper, scrape telegram without api, t.me scraper, telegram +crawler, telegram osint, telegram to csv, telegram to excel, telegram monitor, telegram mcp server, mcp telegram, +claude telegram, ai agent telegram tool, python telegram scraper, اسکرپر تلگرام, استخراج پیام کانال تلگرام -If you have any questions or suggestions, please feel free to contact me at https://t.me/amirspecial +
diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..b97cb6d --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,16 @@ +services: + # Web dashboard on http://localhost:8501 + dashboard: + build: . + entrypoint: ["python", "-m", "streamlit", "run", "/usr/local/lib/python3.12/site-packages/tgscraper/dashboard.py", + "--server.address=0.0.0.0", "--server.port=8501"] + ports: ["8501:8501"] + volumes: ["./data:/data"] + + # Example monitor: forwards new posts to a webhook and archives them in SQLite + watch: + build: . + command: ["watch", "durov", "--interval", "300", "--save", "/data/archive.db"] + volumes: ["./data:/data"] + restart: unless-stopped + profiles: ["watch"] diff --git a/example.py b/example.py index cbc00a7..0ca6770 100644 --- a/example.py +++ b/example.py @@ -1,21 +1,31 @@ -from telegram_scraper import TelegramScraper +"""Tour of the tgscraper API. Run: python example.py""" +import tgscraper as tg + def main(): - # URL کانال - channel_url = 'https://t.me/s/mobydick_crypto' - - # ایجاد نمونه‌ای از TelegramScraper - scraper = TelegramScraper(channel_url, 100) - - # تنظیم پراکسی در صورت نیاز - # scraper.set_proxy('http://your_proxy_address:port') - - # استخراج پیام‌ها - messages = scraper.fetch_messages() - - # نمایش پیام‌ها - for i, message in enumerate(messages, 1): - print(f"{i}: {message}") + # 1) Channel info + info = tg.channel_info("durov") + print(f"{info.title}: {info.subscribers:,} subscribers") + + # 2) Latest posts (newest first) — each one is a rich Message object + posts = tg.scrape("durov", limit=30) + for p in posts[:5]: + print(p.date, p.views, p.url, p.text[:80]) + + # 3) Filters: date range, keywords, media, views... + recent_media = tg.scrape("durov", limit=10, since="2026-01-01", media_only=True) + print(len(recent_media), "recent posts with media") + + # 4) Telegram's own search + hits = tg.search("durov", "privacy", limit=5) + print([h.url for h in hits]) + + # 5) Save anywhere: .json .jsonl .csv .xlsx .db .md + tg.export(posts, "durov.csv") + + # 6) Statistics + print(tg.format_summary(tg.summarize(posts))) + if __name__ == "__main__": main() diff --git a/llms.txt b/llms.txt new file mode 100644 index 0000000..645362b --- /dev/null +++ b/llms.txt @@ -0,0 +1,23 @@ +# Telegram Scraper (tgscraper) + +> Python library, CLI and MCP server that scrapes PUBLIC Telegram channels through the web preview +> (https://t.me/s/) — no API key, login or phone number. Returns structured posts (text, date, views, +> reactions, media, hashtags, links), supports search, date/keyword filters, JSON/CSV/Excel/SQLite export, +> incremental scraping, monitoring with webhooks, media download and analytics. + +Install: `pip install "tgscraper[all] @ git+https://github.com/specialteam/TelegramScraper"` +MCP server: `uvx --from "tgscraper[mcp] @ git+https://github.com/specialteam/TelegramScraper" tgscraper-mcp` + +## Docs +- [README](README.md): features, CLI, Python API, MCP setup for Claude/Cursor/VS Code, Docker, FAQ +- [Agent skill](.claude/skills/telegram-scraper/SKILL.md): task → command table and JSON output schema +- [AGENTS.md](AGENTS.md): how to work on this codebase + +## Key API +- `tgscraper -n 20 --json`, `tgscraper search `, `tgscraper info|stats|media|watch|mcp|dashboard` +- `tgscraper.scrape(channel, limit, since=, until=, keywords=, query=, media_only=, min_views=, incremental=)` +- `tgscraper.scrape_many`, `search`, `channel_info`, `get_message`, `export`, `load`, `summarize`, `download_media`, `watch` +- MCP tools: get_channel_info, get_messages, search_messages, get_message, get_new_messages, analyze_channel, compare_channels, export_messages + +## Limits +- Only public channels with web preview enabled; no private chats, groups, comments or member lists. diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..8bae4c0 --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,49 @@ +[build-system] +requires = ["setuptools>=64"] +build-backend = "setuptools.build_meta" + +[project] +name = "tgscraper" +version = "2.0.0" +description = "Scrape public Telegram channels without API keys or login — Python library, CLI, MCP server for AI agents, and dashboard." +readme = "README.md" +license = { file = "LICENSE" } +requires-python = ">=3.9" +authors = [{ name = "specialteam" }] +keywords = [ + "telegram", "scraper", "telegram-scraper", "telegram-channel", "web-scraping", "osint", + "mcp", "mcp-server", "model-context-protocol", "ai-agent", "claude", "llm-tools", "crawler", "t.me", +] +classifiers = [ + "Programming Language :: Python :: 3", + "License :: OSI Approved :: MIT License", + "Operating System :: OS Independent", + "Topic :: Internet :: WWW/HTTP :: Indexing/Search", + "Topic :: Communications :: Chat", + "Environment :: Console", + "Framework :: AsyncIO", +] +dependencies = ["httpx[socks]>=0.27", "beautifulsoup4>=4.11"] + +[project.optional-dependencies] +excel = ["openpyxl>=3.1"] +mcp = ["mcp>=1.2; python_version>='3.10'"] +dashboard = ["streamlit>=1.35"] +all = ["tgscraper[excel,mcp,dashboard]"] +dev = ["tgscraper[excel,mcp]", "pytest>=7", "pytest-asyncio>=0.23", "respx>=0.21", "flake8"] + +[project.scripts] +tgscraper = "tgscraper.cli:main" +tgscraper-mcp = "tgscraper.mcp_server:main" + +[project.urls] +Homepage = "https://github.com/specialteam/TelegramScraper" +Issues = "https://github.com/specialteam/TelegramScraper/issues" + +[tool.setuptools] +packages = ["tgscraper"] +py-modules = ["telegram_scraper"] + +[tool.pytest.ini_options] +testpaths = ["tests"] +asyncio_mode = "auto" diff --git a/requirements.txt b/requirements.txt index 1190bd8..1853035 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,2 +1,2 @@ -requests -beautifulsoup4 +httpx[socks]>=0.27 +beautifulsoup4>=4.11 diff --git a/telegram_scraper.py b/telegram_scraper.py index 4a7d464..425516a 100644 --- a/telegram_scraper.py +++ b/telegram_scraper.py @@ -1,6 +1,12 @@ -import requests -from bs4 import BeautifulSoup -from urllib.parse import urljoin +"""Backward-compatible wrapper around the old ``TelegramScraper`` class. + +New code should use the ``tgscraper`` package instead:: + + import tgscraper as tg + posts = tg.scrape("mobydick_crypto", limit=100) +""" +from tgscraper import Scraper + class TelegramScraper: def __init__(self, base_url, num_messages=100): @@ -11,49 +17,10 @@ def __init__(self, base_url, num_messages=100): def set_proxy(self, proxy): """تنظیم پراکسی برای درخواست‌ها""" - self.proxy = { - 'http': proxy, - 'https': proxy - } + self.proxy = proxy def fetch_messages(self): - current_url = self.base_url - while len(self.messages) < self.num_messages: - response = self._make_request(current_url) - - if response.status_code != 200: - print(f"Failed to retrieve the page: {response.status_code}") - return self.messages - - soup = BeautifulSoup(response.content, 'html.parser') - self._parse_messages(soup) - - next_url = self._find_next_page(soup) - if next_url: - current_url = urljoin('https://t.me', next_url) - else: - break - + """Return the texts of the latest ``num_messages`` posts (newest first).""" + with Scraper(proxies=self.proxy) as scraper: + self.messages = [m.text for m in scraper.iter_messages(self.base_url, self.num_messages)] return self.messages - - def _make_request(self, url): - """ارسال درخواست HTTP با یا بدون پراکسی""" - if self.proxy: - return requests.get(url, proxies=self.proxy) - else: - return requests.get(url) - - def _parse_messages(self, soup): - message_elements = soup.find_all('div', class_='tgme_widget_message_text') - for element in message_elements: - message_text = element.get_text(strip=True) - if message_text not in self.messages: - self.messages.append(message_text) - if len(self.messages) >= self.num_messages: - break - - def _find_next_page(self, soup): - next_button = soup.find('a', class_='tme_messages_more') - if next_button and next_button.has_attr('href'): - return next_button['href'] - return None diff --git a/test_telegram_scraper.py b/test_telegram_scraper.py deleted file mode 100644 index 6821293..0000000 --- a/test_telegram_scraper.py +++ /dev/null @@ -1,67 +0,0 @@ -import unittest -from unittest.mock import patch, Mock -from telegram_scraper import TelegramScraper - -class TestTelegramScraper(unittest.TestCase): - - @patch('telegram_scraper.requests.get') - def test_fetch_messages(self, mock_get): - # شبیه‌سازی پاسخ HTML - mock_html = """ - - -
Message 1
-
Message 2
- - - - """ - mock_response = Mock() - mock_response.status_code = 200 - mock_response.content = mock_html - mock_get.return_value = mock_response - - # ایجاد نمونه‌ای از TelegramScraper - channel_url = 'https://t.me/s/mobydick_crypto' - scraper = TelegramScraper(channel_url, 2) - - # استخراج پیام‌ها - messages = scraper.fetch_messages() - - # بررسی نتایج - self.assertEqual(len(messages), 2) - self.assertIn('Message 1', messages) - self.assertIn('Message 2', messages) - - @patch('telegram_scraper.requests.get') - def test_set_proxy(self, mock_get): - # ایجاد نمونه‌ای از TelegramScraper - channel_url = 'https://t.me/s/mobydick_crypto' - scraper = TelegramScraper(channel_url, 2) - - # تنظیم پراکسی - scraper.set_proxy('http://127.0.0.1:8080') - - # شبیه‌سازی پاسخ HTML - mock_html = """ - - -
Message 1
- - - """ - mock_response = Mock() - mock_response.status_code = 200 - mock_response.content = mock_html - mock_get.return_value = mock_response - - # استخراج پیام‌ها - messages = scraper.fetch_messages() - - # بررسی اینکه پراکسی تنظیم شده است - mock_get.assert_called_with(channel_url, proxies={'http': 'http://127.0.0.1:8080', 'https': 'http://127.0.0.1:8080'}) - self.assertEqual(len(messages), 1) - self.assertIn('Message 1', messages) - -if __name__ == "__main__": - unittest.main() diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..a597ad0 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,30 @@ +from pathlib import Path + +import pytest + +FIXTURES = Path(__file__).parent / "fixtures" + + +@pytest.fixture +def page1() -> str: + return (FIXTURES / "page1.html").read_text(encoding="utf-8") + + +def make_page(channel, ids, before=None, day=1): + """A minimal t.me/s page with text messages ``ids`` (oldest first).""" + more = (f'' + if before else "") + items = "".join( + f'
' + f'
post {i}
' + f'{i}' + f'' + f"
" + for i in ids + ) + return f"{more}{items}" + + +@pytest.fixture +def page_factory(): + return make_page diff --git a/tests/fixtures/page1.html b/tests/fixtures/page1.html new file mode 100644 index 0000000..b9e6349 --- /dev/null +++ b/tests/fixtures/page1.html @@ -0,0 +1,80 @@ + +Test Channel – Telegram + +
+
+ +
Test Channel
+ +
+
News about crypto
and more
+
+
1.2M subscribers
+
3.4K photos
+
120 videos
+
+
+
+
+ +
+
+
+
Forwarded from Other Channel
+ +
Bitcoin is going up! Great profit 🚀
#BTC #crypto via @someuser link
+
+ 👍1.5K + 🔥320 + 2K + 15 + 7 +
+ +
+
+
+ +
+
+
+ +
Test Channel
+
Bitcoin is going up!
+
+
Warning: market crash, big loss today
+ + +
+ +
+ +
+
+
+ +
+
+
+
سلام دوستان، خبر خوب داریم! رشد عالی
+ +
+
+
+
+ diff --git a/tests/test_client.py b/tests/test_client.py new file mode 100644 index 0000000..6238d5a --- /dev/null +++ b/tests/test_client.py @@ -0,0 +1,116 @@ +import httpx +import pytest +import respx + +import tgscraper as tg +from telegram_scraper import TelegramScraper +from tgscraper import AsyncScraper, ChannelNotFound, Scraper + +URL = "https://t.me/s/chan" + + +def _route_pages(page_factory): + """chan has ids 1..30, 10 per page.""" + def handler(request): + before = int(request.url.params.get("before", 31)) + ids = list(range(max(1, before - 10), before)) + nxt = ids[0] if ids and ids[0] > 1 else None + return httpx.Response(200, text=page_factory("chan", ids, nxt)) + return respx.get(URL).mock(side_effect=handler) + + +@respx.mock +def test_pagination_and_limit(page_factory): + route = _route_pages(page_factory) + with Scraper(delay=0) as s: + msgs = s.get_messages("chan", limit=15) + assert [m.id for m in msgs] == list(range(30, 15, -1)) + assert route.call_count == 2 + + +@respx.mock +def test_whole_history(page_factory): + _route_pages(page_factory) + with Scraper(delay=0) as s: + msgs = s.get_messages("@chan", limit=None) + assert [m.id for m in msgs] == list(range(30, 0, -1)) + + +@respx.mock +def test_min_id_and_before(page_factory): + _route_pages(page_factory) + with Scraper(delay=0) as s: + assert [m.id for m in s.get_messages("chan", limit=None, min_id=25)] == [30, 29, 28, 27, 26] + assert [m.id for m in s.get_messages("chan", limit=3, before=12)] == [11, 10, 9] + + +@respx.mock +def test_filters_and_query(page_factory): + route = _route_pages(page_factory) + with Scraper(delay=0) as s: + msgs = s.get_messages("chan", limit=3, min_views=20, keywords=["post"], query="post") + assert [m.id for m in msgs] == [30, 29, 28] + assert route.calls[0].request.url.params["q"] == "post" + + +@respx.mock +def test_since_stops_paging(page_factory): + def handler(request): + before = int(request.url.params.get("before", 31)) + ids = list(range(before - 10, before)) + # newer pages are on later days + return httpx.Response(200, text=page_factory("chan", ids, ids[0], day=before // 10)) + route = respx.get(URL).mock(side_effect=handler) + with Scraper(delay=0) as s: + msgs = s.get_messages("chan", limit=None, since="2026-01-02") + assert [m.id for m in msgs] == list(range(30, 20, -1)) + list(range(20, 10, -1)) + assert route.call_count == 3 + + +@respx.mock +def test_retry_then_success(page_factory): + respx.get(URL).mock(side_effect=[httpx.Response(503), httpx.Response(200, text=page_factory("chan", [1]))]) + with Scraper(delay=0, retries=2) as s: + s._wait_time = lambda *a: 0 + assert [m.id for m in s.get_messages("chan")] == [1] + + +@respx.mock +def test_channel_not_found(): + respx.get("https://t.me/s/nope").mock(return_value=httpx.Response(302, headers={"Location": "https://t.me/nope"})) + with pytest.raises(ChannelNotFound): + tg.scrape("nope") + + +@respx.mock +def test_get_message_and_info(page1): + respx.get("https://t.me/s/testchan").mock(return_value=httpx.Response(200, text=page1)) + with Scraper(delay=0) as s: + assert s.get_message("testchan", 99).reply_to == 98 + assert s.get_message("testchan", 5) is None + assert s.channel_info("https://t.me/testchan").subscribers == 1_200_000 + + +@respx.mock +def test_incremental(tmp_path, page_factory): + _route_pages(page_factory) + state = tmp_path / "state.json" + assert len(tg.scrape("chan", 5, incremental=True, state_file=state)) == 5 + assert tg.scrape("chan", 5, incremental=True, state_file=state) == [] + + +@respx.mock +async def test_async_scrape_many(page_factory): + _route_pages(page_factory) + respx.get("https://t.me/s/missing").mock(return_value=httpx.Response(404)) + async with AsyncScraper(delay=0) as s: + results = await s.scrape_many(["chan", "missing"], limit=12) + assert [m.id for m in results["chan"]] == list(range(30, 18, -1)) + assert isinstance(results["missing"], ChannelNotFound) + + +@respx.mock +def test_legacy_wrapper(page_factory): + _route_pages(page_factory) + scraper = TelegramScraper("https://t.me/s/chan", 2) + assert scraper.fetch_messages() == ["post 30", "post 29"] diff --git a/tests/test_export_analytics_cli.py b/tests/test_export_analytics_cli.py new file mode 100644 index 0000000..d4f78ea --- /dev/null +++ b/tests/test_export_analytics_cli.py @@ -0,0 +1,93 @@ +import csv +import json + +import httpx +import pytest +import respx + +import tgscraper as tg +from tgscraper.cli import main +from tgscraper.parser import parse_page + + +@pytest.fixture +def messages(page1): + return list(reversed(parse_page(page1, "testchan")[0])) + + +@pytest.mark.parametrize("ext", ["json", "jsonl", "db"]) +def test_export_roundtrip(tmp_path, messages, ext): + path = tg.export(messages, tmp_path / f"out.{ext}") + loaded = tg.load(path) + assert sorted(loaded, key=lambda m: m.id) == sorted(messages, key=lambda m: m.id) + + +def test_sqlite_upsert(tmp_path, messages): + path = tmp_path / "a.db" + tg.export(messages, path) + tg.export(messages[:1], path) + assert len(tg.load(path)) == 3 + + +def test_csv_xlsx_md(tmp_path, messages): + tg.export(messages, tmp_path / "o.csv") + with open(tmp_path / "o.csv", encoding="utf-8-sig") as fh: + rows = list(csv.DictReader(fh)) + assert rows[2]["hashtags"] == "#BTC #crypto" and rows[2]["media_types"] == "photo" + pytest.importorskip("openpyxl") + assert tg.export(messages, tmp_path / "o.xlsx").stat().st_size > 0 + assert "testchan/98" in tg.export(messages, tmp_path / "o.md").read_text(encoding="utf-8") + + +def test_bad_format(tmp_path, messages): + with pytest.raises(ValueError): + tg.export(messages, tmp_path / "out.txt") + + +def test_filter(messages): + f = tg.MessageFilter(since="2026-01-11", until="2026-01-11") + assert [m.id for m in f.apply(messages)] == [99] + assert [m.id for m in tg.MessageFilter(media_types=["video"]).apply(messages)] == [99] + assert [m.id for m in tg.MessageFilter(hashtag="#btc").apply(messages)] == [98] + assert [m.id for m in tg.MessageFilter(regex=r"crash|سلام").apply(messages)] == [100, 99] + + +def test_analytics(messages): + assert tg.sentiment("great profit, bullish!") > 0 + assert tg.sentiment("crash and loss") < 0 + assert tg.sentiment("خبر خوب و رشد عالی") > 0 + stats = tg.summarize(messages) + assert stats["count"] == 3 and stats["total_views"] == 12500 + 980 + 2300 + assert stats["top_posts"][0]["id"] == 98 + assert stats["media_types"] == {"video": 1, "photo": 1} + assert stats["reactions"]["👍"] == 1500 + assert "Top posts" in tg.format_summary(stats) + json.dumps(stats) + assert tg.summarize([]) == {"count": 0} + + +@respx.mock +def test_cli(tmp_path, page1, capsys): + respx.get("https://t.me/s/testchan").mock(return_value=httpx.Response(200, text=page1)) + assert main(["testchan", "-n", "2", "--json"]) == 0 + assert [m["id"] for m in json.loads(capsys.readouterr().out)] == [100, 99] + + out = tmp_path / "x.csv" + assert main(["scrape", "testchan", "-o", str(out), "--media-only"]) == 0 + assert out.exists() + + assert main(["info", "testchan", "--json"]) == 0 + assert json.loads(capsys.readouterr().out)[0]["title"] == "Test Channel" + + assert main(["stats", "testchan"]) == 0 + assert "Top posts" in capsys.readouterr().out + + assert main(["testchan", "-n", "1"]) == 0 + assert "t.me/testchan/100" in capsys.readouterr().out + + +@respx.mock +def test_cli_error(capsys): + respx.get("https://t.me/s/nope").mock(return_value=httpx.Response(404)) + assert main(["nope"]) == 1 + assert "not found" in capsys.readouterr().err diff --git a/tests/test_live.py b/tests/test_live.py new file mode 100644 index 0000000..bf3641f --- /dev/null +++ b/tests/test_live.py @@ -0,0 +1,61 @@ +"""Tests against the real t.me. Skipped unless TGSCRAPER_LIVE=1 (run by the "Live test" workflow).""" +import os + +import pytest + +import tgscraper as tg + +pytestmark = pytest.mark.skipif(os.environ.get("TGSCRAPER_LIVE") != "1", reason="set TGSCRAPER_LIVE=1") + +CHANNEL = os.environ.get("TGSCRAPER_LIVE_CHANNEL", "durov") + + +def test_channel_info(): + info = tg.channel_info(CHANNEL) + print(info) + assert info.title + assert info.subscribers and info.subscribers > 1000 + + +def test_messages_and_paging(): + posts = tg.scrape(CHANNEL, limit=40) # needs at least 2 pages + for p in posts[:3]: + print(p.id, p.date, p.views, p.media_types, p.reactions, repr(p.text[:80])) + assert len(posts) == 40 + ids = [p.id for p in posts] + assert ids == sorted(ids, reverse=True) and len(set(ids)) == 40 + assert all(p.date for p in posts) + assert sum(p.views is not None for p in posts) >= 35 + assert sum(bool(p.text) for p in posts) >= 20 + assert any(p.media for p in posts) + + +def test_get_message(): + latest = tg.scrape(CHANNEL, limit=1)[0] + assert tg.get_message(CHANNEL, latest.id).id == latest.id + + +def test_search(): + posts = tg.search(CHANNEL, "telegram", limit=5) + print([p.url for p in posts]) + assert posts + + +def test_not_found(): + with pytest.raises(tg.ChannelNotFound): + tg.channel_info("this_channel_should_not_exist_1234567") + + +def test_scrape_many(): + results = tg.scrape_many([CHANNEL, "telegram"], limit=5) + for ch, res in results.items(): + assert not isinstance(res, Exception), f"{ch}: {res}" + assert len(res) == 5 + + +def test_reactions_are_split_by_emoji(): + posts = [p for p in tg.scrape(CHANNEL, limit=20) if p.reactions] + print([p.reactions for p in posts[:3]]) + assert posts, "expected some posts with reactions" + assert any(len(p.reactions) > 1 for p in posts) + assert any(not k.startswith("custom") for p in posts for k in p.reactions) diff --git a/tests/test_mcp.py b/tests/test_mcp.py new file mode 100644 index 0000000..2ca74b1 --- /dev/null +++ b/tests/test_mcp.py @@ -0,0 +1,44 @@ +import httpx +import pytest +import respx + +pytest.importorskip("mcp") + +from tgscraper import mcp_server # noqa: E402 + + +@pytest.fixture +def mocked(page1): + with respx.mock: + respx.get("https://t.me/s/testchan").mock(return_value=httpx.Response(200, text=page1)) + respx.get("https://t.me/s/nope").mock(return_value=httpx.Response(404)) + yield + + +async def test_tools_registered(): + names = {t.name for t in await mcp_server.mcp.list_tools()} + assert {"get_channel_info", "get_messages", "search_messages", "get_message", "get_new_messages", + "analyze_channel", "compare_channels", "export_messages"} <= names + + +async def test_tools(mocked, tmp_path): + info = await mcp_server.get_channel_info("testchan") + assert info["subscribers"] == 1_200_000 + + res = await mcp_server.get_messages("@testchan", limit=2) + assert [m["id"] for m in res["messages"]] == [100, 99] + assert "html" not in res["messages"][0] + + assert (await mcp_server.get_message("testchan", 98))["views"] == 12500 + assert (await mcp_server.get_new_messages("testchan", 99))["latest_id"] == 100 + assert (await mcp_server.analyze_channel("testchan"))["stats"]["count"] == 3 + cmp = await mcp_server.compare_channels(["testchan", "nope"]) + assert cmp["testchan"]["subscribers"] == 1_200_000 and "error" in cmp["nope"] + saved = await mcp_server.export_messages("testchan", str(tmp_path / "t.json")) + assert saved["count"] == 3 + assert "error" in await mcp_server.get_channel_info("nope") + + +async def test_call_through_protocol(mocked): + result = await mcp_server.mcp.call_tool("get_messages", {"channel": "testchan", "limit": 1}) + assert "100" in str(result) diff --git a/tests/test_parser.py b/tests/test_parser.py new file mode 100644 index 0000000..94f3559 --- /dev/null +++ b/tests/test_parser.py @@ -0,0 +1,69 @@ +from datetime import datetime, timezone + +import pytest + +from tgscraper.parser import normalize_channel, parse_count, parse_page + + +@pytest.mark.parametrize("raw", ["durov", "@durov", "t.me/durov", "https://t.me/s/durov", "https://t.me/durov/12", + "http://telegram.me/durov", " durov/ "]) +def test_normalize_channel(raw): + assert normalize_channel(raw) == "durov" + + +def test_normalize_channel_invalid(): + with pytest.raises(ValueError): + normalize_channel("not a channel!") + + +@pytest.mark.parametrize("raw,expected", [("1.2K", 1200), ("3,4M", 3_400_000), ("2,300", 2300), ("980", 980), + ("12 345", 12345), ("", None), (None, None), ("abc", None)]) +def test_parse_count(raw, expected): + assert parse_count(raw) == expected + + +def test_parse_page(page1): + messages, info, before = parse_page(page1, "testchan") + assert [m.id for m in messages] == [98, 99, 100] + assert before == 98 + + assert info.title == "Test Channel" + assert info.verified + assert info.subscribers == 1_200_000 + assert info.counters == {"subscribers": 1_200_000, "photos": 3400, "videos": 120} + assert info.description == "News about crypto\nand more" + assert info.photo.endswith("photo.jpg") + + m98, m99, m100 = messages + assert m98.url == "https://t.me/testchan/98" + assert m98.text.startswith("Bitcoin is going up! Great profit 🚀\n#BTC") + assert m98.date == datetime(2026, 1, 10, 9, 30, tzinfo=timezone.utc) + assert m98.views == 12500 + assert m98.forwarded_from == "Other Channel" + assert m98.hashtags == ["BTC", "crypto"] + assert m98.mentions == ["someuser"] + assert m98.links == ["https://example.com/a"] + assert m98.reactions == {"👍": 1500, "🔥": 320, "❤": 2000, "custom:5368324170671202286": 15, "⭐": 7} + assert [(x.type, x.url) for x in m98.media] == [("photo", "https://cdn4.telesco.pe/file/p98.jpg")] + + # the quoted reply text must not leak into the message text + assert m99.text == "Warning: market crash, big loss today" + assert m99.reply_to == 98 + assert m99.edited and m99.author == "Admin" + assert m99.media[0].type == "video" + assert m99.media[0].url.endswith("v99.mp4") + assert m99.media[0].duration == "0:42" + + assert m100.views == 2300 and not m100.media and not m100.edited + + +def test_message_roundtrip(page1): + from tgscraper import Message + messages, _, _ = parse_page(page1, "testchan") + for m in messages: + assert Message.from_dict(m.to_dict()) == m + + +def test_empty_page(): + messages, info, before = parse_page("", "x") + assert messages == [] and before is None and info.username == "x" diff --git a/tgscraper/__init__.py b/tgscraper/__init__.py new file mode 100644 index 0000000..f70839b --- /dev/null +++ b/tgscraper/__init__.py @@ -0,0 +1,93 @@ +"""tgscraper — scrape public Telegram channels without an API key, login or phone number. + +Quick start:: + + import tgscraper as tg + + posts = tg.scrape("durov", limit=50) # list[Message], newest first + tg.export(posts, "durov.csv") # .json .jsonl .csv .xlsx .db .md + print(tg.channel_info("durov").subscribers) + print(tg.format_summary(tg.summarize(posts))) +""" +from __future__ import annotations + +import asyncio +from typing import Dict, List, Optional, Sequence, Union + +from .analytics import format_summary, sentiment, summarize, top_words +from .client import AsyncScraper, ChannelNotFound, Scraper, ScraperError +from .exporters import export, flatten, load, to_json, to_markdown +from .filters import MessageFilter +from .media import download_media +from .models import Channel, Media, Message +from .monitor import telegram_notifier, watch, webhook_notifier +from .parser import normalize_channel +from .state import State + +__version__ = "2.0.0" + +__all__ = [ + "scrape", "scrape_many", "search", "channel_info", "get_message", + "Scraper", "AsyncScraper", "MessageFilter", "State", + "Message", "Media", "Channel", "ScraperError", "ChannelNotFound", + "export", "load", "flatten", "to_json", "to_markdown", "download_media", + "summarize", "format_summary", "sentiment", "top_words", + "watch", "webhook_notifier", "telegram_notifier", "normalize_channel", +] + +Proxy = Union[None, str, Sequence[str]] + + +def scrape( + channel: str, + limit: Optional[int] = 100, + *, + proxy: Proxy = None, + incremental: bool = False, + state_file: Optional[str] = None, + **options, +) -> List[Message]: + """Scrape the latest ``limit`` messages (newest first) of a public channel. + + Filters: ``since``, ``until`` (``"2026-01-01"``), ``keywords=["btc"]``, ``regex``, ``media_only``, + ``media_types=["photo"]``, ``min_views``, ``hashtag``. Server-side search: ``query="text"``. + ``incremental=True`` only returns posts newer than the previous incremental run. + """ + state = State(state_file) if incremental else None + name = normalize_channel(channel) + if state and state.last_id(name) and "min_id" not in options: + options["min_id"] = state.last_id(name) + with Scraper(proxies=proxy) as s: + messages = s.get_messages(name, limit, **options) + if state: + for m in messages: + state.update(name, m.id) + state.save() + return messages + + +def scrape_many( + channels: Sequence[str], limit: Optional[int] = 100, *, proxy: Proxy = None, concurrency: int = 5, **options +) -> Dict[str, Union[List[Message], Exception]]: + """Scrape several channels concurrently. Returns ``{channel: messages or exception}``.""" + async def run(): + async with AsyncScraper(proxies=proxy, concurrency=concurrency) as s: + return await s.scrape_many(channels, limit, **options) + return asyncio.run(run()) + + +def search(channel: str, query: str, limit: Optional[int] = 50, *, proxy: Proxy = None, **options) -> List[Message]: + """Search a channel's history using Telegram's own search.""" + return scrape(channel, limit, proxy=proxy, query=query, **options) + + +def channel_info(channel: str, *, proxy: Proxy = None) -> Channel: + """Title, description, subscriber count and photo of a channel.""" + with Scraper(proxies=proxy) as s: + return s.channel_info(channel) + + +def get_message(channel: str, message_id: int, *, proxy: Proxy = None) -> Optional[Message]: + """One message by id, or ``None``.""" + with Scraper(proxies=proxy) as s: + return s.get_message(channel, message_id) diff --git a/tgscraper/__main__.py b/tgscraper/__main__.py new file mode 100644 index 0000000..dd8a8c9 --- /dev/null +++ b/tgscraper/__main__.py @@ -0,0 +1,5 @@ +import sys + +from .cli import main + +sys.exit(main()) diff --git a/tgscraper/analytics.py b/tgscraper/analytics.py new file mode 100644 index 0000000..b948a09 --- /dev/null +++ b/tgscraper/analytics.py @@ -0,0 +1,130 @@ +"""Quick statistics and a lightweight (lexicon-based) sentiment score. No heavy dependencies.""" +from __future__ import annotations + +import re +from collections import Counter +from typing import Any, Dict, Iterable, List + +from .models import Message + +_WORD_RE = re.compile(r"[^\W\d_]{3,}", re.UNICODE) + +STOPWORDS = set(""" +the and for that this with you your are was were have has had not but from they their them will would there what +when which who how can all any out our about into more some just than then also its it's been being over only other +https http www com org net telegram channel join like very much here now one two new get got via amp +از به با که این را در برای تا آن یک هم و یا اما اگر هر می را ها های بر شد شده است هست بود کرد کند کنید +خواهد باید نیز چه چرا کرده روی بین پس دیگر همه ما شما آنها او من تو ای بعد قبل خود داریم دارد +""".split()) + +POSITIVE = set(""" +good great excellent amazing awesome love best win winning profit gain gains bullish pump moon up rise rising +growth success successful happy strong buy launch launched new record breakout surge rally positive nice congrats +خوب عالی عالیه بهترین سود صعود صعودی رشد موفق موفقیت خوشحال قوی خرید پامپ رکورد افزایش مثبت تبریک فوق‌العاده +""".split()) + +NEGATIVE = set(""" +bad worst terrible awful loss losses lose losing bearish dump crash down fall falling drop scam hack hacked risk +fear sell weak fail failed failure problem warning danger negative sad liquidated rekt fraud +بد بدترین ضرر نزول نزولی ریزش سقوط کلاهبرداری هک ریسک ترس فروش ضعیف شکست مشکل هشدار خطر منفی ناراحت دامپ +""".split()) + + +def sentiment(text: str) -> float: + """Score in [-1, 1] from a small bilingual (English / Persian) word list. Rough, but dependency-free.""" + words = [w.lower() for w in _WORD_RE.findall(text)] + pos = sum(w in POSITIVE for w in words) + neg = sum(w in NEGATIVE for w in words) + return 0.0 if pos + neg == 0 else round((pos - neg) / (pos + neg), 3) + + +def top_words(messages: Iterable[Message], n: int = 20) -> List[tuple]: + counter: Counter = Counter() + for m in messages: + counter.update(w.lower() for w in _WORD_RE.findall(re.sub(r"https?://\S+", "", m.text)) + if w.lower() not in STOPWORDS) + return counter.most_common(n) + + +def summarize(messages: Iterable[Message], top: int = 5) -> Dict[str, Any]: + """Everything you usually want to know about a batch of messages, as a JSON-friendly dict.""" + msgs = list(messages) + if not msgs: + return {"count": 0} + dated = [m for m in msgs if m.date] + viewed = [m for m in msgs if m.views is not None] + per_day: Counter = Counter(m.date.date().isoformat() for m in dated) + per_hour: Counter = Counter(m.date.hour for m in dated) + per_weekday: Counter = Counter(m.date.strftime("%A") for m in dated) + media: Counter = Counter(t for m in msgs for t in m.media_types) + hashtags: Counter = Counter(h.lower() for m in msgs for h in m.hashtags) + mentions: Counter = Counter(x.lower() for m in msgs for x in m.mentions) + reactions: Counter = Counter() + for m in msgs: + reactions.update(m.reactions) + scores = [sentiment(m.text) for m in msgs if m.text] + total_views = sum(m.views for m in viewed) + days = max(len(per_day), 1) + + def brief(m: Message) -> Dict[str, Any]: + return {"id": m.id, "url": m.url, "views": m.views, "date": m.date.isoformat() if m.date else None, + "text": (m.text[:140] + "…") if len(m.text) > 140 else m.text} + + return { + "count": len(msgs), + "channels": sorted({m.channel for m in msgs}), + "first_date": min(m.date for m in dated).isoformat() if dated else None, + "last_date": max(m.date for m in dated).isoformat() if dated else None, + "posts_per_day": round(len(dated) / days, 2), + "total_views": total_views, + "avg_views": round(total_views / len(viewed)) if viewed else None, + "with_media": sum(1 for m in msgs if m.media), + "forwarded": sum(1 for m in msgs if m.forwarded_from), + "media_types": dict(media.most_common()), + "top_posts": [brief(m) for m in sorted(viewed, key=lambda m: m.views or 0, reverse=True)[:top]], + "top_hashtags": hashtags.most_common(top * 2), + "top_mentions": mentions.most_common(top * 2), + "top_words": top_words(msgs, top * 4), + "reactions": dict(reactions.most_common(10)), + "sentiment": { + "average": round(sum(scores) / len(scores), 3) if scores else 0.0, + "positive": sum(s > 0 for s in scores), + "neutral": sum(s == 0 for s in scores), + "negative": sum(s < 0 for s in scores), + }, + "posts_by_day": dict(sorted(per_day.items())), + "posts_by_hour": {h: per_hour.get(h, 0) for h in range(24)}, + "posts_by_weekday": dict(per_weekday.most_common()), + } + + +def format_summary(stats: Dict[str, Any]) -> str: + """Human-readable text version of :func:`summarize`.""" + if not stats.get("count"): + return "No messages." + + def bar(value: int, maximum: int, width: int = 30) -> str: + return "█" * max(1 if value else 0, round(width * value / maximum)) if maximum else "" + + lines = [ + f"📊 {stats['count']} messages from {', '.join(stats['channels'])}", + f" {stats['first_date'] or '?'} → {stats['last_date'] or '?'} ({stats['posts_per_day']} posts/day)", + f"👁 total views {stats['total_views']:,} · average {stats['avg_views'] or 0:,}", + f"📎 with media {stats['with_media']} {stats['media_types']} · forwarded {stats['forwarded']}", + f"🙂 sentiment avg {stats['sentiment']['average']:+} " + f"(+{stats['sentiment']['positive']} / ={stats['sentiment']['neutral']} / -{stats['sentiment']['negative']})", + "", + "🔥 Top posts:", + ] + for p in stats["top_posts"]: + lines.append(f" {p['views'] or 0:>9,} {p['url']} {p['text'][:60]!r}") + if stats["top_hashtags"]: + lines += ["", "#️⃣ " + " ".join(f"#{h}({c})" for h, c in stats["top_hashtags"])] + if stats["top_words"]: + lines += ["🔤 " + " ".join(f"{w}({c})" for w, c in stats["top_words"][:15])] + hours = stats["posts_by_hour"] + peak = max(hours.values()) if hours else 0 + if peak: + lines += ["", "🕒 Posts by hour (UTC):"] + lines += [f" {h:02d} {bar(c, peak)} {c}" for h, c in hours.items() if c] + return "\n".join(lines) diff --git a/tgscraper/cli.py b/tgscraper/cli.py new file mode 100644 index 0000000..c463e8d --- /dev/null +++ b/tgscraper/cli.py @@ -0,0 +1,265 @@ +"""Command line interface: ``tgscraper durov -n 50 -o durov.csv``.""" +from __future__ import annotations + +import argparse +import json +import logging +import sys +from pathlib import Path +from typing import List, Optional, Sequence + +from . import __version__ +from .analytics import format_summary, summarize +from .client import ScraperError +from .exporters import FORMATS, export, load, to_json, to_markdown +from .media import download_media +from .models import Message + +COMMANDS = {"scrape", "info", "search", "stats", "watch", "media", "mcp", "dashboard"} + +EPILOG = """examples: + tgscraper durov # latest 20 posts, pretty in the terminal + tgscraper durov -n 500 -o durov.csv # save as CSV (.json .jsonl .xlsx .db .md too) + tgscraper durov telegram -n 100 -o all.db + tgscraper durov --since 2026-01-01 --keyword ton --media-only + tgscraper search durov "privacy" -n 30 + tgscraper info durov + tgscraper stats durov -n 300 # or: tgscraper stats durov.json + tgscraper media durov -n 50 -d ./media + tgscraper watch durov --webhook https://example.com/hook + tgscraper mcp # run the MCP server for AI assistants + tgscraper durov --json | jq '.[0].text' # machine-readable output +""" + + +def _add_scrape_options(p: argparse.ArgumentParser, many: bool = True) -> None: + if many: + p.add_argument("channels", nargs="+", help="channel usernames or links (durov, @durov, t.me/durov)") + p.add_argument("-n", "--limit", type=int, default=20, help="max messages per channel, 0 = all (default 20)") + p.add_argument("-o", "--output", help="save to file; format from extension: " + ", ".join(FORMATS)) + p.add_argument("-f", "--format", choices=FORMATS, help="force output format") + p.add_argument("--json", action="store_true", help="print JSON to stdout (for scripts and AI agents)") + p.add_argument("--since", help="only posts on/after this date (YYYY-MM-DD)") + p.add_argument("--until", help="only posts on/before this date (YYYY-MM-DD)") + p.add_argument("-k", "--keyword", action="append", default=[], help="keep posts containing a keyword (repeatable)") + p.add_argument("--regex", help="keep posts matching a regular expression") + p.add_argument("--hashtag", help="keep posts with this hashtag") + p.add_argument("--media-only", action="store_true", help="only posts with media") + p.add_argument("--media-type", action="append", default=[], + help="photo, video, voice, document, round_video, sticker, link_preview (repeatable)") + p.add_argument("--min-views", type=int, help="only posts with at least N views") + p.add_argument("-q", "--query", help="Telegram server-side search") + p.add_argument("--incremental", action="store_true", help="only posts newer than the last --incremental run") + p.add_argument("--download-media", metavar="DIR", help="also download photos/videos into DIR") + _add_common(p) + + +def _add_common(p: argparse.ArgumentParser) -> None: + p.add_argument("-p", "--proxy", action="append", help="http/socks5 proxy URL (repeat to rotate)") + p.add_argument("-v", "--verbose", action="store_true") + + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser( + prog="tgscraper", + description="Scrape public Telegram channels — no API key, no login.", + epilog=EPILOG, + formatter_class=argparse.RawDescriptionHelpFormatter, + ) + parser.add_argument("--version", action="version", version=f"tgscraper {__version__}") + sub = parser.add_subparsers(dest="command") + + _add_scrape_options(sub.add_parser("scrape", help="scrape messages (default command)")) + + p = sub.add_parser("search", help="search inside a channel") + p.add_argument("channels", nargs=1, metavar="channel") + p.add_argument("search_query", metavar="query") + _add_scrape_options(p, many=False) + + p = sub.add_parser("info", help="channel title, description, subscribers") + p.add_argument("channels", nargs="+") + p.add_argument("--json", action="store_true") + _add_common(p) + + p = sub.add_parser("stats", help="statistics for a channel or a saved .json/.jsonl/.db file") + p.add_argument("source", nargs="+", help="channel(s) or file") + p.add_argument("-n", "--limit", type=int, default=200) + p.add_argument("--json", action="store_true") + _add_common(p) + + p = sub.add_parser("media", help="download photos/videos of recent posts") + p.add_argument("channels", nargs="+") + p.add_argument("-n", "--limit", type=int, default=20) + p.add_argument("-d", "--dir", default="media") + p.add_argument("-t", "--type", action="append", default=[], help="photo, video, voice... (repeatable)") + _add_common(p) + + p = sub.add_parser("watch", help="notify on new posts (Ctrl+C to stop)") + p.add_argument("channels", nargs="+") + p.add_argument("-i", "--interval", type=float, default=60, help="seconds between checks (default 60)") + p.add_argument("-k", "--keyword", action="append", default=[], help="only notify on these keywords") + p.add_argument("--webhook", help="POST each new post as JSON to this URL") + p.add_argument("--bot-token", help="Telegram bot token for notifications") + p.add_argument("--chat-id", help="chat id that receives bot notifications") + p.add_argument("--save", help="append new posts to this .jsonl/.db file") + p.add_argument("--backfill", type=int, default=0, help="emit N existing posts on start") + _add_common(p) + + p = sub.add_parser("mcp", help="run the MCP server (stdio) so AI assistants can use tgscraper") + p.add_argument("--transport", default="stdio", choices=["stdio", "sse", "streamable-http"]) + + sub.add_parser("dashboard", help="open the web dashboard (needs tgscraper[dashboard])") + return parser + + +def _filters(args) -> dict: + return { + "since": args.since, "until": args.until, "keywords": args.keyword, "regex": args.regex, + "media_only": args.media_only, "media_types": args.media_type, "min_views": args.min_views, + "hashtag": args.hashtag, + } + + +def _print_messages(messages: Sequence[Message]) -> None: + for m in messages: + meta = [m.date.strftime("%Y-%m-%d %H:%M") if m.date else "", f"👁 {m.views:,}" if m.views is not None else ""] + if m.media: + meta.append("📎 " + ",".join(m.media_types)) + if m.forwarded_from: + meta.append(f"↪ {m.forwarded_from}") + print(f"\033[1m{m.url}\033[0m " + " ".join(x for x in meta if x)) + print(m.text or "(no text)") + print() + + +def _scrape(args) -> List[Message]: + from . import scrape + limit = None if args.limit == 0 else args.limit + query = getattr(args, "search_query", None) or args.query + out: List[Message] = [] + for ch in args.channels: + out += scrape(ch, limit, proxy=args.proxy, incremental=args.incremental, query=query, **_filters(args)) + return out + + +def _emit(messages: List[Message], args) -> None: + if args.output: + path = export(messages, args.output, args.format) + print(f"✅ Saved {len(messages)} messages to {path}", file=sys.stderr) + elif args.json or args.format == "json": + print(to_json(messages)) + elif args.format == "md": + print(to_markdown(messages)) + elif args.format: + print("Use -o FILE with --format " + args.format, file=sys.stderr) + else: + _print_messages(messages) + print(f"— {len(messages)} messages", file=sys.stderr) + if getattr(args, "download_media", None): + files = download_media(messages, args.download_media, proxy=(args.proxy or [None])[0]) + print(f"🖼 Downloaded {len(files)} files to {args.download_media}", file=sys.stderr) + + +def main(argv: Optional[Sequence[str]] = None) -> int: + argv = list(sys.argv[1:] if argv is None else argv) + # `tgscraper durov` == `tgscraper scrape durov` + if argv and not argv[0].startswith("-") and argv[0] not in COMMANDS: + argv.insert(0, "scrape") + parser = build_parser() + args = parser.parse_args(argv) + if not args.command: + parser.print_help() + return 0 + logging.basicConfig(level=logging.INFO if getattr(args, "verbose", False) else logging.WARNING, + format="%(levelname)s %(message)s") + try: + return _run(args) + except ScraperError as exc: + print(f"❌ {exc}", file=sys.stderr) + return 1 + except KeyboardInterrupt: + return 130 + + +def _run(args) -> int: + from . import channel_info, scrape + + if args.command in ("scrape", "search"): + _emit(_scrape(args), args) + elif args.command == "info": + infos = [channel_info(ch, proxy=args.proxy) for ch in args.channels] + if args.json: + print(json.dumps([i.to_dict() for i in infos], ensure_ascii=False, indent=2)) + for i in infos if not args.json else []: + print(f"\033[1m{i.title}\033[0m ({i.url}){' ✔' if i.verified else ''}") + print(f"👥 {i.subscribers:,} subscribers" if i.subscribers is not None else "👥 ?") + if i.counters: + print(" " + " · ".join(f"{k}: {v:,}" for k, v in i.counters.items())) + if i.description: + print(i.description) + print() + elif args.command == "stats": + messages: List[Message] = [] + for src in args.source: + if Path(src).is_file(): + messages += load(src) + else: + messages += scrape(src, args.limit or None, proxy=args.proxy) + stats = summarize(messages) + print(json.dumps(stats, ensure_ascii=False, indent=2, default=str) if args.json else format_summary(stats)) + elif args.command == "media": + total = 0 + for ch in args.channels: + msgs = scrape(ch, args.limit, proxy=args.proxy, media_only=True) + total += len(download_media(msgs, args.dir, types=args.type or None, proxy=(args.proxy or [None])[0])) + print(f"🖼 {total} files in {args.dir}") + elif args.command == "watch": + _watch(args) + elif args.command == "mcp": + from .mcp_server import main as mcp_main + mcp_main(args.transport) + elif args.command == "dashboard": + _dashboard() + return 0 + + +def _watch(args) -> None: + from .client import Scraper + from .monitor import telegram_notifier, watch, webhook_notifier + + def printer(m: Message) -> None: + _print_messages([m]) + sys.stdout.flush() + + handlers = [printer] + if args.webhook: + handlers.append(webhook_notifier(args.webhook)) + if args.bot_token and args.chat_id: + handlers.append(telegram_notifier(args.bot_token, args.chat_id)) + if args.save: + def saver(m: Message) -> None: + if args.save.endswith((".db", ".sqlite", ".sqlite3")): + export([m], args.save) # sqlite export is an upsert + else: + with open(args.save, "a", encoding="utf-8") as fh: + fh.write(json.dumps(m.to_dict(), ensure_ascii=False) + "\n") + handlers.append(saver) + print(f"👀 Watching {', '.join(args.channels)} every {args.interval:g}s — Ctrl+C to stop", file=sys.stderr) + with Scraper(proxies=args.proxy) as s: + watch(args.channels, handlers, interval=args.interval, keywords=args.keyword or None, + scraper=s, backfill=args.backfill) + + +def _dashboard() -> None: + import subprocess + try: + import streamlit # noqa: F401 + except ImportError: + print("Dashboard needs streamlit: pip install 'tgscraper[dashboard]'", file=sys.stderr) + raise SystemExit(1) + app = Path(__file__).with_name("dashboard.py") + raise SystemExit(subprocess.call([sys.executable, "-m", "streamlit", "run", str(app)])) + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tgscraper/client.py b/tgscraper/client.py new file mode 100644 index 0000000..3737d16 --- /dev/null +++ b/tgscraper/client.py @@ -0,0 +1,318 @@ +"""Sync and async HTTP clients that page through ``https://t.me/s/``.""" +from __future__ import annotations + +import asyncio +import itertools +import logging +import random +import time +from dataclasses import dataclass, field +from typing import AsyncIterator, Dict, Iterator, List, Optional, Sequence, Set, Tuple, Union + +import httpx + +from .filters import MessageFilter +from .models import Channel, Message +from .parser import normalize_channel, parse_page + +log = logging.getLogger("tgscraper") + +BASE_URL = "https://t.me/s/" +DEFAULT_HEADERS = { + "User-Agent": ( + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 " + "(KHTML, like Gecko) Chrome/128.0 Safari/537.36" + ), + "Accept-Language": "en-US,en;q=0.9", +} +RETRY_STATUSES = {429, 500, 502, 503, 504} + + +class ScraperError(Exception): + """Base error of this package.""" + + +class ChannelNotFound(ScraperError): + """The channel does not exist, is private, or has no public web preview.""" + + +Page = Tuple[List[Message], Channel, Optional[int]] + + +@dataclass +class _Walk: + """Pagination state shared by the sync and async clients.""" + + limit: Optional[int] + flt: MessageFilter + min_id: Optional[int] = None + yielded: int = 0 + seen: Set[int] = field(default_factory=set) + done: bool = False + + def feed(self, messages: List[Message], before: Optional[int]) -> List[Message]: + """Take one page (oldest->newest); return matches newest-first and update ``done``.""" + out = [] + new = [m for m in reversed(messages) if m.id not in self.seen] + if not new: + self.done = True + return out + for msg in new: + self.seen.add(msg.id) + if self.min_id is not None and msg.id <= self.min_id: + self.done = True + break + if self.flt.since and msg.date and msg.date < self.flt.since: + self.done = True + break + if self.flt.matches(msg): + out.append(msg) + self.yielded += 1 + if self.limit is not None and self.yielded >= self.limit: + self.done = True + break + if before is None: + self.done = True + return out + + +def _build_filter(flt: Optional[MessageFilter], kwargs: Dict) -> MessageFilter: + if flt is not None and kwargs: + raise TypeError("Pass either filter=MessageFilter(...) or filter keyword arguments, not both") + return flt or MessageFilter(**kwargs) + + +def _page_params(before: Optional[int], query: Optional[str]) -> Dict[str, Union[str, int]]: + params: Dict[str, Union[str, int]] = {} + if before: + params["before"] = before + if query: + params["q"] = query + return params + + +class _Base: + def __init__( + self, + proxies: Union[None, str, Sequence[str]] = None, + timeout: float = 20.0, + retries: int = 3, + delay: float = 0.5, + headers: Optional[Dict[str, str]] = None, + ) -> None: + if isinstance(proxies, str): + proxies = [proxies] + self.proxies: List[Optional[str]] = list(proxies) if proxies else [None] + self.timeout = timeout + self.retries = retries + self.delay = delay + self.headers = {**DEFAULT_HEADERS, **(headers or {})} + + def _wait_time(self, attempt: int, response: Optional[httpx.Response]) -> float: + if response is not None and response.headers.get("Retry-After", "").isdigit(): + return float(response.headers["Retry-After"]) + return min(30.0, (2 ** attempt) * max(self.delay, 0.5)) + random.uniform(0, 0.3) + + @staticmethod + def _check(response: httpx.Response, channel: str) -> None: + if response.status_code == 404 or response.is_redirect: + raise ChannelNotFound( + f"Channel '{channel}' not found or has no public preview (https://t.me/s/{channel})" + ) + + +class Scraper(_Base): + """Synchronous scraper. + + >>> with Scraper() as s: + ... for msg in s.iter_messages("durov", limit=10): + ... print(msg.date, msg.text) + """ + + def __init__(self, *args, **kwargs) -> None: + super().__init__(*args, **kwargs) + self._clients = [ + httpx.Client(proxy=p, timeout=self.timeout, headers=self.headers, follow_redirects=False) + for p in self.proxies + ] + self._cycle = itertools.cycle(self._clients) + self._last_request = 0.0 + + # --- plumbing ------------------------------------------------------- + def _get(self, channel: str, params: Dict) -> str: + url = BASE_URL + channel + response = None + for attempt in range(self.retries + 1): + wait = self.delay - (time.monotonic() - self._last_request) + if wait > 0: + time.sleep(wait) + self._last_request = time.monotonic() + try: + response = next(self._cycle).get(url, params=params) + except httpx.TransportError as exc: + if attempt >= self.retries: + raise ScraperError(f"Network error for {url}: {exc}") from exc + log.warning("Network error (%s), retrying...", exc) + time.sleep(self._wait_time(attempt, None)) + continue + self._check(response, channel) + if response.status_code in RETRY_STATUSES and attempt < self.retries: + log.warning("HTTP %s for %s, retrying...", response.status_code, url) + time.sleep(self._wait_time(attempt, response)) + continue + if response.status_code != 200: + raise ScraperError(f"HTTP {response.status_code} for {url}") + return response.text + raise ScraperError(f"Giving up on {url}") + + def get_page(self, channel: str, before: Optional[int] = None, query: Optional[str] = None) -> Page: + channel = normalize_channel(channel) + return parse_page(self._get(channel, _page_params(before, query)), channel) + + # --- public API ----------------------------------------------------- + def channel_info(self, channel: str) -> Channel: + """Title, description, subscriber count, photo... of a channel.""" + return self.get_page(channel)[1] + + def iter_messages( + self, + channel: str, + limit: Optional[int] = 100, + *, + query: Optional[str] = None, + min_id: Optional[int] = None, + before: Optional[int] = None, + max_pages: Optional[int] = None, + filter: Optional[MessageFilter] = None, + **filter_kwargs, + ) -> Iterator[Message]: + """Yield messages newest-first. + + ``query`` is a server-side search. Filters (``since``, ``until``, ``keywords``, ``regex``, + ``media_only``, ``media_types``, ``min_views``, ``hashtag``) are applied client-side. + ``min_id`` stops at messages older than or equal to that id (used for incremental scraping). + ``before`` starts from messages older than that id. ``limit=None`` scrapes the whole history. + """ + channel = normalize_channel(channel) + walk = _Walk(limit=limit, flt=_build_filter(filter, filter_kwargs), min_id=min_id) + for page_no in itertools.count(1): + messages, _, before = self.get_page(channel, before, query) + yield from walk.feed(messages, before) + if walk.done or (max_pages and page_no >= max_pages): + return + + def get_messages(self, channel: str, limit: Optional[int] = 100, **kwargs) -> List[Message]: + """Same as :meth:`iter_messages` but returns a list.""" + return list(self.iter_messages(channel, limit, **kwargs)) + + def get_message(self, channel: str, message_id: int) -> Optional[Message]: + """Fetch a single message by id (``None`` if deleted).""" + messages, _, _ = self.get_page(channel, before=message_id + 1) + return next((m for m in messages if m.id == message_id), None) + + def close(self) -> None: + for client in self._clients: + client.close() + + def __enter__(self) -> "Scraper": + return self + + def __exit__(self, *exc) -> None: + self.close() + + +class AsyncScraper(_Base): + """Asynchronous scraper, ideal for many channels at once. + + >>> async with AsyncScraper() as s: + ... results = await s.scrape_many(["durov", "telegram"], limit=50) + """ + + def __init__(self, *args, concurrency: int = 5, **kwargs) -> None: + super().__init__(*args, **kwargs) + self._clients = [ + httpx.AsyncClient(proxy=p, timeout=self.timeout, headers=self.headers, follow_redirects=False) + for p in self.proxies + ] + self._cycle = itertools.cycle(self._clients) + self._semaphore = asyncio.Semaphore(concurrency) + + async def _get(self, channel: str, params: Dict) -> str: + url = BASE_URL + channel + for attempt in range(self.retries + 1): + async with self._semaphore: + if self.delay: + await asyncio.sleep(self.delay) + try: + response = await next(self._cycle).get(url, params=params) + except httpx.TransportError as exc: + if attempt >= self.retries: + raise ScraperError(f"Network error for {url}: {exc}") from exc + await asyncio.sleep(self._wait_time(attempt, None)) + continue + self._check(response, channel) + if response.status_code in RETRY_STATUSES and attempt < self.retries: + await asyncio.sleep(self._wait_time(attempt, response)) + continue + if response.status_code != 200: + raise ScraperError(f"HTTP {response.status_code} for {url}") + return response.text + raise ScraperError(f"Giving up on {url}") + + async def get_page(self, channel: str, before: Optional[int] = None, query: Optional[str] = None) -> Page: + channel = normalize_channel(channel) + return parse_page(await self._get(channel, _page_params(before, query)), channel) + + async def channel_info(self, channel: str) -> Channel: + return (await self.get_page(channel))[1] + + async def iter_messages( + self, + channel: str, + limit: Optional[int] = 100, + *, + query: Optional[str] = None, + min_id: Optional[int] = None, + before: Optional[int] = None, + max_pages: Optional[int] = None, + filter: Optional[MessageFilter] = None, + **filter_kwargs, + ) -> AsyncIterator[Message]: + channel = normalize_channel(channel) + walk = _Walk(limit=limit, flt=_build_filter(filter, filter_kwargs), min_id=min_id) + for page_no in itertools.count(1): + messages, _, before = await self.get_page(channel, before, query) + for msg in walk.feed(messages, before): + yield msg + if walk.done or (max_pages and page_no >= max_pages): + return + + async def get_messages(self, channel: str, limit: Optional[int] = 100, **kwargs) -> List[Message]: + return [m async for m in self.iter_messages(channel, limit, **kwargs)] + + async def get_message(self, channel: str, message_id: int) -> Optional[Message]: + messages, _, _ = await self.get_page(channel, before=message_id + 1) + return next((m for m in messages if m.id == message_id), None) + + async def scrape_many( + self, channels: Sequence[str], limit: Optional[int] = 100, **kwargs + ) -> Dict[str, Union[List[Message], Exception]]: + """Scrape several channels concurrently. Failed channels map to their exception.""" + async def one(ch: str): + try: + return await self.get_messages(ch, limit, **kwargs) + except Exception as exc: # noqa: BLE001 - reported per channel + return exc + + results = await asyncio.gather(*(one(c) for c in channels)) + return {normalize_channel(c): r for c, r in zip(channels, results)} + + async def aclose(self) -> None: + for client in self._clients: + await client.aclose() + + async def __aenter__(self) -> "AsyncScraper": + return self + + async def __aexit__(self, *exc) -> None: + await self.aclose() diff --git a/tgscraper/dashboard.py b/tgscraper/dashboard.py new file mode 100644 index 0000000..0381556 --- /dev/null +++ b/tgscraper/dashboard.py @@ -0,0 +1,112 @@ +"""Web dashboard. Run with ``tgscraper dashboard`` (needs ``pip install 'tgscraper[dashboard]'``).""" +from __future__ import annotations + +import datetime as dt +import io +import json + +import streamlit as st + +import tgscraper as tg +from tgscraper.exporters import FLAT_COLUMNS, flatten + +st.set_page_config(page_title="Telegram Scraper", page_icon="📡", layout="wide") +st.title("📡 Telegram Scraper") +st.caption("Read any public Telegram channel — no API key, no login.") + +with st.sidebar: + channels_text = st.text_input("Channel(s)", "durov", help="Comma separated: durov, telegram, t.me/xyz") + limit = st.slider("Messages per channel", 10, 1000, 100, step=10) + query = st.text_input("Telegram search (optional)") + keywords = st.text_input("Keywords filter (comma separated)") + use_dates = st.checkbox("Date range") + since = until = None + if use_dates: + since = st.date_input("Since", dt.date.today() - dt.timedelta(days=30)) + until = st.date_input("Until", dt.date.today()) + media_only = st.checkbox("Only posts with media") + proxy = st.text_input("Proxy (optional)", placeholder="socks5://127.0.0.1:1080") + go = st.button("🚀 Scrape", type="primary", use_container_width=True) + + +@st.cache_data(ttl=300, show_spinner=False) +def fetch(channel, limit, query, keywords, since, until, media_only, proxy): + info = tg.channel_info(channel, proxy=proxy or None) + msgs = tg.scrape(channel, limit, proxy=proxy or None, query=query or None, + keywords=keywords, since=since, until=until, media_only=media_only) + return info.to_dict(), [m.to_dict() for m in msgs] + + +if go: + st.session_state["channels"] = [c.strip() for c in channels_text.split(",") if c.strip()] + +for channel in st.session_state.get("channels", []): + try: + with st.spinner(f"Scraping {channel}…"): + info, raw = fetch(channel, limit, query, [k.strip() for k in keywords.split(",") if k.strip()], + since, until, media_only, proxy) + except Exception as exc: # noqa: BLE001 + st.error(f"{channel}: {exc}") + continue + messages = [tg.Message.from_dict(d) for d in raw] + stats = tg.summarize(messages) + + left, right = st.columns([1, 5]) + if info.get("photo"): + left.image(info["photo"], width=96) + right.subheader(f"{info.get('title') or channel} · [t.me/{info['username']}]({info['url']})") + if info.get("description"): + right.caption(info["description"]) + + c = st.columns(5) + c[0].metric("Subscribers", f"{info.get('subscribers') or 0:,}") + c[1].metric("Messages", stats.get("count", 0)) + c[2].metric("Avg views", f"{stats.get('avg_views') or 0:,}") + c[3].metric("Posts / day", stats.get("posts_per_day", 0)) + c[4].metric("Sentiment", f"{stats.get('sentiment', {}).get('average', 0):+.2f}") + if not messages: + st.info("No messages matched.") + continue + + rows = [flatten(m) for m in messages] + tab_posts, tab_charts, tab_words, tab_export = st.tabs(["📰 Posts", "📈 Charts", "🔤 Words & tags", "💾 Export"]) + with tab_posts: + st.dataframe(rows, use_container_width=True, hide_index=True, column_config={ + "url": st.column_config.LinkColumn("url"), "text": st.column_config.TextColumn("text", width="large")}) + with tab_charts: + a, b = st.columns(2) + a.markdown("**Posts per day**") + a.bar_chart(stats["posts_by_day"]) + b.markdown("**Posts by hour (UTC)**") + b.bar_chart({str(k): v for k, v in stats["posts_by_hour"].items()}) + views = {m.date.isoformat(): m.views for m in messages if m.date and m.views is not None} + if views: + st.markdown("**Views per post**") + st.line_chart(dict(sorted(views.items()))) + with tab_words: + a, b = st.columns(2) + a.markdown("**Top words**") + a.bar_chart(dict(stats["top_words"])) + b.markdown("**Top hashtags**") + if stats["top_hashtags"]: + b.bar_chart(dict(stats["top_hashtags"])) + else: + b.caption("No hashtags") + st.markdown("**🔥 Most viewed**") + for p in stats["top_posts"]: + st.markdown(f"- **{p['views'] or 0:,}** views — [{p['url']}]({p['url']}): {p['text']}") + with tab_export: + buf = io.StringIO() + import csv + writer = csv.DictWriter(buf, fieldnames=FLAT_COLUMNS) + writer.writeheader() + writer.writerows(rows) + d1, d2, d3 = st.columns(3) + d1.download_button("⬇️ CSV", buf.getvalue().encode("utf-8-sig"), f"{channel}.csv", "text/csv") + d2.download_button("⬇️ JSON", json.dumps(raw, ensure_ascii=False, indent=2), f"{channel}.json", + "application/json") + d3.download_button("⬇️ Markdown", tg.to_markdown(messages), f"{channel}.md", "text/markdown") + st.divider() + +if not st.session_state.get("channels"): + st.info("👈 Enter a channel and press **Scrape**.") diff --git a/tgscraper/exporters.py b/tgscraper/exporters.py new file mode 100644 index 0000000..45200e5 --- /dev/null +++ b/tgscraper/exporters.py @@ -0,0 +1,163 @@ +"""Save / load messages as JSON, JSON Lines, CSV, Excel, SQLite or Markdown.""" +from __future__ import annotations + +import csv +import json +import sqlite3 +from pathlib import Path +from typing import Any, Dict, Iterable, List, Optional, Union + +from .models import Message + +FORMATS = ("json", "jsonl", "csv", "xlsx", "sqlite", "md") +_EXT = { + ".json": "json", ".jsonl": "jsonl", ".ndjson": "jsonl", ".csv": "csv", ".xlsx": "xlsx", + ".db": "sqlite", ".sqlite": "sqlite", ".sqlite3": "sqlite", ".md": "md", +} +FLAT_COLUMNS = [ + "channel", "id", "url", "date", "text", "views", "author", "edited", "forwarded_from", "reply_to", + "media_types", "media_urls", "reactions", "reactions_total", "hashtags", "mentions", "links", +] +_SQL_COLUMNS = { + "channel": "TEXT", "id": "INTEGER", "url": "TEXT", "date": "TEXT", "text": "TEXT", "html": "TEXT", + "views": "INTEGER", "author": "TEXT", "edited": "INTEGER", "forwarded_from": "TEXT", "reply_to": "INTEGER", + "media": "TEXT", "reactions": "TEXT", "hashtags": "TEXT", "mentions": "TEXT", "links": "TEXT", +} + + +def detect_format(path: Union[str, Path], fmt: Optional[str] = None) -> str: + if fmt: + fmt = fmt.lower() + if fmt not in FORMATS: + raise ValueError(f"Unknown format {fmt!r}. Choose one of: {', '.join(FORMATS)}") + return fmt + ext = Path(path).suffix.lower() + if ext not in _EXT: + raise ValueError(f"Cannot guess format from {str(path)!r}; use one of {', '.join(_EXT)} or pass format=") + return _EXT[ext] + + +def flatten(msg: Message) -> Dict[str, Any]: + """One flat row per message, friendly for CSV / Excel / pandas.""" + return { + "channel": msg.channel, + "id": msg.id, + "url": msg.url, + "date": msg.date.isoformat() if msg.date else "", + "text": msg.text, + "views": msg.views, + "author": msg.author or "", + "edited": msg.edited, + "forwarded_from": msg.forwarded_from or "", + "reply_to": msg.reply_to, + "media_types": ",".join(msg.media_types), + "media_urls": " ".join(m.url for m in msg.media if m.url), + "reactions": json.dumps(msg.reactions, ensure_ascii=False) if msg.reactions else "", + "reactions_total": sum(msg.reactions.values()), + "hashtags": " ".join("#" + h for h in msg.hashtags), + "mentions": " ".join("@" + m for m in msg.mentions), + "links": " ".join(msg.links), + } + + +def to_json(messages: Iterable[Message], indent: Optional[int] = 2) -> str: + return json.dumps([m.to_dict() for m in messages], ensure_ascii=False, indent=indent) + + +def to_markdown(messages: Iterable[Message]) -> str: + parts = [] + for m in messages: + head = f"### [{m.channel}/{m.id}]({m.url})" + meta = " · ".join(x for x in [ + m.date.strftime("%Y-%m-%d %H:%M") if m.date else "", + f"👁 {m.views:,}" if m.views is not None else "", + ("📎 " + ", ".join(m.media_types)) if m.media else "", + ] if x) + parts.append(f"{head}\n_{meta}_\n\n{m.text or '(no text)'}\n") + return "\n---\n\n".join(parts) + + +def export(messages: Iterable[Message], path: Union[str, Path], format: Optional[str] = None) -> Path: + """Write messages to ``path``; the format is guessed from the extension. + + SQLite exports are *upserts*, so running the same export repeatedly builds a growing archive. + """ + path = Path(path) + fmt = detect_format(path, format) + messages = list(messages) + if path.parent and not path.parent.exists(): + path.parent.mkdir(parents=True, exist_ok=True) + + if fmt == "json": + path.write_text(to_json(messages), encoding="utf-8") + elif fmt == "jsonl": + with path.open("w", encoding="utf-8") as fh: + for m in messages: + fh.write(json.dumps(m.to_dict(), ensure_ascii=False) + "\n") + elif fmt == "csv": + # utf-8-sig so Excel opens Persian/Arabic/emoji text correctly + with path.open("w", encoding="utf-8-sig", newline="") as fh: + writer = csv.DictWriter(fh, fieldnames=FLAT_COLUMNS) + writer.writeheader() + writer.writerows(flatten(m) for m in messages) + elif fmt == "xlsx": + try: + from openpyxl import Workbook + except ImportError as exc: # pragma: no cover + raise ImportError("Excel export needs openpyxl: pip install 'tgscraper[excel]'") from exc + wb = Workbook() + ws = wb.active + ws.title = "messages" + ws.append(FLAT_COLUMNS) + for m in messages: + ws.append([flatten(m)[c] for c in FLAT_COLUMNS]) + wb.save(path) + elif fmt == "sqlite": + _to_sqlite(messages, path) + elif fmt == "md": + path.write_text(to_markdown(messages), encoding="utf-8") + return path + + +def _to_sqlite(messages: List[Message], path: Path) -> None: + cols = list(_SQL_COLUMNS) + with sqlite3.connect(path) as db: + db.execute( + "CREATE TABLE IF NOT EXISTS messages (" + + ", ".join(f"{c} {t}" for c, t in _SQL_COLUMNS.items()) + + ", PRIMARY KEY (channel, id))" + ) + db.execute("CREATE INDEX IF NOT EXISTS idx_messages_date ON messages(date)") + rows = [] + for m in messages: + d = m.to_dict() + for key in ("media", "reactions", "hashtags", "mentions", "links"): + d[key] = json.dumps(d[key], ensure_ascii=False) + d["edited"] = int(m.edited) + rows.append([d[c] for c in cols]) + db.executemany( + f"INSERT OR REPLACE INTO messages ({', '.join(cols)}) VALUES ({', '.join('?' * len(cols))})", rows + ) + + +def load(path: Union[str, Path], format: Optional[str] = None) -> List[Message]: + """Load messages previously saved as json / jsonl / sqlite.""" + path = Path(path) + fmt = detect_format(path, format) + if fmt == "json": + return [Message.from_dict(d) for d in json.loads(path.read_text(encoding="utf-8"))] + if fmt == "jsonl": + with path.open(encoding="utf-8") as fh: + return [Message.from_dict(json.loads(line)) for line in fh if line.strip()] + if fmt == "sqlite": + with sqlite3.connect(path) as db: + db.row_factory = sqlite3.Row + out = [] + for row in db.execute("SELECT * FROM messages ORDER BY date DESC, id DESC"): + d = dict(row) + for key in ("media", "reactions", "hashtags", "mentions", "links"): + d[key] = json.loads(d[key]) if d[key] else ([] if key != "reactions" else {}) + d["edited"] = bool(d["edited"]) + out.append(Message.from_dict(d)) + return out + raise ValueError(f"Loading {fmt} is not supported; use json, jsonl or sqlite") diff --git a/tgscraper/filters.py b/tgscraper/filters.py new file mode 100644 index 0000000..56215c2 --- /dev/null +++ b/tgscraper/filters.py @@ -0,0 +1,80 @@ +"""Client-side message filters.""" +from __future__ import annotations + +import re +from dataclasses import dataclass, field +from datetime import date, datetime, time, timezone +from typing import Iterable, List, Optional, Pattern, Sequence, Union + +from .models import Message + +DateLike = Union[str, date, datetime, None] + + +def to_datetime(value: DateLike, end_of_day: bool = False) -> Optional[datetime]: + """Parse ``"2026-01-31"``, ``"2026-01-31T10:00"``, ``date`` or ``datetime`` into an aware UTC datetime.""" + if value is None or value == "": + return None + if isinstance(value, str): + value = value.strip() + parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) + if len(value) == 10 and end_of_day: + parsed = datetime.combine(parsed.date(), time.max) + value = parsed + elif not isinstance(value, datetime): + value = datetime.combine(value, time.max if end_of_day else time.min) + if value.tzinfo is None: + value = value.replace(tzinfo=timezone.utc) + return value + + +@dataclass +class MessageFilter: + """Filters applied to every scraped message. All conditions must match.""" + + since: DateLike = None + until: DateLike = None + keywords: Sequence[str] = field(default_factory=list) # any of them, case-insensitive + regex: Optional[Union[str, Pattern]] = None + media_only: bool = False + media_types: Sequence[str] = field(default_factory=list) # photo, video, document, voice, ... + min_views: Optional[int] = None + hashtag: Optional[str] = None + + def __post_init__(self) -> None: + self.since = to_datetime(self.since) + self.until = to_datetime(self.until, end_of_day=True) + if isinstance(self.keywords, str): + self.keywords = [self.keywords] + self.keywords = [k.lower() for k in self.keywords if k] + if isinstance(self.media_types, str): + self.media_types = [self.media_types] + if isinstance(self.regex, str) and self.regex: + self.regex = re.compile(self.regex, re.IGNORECASE) + if self.hashtag: + self.hashtag = self.hashtag.lstrip("#").lower() + + def matches(self, msg: Message) -> bool: + date_ = msg.date + if date_ is not None and date_.tzinfo is None: + date_ = date_.replace(tzinfo=timezone.utc) + if self.since and date_ and date_ < self.since: + return False + if self.until and date_ and date_ > self.until: + return False + if self.keywords and not any(k in msg.text.lower() for k in self.keywords): + return False + if self.regex and not self.regex.search(msg.text): + return False + if self.media_only and not msg.media: + return False + if self.media_types and not set(self.media_types) & set(msg.media_types): + return False + if self.min_views is not None and (msg.views or 0) < self.min_views: + return False + if self.hashtag and self.hashtag not in (h.lower() for h in msg.hashtags): + return False + return True + + def apply(self, messages: Iterable[Message]) -> List[Message]: + return [m for m in messages if self.matches(m)] diff --git a/tgscraper/mcp_server.py b/tgscraper/mcp_server.py new file mode 100644 index 0000000..6089e17 --- /dev/null +++ b/tgscraper/mcp_server.py @@ -0,0 +1,220 @@ +"""MCP (Model Context Protocol) server so AI assistants (Claude, Cursor, ChatGPT...) can read Telegram channels. + +Run: ``tgscraper-mcp`` (or ``tgscraper mcp``, or ``python -m tgscraper.mcp_server``) +""" +from __future__ import annotations + +import os +from typing import Any, Dict, List, Optional + +try: # mcp >= 2 + from mcp.server.mcpserver import MCPServer as FastMCP +except ImportError: + try: # mcp 1.x + from mcp.server.fastmcp import FastMCP + except ImportError as exc: # pragma: no cover + raise ImportError("The MCP server needs the 'mcp' package: pip install 'tgscraper[mcp]'") from exc + +from .analytics import summarize +from .client import AsyncScraper, ScraperError +from .exporters import export +from .models import Message + +PROXY = os.environ.get("TGSCRAPER_PROXY") or None +MAX_LIMIT = int(os.environ.get("TGSCRAPER_MCP_MAX_LIMIT", "500")) +MAX_TEXT = 4000 + +mcp = FastMCP( + "telegram-scraper", + instructions=( + "Read public Telegram channels (t.me/s/) without any API key. " + "Channel arguments accept 'durov', '@durov' or 'https://t.me/durov'. " + "Messages are returned newest first. Start with get_channel_info or get_messages; " + "use search_messages for topics and analyze_channel for statistics. " + "Only public channels with web preview enabled can be read." + ), +) + + +def _scraper() -> AsyncScraper: + return AsyncScraper(proxies=PROXY, delay=0.3) + + +def _compact(m: Message) -> Dict[str, Any]: + d = m.to_dict() + d.pop("html", None) + if len(d["text"]) > MAX_TEXT: + d["text"] = d["text"][:MAX_TEXT] + "… [truncated]" + d["media"] = [{k: v for k, v in x.items() if v} for x in d["media"]] + return {k: v for k, v in d.items() if v not in (None, [], {}, "", False) or k in ("id", "text")} + + +def _limit(n: Optional[int]) -> int: + return max(1, min(int(n or 20), MAX_LIMIT)) + + +def _error(exc: Exception) -> Dict[str, Any]: + return {"error": str(exc)} + + +@mcp.tool() +async def get_channel_info(channel: str) -> Dict[str, Any]: + """Get a public Telegram channel's title, description, subscriber count, photo and media counters.""" + try: + async with _scraper() as s: + return (await s.channel_info(channel)).to_dict() + except ScraperError as exc: + return _error(exc) + + +@mcp.tool() +async def get_messages( + channel: str, + limit: int = 20, + since: Optional[str] = None, + until: Optional[str] = None, + keywords: Optional[List[str]] = None, + hashtag: Optional[str] = None, + media_only: bool = False, + min_views: Optional[int] = None, + before_id: Optional[int] = None, +) -> Dict[str, Any]: + """Get recent posts of a public Telegram channel, newest first. + + Args: + channel: username or link, e.g. "durov" or "https://t.me/durov". + limit: number of posts to return (1-500, default 20). + since / until: ISO dates like "2026-01-31" to restrict the time range. + keywords: keep only posts containing any of these words (case-insensitive). + hashtag: keep only posts with this hashtag. + media_only: keep only posts with photos/videos/files. + min_views: keep only posts with at least this many views. + before_id: return posts older than this message id (for paging backwards). + """ + try: + async with _scraper() as s: + kwargs: Dict[str, Any] = dict(since=since, until=until, keywords=keywords or [], hashtag=hashtag, + media_only=media_only, min_views=min_views) + msgs = await s.get_messages(channel, _limit(limit), before=before_id, max_pages=100, **kwargs) + return {"channel": channel, "count": len(msgs), "messages": [_compact(m) for m in msgs]} + except (ScraperError, ValueError) as exc: + return _error(exc) + + +@mcp.tool() +async def search_messages(channel: str, query: str, limit: int = 20) -> Dict[str, Any]: + """Search the whole history of a public Telegram channel using Telegram's own search.""" + try: + async with _scraper() as s: + msgs = await s.get_messages(channel, _limit(limit), query=query, max_pages=50) + return {"channel": channel, "query": query, "count": len(msgs), "messages": [_compact(m) for m in msgs]} + except (ScraperError, ValueError) as exc: + return _error(exc) + + +@mcp.tool() +async def get_message(channel: str, message_id: int) -> Dict[str, Any]: + """Get one specific post, e.g. for a link like https://t.me/durov/123 use channel="durov", message_id=123.""" + try: + async with _scraper() as s: + msg = await s.get_message(channel, message_id) + return _compact(msg) if msg else {"error": f"Message {message_id} not found in {channel}"} + except (ScraperError, ValueError) as exc: + return _error(exc) + + +@mcp.tool() +async def get_new_messages(channel: str, after_id: int, limit: int = 100) -> Dict[str, Any]: + """Get only posts newer than ``after_id`` — use the highest id you saw before to follow a channel.""" + try: + async with _scraper() as s: + msgs = await s.get_messages(channel, _limit(limit), min_id=after_id, max_pages=20) + return {"channel": channel, "count": len(msgs), "latest_id": max([m.id for m in msgs] or [after_id]), + "messages": [_compact(m) for m in msgs]} + except (ScraperError, ValueError) as exc: + return _error(exc) + + +@mcp.tool() +async def analyze_channel(channel: str, limit: int = 200) -> Dict[str, Any]: + """Statistics for a channel's recent posts: views, top posts, posting frequency and hours, hashtags, + frequent words, media mix, reactions and a rough sentiment score. Also includes channel info.""" + try: + async with _scraper() as s: + info = await s.channel_info(channel) + msgs = await s.get_messages(channel, _limit(limit)) + stats = summarize(msgs) + stats.pop("posts_by_day", None) + return {"channel": info.to_dict(), "stats": stats} + except (ScraperError, ValueError) as exc: + return _error(exc) + + +@mcp.tool() +async def compare_channels(channels: List[str], limit: int = 100) -> Dict[str, Any]: + """Compare several public channels side by side (subscribers, activity, average views, sentiment).""" + out: Dict[str, Any] = {} + async with _scraper() as s: + results = await s.scrape_many(channels[:10], _limit(limit)) + for ch, res in results.items(): + if isinstance(res, Exception): + out[ch] = {"error": str(res)} + continue + try: + info = await s.channel_info(ch) + except ScraperError as exc: + out[ch] = {"error": str(exc)} + continue + st = summarize(res) + out[ch] = { + "title": info.title, "subscribers": info.subscribers, "posts_analyzed": st["count"], + "posts_per_day": st.get("posts_per_day"), "avg_views": st.get("avg_views"), + "engagement_rate": round(st["avg_views"] / info.subscribers, 3) + if st.get("avg_views") and info.subscribers else None, + "with_media": st.get("with_media"), "sentiment": st.get("sentiment", {}).get("average"), + "top_hashtags": st.get("top_hashtags", [])[:5], + } + return out + + +@mcp.tool() +async def export_messages(channel: str, path: str, limit: int = 100) -> Dict[str, Any]: + """Save a channel's posts to a local file. Format from the extension: .json .jsonl .csv .xlsx .db .md""" + try: + async with _scraper() as s: + msgs = await s.get_messages(channel, _limit(limit)) + saved = export(msgs, os.path.expanduser(path)) + return {"saved": str(saved.resolve()), "count": len(msgs)} + except (ScraperError, ValueError, ImportError, OSError) as exc: + return _error(exc) + + +@mcp.resource("telegram://channel/{channel}") +async def channel_resource(channel: str) -> str: + """The 20 latest posts of a channel as Markdown.""" + from .exporters import to_markdown + async with _scraper() as s: + return to_markdown(await s.get_messages(channel, 20)) + + +@mcp.prompt() +def summarize_channel(channel: str, days: int = 7) -> str: + """Summarize what a Telegram channel posted recently.""" + return (f"Use get_channel_info and get_messages (since = {days} days ago) for the Telegram channel " + f"'{channel}'. Summarize the main topics, the most viewed posts (with links), notable announcements, " + f"and the overall tone. Answer in the user's language.") + + +@mcp.prompt() +def track_topic(channels: str, topic: str) -> str: + """Find what several channels say about a topic.""" + return (f"For each of these Telegram channels: {channels} — call search_messages with query '{topic}'. " + f"Compare what each channel says about '{topic}', cite post links, and note dates.") + + +def main(transport: str = "stdio") -> None: + mcp.run(transport=transport) + + +if __name__ == "__main__": + main() diff --git a/tgscraper/media.py b/tgscraper/media.py new file mode 100644 index 0000000..e1eed92 --- /dev/null +++ b/tgscraper/media.py @@ -0,0 +1,58 @@ +"""Download photos / videos / voice notes attached to messages.""" +from __future__ import annotations + +import logging +import mimetypes +from pathlib import Path +from typing import Iterable, List, Optional, Sequence, Union +from urllib.parse import urlparse + +import httpx + +from .client import DEFAULT_HEADERS +from .models import Message + +log = logging.getLogger("tgscraper") + +DOWNLOADABLE = ("photo", "video", "round_video", "voice", "sticker") + + +def download_media( + messages: Iterable[Message], + folder: Union[str, Path] = "media", + types: Optional[Sequence[str]] = None, + proxy: Optional[str] = None, + overwrite: bool = False, +) -> List[Path]: + """Download media files to ``folder`` as ``__.``. Returns the saved paths. + + Documents are not downloadable from the public preview (they link back to the Telegram app). + """ + folder = Path(folder) + folder.mkdir(parents=True, exist_ok=True) + wanted = set(types or DOWNLOADABLE) & set(DOWNLOADABLE) + saved: List[Path] = [] + with httpx.Client(proxy=proxy, headers=DEFAULT_HEADERS, timeout=60, follow_redirects=True) as client: + for msg in messages: + for n, media in enumerate(msg.media, 1): + if media.type not in wanted or not media.url or not media.url.startswith("http"): + continue + ext = Path(urlparse(media.url).path).suffix + stem = folder / f"{msg.channel}_{msg.id}_{n}" + existing = list(folder.glob(stem.name + ".*")) + if existing and not overwrite: + saved.append(existing[0]) + continue + try: + response = client.get(media.url) + response.raise_for_status() + except httpx.HTTPError as exc: + log.warning("Could not download %s: %s", media.url, exc) + continue + if not ext: + content_type = response.headers.get("content-type", "").split(";")[0] + ext = mimetypes.guess_extension(content_type) or ".bin" + target = stem.with_suffix(ext) + target.write_bytes(response.content) + saved.append(target) + return saved diff --git a/tgscraper/models.py b/tgscraper/models.py new file mode 100644 index 0000000..6dc2e18 --- /dev/null +++ b/tgscraper/models.py @@ -0,0 +1,84 @@ +"""Data models returned by the scraper.""" +from __future__ import annotations + +from dataclasses import asdict, dataclass, field +from datetime import datetime +from typing import Any, Dict, List, Optional + + +@dataclass +class Media: + """A photo, video, document, voice note, etc. attached to a message.""" + + type: str # photo | video | round_video | document | voice | sticker | link_preview + url: Optional[str] = None + thumbnail: Optional[str] = None + duration: Optional[str] = None + title: Optional[str] = None + + def to_dict(self) -> Dict[str, Any]: + return asdict(self) + + +@dataclass +class Message: + """One post of a public Telegram channel.""" + + id: int + channel: str + url: str + date: Optional[datetime] = None + text: str = "" + html: str = "" + views: Optional[int] = None + author: Optional[str] = None + edited: bool = False + forwarded_from: Optional[str] = None + reply_to: Optional[int] = None + media: List[Media] = field(default_factory=list) + reactions: Dict[str, int] = field(default_factory=dict) + hashtags: List[str] = field(default_factory=list) + mentions: List[str] = field(default_factory=list) + links: List[str] = field(default_factory=list) + + @property + def has_media(self) -> bool: + return bool(self.media) + + @property + def media_types(self) -> List[str]: + return [m.type for m in self.media] + + def to_dict(self) -> Dict[str, Any]: + data = asdict(self) + data["date"] = self.date.isoformat() if self.date else None + return data + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "Message": + data = dict(data) + if isinstance(data.get("date"), str) and data["date"]: + data["date"] = datetime.fromisoformat(data["date"]) + data["media"] = [m if isinstance(m, Media) else Media(**m) for m in data.get("media") or []] + known = cls.__dataclass_fields__.keys() + return cls(**{k: v for k, v in data.items() if k in known}) + + def __str__(self) -> str: + return self.text + + +@dataclass +class Channel: + """Public information about a channel.""" + + username: str + url: str + title: Optional[str] = None + description: Optional[str] = None + photo: Optional[str] = None + subscribers: Optional[int] = None + counters: Dict[str, int] = field(default_factory=dict) # subscribers, photos, videos, links, files... + verified: bool = False + + def to_dict(self) -> Dict[str, Any]: + return asdict(self) diff --git a/tgscraper/monitor.py b/tgscraper/monitor.py new file mode 100644 index 0000000..64738cb --- /dev/null +++ b/tgscraper/monitor.py @@ -0,0 +1,93 @@ +"""Watch channels and push new posts to a callback, a webhook or a Telegram bot.""" +from __future__ import annotations + +import logging +import time +from typing import Callable, List, Optional, Sequence + +import httpx + +from .client import Scraper, ScraperError +from .models import Message +from .parser import normalize_channel + +log = logging.getLogger("tgscraper") + +Handler = Callable[[Message], None] + + +def webhook_notifier(url: str, timeout: float = 15) -> Handler: + """POST every new message as JSON to ``url`` (Slack/Discord/n8n/Zapier/your API...).""" + def send(msg: Message) -> None: + payload = msg.to_dict() + # Slack & Discord read "text"/"content"; everyone else gets the full message. + payload.update({"content": f"{msg.url}\n{msg.text}"[:1900]}) + try: + httpx.post(url, json=payload, timeout=timeout).raise_for_status() + except httpx.HTTPError as exc: + log.warning("Webhook failed: %s", exc) + return send + + +def telegram_notifier(bot_token: str, chat_id: str, timeout: float = 15) -> Handler: + """Forward new posts to a chat through a Telegram bot (create one with @BotFather).""" + api = f"https://api.telegram.org/bot{bot_token}/sendMessage" + + def send(msg: Message) -> None: + text = f"📢 {msg.channel}\n\n{msg.text}\n\n{msg.url}" + try: + httpx.post(api, json={"chat_id": chat_id, "text": text[:4096]}, timeout=timeout).raise_for_status() + except httpx.HTTPError as exc: + log.warning("Telegram bot notification failed: %s", exc) + return send + + +def watch( + channels: Sequence[str], + handlers: Sequence[Handler], + interval: float = 60, + keywords: Optional[Sequence[str]] = None, + scraper: Optional[Scraper] = None, + backfill: int = 0, + iterations: Optional[int] = None, +) -> None: + """Poll ``channels`` every ``interval`` seconds and call every handler for each new message. + + ``backfill`` = how many existing posts to emit on start (0 = only brand-new posts). + ``keywords`` = only notify when a post contains one of them. + ``iterations`` = stop after N polls (``None`` = forever, Ctrl+C to quit). + """ + own = scraper is None + scraper = scraper or Scraper() + names: List[str] = [normalize_channel(c) for c in channels] + last: dict = {} + try: + for ch in names: + recent = scraper.get_messages(ch, limit=max(backfill, 1)) + last[ch] = max((m.id for m in recent), default=0) + for msg in reversed(recent[:backfill]): + _dispatch(msg, handlers, keywords) + log.info("Watching %s (last id %s)", ch, last[ch]) + polls = 0 + while iterations is None or polls < iterations: + time.sleep(interval) + polls += 1 + for ch in names: + try: + new = scraper.get_messages(ch, limit=None, min_id=last[ch], max_pages=5) + except ScraperError as exc: + log.warning("Polling %s failed: %s", ch, exc) + continue + for msg in reversed(new): # oldest first + _dispatch(msg, handlers, keywords) + last[ch] = max(last[ch], msg.id) + finally: + if own: + scraper.close() + + +def _dispatch(msg: Message, handlers: Sequence[Handler], keywords: Optional[Sequence[str]]) -> None: + if keywords and not any(k.lower() in msg.text.lower() for k in keywords): + return + for handler in handlers: + handler(msg) diff --git a/tgscraper/parser.py b/tgscraper/parser.py new file mode 100644 index 0000000..e495732 --- /dev/null +++ b/tgscraper/parser.py @@ -0,0 +1,240 @@ +"""HTML parsing for the public web preview of a channel (https://t.me/s/).""" +from __future__ import annotations + +import re +from datetime import datetime +from typing import Dict, List, Optional, Tuple + +from bs4 import BeautifulSoup, Tag + +from .models import Channel, Media, Message + +_CHANNEL_RE = re.compile(r"^(?:https?://)?(?:www\.)?(?:t\.me|telegram\.me)/(?:s/)?([A-Za-z0-9_]+)", re.I) +_BG_URL_RE = re.compile(r"url\(['\"]?([^'\")]+)['\"]?\)") +_HASHTAG_RE = re.compile(r"(? str: + """Accept ``durov``, ``@durov``, ``t.me/durov``, ``https://t.me/s/durov`` ... and return ``durov``.""" + channel = channel.strip() + match = _CHANNEL_RE.match(channel) + if match: + return match.group(1) + channel = channel.lstrip("@").strip("/") + if not re.fullmatch(r"[A-Za-z0-9_]+", channel): + raise ValueError(f"Invalid channel name or URL: {channel!r}") + return channel + + +def parse_count(value: Optional[str]) -> Optional[int]: + """Convert Telegram counters like ``1.2K``, ``3,4M`` or ``12 345`` to ints.""" + if not value: + return None + value = value.strip().replace(" ", "").replace(" ", "") + multiplier = 1 + if value[-1:].upper() in ("K", "M", "B"): + multiplier = {"K": 1_000, "M": 1_000_000, "B": 1_000_000_000}[value[-1].upper()] + value = value[:-1].replace(",", ".") + else: + value = value.replace(",", "") + try: + return int(round(float(value) * multiplier)) + except ValueError: + return None + + +def _bg_url(tag: Tag) -> Optional[str]: + match = _BG_URL_RE.search(tag.get("style", "") or "") + return match.group(1) if match else None + + +def _text_of(tag: Tag) -> str: + for br in tag.find_all("br"): + br.replace_with("\n") + return tag.get_text().strip() + + +def _own(root: Tag, selector: str) -> List[Tag]: + """Elements matching ``selector`` that are not inside a quoted reply / link preview.""" + out = [] + for el in root.select(selector): + if el.find_parent(class_=["tgme_widget_message_reply", "link_preview_wrap"]) is not None: + continue + out.append(el) + return out + + +def _parse_media(node: Tag) -> List[Media]: + media: List[Media] = [] + for el in _own(node, "a.tgme_widget_message_photo_wrap"): + media.append(Media("photo", url=_bg_url(el), thumbnail=_bg_url(el))) + for el in _own(node, ".tgme_widget_message_video_player"): + video = el.select_one("video") + thumb = el.select_one(".tgme_widget_message_video_thumb") + duration = el.select_one(".message_video_duration") + kind = "round_video" if "tgme_widget_message_roundvideo_player" in (el.get("class") or []) else "video" + media.append(Media( + kind, + url=video.get("src") if video else None, + thumbnail=_bg_url(thumb) if thumb else None, + duration=duration.get_text(strip=True) if duration else None, + )) + for el in _own(node, "audio.tgme_widget_message_voice"): + duration = node.select_one(".tgme_widget_message_voice_duration") + media.append(Media("voice", url=el.get("src"), duration=duration.get_text(strip=True) if duration else None)) + for el in _own(node, ".tgme_widget_message_document_wrap"): + title = el.select_one(".tgme_widget_message_document_title") + kind = "audio" if el.select_one(".audio") else "document" + media.append(Media(kind, url=el.get("href"), title=title.get_text(strip=True) if title else None)) + for el in _own(node, ".tgme_widget_message_sticker_wrap"): + img = el.select_one("img, video") + media.append(Media("sticker", url=img.get("src") if img else _bg_url(el))) + for el in node.select("a.tgme_widget_message_link_preview"): + title = el.select_one(".link_preview_title") or el.select_one(".link_preview_site_name") + image = el.select_one(".link_preview_image, .link_preview_right_image") + media.append(Media("link_preview", url=el.get("href"), + thumbnail=_bg_url(image) if image else None, + title=title.get_text(strip=True) if title else None)) + return media + + +_EMOJI_IMG_RE = re.compile(r"/emoji/\d+/([0-9A-Fa-f]+)\.png") + + +def _reaction_key(el: Tag) -> str: + """Emoji of a reaction: inline text, the hex-encoded emoji image name, or a custom emoji id.""" + bold = el.select_one("b") + if bold and bold.get_text(strip=True): + return bold.get_text(strip=True) + for styled in [el] + el.find_all(style=True): + match = _EMOJI_IMG_RE.search(styled.get("style", "") or "") + if match: + try: + return bytes.fromhex(match.group(1)).decode("utf-8") + except ValueError: + pass + if "tgme_reaction_paid" in (el.get("class") or []): + return "⭐" + custom = el.select_one("[emoji-id]") + if custom is not None: + return f"custom:{custom['emoji-id']}" + return el.get("data-emoji") or "custom" + + +def _parse_reactions(node: Tag) -> Dict[str, int]: + reactions: Dict[str, int] = {} + for el in node.select(".tgme_reaction"): + emoji = _reaction_key(el) + for child in el.find_all(["i", "tg-emoji", "b"]): + child.extract() + count = parse_count(el.get_text(strip=True)) + reactions[emoji] = reactions.get(emoji, 0) + (count or 0) + return reactions + + +def parse_message(node: Tag, channel: str) -> Optional[Message]: + post = node.get("data-post") or "" + if "/" not in post: + return None + post_channel, _, post_id = post.rpartition("/") + try: + msg_id = int(post_id) + except ValueError: + return None + + text_nodes = _own(node, ".tgme_widget_message_text") + text_el = text_nodes[0] if text_nodes else None + html = text_el.decode_contents() if text_el else "" + links: List[str] = [] + if text_el: + for a in text_el.find_all("a", href=True): + href = a["href"] + if not href.startswith("?q=") and href not in links: + links.append(href) + text = _text_of(text_el) if text_el else "" + + date = None + time_el = node.select_one(".tgme_widget_message_date time[datetime]") or node.select_one("time[datetime]") + if time_el: + try: + date = datetime.fromisoformat(time_el["datetime"]) + except ValueError: + date = None + + views_el = node.select_one(".tgme_widget_message_views") + author_el = node.select_one(".tgme_widget_message_from_author") + meta_el = node.select_one(".tgme_widget_message_meta") + fwd_el = node.select_one(".tgme_widget_message_forwarded_from_name") + reply_el = node.select_one("a.tgme_widget_message_reply") + reply_to = None + if reply_el and reply_el.get("href"): + match = re.search(r"/(\d+)(?:\?|$)", reply_el["href"]) + reply_to = int(match.group(1)) if match else None + + return Message( + id=msg_id, + channel=post_channel or channel, + url=f"https://t.me/{post_channel or channel}/{msg_id}", + date=date, + text=text, + html=html, + views=parse_count(views_el.get_text()) if views_el else None, + author=author_el.get_text(strip=True) if author_el else None, + edited=bool(meta_el and "edited" in meta_el.get_text().lower()), + forwarded_from=fwd_el.get_text(strip=True) if fwd_el else None, + reply_to=reply_to, + media=_parse_media(node), + reactions=_parse_reactions(node), + hashtags=list(dict.fromkeys(_HASHTAG_RE.findall(text))), + mentions=list(dict.fromkeys(_MENTION_RE.findall(text))), + links=links, + ) + + +def parse_channel_info(soup: BeautifulSoup, channel: str) -> Channel: + info = Channel(username=channel, url=f"https://t.me/{channel}") + title = soup.select_one(".tgme_channel_info_header_title") + if title: + info.title = title.get_text(strip=True) + info.verified = title.select_one(".verified-icon") is not None + desc = soup.select_one(".tgme_channel_info_description") + if desc: + info.description = _text_of(desc) + photo = soup.select_one(".tgme_channel_info_header .tgme_page_photo_image img") + if photo: + info.photo = photo.get("src") + for counter in soup.select(".tgme_channel_info_counter"): + value = counter.select_one(".counter_value") + kind = counter.select_one(".counter_type") + if value and kind: + info.counters[kind.get_text(strip=True).lower()] = parse_count(value.get_text()) or 0 + for key in ("subscribers", "subscriber", "members", "member"): + if key in info.counters: + info.subscribers = info.counters[key] + break + return info + + +def parse_page(html: str, channel: str) -> Tuple[List[Message], Channel, Optional[int]]: + """Parse one ``t.me/s/`` page. + + Returns ``(messages oldest->newest, channel info, before_id for the next older page or None)``. + """ + soup = BeautifulSoup(html, "html.parser") + messages = [] + for node in soup.select(".tgme_widget_message[data-post]"): + msg = parse_message(node, channel) + if msg: + messages.append(msg) + before = None + more = soup.select_one("a.tme_messages_more[data-before]") + if more: + try: + before = int(more["data-before"]) + except (KeyError, ValueError): + before = None + if before is None and more is not None and more.get("href"): + match = re.search(r"before=(\d+)", more["href"]) + before = int(match.group(1)) if match else None + return messages, parse_channel_info(soup, channel), before diff --git a/tgscraper/state.py b/tgscraper/state.py new file mode 100644 index 0000000..f774aea --- /dev/null +++ b/tgscraper/state.py @@ -0,0 +1,38 @@ +"""Remember the last scraped message per channel so later runs only fetch new posts.""" +from __future__ import annotations + +import json +import os +from pathlib import Path +from typing import Dict, Optional, Union + +DEFAULT_STATE_FILE = Path(os.environ.get("TGSCRAPER_STATE", Path.home() / ".tgscraper" / "state.json")) + + +class State: + def __init__(self, path: Union[str, Path, None] = None) -> None: + self.path = Path(path) if path else DEFAULT_STATE_FILE + self._data: Dict[str, int] = {} + if self.path.exists(): + try: + self._data = {k: int(v) for k, v in json.loads(self.path.read_text()).items()} + except (ValueError, OSError): + self._data = {} + + def last_id(self, channel: str) -> Optional[int]: + return self._data.get(channel.lower()) + + def update(self, channel: str, message_id: int) -> None: + key = channel.lower() + if message_id > self._data.get(key, 0): + self._data[key] = message_id + + def save(self) -> None: + self.path.parent.mkdir(parents=True, exist_ok=True) + self.path.write_text(json.dumps(self._data, indent=2, sort_keys=True)) + + def reset(self, channel: Optional[str] = None) -> None: + if channel: + self._data.pop(channel.lower(), None) + else: + self._data.clear()