Compare commits

...
19 Commits
Author SHA1 Message Date
YoursFunny f7cb809e5a bump version to 1.0.6 2026-08-05 02:07:34 +08:00
YoursFunny 51b40cdb42 feat: cache sent media file ids for instant repeat sends 2026-08-05 02:06:30 +08:00
YoursFunny ab2306002a perf: handle batch-forwarded URLs concurrently with queue workers 2026-08-05 01:26:24 +08:00
YoursFunny a92b12f633 bump version to 1.0.5 2026-08-05 00:16:45 +08:00
YoursFunny 74e3b7593c docs: note fresh twitter query ids and TID scope in auth fallback 2026-08-05 00:07:30 +08:00
YoursFunny 2faccaac42 fix: cut tweet text by code points, not UTF-16 units 2026-08-04 23:45:54 +08:00
YoursFunny f4e60d8946 fix: trim CRLF from TWITTER_AUTH_TOKEN 2026-08-04 23:45:54 +08:00
YoursFunny 3006dcd98c feat: fetch NSFW tweets via authenticated twitter API fallback 2026-08-04 23:45:53 +08:00
YoursFunny 8e8acdd859 deploy: add container names and startup order to compose 2026-08-04 23:11:04 +08:00
YoursFunny 8b9dd963e1 bump version to 1.0.4 2026-08-04 22:10:38 +08:00
YoursFunny e83d48f1f7 fix: handle SIGTERM for graceful shutdown on docker stop 2026-08-04 22:09:45 +08:00
YoursFunny 7c55b26731 deploy: add nginx-proxy reverse proxy for webhook TLS 2026-08-04 21:50:19 +08:00
YoursFunny 950db48a13 ci: cache buildkit layers across runs 2026-08-04 18:52:28 +08:00
YoursFunny 32ea8ec6ca fix docker build: fetch ffmpeg from martin-riedl.de 2026-08-04 18:52:27 +08:00
YoursFunny 62dc033452 docs: add AGENTS.md with repository guidelines 2026-08-04 18:52:21 +08:00
YoursFunny 2801aaa39c bump version to 1.0.3 2026-08-04 16:24:16 +08:00
YoursFunny 93d47752cf fallback to smaller media when file too large 2026-08-04 16:22:49 +08:00
YoursFunny 6fd4edb3c5 webhook: drop duplicate set_webhook, ignore empty env vars 2026-08-04 15:58:32 +08:00
YoursFunny 7c5afce0b4 ci: build once on tagged commits, bump docker actions 2026-08-04 01:58:19 +08:00
23 changed files with 1802 additions and 215 deletions
+4
View File
@@ -28,5 +28,9 @@ LICENSE
README.md
data/
cert/
nginx-certs/
nginx-vhost.d/
nginx-html/
nginx-acme/
**/target/
.idea/
+37 -6
View File
@@ -12,12 +12,35 @@ env:
DOCKERHUB_REPO: yoursfunny/telegram-twitter-media-bot
jobs:
# A tag push and a branch push to the same commit fire two workflow runs;
# build only once. Tag runs always build; master runs build only when the
# pushed commit is not already tagged (the tag run covers it).
should-build:
runs-on: ubuntu-latest
outputs:
build: ${{ steps.check.outputs.build }}
steps:
- uses: actions/checkout@v7
with:
fetch-depth: 0
- id: check
shell: bash
run: |
if [ "$GITHUB_REF_TYPE" = "branch" ] && git tag --points-at "$GITHUB_SHA" | grep -q .; then
echo "commit already tagged; the tag run builds the image"
echo "build=false" >> "$GITHUB_OUTPUT"
else
echo "build=true" >> "$GITHUB_OUTPUT"
fi
docker:
needs: should-build
if: needs.should-build.outputs.build == 'true'
runs-on: ubuntu-latest
steps:
- name: Docker meta
id: meta
uses: docker/metadata-action@v5
uses: docker/metadata-action@v6
with:
images: ${{ env.DOCKERHUB_REPO }}
tags: |
@@ -28,22 +51,30 @@ jobs:
type=sha
-
name: Set up QEMU
uses: docker/setup-qemu-action@v3
uses: docker/setup-qemu-action@v4
-
name: Set up Docker Buildx
uses: docker/setup-buildx-action@v3
uses: docker/setup-buildx-action@v4
-
name: Login to Docker Hub
uses: docker/login-action@v3
uses: docker/login-action@v4
with:
username: ${{ secrets.DOCKERHUB_USERNAME }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
# Buildkit cache via the GitHub Actions cache backend (uses the
# automatic GITHUB_TOKEN, no extra secrets). mode=max keeps every
# stage's layers so the cargo-deps and ffmpeg layers are restored
# instead of re-downloaded/recompiled. The scope must be pinned to a
# fixed string: the gha backend defaults to the current git ref, which
# would give every new tag a cold cache on release builds.
-
name: Build and push
uses: docker/build-push-action@v5
uses: docker/build-push-action@v7
with:
push: true
build-args: |
APP_NAME=${{ env.APP_NAME }}
tags: ${{ steps.meta.outputs.tags }}
labels: ${{ steps.meta.outputs.labels }}
labels: ${{ steps.meta.outputs.labels }}
cache-from: type=gha,scope=tgxmb-build
cache-to: type=gha,mode=max,scope=tgxmb-build
+4
View File
@@ -2,6 +2,10 @@
__pycache__/
cert/
data/
nginx-certs/
nginx-vhost.d/
nginx-html/
nginx-acme/
docker-compose.yml
.env
+98
View File
@@ -0,0 +1,98 @@
# Repository Guidelines
## Project Overview
Telegram bot (teloxide) that turns post links from X/Twitter, Pixiv, and Bluesky into media messages (images, video, GIF) with the post's title, author, and tags. It supports batch media splitting, retry with persistence, inline queries, forward-channel rebinding with caption templates, and Pixiv ugoira→MP4 transcoding. README and user-facing strings are in Chinese. The project is a Rust port of a Python predecessor (see `queue.rs` comments referencing `utils/task_queue.py`).
Two-crate Cargo workspace (both v1.0.3, edition 2024, resolver 3):
- **`crates/x-media`** — library that fetches and normalizes media from the three sites. Pure, no Telegram knowledge.
- **`crates/xmedia-bot`** — the bot binary: teloxide dispatcher, SQLite-backed chat state, persistent task queue.
## Architecture & Data Flow
```
Telegram update → Dispatcher (polling or axum webhook) → dptree branches
├─ message → commands (any chat) / URL links (private chat only)
├─ inline_query → InlineQueryResult Photo/Video/Mpeg4Gif
└─ callback_query → "forward" (copy to channel) / "template|<name>" (apply caption template)
```
Message flow: `message_handler` extracts URLs (from `url`/`text_link` entities, text + caption, deduped) → `x_media::site::fetch(url)``Fetched` → builds a `Task``send::send_media_sequence` (media groups ≤ 9, caption on first item) or `send::send_animation`. On Telegram URL-fetch failure or size error (`send_batch_via_upload`): download via `x_media::site::download_media` to a temp file (≤ 10 MiB), sniff magic bytes (`sniff_ext`), upload via multipart; oversized items fall back to `fallback_url`. On failure: `enqueue_retry` persists resume-state `Task` into the SQLite queue → single worker leases (120 s lock TTL) → retry with exponential backoff (≤ 30 s, `MAX_RETRIES = 2`) → dead-letter → `notify_failure`. Success → `post_send_actions`: edit-before-forward prompt with inline buttons, or `copy_messages` to the bound forward channel.
The `x-media` library: `site::fetch(url)` dispatches (in order) twitter → bsky → pixiv via per-site regex `PATTERN` and returns `Ok(None)` for unmatched URLs. `Fetched { source_url, caption, title, media: Vec<Media>, sensitive, … }`; `caption_with(format)` substitutes `{url} {author} {author_url} {title} {tags}`.
## Key Directories
| Path | Purpose |
|---|---|
| `crates/x-media/src/` | Fetch library. `site/mod.rs` = dispatcher + `Fetched`/`FetchError`/`download_media`/`media_size`; `media.rs` = `Media` enum; `examples/fetch.rs` = end-to-end usage sample |
| `crates/x-media/src/site/<twitter\|pixiv\|bsky>/` | One directory per site: `mod.rs` (re-exports), `interface.rs` (PATTERN, `enabled()`, `fetch_from_url()`, site struct, `From<SiteStruct> for Fetched`), `model.rs` (serde DTOs). Pixiv adds `api.rs` (auth + transport); twitter adds `auth.rs` (logged-in GraphQL `TweetDetail` fallback for NSFW tweets, gated on `TWITTER_AUTH_TOKEN`) |
| `crates/xmedia-bot/src/main.rs` | Entry point: env/log init, queue worker start, pixiv validation, 300 s edit-expiry sweep, dptree handler tree, webhook vs polling dispatch |
| `crates/xmedia-bot/src/config.rs` | Manual env parsing into `Config` |
| `crates/xmedia-bot/src/handlers.rs` | `Command` enum (teloxide `BotCommands`), message/inline/callback handlers, URL extraction, global statics; per-URL work spawned with a `Semaphore(8)` cap (teloxide's per-chat workers are sequential — batch-forwards need concurrency) |
| `crates/xmedia-bot/src/state.rs` | `ChatStore`: parking_lot `Mutex<HashMap>` cache + SQLite write-through (`chat_state` table) |
| `crates/xmedia-bot/src/link_cache.rs` | `LinkCache`: SQLite-backed cache (`link_cache` table) of successfully sent posts — raw caption fields + Telegram `file_id`s; repeat links re-send locally (no fetch/upload), TTL + prune, invalidated on permanent send failure |
| `crates/xmedia-bot/src/queue.rs` | `PersistentTaskQueue`: SQLite-backed queue (`tasks` table), `QUEUE_WORKERS = 4` concurrent workers (lease via `BEGIN IMMEDIATE` + `locked_until` TTL), retry→dead-letter, `Notify::notify_waiters` wakeup, `busy_timeout` on all connections |
| `crates/xmedia-bot/src/send.rs` | Media senders, upload fallback, error classification, queue task handlers |
## Development Commands
```bash
export TELOXIDE_TOKEN=<token> # required; PIXIV_REFRESH_TOKEN optional (Pixiv disabled without it)
cargo run -p xmedia-bot # run the bot (polling by default)
cargo run -p x-media --example fetch -- <url> # test a link through the fetch library
cargo test --workspace # full test suite (no CI test step exists — run locally)
cargo build --release -p xmedia-bot # release build (Dockerfile does this)
cargo clippy --workspace --all-targets # lint (Clippy is the configured IDE linter)
cargo fmt --check # formatting
```
Docker: `docker build -t tgxmb .` then `docker run --rm -d --name tgxmb --env-file .env -v ./data:/app/data tgxmb`. Runtime requires **ffmpeg** (built into the image).
## Code Conventions & Common Patterns
- **No anyhow/thiserror.** Errors are hand-rolled enums with manual `Display`/`source()`/`From` impls: `QueueError` (`Retryable { delay_seconds, payload }` / `Permanent`), `SendError` (Retryable/Permanent), `FetchError` (`Http`/`Json`/`Pixiv`/`NotFound`/`Blocked`), `PixivError`, `Classification`. New errors should follow this pattern.
- **Global state via `std::sync::LazyLock` statics**, not DI: `CONFIG`, `CHAT_STORE`, `TASK_QUEUE` in `handlers.rs`; shared reqwest `CLIENT` in `x-media/src/site/mod.rs`. `Bot` is passed/cloned into handlers; queue workers rebuild `Bot::from_env()`.
- **Async**: tokio multi-thread runtime (`#[tokio::main]` default). All rusqlite I/O inside `tokio::task::spawn_blocking`. Long loops use `tokio::select!` with `tokio::sync::{watch, Notify}` stop/wake channels. No streams.
- **Blocking sync primitives**: `parking_lot::Mutex` for hot caches, `tokio::sync::Mutex` for async-shared state (pixiv token cache), `AtomicBool` for feature gates.
- **Site adapter convention** (no trait, no enum dispatch — follow the existing convention): each site module exports `PATTERN: LazyLock<Regex>`, `enabled() -> bool`, `fetch_from_url(url) -> Result<Fetched, FetchError>`; `site/mod.rs` re-exports the site struct and `fetch_once` adds one guarded if-branch. Adding a site = new `site/<name>/{mod.rs,interface.rs,model.rs}` + one branch in `fetch_once`.
- **Serde**: per-site `model.rs` are pure `Deserialize` DTOs mirroring API JSON; site structs in `interface.rs` have private fields, a `caption()` builder, and `impl From<SiteStruct> for Fetched`. Persisted payloads use internally-tagged enums (`#[serde(tag = "kind")]` / `type`).
- **Naming**: module-per-concern, snake_case files, `CamelCase` types, `snake_case` fns. `//!` module docs and `///` docs on non-obvious logic (syndication token, ugoira encoding, `display_text_range`).
- **Retries**: only `x-media::site::fetch` retries (3 attempts, `1 << attempt` backoff, HTTP errors only). Queue retries are explicit `QueueError::Retryable` with computed delay (`retry_delay_seconds`).
- Logging via `log` macros (`pretty_env_logger`, level from `RUST_LOG`).
## Important Files
| File | Why it matters |
|---|---|
| `crates/xmedia-bot/src/main.rs` | Startup sequence, webhook vs polling, graceful shutdown (SIGINT via teloxide ctrlc / SIGTERM via `stop_token` for docker, → sweep stop → admin msg → queue stop) |
| `crates/xmedia-bot/src/handlers.rs` | `CHAT_STORE`/`TASK_QUEUE`/`CONFIG` singletons (open `data/task_queue.db` **relative to CWD**); command dispatch; URL extraction; retry enqueue |
| `crates/xmedia-bot/src/send.rs` | Constants `MAX_MEDIA_GROUP = 9`, `MAX_UPLOAD_BYTES = 10 MiB`; fallback chain; `classify_request_error` |
| `crates/x-media/src/site/mod.rs` | Dispatcher, `Fetched`/`FetchError`, shared `CLIENT`, `download_media` (adds `Referer: https://www.pixiv.net/` for `pximg.net` hotlink protection) |
| `crates/x-media/src/site/pixiv/api.rs` | OAuth token exchange (hardcoded app client id/secret), access-token cache, ugoira zip→MP4 via ffmpeg in `spawn_blocking` |
| `Dockerfile` | Multi-stage: cached dep layer via stub sources + `touch *.rs` mtime hack, static ffmpeg from ffmpeg.martin-riedl.de (`FFMPEG_URL` arg, `unzip -t` integrity check), `debian:bookworm-slim` runtime, entrypoint |
| `docker-entrypoint.sh` | Privilege drop: `useradd` with `LOCAL_USER_ID` (default 9001) + `setpriv` (no gosu on bookworm-slim) |
| `docker-compose.yml.example` | Deployment env reference (real `docker-compose.yml` is gitignored). Ships nginx-proxy + acme-companion: webhook mode needs TLS termination in front (teloxide's axum listener is HTTP-only; `WEBHOOK_CERT` only feeds `set_webhook`), bot exposes `VIRTUAL_HOST`/`VIRTUAL_PORT` on the shared `proxy` network, no host port; container names `nginx-proxy`/`acme-companion`/`tgxmb`, start order via `depends_on` (proxy → acme → bot) |
| `.github/workflows/docker.yml` | CI: build+push to Docker Hub on tag `v*`/master; **no test step**; buildx gha cache (`cache-from`/`cache-to`, scope `tgxmb-build`, `mode=max`) so cargo deps + ffmpeg layers are restored across runs |
| `README.md` | Feature docs + command table (Chinese) |
## Runtime/Tooling Preferences
- **Rust, stable, edition 2024**, workspace resolver 3. No `rust-version`/MSRV pin, no `rust-toolchain.toml` — recent stable is assumed. No nightly features.
- Package manager: **Cargo** (workspace with path dep `x-media``xmedia-bot`). No `[workspace.package]`/shared deps — each crate lists deps independently.
- **Two reqwest versions coexist in the lock** (0.12.28 via teloxide, 0.13.3 in x-media) — don't unify casually.
- Config is **environment-variable driven** (dotenv loads `.env`, gitignored; no `.env.example` exists). Key vars: `TELOXIDE_TOKEN` (required), `PIXIV_REFRESH_TOKEN`, `TWITTER_AUTH_TOKEN` (optional; x.com `auth_token` cookie — enables the logged-in GraphQL fallback that fetches NSFW tweets syndication withholds), `BOT_ADMIN` (comma-separated ids), `EDIT_MESSAGE_TTL_SECONDS` (default 86400), `LINK_CACHE_TTL_SECONDS` (default 604800), `WEBHOOK`/`WEBHOOK_URL`/`WEBHOOK_LISTEN`/`WEBHOOK_PORT`/`WEBHOOK_CERT`/`WEBHOOK_SECRET_TOKEN` (webhook mode requires URL/listen/port, `.expect`ed; `WEBHOOK_CERT` is Telegram-facing self-signed validation only — TLS must be terminated by a reverse proxy), `RUST_LOG`, `TELOXIDE_PROXY`, `LOCAL_USER_ID` (entrypoint only).
- SQLite via `rusqlite` with `bundled` feature (no system libsqlite needed). DB file `data/task_queue.db` is CWD-relative — run from the workspace root, or `/app` in Docker. Mount `./data` and `./cert` volumes.
- `.gitattributes` enforces LF for `*.sh` (CRLF breaks shebangs in containers). `.gitignore`: `.env`, `data/`, `cert/`, `docker-compose.yml`, `/target`, `.idea/`.
- Docs are in Chinese; user-facing bot strings too. Keep that convention when editing captions/templates/docs.
## Testing & QA
- **~51 tests, all inline `#[cfg(test)] mod tests`** — no `tests/` integration directories. Framework: built-in Rust test + `#[tokio::test]` (dev-deps only in `x-media`: tokio macros/rt-multi-thread, dotenv).
- No mocking framework anywhere (no mockito/wiremock/mockall). Conventions: pure-function units (regex parsing, serde round-trips, chunking, retry math) tested synchronously; async tests use real dependencies — file-backed SQLite via `tempfile` (`queue.rs::new_queue()` helper), live network fetches.
- Live-network tests exist in `site/twitter/interface.rs` (3), `site/bsky/interface.rs` (2), `site/pixiv/interface.rs`/`api.rs` (env-gated on `PIXIV_REFRESH_TOKEN`/dotenv, skip by early return). Run the full suite with `cargo test --workspace`.
- Fixtures are inline `serde_json::json!` builder fns (`fixture()`, `thread_json()`, `illust_json()`), not files. The shared `CLIENT` sets `pool_max_idle_per_host(0)` under `#[cfg(test)]` to avoid cross-runtime `DispatchGone`.
- **CI runs no tests** — `.github/workflows/docker.yml` only builds/pushes the image; verification is a local responsibility.
- Untested and hard to test without a mock seam: `handlers.rs` (depends directly on teloxide `Bot`); `main.rs`, `config.rs`, `state.rs`; `media.rs`, `lib.rs`, all `model.rs`.
- No coverage tracking, no lint gate in CI.
Generated
+3 -2
View File
@@ -3296,12 +3296,13 @@ checksum = "1ffae5123b2d3fc086436f8834ae3ab053a283cfac8fe0a0b8eaae044768a4c4"
[[package]]
name = "x-media"
version = "1.0.2"
version = "1.0.6"
dependencies = [
"bytes",
"dotenv",
"html-escape",
"log",
"rand 0.8.6",
"regex",
"reqwest 0.13.3",
"serde",
@@ -3314,7 +3315,7 @@ dependencies = [
[[package]]
name = "xmedia-bot"
version = "1.0.2"
version = "1.0.6"
dependencies = [
"dotenv",
"html-escape",
+16 -9
View File
@@ -1,12 +1,16 @@
# ---------- build stage ----------
# rust:1-bookworm (full, not slim) ships the C toolchain needed by
# rusqlite's bundled SQLite, plus wget/xz for the ffmpeg download.
# rusqlite's bundled SQLite, plus wget/unzip for the ffmpeg download.
FROM rust:1-bookworm AS builder
ARG APP_NAME=telegram-twitter-media-bot
# Statically compiled ffmpeg (ugoira MP4 encoding). amd64 by default; override
# for other platforms or pin a different johnvansickle build.
ARG FFMPEG_URL=https://johnvansickle.com/ffmpeg/releases/ffmpeg-7.0.2-amd64-static.tar.xz
# Prebuilt static ffmpeg (glibc-linked, includes libx264) for ugoira MP4
# encoding. Served from https://ffmpeg.martin-riedl.de (Cloudflare CDN,
# built on Debian 12 — glibc-compatible with the bookworm-slim runtime).
# johnvansickle.com throttles datacenter IPs and served garbage from GitHub
# runners. `/redirect/latest/` floats to the newest release build; each build
# also ships a .sha256. Swap `amd64` for `arm64` when building arm64 images.
ARG FFMPEG_URL=https://ffmpeg.martin-riedl.de/redirect/latest/linux/amd64/release/ffmpeg.zip
WORKDIR /build
@@ -22,11 +26,14 @@ RUN mkdir -p crates/x-media/src crates/xmedia-bot/src \
&& cargo build --release -p xmedia-bot
# 2. Static ffmpeg next (cached unless FFMPEG_URL changes), so source edits
# never re-download it. The johnvansickle tarball has a
# `{build}/ffmpeg` layout, so strip one path component.
RUN wget -q -O /tmp/ffmpeg.tar.xz "$FFMPEG_URL" \
&& tar -xJf /tmp/ffmpeg.tar.xz -C /usr/local/bin --strip-components=1 --wildcards '*/ffmpeg' \
&& rm /tmp/ffmpeg.tar.xz \
# never re-download it. The zip contains a single `ffmpeg` binary at the
# root. `unzip -t` verifies the archive before extraction so a bad
# download fails loudly here instead of a cryptic later error.
RUN wget -q -O /tmp/ffmpeg.zip "$FFMPEG_URL" \
&& unzip -tq /tmp/ffmpeg.zip \
&& unzip -q /tmp/ffmpeg.zip -d /usr/local/bin \
&& chmod +x /usr/local/bin/ffmpeg \
&& rm /tmp/ffmpeg.zip \
&& /usr/local/bin/ffmpeg -version >/dev/null
# 3. Real sources last: only our crates recompile on source changes. The
+48 -1
View File
@@ -10,6 +10,7 @@ Telegram 机器人,将 X / Twitter、Pixiv、Bluesky 的帖子链接转换为
- 可绑定转发频道自动转发;支持转发前编辑 caption 与自定义模板
- 发送失败自动重试并持久化,重试耗尽后通知用户
- Pixiv ugoira 动图自动转码为 MP4
- 链接结果本地缓存:成功发送后缓存 Telegram file id 与 caption 等,再次收到相同链接直接本地重发,不再请求源站、不保存媒体文件(`LINK_CACHE_TTL_SECONDS` 控制过期,默认 7 天)
## 快速开始
@@ -28,7 +29,53 @@ docker build -t tgxmb .
docker run --rm -d --name tgxmb --env-file .env -v ./data:/app/data tgxmb
```
环境变量:`TELOXIDE_TOKEN`(必填)、`PIXIV_REFRESH_TOKEN``BOT_ADMIN``EDIT_MESSAGE_TTL_SECONDS``RUST_LOG``WEBHOOK*`
环境变量:`TELOXIDE_TOKEN`(必填)、`PIXIV_REFRESH_TOKEN``BOT_ADMIN``EDIT_MESSAGE_TTL_SECONDS``LINK_CACHE_TTL_SECONDS``RUST_LOG``WEBHOOK*``TWITTER_AUTH_TOKEN`(可选)
NSFW 推文:公开的 syndication 接口不返回敏感内容。设置 `TWITTER_AUTH_TOKEN`(登录 x.com 后浏览器 Cookie 里的 `auth_token` 值)后,bot 会仅在遇到 NSFW 推文时以登录态获取媒体;未设置则提示无媒体。
### Webhook 部署(需要反向代理)
`docker-compose.yml.example` 内置了 [nginx-proxy](https://github.com/nginx-proxy/nginx-proxy) + [acme-companion](https://github.com/nginx-proxy/acme-companion) 反向代理编排,按部署环境二选一:
**有域名**
1. DNS A 记录指向服务器
2. compose 里设 `VIRTUAL_HOST``LETSENCRYPT_HOST` 为域名,`WEBHOOK_URL` 设为 `https://域名/`
3. 证书自动签发与续期,无需手动处理
**只有 IP**
1. 生成自签证书(PEM 格式,见第 3 步):
`openssl req -x509 -newkey rsa:2048 -nodes -days 365 -keyout nginx-certs/default.key -out nginx-certs/default.crt`
2. compose 里 nginx-proxy 设 `DEFAULT_HOST`bot 设 `WEBHOOK_CERT: './cert/cert.pem'`(须与代理所服务的为同一张证书)
3. 证书必须是 PEM 编码(ASCII BASE64,以 `-----BEGIN CERTIFICATE-----` 开头)—— Telegram 只接受该格式;若现有证书是 DER 二进制,转换:
`openssl x509 -in cert.der -inform DER -out cert.pem -outform PEM`
(私钥同理:`openssl rsa -in key.der -inform DER -out key.pem -outform PEM`
Telegram 只接受 443/80/88/8443 端口。
<details>
<summary>环境变量说明</summary>
| 变量 | 说明 |
|---|---|
| `TELOXIDE_TOKEN` | Bot token(必填) |
| `PIXIV_REFRESH_TOKEN` | Pixiv 刷新令牌;未设置则禁用 Pixiv |
| `BOT_ADMIN` | 管理员聊天 ID,逗号分隔;接收启动/停止通知 |
| `EDIT_MESSAGE_TTL_SECONDS` | 转发前编辑记录过期秒数,默认 86400 |
| `LINK_CACHE_TTL_SECONDS` | 链接结果缓存过期秒数,默认 604800(7 天) |
| `RUST_LOG` | 日志级别 |
| `LOCAL_USER_ID` | 容器内运行用户 UID,默认 9001 |
| `VIRTUAL_HOST` | 对外域名或 IPnginx-proxy 按此路由 |
| `VIRTUAL_PORT` | bot 容器内监听端口,nginx-proxy 的转发目标 |
| `LETSENCRYPT_HOST` | 设为域名时由 acme-companion 自动签发/续期证书 |
| `DEFAULT_HOST` | nginx-proxy 将未知 Host 的请求路由到该 vhost(IP 访问时需要) |
| `DEFAULT_EMAIL` | acme-companion 证书通知邮箱 |
| `WEBHOOK` | `true` 启用 webhook 模式(默认轮询) |
| `WEBHOOK_LISTEN` / `WEBHOOK_PORT` | bot 容器内监听地址/端口 |
| `WEBHOOK_URL` | 对外公网 HTTPS 地址(`https://域名/` |
| `WEBHOOK_CERT` | 自签证书路径(仅 IP 路径需要,须为 PEM 且与代理所服务的一致) |
| `WEBHOOK_SECRET_TOKEN` | 更新校验令牌(`X-Telegram-Bot-Api-Secret-Token` |
</details>
## 命令
+2 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "x-media"
version = "1.0.2"
version = "1.0.6"
edition = "2024"
[dependencies]
@@ -13,6 +13,7 @@ url = "2.5.2"
bytes = "1"
zip = "2"
tempfile = "3"
rand = "0.8"
log = "0.4"
tokio = { version = "1.40", features = ["time"] }
+18
View File
@@ -14,6 +14,24 @@ impl Media {
Media::Animated { thumbnail_url, .. } => Some(thumbnail_url),
}
}
/// A smaller variant of this media's file (used as the fallback when the
/// primary URL or upload exceeds Telegram's size limits). None when no
/// smaller variant exists (videos, animated gifs).
pub fn smaller_url(&self) -> Option<&str> {
match self {
Media::Illustration {
url,
fallback_url,
thumbnail_url,
..
} => fallback_url
.as_deref()
.or(thumbnail_url.as_deref())
.filter(|smaller| *smaller != url),
Media::Video { .. } | Media::Animated { .. } => None,
}
}
}
#[derive(Debug)]
+131 -10
View File
@@ -68,18 +68,73 @@ impl Fetched {
/// caption.
pub fn caption_with(&self, format: &str) -> String {
match (&self.render_data, format.is_empty()) {
(Some(data), false) => {
let escaped = html_escape::encode_text(format).into_owned();
escaped
.replace("{url}", &data.url)
.replace("{author}", &data.author)
.replace("{author_url}", &data.author_url)
.replace("{title}", &data.title)
.replace("{tags}", &data.tags)
}
(Some(data), false) => caption_from_fields(
format,
"",
&data.url,
&data.author,
&data.author_url,
&data.title,
&data.tags,
),
_ => self.caption.clone(),
}
}
/// The pre-escaped placeholder values (author, author_url, title, tags)
/// a caller needs to rebuild a caption later, e.g. for a cached post
/// where the [`Fetched`] is no longer available.
pub fn render_fields(&self) -> Option<(&str, &str, &str, &str)> {
self.render_data.as_ref().map(|d| {
(
d.author.as_str(),
d.author_url.as_str(),
d.title.as_str(),
d.tags.as_str(),
)
})
}
}
/// Renders a user-supplied caption format from raw (already-escaped) field
/// values with the same escaping/substitution rules as
/// [`Fetched::caption_with`]. An empty format returns `built_in` unchanged.
pub fn caption_from_fields(
format: &str,
built_in: &str,
url: &str,
author: &str,
author_url: &str,
title: &str,
tags: &str,
) -> String {
if format.is_empty() {
return built_in.to_string();
}
let escaped = html_escape::encode_text(format).into_owned();
escaped
.replace("{url}", url)
.replace("{author}", author)
.replace("{author_url}", author_url)
.replace("{title}", title)
.replace("{tags}", tags)
}
/// Stable per-post cache key derived from any supported URL, so variant
/// domains (x.com / twitter.com / fxtwitter.com, mobile, `/photo/N`
/// suffixes) map to the same post. Returns `"twitter:<id>"`,
/// `"pixiv:<id>"` or `"bsky:<handle>/<rkey>"`.
pub fn cache_key(url: &str) -> Option<String> {
if let Some(caps) = twitter::PATTERN.captures(url) {
return Some(format!("twitter:{}", &caps[1]));
}
if let Some(caps) = pixiv::PATTERN.captures(url) {
return Some(format!("pixiv:{}", &caps[1]));
}
if let Some(caps) = bsky::PATTERN.captures(url) {
return Some(format!("bsky:{}/{}", &caps[1], &caps[2]));
}
None
}
#[derive(Debug)]
@@ -89,6 +144,9 @@ pub enum FetchError {
Pixiv(PixivError),
NotFound,
Blocked,
/// The post exists but its content is withheld (twitter NSFW /
/// age-restricted tweets come back as an empty `{}` from syndication).
Sensitive,
}
impl fmt::Display for FetchError {
@@ -99,6 +157,7 @@ impl fmt::Display for FetchError {
FetchError::Pixiv(e) => write!(f, "pixiv error: {e}"),
FetchError::NotFound => write!(f, "not found"),
FetchError::Blocked => write!(f, "blocked"),
FetchError::Sensitive => write!(f, "content withheld (sensitive)"),
}
}
}
@@ -109,7 +168,7 @@ impl std::error::Error for FetchError {
FetchError::Http(e) => Some(e),
FetchError::Json(e) => Some(e),
FetchError::Pixiv(e) => Some(e),
FetchError::NotFound | FetchError::Blocked => None,
FetchError::NotFound | FetchError::Blocked | FetchError::Sensitive => None,
}
}
}
@@ -195,6 +254,19 @@ async fn fetch_once(url: &str) -> Result<Option<Fetched>, FetchError> {
/// fetch of a media URL is blocked (hotlink protection), the bot downloads
/// the file itself and uploads it via multipart. Site-appropriate headers:
/// pixiv image hosts need the `Referer` header.
/// Returns the Content-Length of a media URL, or `None` when the server does
/// not report one. Used to check whether a file fits Telegram's size limits
/// before downloading/uploading it.
pub async fn media_size(url: &str) -> Result<Option<u64>, FetchError> {
let mut request = CLIENT.get(url);
let lower = url.to_ascii_lowercase();
if lower.contains("pximg.net") {
request = request.header("Referer", "https://www.pixiv.net/");
}
let response = request.send().await?;
Ok(response.content_length())
}
pub async fn download_media(url: &str) -> Result<bytes::Bytes, FetchError> {
let mut request = CLIENT.get(url);
let lower = url.to_ascii_lowercase();
@@ -209,6 +281,55 @@ pub async fn download_media(url: &str) -> Result<bytes::Bytes, FetchError> {
mod tests {
use super::*;
#[test]
fn cache_key_normalizes_domain_variants() {
assert_eq!(
cache_key("https://x.com/user/status/1234567890/photo/1"),
Some("twitter:1234567890".into())
);
assert_eq!(
cache_key("https://mobile.twitter.com/user/status/1234567890"),
Some("twitter:1234567890".into())
);
assert_eq!(
cache_key("https://fxtwitter.com/user/status/1234567890"),
Some("twitter:1234567890".into())
);
assert_eq!(
cache_key("https://www.pixiv.net/artworks/123456"),
Some("pixiv:123456".into())
);
assert_eq!(
cache_key("https://bsky.app/profile/handle.example/post/3lorem"),
Some("bsky:handle.example/3lorem".into())
);
assert_eq!(cache_key("https://example.com/not-a-post"), None);
}
#[test]
fn caption_from_fields_substitutes_and_escapes() {
// The format string is escaped, the field values are substituted
// verbatim (callers pass the already-escaped render data).
let out = caption_from_fields(
"see {author} at {url} — {title}",
"",
"https://x.com/u/status/1",
"A &amp; B",
"https://x.com/u",
"hello <world>",
"",
);
assert_eq!(
out,
"see A &amp; B at https://x.com/u/status/1 — hello <world>"
);
// Empty format keeps the built-in caption untouched.
assert_eq!(
caption_from_fields("", "built-in", "u", "a", "au", "t", "g"),
"built-in"
);
}
#[tokio::test]
async fn unsupported_url_returns_none() {
let result = fetch("https://example.com/some/article").await;
+376
View File
@@ -0,0 +1,376 @@
//! Authenticated fallback for tweets the public syndication endpoint refuses
//! to serve (NSFW / age-restricted tweets come back as an empty `{}`).
//!
//! Mirrors nazurin's web API client ([`web.py`]) and is used *only* when
//! syndication reports [`FetchError::Sensitive`]: the private GraphQL
//! `TweetDetail` endpoint, authenticated with a browser session cookie from
//! `TWITTER_AUTH_TOKEN` (the `auth_token` cookie value of a logged-in x.com
//! session). A fresh random `ct0` is generated per call; X checks that the
//! `x-csrf-token` header matches the cookie, not that it issued the value.
//!
//! [`web.py`]: https://github.com/y-young/nazurin/blob/master/nazurin/sites/twitter/api/web.py
//!
//! # Caveats
//! - X rotates the GraphQL query id when it rolls the web app; if requests
//! start failing, update [`TWEET_DETAIL_QUERY_ID`]. Fresh references from
//! the actively maintained FxEmbed/FxEmbed: TweetDetail
//! `R9IzzyzQBV87-DOWpcvDmw`, TweetResultByRestId `f2sagi1jweVHFkTUIHzmMQ`
//! (the latter is anonymous and surfaces NSFW tweets as
//! `reason: NsfwLoggedOut`).
//! - `x-client-transaction-id` is only required for `SearchTimeline`
//! (verified against FxEmbed's `proxy/allowlist.ts`) — TweetDetail works
//! without it; no need for the nazurin home-page/JS-bundle derivation.
use std::sync::LazyLock;
use serde_json::{json, Value};
use crate::site::FetchError;
use super::interface::Tweet;
/// `auth_token` cookie of a logged-in x.com session; enables the fallback.
/// Trimmed: a CRLF `.env` (Windows) leaves a trailing `\r` on the value,
/// which would make the Cookie header invalid.
static AUTH_TOKEN: LazyLock<Option<String>> = LazyLock::new(|| {
std::env::var("TWITTER_AUTH_TOKEN")
.ok()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
});
/// Public "logged in" client token used by the x.com web app.
const LOGGED_IN_BEARER: &str =
"Bearer AAAAAAAAAAAAAAAAAAAAANRILgAAAAAAnNwIzUejRCOuH5E6I8xnZz4puTs%3D1Zv7ttfk8LF81IUq16cHjhLTvJu4FA33AGWWjCpTnA";
/// `TweetDetail` query id (from nazurin; still valid as of 2026-08,
/// corroborated by the current FxEmbed build — see module caveats).
const TWEET_DETAIL_QUERY_ID: &str = "_8aYOgEDz35BrBcBal1-_w";
fn variables(id: &str) -> Value {
json!({
"focalTweetId": id,
"with_rux_injections": false,
"includePromotedContent": false,
"withCommunity": true,
"withQuickPromoteEligibilityTweetFields": false,
"withBirdwatchNotes": false,
"withVoice": true,
})
}
fn features() -> Value {
json!({
"rweb_video_screen_enabled": false,
"profile_label_improvements_pcf_label_in_post_enabled": true,
"rweb_tipjar_consumption_enabled": true,
"verified_phone_label_enabled": false,
"creator_subscriptions_tweet_preview_api_enabled": true,
"responsive_web_graphql_timeline_navigation_enabled": true,
"responsive_web_graphql_skip_user_profile_image_extensions_enabled": false,
"premium_content_api_read_enabled": false,
"communities_web_enable_tweet_community_results_fetch": true,
"c9s_tweet_anatomy_moderator_badge_enabled": true,
"responsive_web_grok_analyze_button_fetch_trends_enabled": false,
"responsive_web_grok_analyze_post_followups_enabled": true,
"responsive_web_jetfuel_frame": false,
"responsive_web_grok_share_attachment_enabled": true,
"articles_preview_enabled": true,
"responsive_web_edit_tweet_api_enabled": true,
"graphql_is_translatable_rweb_tweet_is_translatable_enabled": true,
"view_counts_everywhere_api_enabled": true,
"longform_notetweets_consumption_enabled": true,
"responsive_web_twitter_article_tweet_consumption_enabled": true,
"tweet_awards_web_tipping_enabled": false,
"responsive_web_grok_show_grok_translated_post": false,
"responsive_web_grok_analysis_button_from_backend": true,
"creator_subscriptions_quote_tweet_preview_enabled": false,
"freedom_of_speech_not_reach_fetch_enabled": true,
"standardized_nudges_misinfo": true,
"tweet_with_visibility_results_prefer_gql_limited_actions_policy_enabled": true,
"longform_notetweets_rich_text_read_enabled": true,
"longform_notetweets_inline_media_enabled": true,
"responsive_web_grok_image_annotation_enabled": true,
"responsive_web_enhance_cards_enabled": false,
})
}
/// Whether the authenticated fallback is available.
pub fn enabled() -> bool {
AUTH_TOKEN.is_some()
}
/// Fetches a tweet as the logged-in user via the private GraphQL API.
/// Returns the syndication-shaped [`Tweet`] (media included for NSFW posts).
pub async fn fetch(id: &str) -> Result<Tweet, FetchError> {
let token = AUTH_TOKEN
.as_deref()
.ok_or(FetchError::Sensitive)?;
// 16 random bytes as 32 hex chars: X rejects ct0 values of any other
// length with 403 code 353 ("matching csrf cookie and header").
let ct0: String = (0..16)
.map(|_| format!("{:02x}", rand::random::<u8>()))
.collect();
let response = crate::site::CLIENT
.get(format!(
"https://x.com/i/api/graphql/{TWEET_DETAIL_QUERY_ID}/TweetDetail"
))
.query(&[
("variables", variables(id).to_string()),
("features", features().to_string()),
])
.header("authorization", LOGGED_IN_BEARER)
.header("x-csrf-token", &ct0)
.header("x-twitter-auth-type", "OAuth2Session")
.header("cookie", format!("auth_token={token}; ct0={ct0}"))
.header("x-twitter-client-language", "en")
.header("x-twitter-active-user", "yes")
.header("referer", "https://x.com/")
.send()
.await?;
if !response.status().is_success() {
log::warn!("twitter auth fetch {id}: HTTP {}", response.status());
return Err(FetchError::NotFound);
}
let text = response.text().await?;
let json: Value = serde_json::from_str(&text)?;
let result = parse_tweet_result(&json, id)?;
let syndication_shape = to_syndication_shape(&result)
.ok_or_else(|| FetchError::Json(serde_json::Error::io(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"missing tweet fields in GraphQL response",
))))?;
Tweet::from_syndication_json(&syndication_shape.to_string()).map_err(FetchError::Json)
}
/// Locates the tweet for `id` in a `TweetDetail` response and unwraps
/// visibility wrappers / retweets, mirroring nazurin's `_process_response`.
fn parse_tweet_result(json: &Value, id: &str) -> Result<Value, FetchError> {
if let Some(errors) = json.get("errors").and_then(|e| e.as_array()) {
let messages: Vec<&str> = errors
.iter()
.filter_map(|e| e.get("message").and_then(|m| m.as_str()))
.collect();
log::warn!("twitter auth fetch {id} failed: {}", messages.join("; "));
return Err(FetchError::NotFound);
}
let instructions = json
.pointer("/data/threaded_conversation_with_injections_v2/instructions")
.and_then(|v| v.as_array())
.ok_or(FetchError::NotFound)?;
for instruction in instructions {
if instruction.get("type").and_then(|t| t.as_str()) != Some("TimelineAddEntries") {
continue;
}
let entries = instruction
.get("entries")
.and_then(|e| e.as_array())
.ok_or(FetchError::NotFound)?;
let wanted = format!("tweet-{id}");
for entry in entries {
if entry.get("entryId").and_then(|i| i.as_str()) == Some(wanted.as_str()) {
let result = entry
.pointer("/content/itemContent/tweet_results/result")
.ok_or(FetchError::NotFound)?;
return normalize_tweet_result(result);
}
}
}
Err(FetchError::NotFound)
}
/// Unwraps TweetTombstone/TweetUnavailable errors, the
/// TweetWithVisibilityResults wrapper and retweets, returning the
/// `{core, legacy, ...}` tweet object.
fn normalize_tweet_result(result: &Value) -> Result<Value, FetchError> {
match result.get("__typename").and_then(|t| t.as_str()) {
Some("TweetTombstone") => {
let text = result
.pointer("/tombstone/text/text")
.and_then(|t| t.as_str())
.unwrap_or("tweet is unavailable");
log::warn!("twitter auth fetch: tombstone: {text}");
return Err(FetchError::NotFound);
}
Some("TweetUnavailable") => {
let reason = result
.get("reason")
.and_then(|r| r.as_str())
.unwrap_or("unknown");
log::warn!("twitter auth fetch: tweet unavailable: {reason}");
return Err(FetchError::NotFound);
}
_ => {}
}
// TweetWithVisibilityResults (e.g. limited replies) nests the real tweet.
let tweet = result.get("tweet").unwrap_or(result);
// A retweet's media lives on the original tweet.
if let Some(original) = tweet.pointer("/legacy/retweeted_status_result/result") {
return Ok(original.clone());
}
Ok(tweet.clone())
}
/// Maps a GraphQL `{core, legacy, ...}` tweet onto the syndication JSON
/// shape [`Tweet::from_syndication_json`] parses, so the existing text /
/// media handling (t.co expansion, `name=orig`, mp4 variant) is reused.
fn to_syndication_shape(tweet: &Value) -> Option<Value> {
let legacy = tweet.get("legacy")?;
let user = tweet.pointer("/core/user_results/result/legacy")?;
Some(json!({
"id_str": legacy.get("id_str"),
"text": legacy.get("full_text"),
"user": {
"name": user.get("name"),
"screen_name": user.get("screen_name"),
},
"possibly_sensitive": legacy.get("possibly_sensitive"),
"display_text_range": legacy.get("display_text_range"),
"entities": legacy.get("entities"),
"mediaDetails": legacy.pointer("/extended_entities/media"),
}))
}
#[cfg(test)]
mod tests {
use super::*;
fn tweet_result() -> Value {
json!({
"__typename": "Tweet",
"core": {
"user_results": {
"result": {
"legacy": { "name": "Display Name", "screen_name": "nsfw_author" }
}
}
},
"legacy": {
"id_str": "2083868672721039569",
"full_text": "nsfw content https://t.co/abc123",
"display_text_range": [0, 12],
"possibly_sensitive": true,
"entities": {
"urls": [
{ "url": "https://t.co/abc123", "expanded_url": "https://example.com/x" }
]
},
"extended_entities": {
"media": [
{
"type": "photo",
"media_url_https": "https://pbs.twimg.com/media/nsfw.jpg",
"original_info": { "width": 1200, "height": 800 }
},
{
"type": "video",
"media_url_https": "https://pbs.twimg.com/thumb.jpg",
"video_info": {
"variants": [
{ "content_type": "application/x-mpegURL", "url": "https://x.com/pl.m3u8" },
{ "content_type": "video/mp4", "url": "https://video.twimg.com/nsfw.mp4" }
]
}
}
]
}
}
})
}
fn conversation(tweet: Value) -> Value {
json!({
"data": {
"threaded_conversation_with_injections_v2": {
"instructions": [
{ "type": "TimelineAddEntries", "entries": [
{ "entryId": "tweet-2083868672721039569",
"content": { "itemContent": { "tweet_results": { "result": tweet } } } }
]}
]
}
}
})
}
#[test]
fn parses_graphql_tweet_into_fetched() {
let json = conversation(tweet_result());
let result = parse_tweet_result(&json, "2083868672721039569").unwrap();
let shape = to_syndication_shape(&result).unwrap();
let tweet = Tweet::from_syndication_json(&shape.to_string()).unwrap();
let fetched: crate::site::Fetched = tweet.into();
assert!(fetched.sensitive);
assert_eq!(fetched.media.len(), 2);
match &fetched.media[0] {
crate::media::Media::Illustration { url, .. } => {
assert_eq!(url, "https://pbs.twimg.com/media/nsfw.jpg?name=orig");
}
other => panic!("expected illustration, got {other:?}"),
}
match &fetched.media[1] {
crate::media::Media::Video { url, .. } => {
assert_eq!(url, "https://video.twimg.com/nsfw.mp4");
}
other => panic!("expected video, got {other:?}"),
}
assert_eq!(fetched.source_url, "https://x.com/nsfw_author/status/2083868672721039569");
// display_text_range cuts the trailing t.co link.
assert_eq!(fetched.title, "nsfw content");
}
#[test]
fn unwraps_retweet_to_original() {
let original = tweet_result();
let mut rt = tweet_result();
rt["legacy"]["retweeted_status_result"] = json!({ "result": original });
let json = conversation(rt);
let result = parse_tweet_result(&json, "2083868672721039569").unwrap();
assert!(result.pointer("/legacy/retweeted_status_result").is_none());
assert_eq!(result.pointer("/legacy/id_str").unwrap(), "2083868672721039569");
}
#[test]
fn error_response_maps_to_not_found() {
let json = json!({ "errors": [{ "message": "NsfwLoggedOut" }] });
assert!(matches!(
parse_tweet_result(&json, "1"),
Err(FetchError::NotFound)
));
}
#[test]
fn missing_entry_maps_to_not_found() {
let json = conversation(json!({ "__typename": "Tweet" }));
assert!(matches!(
parse_tweet_result(&json, "999"),
Err(FetchError::NotFound)
));
}
#[test]
fn tombstone_maps_to_not_found() {
let tombstone = json!({
"__typename": "TweetTombstone",
"tombstone": { "text": { "text": "Age-restricted adult content" } }
});
let json = conversation(tombstone);
assert!(matches!(
parse_tweet_result(&json, "2083868672721039569"),
Err(FetchError::NotFound)
));
}
#[test]
fn visibility_wrapper_unwraps() {
let inner = tweet_result();
let wrapped = json!({ "__typename": "TweetWithVisibilityResults", "tweet": inner });
let json = conversation(wrapped);
let result = parse_tweet_result(&json, "2083868672721039569").unwrap();
assert_eq!(result.get("__typename").unwrap(), "Tweet");
}
}
+59 -15
View File
@@ -19,7 +19,44 @@ pub async fn fetch_from_url(url: &str) -> Result<Fetched, FetchError> {
.and_then(|caps| caps.get(1))
.map(|m| m.as_str())
.ok_or(FetchError::NotFound)?;
Ok(fetch(id).await?.into())
match fetch(id).await {
Ok(tweet) => Ok(tweet.into()),
// Syndication withholds NSFW/age-restricted tweets (empty `{}`).
// Retry as the logged-in user when TWITTER_AUTH_TOKEN is set;
// otherwise degrade to an empty result (the bot replies
// "No media found").
Err(FetchError::Sensitive) => {
if super::auth::enabled() {
match super::auth::fetch(id).await {
Ok(tweet) => Ok(tweet.into()),
Err(e) => {
log::warn!("twitter auth fallback failed for {id}: {e}");
Ok(empty_fetched(url))
}
}
} else {
log::info!(
"tweet {id} is sensitive; set TWITTER_AUTH_TOKEN to fetch NSFW media"
);
Ok(empty_fetched(url))
}
}
Err(e) => Err(e),
}
}
/// A Fetched with no media for withheld tweets: the bot replies
/// "No media found" and moves on instead of erroring.
fn empty_fetched(url: &str) -> Fetched {
Fetched {
source_url: url.to_string(),
caption: url.to_string(),
title: String::new(),
media: vec![],
sensitive: true,
render_data: None,
_keep_alive: None,
}
}
/// Fetches a tweet from the syndication endpoint. Deleted/blocked tweets
@@ -44,7 +81,16 @@ pub async fn fetch(id: &str) -> Result<Tweet, FetchError> {
{
return Err(FetchError::NotFound);
}
Ok(Tweet::from_syndication_json(&text).map_err(FetchError::Json)?)
// NSFW / age-restricted tweets exist but are served as an empty `{}` —
// they surface as FetchError::Sensitive so the caller can retry as a
// logged-in user.
if serde_json::from_str::<serde_json::Value>(&text)
.map(|v| v.get("id_str").is_none())
.unwrap_or(false)
{
return Err(FetchError::Sensitive);
}
Tweet::from_syndication_json(&text).map_err(FetchError::Json)
}
/// The syndication token: JS `((id / 1e15) * PI).toString(36)` (the
@@ -129,7 +175,9 @@ impl Tweet {
title: None,
url: original_twimg_url(&item.media_url_https),
thumbnail_url: None,
fallback_url: None,
// The param-less base URL is a reduced-size variant;
// used as the fallback when the original is too large.
fallback_url: Some(item.media_url_https.clone()),
}),
"video" => media.push(Media::Video {
title: None,
@@ -157,21 +205,17 @@ impl Tweet {
}
/// The raw syndication `text` ends with the appended media short link
/// (" https://t.co/wmI8McgXul"). `display_text_range` (UTF-16 indices) marks
/// the visible text; a regex strips any remaining trailing t.co link when the
/// range is absent or a tweet ends in a URL short link.
/// (" https://t.co/wmI8McgXul"). `display_text_range` marks the visible text;
/// a regex strips any remaining trailing t.co link when the range is absent
/// or a tweet ends in a URL short link.
///
/// X reports these indices in Unicode **code points**, not UTF-16 units
/// (verified against GraphQL responses containing emoji: cutting an emoji
/// tweet by UTF-16 units silently drops the character after the emoji).
fn strip_trailing_short_links(text: &str, display_text_range: Option<[usize; 2]>) -> String {
let mut out = match display_text_range {
Some([start, end]) if start < end => {
let units: Vec<u16> = text
.encode_utf16()
.skip(start)
.take(end - start)
.collect();
// Drop the replacement char that a surrogate cut at the boundary
// would produce (the range end is a valid UTF-16 boundary in
// practice, so this is just a safety net).
String::from_utf16_lossy(&units).replace('\u{FFFD}', "")
text.chars().skip(start).take(end - start).collect()
}
_ => text.to_string(),
};
+1
View File
@@ -1,3 +1,4 @@
mod auth;
mod interface;
mod model;
+1 -1
View File
@@ -10,7 +10,7 @@ pub struct SyndicationTweet {
#[serde(default)]
pub possibly_sensitive: Option<bool>,
/// Visible-text span; the raw `text` field has the appended media short
/// link after it. Indices are UTF-16 code units.
/// link after it. Indices are Unicode code points (not UTF-16 units).
#[serde(default, rename = "display_text_range")]
pub display_text_range: Option<[usize; 2]>,
#[serde(default)]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "xmedia-bot"
version = "1.0.2"
version = "1.0.6"
edition = "2024"
[dependencies]
+17 -2
View File
@@ -10,6 +10,8 @@ pub struct Config {
pub admin_ids: Vec<i64>,
/// EDIT_MESSAGE_TTL_SECONDS, default 86400 (24h).
pub edit_message_ttl: Duration,
/// LINK_CACHE_TTL_SECONDS, default 604800 (7 days).
pub link_cache_ttl: Duration,
// Webhook settings (moved out of main; names/defaults unchanged).
pub webhook_enabled: bool,
pub webhook_url: Option<url::Url>,
@@ -36,17 +38,30 @@ impl Config {
.map(Duration::from_secs)
.unwrap_or(Duration::from_secs(86400));
let link_cache_ttl = env::var("LINK_CACHE_TTL_SECONDS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.map(Duration::from_secs)
.unwrap_or(Duration::from_secs(7 * 24 * 3600));
let webhook_enabled = env::var("WEBHOOK")
.is_ok_and(|v| matches!(v.to_lowercase().as_str(), "true" | "yes" | "1"));
let webhook_url = env::var("WEBHOOK_URL").ok().and_then(|s| s.parse().ok());
let webhook_listen = env::var("WEBHOOK_LISTEN").ok().and_then(|s| s.parse().ok());
let webhook_port = env::var("WEBHOOK_PORT").ok().and_then(|s| s.parse().ok());
let webhook_cert = env::var("WEBHOOK_CERT").ok();
let webhook_secret_token = env::var("WEBHOOK_SECRET_TOKEN").ok();
// Empty strings count as unset (e.g. `-e WEBHOOK_CERT=` to disable a
// value that would otherwise come from `.env`).
let webhook_cert = env::var("WEBHOOK_CERT")
.ok()
.filter(|s| !s.is_empty());
let webhook_secret_token = env::var("WEBHOOK_SECRET_TOKEN")
.ok()
.filter(|s| !s.is_empty());
Config {
admin_ids,
edit_message_ttl,
link_cache_ttl,
webhook_enabled,
webhook_url,
webhook_listen,
+204 -66
View File
@@ -1,10 +1,12 @@
use crate::config::Config;
use crate::link_cache::{CachedMediaKind, CachedPost, LinkCache};
use crate::queue::PersistentTaskQueue;
use crate::send::{self, MediaItemPayload, Task};
use crate::state::{ChatStore, unix_now};
use crate::state::{ChatData, ChatStore, unix_now};
use std::collections::HashSet;
use std::sync::LazyLock;
use teloxide::prelude::*;
use tokio::sync::Semaphore;
use teloxide::types::{
CallbackQuery, ChatAction, ChatId, ChatKind, InlineQuery, InlineQueryResult,
InlineQueryResultMpeg4Gif, InlineQueryResultPhoto, InlineQueryResultVideo, Message,
@@ -19,8 +21,18 @@ pub static CHAT_STORE: LazyLock<ChatStore> = LazyLock::new(|| {
});
pub static TASK_QUEUE: LazyLock<PersistentTaskQueue> =
LazyLock::new(|| PersistentTaskQueue::new("data/task_queue.db"));
pub static LINK_CACHE: LazyLock<LinkCache> =
LazyLock::new(|| LinkCache::open("data/task_queue.db"));
pub static CONFIG: LazyLock<Config> = LazyLock::new(Config::load);
/// Cap on concurrent per-URL processing. teloxide dispatches updates to a
/// per-chat worker that handles them sequentially, so a batch-forward of many
/// messages would otherwise be processed one at a time (fetch + send each,
/// roughly a second per message). Moving the work into spawned tasks trades
/// per-chat reply ordering for throughput; the semaphore bounds how many run
/// at once so a big burst cannot hammer Telegram's rate limits.
static URL_TASKS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(8));
#[derive(BotCommands, Clone)]
#[command(rename_rule = "snake_case", description = "")]
enum Command {
@@ -322,22 +334,29 @@ fn thumbnail_for(media: &Media) -> Option<String> {
}
fn media_to_payload(media: &Media, sensitive: bool) -> MediaItemPayload {
let fallback_url = media.smaller_url().map(str::to_string);
match media {
// A gif inside a group becomes a video item; a lone gif takes the
// animation path (see url_media).
Media::Illustration { .. } => MediaItemPayload::Photo {
media: media.url().to_string(),
has_spoiler: sensitive,
fallback_url,
file_id: false,
},
Media::Video { .. } => MediaItemPayload::Video {
media: media.url().to_string(),
has_spoiler: sensitive,
thumbnail: thumbnail_for(media),
fallback_url,
file_id: false,
},
Media::Animated { .. } => MediaItemPayload::Video {
media: media.url().to_string(),
has_spoiler: sensitive,
thumbnail: thumbnail_for(media),
fallback_url,
file_id: false,
},
}
}
@@ -350,11 +369,148 @@ async fn enqueue_retry(task: Task, delay_seconds: f64) {
}
}
/// Sends a task and handles the outcome: post-send actions on success, retry
/// enqueue on retryable failure, reply + link-cache invalidation on
/// permanent failure (a stale cached file id must not repeat forever).
async fn dispatch_send(bot: Bot, message: &Message, task: &Task, url: &str) {
let result = match task {
Task::SendAnimation { .. } => send::send_animation(&bot, task).await,
Task::SendMediaSequence { .. } => send::send_media_sequence(&bot, task).await,
Task::ForwardMessages { .. } => unreachable!(),
};
match result {
Ok(message_ids) => {
log::info!("sent {} message(s) for {url}", message_ids.len());
send::post_send_actions(&bot, task, message_ids).await;
}
Err(send::SendError::Retryable { delay_seconds, task }) => {
log::info!("send for {url} failed, queued for retry in {delay_seconds:.1}s");
enqueue_retry(task, delay_seconds).await;
let _ = reply(bot, message.clone(), "Send failed. Task queued for retry.").await;
}
Err(send::SendError::Permanent {
message: err_message,
task,
}) => {
send::invalidate_cache(&task).await;
log::error!("send for {url} failed permanently: {err_message}");
let _ = reply(bot, message.clone(), format!("Send failed: {err_message}")).await;
}
}
}
/// Builds the send task from ready-made items, sharing the payload shape
/// between the fresh-fetch and link-cache paths.
#[allow(clippy::too_many_arguments)]
fn build_send_task(
chat_data: &ChatData,
message: &Message,
source_url: String,
caption: String,
items: Vec<MediaItemPayload>,
cache_data: Option<CachedPost>,
) -> Task {
let chat_id = message.chat.id.0;
if items.len() == 1 && matches!(items[0], MediaItemPayload::Animation { .. }) {
Task::SendAnimation {
chat_id,
reply_to_message_id: message.id.0 as i64,
caption,
animation: items.into_iter().next().unwrap(),
source_url,
edit_before_forward: chat_data.edit_before_forward,
forward_channel_id: chat_data.forward_channel_id,
notify_chat_id: Some(chat_id),
notify_message_id: Some(message.id.0 as i64),
cache_data,
}
} else {
Task::SendMediaSequence {
chat_id,
reply_to_message_id: message.id.0 as i64,
caption,
media_batches: send::chunk_media_items(items),
batch_index: 0,
sent_message_ids: vec![],
source_url,
edit_before_forward: chat_data.edit_before_forward,
forward_channel_id: chat_data.forward_channel_id,
notify_chat_id: Some(chat_id),
notify_message_id: Some(message.id.0 as i64),
cache_data,
}
}
}
async fn url_media(bot: Bot, message: &Message, url: &str) {
let chat_id = message.chat.id.0;
if let Err(e) = bot.send_chat_action(ChatId(chat_id), ChatAction::Typing).await {
log::error!("send_chat_action failed: {e}");
}
// Link cache: a post sent before is re-sent from Telegram file ids —
// no source-site request, no download, no upload. Keyed by the
// normalized post id so x.com / fxtwitter / /photo/N variants collide.
if let Some(key) = x_media::site::cache_key(url)
&& let Some(cached) = LINK_CACHE.get(&key, CONFIG.link_cache_ttl).await
{
log::info!("link cache hit for {url}");
let chat_data = CHAT_STORE.get(chat_id).await;
let site = key.split(':').next().unwrap_or("unknown");
let format = chat_data
.message_format
.get(site)
.cloned()
.unwrap_or_default();
let caption = if format.is_empty() {
cached.caption.clone()
} else {
x_media::site::caption_from_fields(
&format,
"",
&cached.url,
&cached.author,
&cached.author_url,
&cached.title,
&cached.tags,
)
};
let items: Vec<MediaItemPayload> = cached
.media
.iter()
.map(|m| match m.kind {
CachedMediaKind::Photo => MediaItemPayload::Photo {
media: m.file_id.clone(),
has_spoiler: cached.sensitive,
fallback_url: None,
file_id: true,
},
CachedMediaKind::Video => MediaItemPayload::Video {
media: m.file_id.clone(),
has_spoiler: cached.sensitive,
thumbnail: None,
fallback_url: None,
file_id: true,
},
CachedMediaKind::Animation => MediaItemPayload::Animation {
media: m.file_id.clone(),
has_spoiler: cached.sensitive,
file_id: true,
},
})
.collect();
let task = build_send_task(
&chat_data,
message,
cached.url.clone(),
caption,
items,
Some(cached),
);
dispatch_send(bot, message, &task, url).await;
return;
}
log::info!("fetching {url}");
match x_media::site::fetch(url).await {
// Unsupported links are ignored silently (Python parity).
@@ -384,66 +540,34 @@ async fn url_media(bot: Bot, message: &Message, url: &str) {
.cloned()
.unwrap_or_default();
let caption = fetched.caption_with(&format);
let task = if fetched.media.len() == 1
&& matches!(fetched.media[0], Media::Animated { .. })
{
Task::SendAnimation {
chat_id,
reply_to_message_id: message.id.0 as i64,
caption: caption.clone(),
animation: MediaItemPayload::Animation {
media: fetched.media[0].url().to_string(),
has_spoiler: fetched.sensitive,
},
source_url: fetched.source_url.clone(),
edit_before_forward: chat_data.edit_before_forward,
forward_channel_id: chat_data.forward_channel_id,
notify_chat_id: Some(chat_id),
notify_message_id: Some(message.id.0 as i64),
// Raw render data for the link cache; the send fills in the
// Telegram file ids and persists the entry.
let cache_data = fetched.render_fields().map(|(author, author_url, title, tags)| {
CachedPost {
url: fetched.source_url.clone(),
caption: fetched.caption.clone(),
title: title.to_string(),
author: author.to_string(),
author_url: author_url.to_string(),
tags: tags.to_string(),
sensitive: fetched.sensitive,
media: vec![],
}
} else {
let items: Vec<MediaItemPayload> = fetched
.media
.iter()
.map(|media| media_to_payload(media, fetched.sensitive))
.collect();
Task::SendMediaSequence {
chat_id,
reply_to_message_id: message.id.0 as i64,
caption: caption.clone(),
media_batches: send::chunk_media_items(items),
batch_index: 0,
sent_message_ids: vec![],
source_url: fetched.source_url.clone(),
edit_before_forward: chat_data.edit_before_forward,
forward_channel_id: chat_data.forward_channel_id,
notify_chat_id: Some(chat_id),
notify_message_id: Some(message.id.0 as i64),
}
};
let result = match &task {
Task::SendAnimation { .. } => send::send_animation(&bot, &task).await,
Task::SendMediaSequence { .. } => send::send_media_sequence(&bot, &task).await,
Task::ForwardMessages { .. } => unreachable!(),
};
match result {
Ok(message_ids) => {
log::info!("sent {} message(s) for {url}", message_ids.len());
send::post_send_actions(&bot, &task, message_ids).await;
}
Err(send::SendError::Retryable { delay_seconds, task }) => {
log::info!("send for {url} failed, queued for retry in {delay_seconds:.1}s");
enqueue_retry(task, delay_seconds).await;
let _ = reply(bot, message.clone(), "Send failed. Task queued for retry.").await;
}
Err(send::SendError::Permanent {
message: err_message,
..
}) => {
log::error!("send for {url} failed permanently: {err_message}");
let _ = reply(bot, message.clone(), format!("Send failed: {err_message}")).await;
}
}
});
let items: Vec<MediaItemPayload> = fetched
.media
.iter()
.map(|media| media_to_payload(media, fetched.sensitive))
.collect();
let task = build_send_task(
&chat_data,
message,
fetched.source_url.clone(),
caption,
items,
cache_data,
);
dispatch_send(bot, message, &task, url).await;
}
}
}
@@ -477,7 +601,13 @@ pub async fn message_handler(bot: Bot, message: Message) -> Result<(), RequestEr
log::info!("extracted {} URL(s): {urls:?}", urls.len());
}
for url in urls {
url_media(bot.clone(), &message, &url).await;
let bot = bot.clone();
let message = message.clone();
tokio::spawn(async move {
// Held for the whole task; the semaphore is never closed.
let _permit = URL_TASKS.acquire().await.expect("URL semaphore closed");
url_media(bot, &message, &url).await;
});
}
}
respond(())
@@ -502,11 +632,19 @@ pub async fn inline_query_handler(bot: Bot, query: InlineQuery) -> Result<(), Re
.unwrap_or_else(|| url.clone());
let caption = fetched.caption.clone();
let result = match media {
Media::Illustration { .. } => InlineQueryResult::Photo(
InlineQueryResultPhoto::new(id, url, thumbnail)
.caption(caption)
.parse_mode(ParseMode::Html),
),
Media::Illustration { .. } => {
// Inline photo results have their own (smaller) size
// cap; use the reduced variant when one exists.
let photo_url = media
.smaller_url()
.and_then(|u| url::Url::parse(u).ok())
.unwrap_or_else(|| url.clone());
InlineQueryResult::Photo(
InlineQueryResultPhoto::new(id, photo_url, thumbnail)
.caption(caption)
.parse_mode(ParseMode::Html),
)
}
Media::Video { .. } => InlineQueryResult::Video(
InlineQueryResultVideo::new(
id,
+230
View File
@@ -0,0 +1,230 @@
//! Persistent cache of successfully sent posts.
//!
//! After a media send succeeds, the raw render data plus the Telegram
//! `file_id`s of the sent items are stored keyed by [`crate::site` cache
//! key]. A repeated link is then answered entirely from local state — no
//! re-fetch of the source site, no re-upload — and no media file is stored
//! on disk (the file ids point at Telegram's servers). Entries expire after
//! [`Config::link_cache_ttl`]; a stale entry is dropped lazily on read and
//! by the periodic prune in `main`.
use rusqlite::{params, Connection};
use serde::{Deserialize, Serialize};
use std::time::Duration;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum CachedMediaKind {
Photo,
Video,
Animation,
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct CachedMedia {
pub kind: CachedMediaKind,
pub file_id: String,
}
/// Everything needed to re-send a post without touching the source site:
/// the canonical URL, pre-escaped caption fields, and the file ids produced
/// by the original successful send.
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct CachedPost {
pub url: String,
/// The site's built-in caption (used when the chat has no format
/// override).
pub caption: String,
pub title: String,
pub author: String,
pub author_url: String,
pub tags: String,
pub sensitive: bool,
pub media: Vec<CachedMedia>,
}
/// SQLite-backed cache sharing `data/task_queue.db` with the queue and chat
/// state (same `open_db` pattern: busy timeout, `spawn_blocking` I/O).
pub struct LinkCache {
db_path: String,
}
fn open_db(path: &str) -> rusqlite::Result<Connection> {
let conn = Connection::open(path)?;
conn.busy_timeout(Duration::from_secs(5))?;
Ok(conn)
}
impl LinkCache {
pub fn open(db_path: &str) -> Self {
if let Ok(conn) = Connection::open(db_path)
&& let Err(e) = conn.execute_batch(
"CREATE TABLE IF NOT EXISTS link_cache (url TEXT PRIMARY KEY, \
payload TEXT NOT NULL, created_at REAL NOT NULL);",
)
{
log::error!("failed to initialize link cache schema: {e}");
}
Self {
db_path: db_path.to_string(),
}
}
/// Returns the cached post if present and not expired; a stale entry is
/// removed on the spot.
pub async fn get(&self, key: &str, ttl: Duration) -> Option<CachedPost> {
let db_path = self.db_path.clone();
let key = key.to_string();
let ttl = ttl.as_secs_f64();
tokio::task::spawn_blocking(move || -> rusqlite::Result<Option<CachedPost>> {
let conn = open_db(&db_path)?;
let mut stmt =
conn.prepare("SELECT payload, created_at FROM link_cache WHERE url = ?1")?;
let mut rows = stmt.query(params![key])?;
let Some(row) = rows.next()? else {
return Ok(None);
};
let payload: String = row.get(0)?;
let created_at: f64 = row.get(1)?;
if now_f64() - created_at > ttl {
conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key])?;
return Ok(None);
}
serde_json::from_str(&payload).map(Some).map_err(|e| {
rusqlite::Error::ToSqlConversionFailure(Box::new(e))
})
})
.await
.expect("link cache read worker panicked")
.unwrap_or_else(|e| {
log::error!("link cache read failed: {e}");
None
})
}
pub async fn put(&self, key: &str, post: &CachedPost) {
let db_path = self.db_path.clone();
let key = key.to_string();
let payload = serde_json::to_string(post).expect("cached post serializes");
tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = open_db(&db_path)?;
conn.execute(
"INSERT OR REPLACE INTO link_cache (url, payload, created_at) VALUES (?1, ?2, ?3)",
params![key, payload, now_f64()],
)?;
Ok(())
})
.await
.expect("link cache write worker panicked")
.unwrap_or_else(|e| log::error!("link cache write failed: {e}"));
}
/// Drops an entry (e.g. a cached file id that turned out invalid).
pub async fn remove(&self, key: &str) {
let db_path = self.db_path.clone();
let key = key.to_string();
tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = open_db(&db_path)?;
conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key])?;
Ok(())
})
.await
.expect("link cache delete worker panicked")
.unwrap_or_else(|e| log::error!("link cache delete failed: {e}"));
}
/// Removes expired entries; returns how many were deleted.
pub async fn prune(&self, ttl: Duration) -> usize {
let db_path = self.db_path.clone();
let cutoff = now_f64() - ttl.as_secs_f64();
tokio::task::spawn_blocking(move || -> rusqlite::Result<usize> {
let conn = open_db(&db_path)?;
conn.execute(
"DELETE FROM link_cache WHERE created_at < ?1",
params![cutoff],
)
})
.await
.expect("link cache prune worker panicked")
.unwrap_or_else(|e| {
log::error!("link cache prune failed: {e}");
0
})
}
}
fn now_f64() -> f64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs_f64())
.unwrap_or(0.0)
}
#[cfg(test)]
mod tests {
use super::*;
fn entry() -> CachedPost {
CachedPost {
url: "https://x.com/u/status/1".into(),
caption: "cap".into(),
title: "t".into(),
author: "a".into(),
author_url: "au".into(),
tags: "".into(),
sensitive: true,
media: vec![CachedMedia {
kind: CachedMediaKind::Photo,
file_id: "AgAC...".into(),
}],
}
}
#[tokio::test]
async fn put_get_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let cache = LinkCache::open(dir.path().join("c.db").to_str().unwrap());
cache.put("twitter:1", &entry()).await;
let got = cache.get("twitter:1", Duration::from_secs(3600)).await;
assert!(got.is_some());
let got = got.unwrap();
assert_eq!(got.url, "https://x.com/u/status/1");
assert_eq!(got.media[0].file_id, "AgAC...");
}
#[tokio::test]
async fn expired_entry_removed_on_read() {
let dir = tempfile::tempdir().unwrap();
let cache = LinkCache::open(dir.path().join("c.db").to_str().unwrap());
cache.put("twitter:1", &entry()).await;
// Force the row into the past so a 1s TTL expires it.
{
let conn = Connection::open(dir.path().join("c.db")).unwrap();
conn.execute(
"UPDATE link_cache SET created_at = created_at - 100",
[],
)
.unwrap();
}
assert!(cache.get("twitter:1", Duration::from_secs(1)).await.is_none());
assert!(cache.get("twitter:1", Duration::from_secs(3600)).await.is_none());
}
#[tokio::test]
async fn remove_and_prune() {
let dir = tempfile::tempdir().unwrap();
let cache = LinkCache::open(dir.path().join("c.db").to_str().unwrap());
cache.put("twitter:1", &entry()).await;
cache.put("pixiv:2", &entry()).await;
cache.remove("twitter:1").await;
assert!(cache.get("twitter:1", Duration::from_secs(3600)).await.is_none());
assert!(cache.get("pixiv:2", Duration::from_secs(3600)).await.is_some());
{
let conn = Connection::open(dir.path().join("c.db")).unwrap();
conn.execute("UPDATE link_cache SET created_at = created_at - 100", [])
.unwrap();
}
assert_eq!(cache.prune(Duration::from_secs(1)).await, 1);
assert!(cache.get("pixiv:2", Duration::from_secs(3600)).await.is_none());
}
}
+49 -8
View File
@@ -1,18 +1,38 @@
use dotenv::dotenv;
use teloxide::dptree::endpoint;
use teloxide::stop::StopToken;
use teloxide::types::{ChatId, InputFile, MessageId};
use teloxide::update_listeners::webhooks;
use teloxide::update_listeners::{self, webhooks, UpdateListener};
use teloxide::prelude::*;
use tokio::sync::watch;
use x_media::site;
mod config;
mod handlers;
mod link_cache;
mod queue;
mod send;
mod state;
use handlers::{CHAT_STORE, CONFIG, TASK_QUEUE};
use handlers::{CHAT_STORE, CONFIG, LINK_CACHE, TASK_QUEUE};
/// Docker `stop` / `compose down` delivers SIGTERM, which teloxide's ctrlc
/// handler (SIGINT only) never sees — without this the process would die
/// before the graceful shutdown below (admin notice, queue drain). Stopping
/// the token unwinds the dispatcher exactly like Ctrl+C does.
#[cfg(unix)]
fn spawn_sigterm_handler(stop_token: StopToken) {
tokio::spawn(async move {
let mut sigterm = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("failed to install SIGTERM handler");
sigterm.recv().await;
log::info!("SIGTERM received, stopping the dispatcher");
stop_token.stop();
});
}
#[cfg(not(unix))]
fn spawn_sigterm_handler(_stop_token: StopToken) {}
#[tokio::main]
async fn main() {
@@ -66,6 +86,10 @@ async fn main() {
}
let ttl = CONFIG.edit_message_ttl;
let removed = CHAT_STORE.prune_expired(ttl).await;
let pruned = LINK_CACHE.prune(CONFIG.link_cache_ttl).await;
if pruned > 0 {
log::info!("link cache: pruned {pruned} expired entr(ies)");
}
for (chat_id, prompt_message_id) in removed {
// If the prompt was already deleted, this fails with a
// 400 "message to edit not found" — log and ignore.
@@ -96,7 +120,8 @@ async fn main() {
.webhook_url
.clone()
.expect("WEBHOOK_URL is not set");
bot.set_webhook(url.clone()).await.unwrap();
// `webhooks::axum` calls set_webhook itself (with the full options,
// secret token included) — no explicit registration here.
let listen = CONFIG.webhook_listen.expect("WEBHOOK_LISTEN is not set");
let port = CONFIG.webhook_port.expect("WEBHOOK_PORT is not set");
let mut options = webhooks::Options::new((listen, port).into(), url);
@@ -107,20 +132,36 @@ async fn main() {
options = options.secret_token(secret.clone());
}
let mut listener = webhooks::axum(bot.clone(), options)
.await
.expect("Failed to create webhook listener");
let stop_token = listener.stop_token();
spawn_sigterm_handler(stop_token);
dispatcher
.dispatch_with_listener(
webhooks::axum(bot.clone(), options)
.await
.expect("Failed to create webhook listener"),
listener,
LoggingErrorHandler::with_custom_text("Error from update listener"),
)
.await;
} else {
log::info!("running in polling mode");
dispatcher.dispatch().await;
// Same listener `dispatch()` builds internally — using
// `dispatch_with_listener` just exposes its stop token so SIGTERM can
// unwind the dispatcher before the graceful shutdown below.
let mut listener = update_listeners::polling_default(bot.clone()).await;
let stop_token = listener.stop_token();
spawn_sigterm_handler(stop_token);
dispatcher
.dispatch_with_listener(
listener,
LoggingErrorHandler::with_custom_text("Error from update listener"),
)
.await;
}
// Graceful stop (Ctrl+C): stop the sweep, notify the admin, drain the queue.
// Graceful stop (Ctrl+C / SIGTERM): stop the sweep, notify the admin, drain the queue.
log::info!("Stopping bot");
let _ = stop_tx.send(true);
if let Some(admin) = CONFIG.admin_ids.first() {
+52 -27
View File
@@ -6,7 +6,7 @@
//! replaced by dedicated columns.
use parking_lot::Mutex;
use rusqlite::{params, Connection};
use rusqlite::{params, Connection, TransactionBehavior};
use serde_json::Value;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
@@ -18,6 +18,12 @@ use tokio::task::JoinHandle;
pub const MAX_RETRIES: u32 = 2;
pub const LOCK_TTL_SECONDS: f64 = 120.0;
/// Number of concurrent worker loops. Tasks are independent (retries and
/// forward resumes); leases serialize row claims via SQLite transactions, so
/// extra workers drain backlogs faster. Each worker can be mid-send to
/// Telegram at the same time as handler tasks, so keep this modest.
const QUEUE_WORKERS: usize = 4;
/// What a handler returns instead of throwing. The payload it carries is the
/// (possibly updated) task state to persist for the next attempt.
pub enum QueueError {
@@ -42,7 +48,7 @@ pub struct PersistentTaskQueue {
db_path: String,
notify: Arc<Notify>,
stop: Arc<AtomicBool>,
worker: Mutex<Option<JoinHandle<()>>>,
worker: Mutex<Vec<JoinHandle<()>>>,
counter: AtomicU64,
}
@@ -68,6 +74,15 @@ fn now_f64() -> f64 {
.unwrap_or(0.0)
}
/// Opens the queue DB with a busy timeout. Handler tasks enqueue while
/// workers lease/update rows concurrently; without the timeout a concurrent
/// write fails immediately with SQLITE_BUSY and the operation is lost.
fn open_db(path: &str) -> rusqlite::Result<Connection> {
let conn = Connection::open(path)?;
conn.busy_timeout(Duration::from_secs(5))?;
Ok(conn)
}
fn ensure_schema(conn: &Connection) -> rusqlite::Result<()> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS tasks (id TEXT PRIMARY KEY, payload TEXT NOT NULL, \
@@ -96,12 +111,12 @@ impl PersistentTaskQueue {
db_path: db_path.to_string(),
notify: Arc::new(Notify::new()),
stop: Arc::new(AtomicBool::new(false)),
worker: Mutex::new(None),
worker: Mutex::new(Vec::new()),
counter: AtomicU64::new(0),
}
}
/// Starts the worker loop. Also recovers rows left `in_progress` by a
/// Starts the worker loops. Also recovers rows left `in_progress` by a
/// previous process (lease expired).
pub async fn start<H, F, D, G>(&self, handler: H, dead_letter: D)
where
@@ -114,21 +129,25 @@ impl PersistentTaskQueue {
let dead_letter: Arc<DeadLetter> =
Arc::new(move |payload, message| Box::pin(dead_letter(payload, message)));
self.recover_stale().await;
let worker = QueueWorker {
db_path: self.db_path.clone(),
notify: Arc::clone(&self.notify),
stop: Arc::clone(&self.stop),
handler,
dead_letter,
};
let worker = tokio::spawn(worker.run_loop());
*self.worker.lock() = Some(worker);
let mut handles = Vec::with_capacity(QUEUE_WORKERS);
for _ in 0..QUEUE_WORKERS {
let worker = QueueWorker {
db_path: self.db_path.clone(),
notify: Arc::clone(&self.notify),
stop: Arc::clone(&self.stop),
handler: Arc::clone(&handler),
dead_letter: Arc::clone(&dead_letter),
};
handles.push(tokio::spawn(worker.run_loop()));
}
*self.worker.lock() = handles;
}
pub async fn stop(&self) {
self.stop.store(true, Ordering::Relaxed);
self.notify.notify_one();
if let Some(handle) = self.worker.lock().take() {
self.notify.notify_waiters();
let handles = std::mem::take(&mut *self.worker.lock());
for handle in handles {
let _ = handle.await;
}
}
@@ -145,8 +164,8 @@ impl PersistentTaskQueue {
let payload = payload.to_string();
let db_path = self.db_path.clone();
log::info!("enqueued {id} (run_after {run_after:.1})");
let result = tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = Connection::open(&db_path)?;
tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = open_db(&db_path)?;
conn.execute(
"INSERT OR REPLACE INTO tasks (id, payload, run_after, attempts, status, locked_until, created_at) \
VALUES (?1, ?2, ?3, 0, 'pending', 0, ?4)",
@@ -156,14 +175,16 @@ impl PersistentTaskQueue {
})
.await
.expect("queue insert worker panicked")?;
self.notify.notify_one();
Ok(result)
// Wake every sleeping worker: with several workers the one that finds
// nothing due must not starve the newly inserted row.
self.notify.notify_waiters();
Ok(())
}
async fn recover_stale(&self) {
let db_path = self.db_path.clone();
tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = Connection::open(&db_path)?;
let conn = open_db(&db_path)?;
conn.execute(
"UPDATE tasks SET status='pending', locked_until=0 WHERE status='in_progress' AND locked_until < ?1",
params![now_f64()],
@@ -206,11 +227,15 @@ impl QueueWorker {
async fn lease_next(&self) -> Option<LeasedRow> {
let db_path = self.db_path.clone();
tokio::task::spawn_blocking(move || -> rusqlite::Result<Option<LeasedRow>> {
let mut conn = Connection::open(&db_path)?;
let tx = conn.transaction()?;
let mut conn = open_db(&db_path)?;
// BEGIN IMMEDIATE: with several workers, a deferred transaction
// that read before another worker's lease commit would fail with
// SQLITE_BUSY_SNAPSHOT. Taking the write lock up front serializes
// leases and re-reads the freshest committed state.
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let now = now_f64();
let row = tx.query_row(
"SELECT id, payload, attempts FROM tasks WHERE status='pending' AND run_after <= ?1 \
"SELECT id, payload, attempts FROM tasks WHERE status='pending' AND run_after <= ?1 AND locked_until <= ?1 \
ORDER BY run_after LIMIT 1",
params![now],
|r| {
@@ -251,7 +276,7 @@ impl QueueWorker {
async fn earliest_run_after(&self) -> Option<f64> {
let db_path = self.db_path.clone();
tokio::task::spawn_blocking(move || -> rusqlite::Result<Option<f64>> {
let conn = Connection::open(&db_path)?;
let conn = open_db(&db_path)?;
let mut stmt = conn.prepare("SELECT MIN(run_after) FROM tasks WHERE status='pending'")?;
let mut rows = stmt.query([])?;
match rows.next()? {
@@ -314,7 +339,7 @@ impl QueueWorker {
let db_path = self.db_path.clone();
let id = id.to_string();
tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = Connection::open(&db_path)?;
let conn = open_db(&db_path)?;
conn.execute("DELETE FROM tasks WHERE id = ?1", params![id])?;
Ok(())
})
@@ -328,7 +353,7 @@ impl QueueWorker {
let id = id.to_string();
let payload = payload.to_string();
tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = Connection::open(&db_path)?;
let conn = open_db(&db_path)?;
conn.execute(
"UPDATE tasks SET payload=?1, run_after=?2, attempts=?3, status='pending', locked_until=0 WHERE id=?4",
params![payload, now_f64() + delay_seconds, attempts, id],
@@ -338,7 +363,7 @@ impl QueueWorker {
.await
.expect("queue reschedule worker panicked")
.unwrap_or_else(|e| log::error!("queue reschedule failed: {e}"));
self.notify.notify_one();
self.notify.notify_waiters();
}
}
+392 -52
View File
@@ -3,7 +3,8 @@
//! URL is blocked by hotlink protection; the bot downloads the file itself
//! and uploads it via multipart).
use crate::handlers::{CHAT_STORE, TASK_QUEUE};
use crate::handlers::{CHAT_STORE, LINK_CACHE, TASK_QUEUE};
use crate::link_cache::{CachedMedia, CachedMediaKind, CachedPost};
use crate::queue::QueueError;
use crate::state::{EditMessage, unix_now};
use rand::Rng;
@@ -25,18 +26,43 @@ pub enum MediaItemPayload {
Photo {
media: String,
has_spoiler: bool,
/// Smaller variant used when the primary media exceeds Telegram's
/// size limits.
#[serde(default)]
fallback_url: Option<String>,
/// `media` is a Telegram file id (link-cache hit), not a URL.
#[serde(default)]
file_id: bool,
},
Video {
media: String,
has_spoiler: bool,
thumbnail: Option<String>,
#[serde(default)]
fallback_url: Option<String>,
/// `media` is a Telegram file id (link-cache hit), not a URL.
#[serde(default)]
file_id: bool,
},
Animation {
media: String,
has_spoiler: bool,
/// `media` is a Telegram file id (link-cache hit), not a URL.
#[serde(default)]
file_id: bool,
},
}
impl MediaItemPayload {
fn fallback_url(&self) -> Option<&str> {
match self {
MediaItemPayload::Photo { fallback_url, .. }
| MediaItemPayload::Video { fallback_url, .. } => fallback_url.as_deref(),
MediaItemPayload::Animation { .. } => None,
}
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum Task {
@@ -52,6 +78,10 @@ pub enum Task {
forward_channel_id: Option<i64>,
notify_chat_id: Option<i64>,
notify_message_id: Option<i64>,
/// Raw render data captured on a cache miss; the send fills in the
/// Telegram file ids and persists the entry (see `link_cache`).
#[serde(default)]
cache_data: Option<CachedPost>,
},
SendAnimation {
chat_id: i64,
@@ -63,6 +93,10 @@ pub enum Task {
forward_channel_id: Option<i64>,
notify_chat_id: Option<i64>,
notify_message_id: Option<i64>,
/// Raw render data captured on a cache miss; the send fills in the
/// Telegram file id and persists the entry (see `link_cache`).
#[serde(default)]
cache_data: Option<CachedPost>,
},
ForwardMessages {
from_chat_id: i64,
@@ -73,8 +107,112 @@ pub enum Task {
},
}
impl Task {
fn cache_data(&self) -> Option<&CachedPost> {
match self {
Task::SendMediaSequence { cache_data, .. }
| Task::SendAnimation { cache_data, .. } => cache_data.as_ref(),
Task::ForwardMessages { .. } => None,
}
}
fn source_url(&self) -> Option<&str> {
match self {
Task::SendMediaSequence { source_url, .. }
| Task::SendAnimation { source_url, .. } => Some(source_url),
Task::ForwardMessages { .. } => None,
}
}
/// True when the media payloads are Telegram file ids from the link cache
/// (a cached file id that goes permanently bad should be dropped so the
/// next request re-fetches).
fn is_cached_send(&self) -> bool {
self.cache_data().is_some_and(|c| !c.media.is_empty())
}
}
/// Telegram file id of the message's media, matched to the payload kind.
fn file_id_of_message(message: &Message, item: &MediaItemPayload) -> Option<String> {
match item {
// `photo()` returns all sizes, smallest first — the largest carries
// the file id of the sent media.
MediaItemPayload::Photo { .. } => {
message.photo().and_then(|sizes| sizes.last()).map(|p| p.file.id.to_string())
}
MediaItemPayload::Video { .. } => message.video().map(|v| v.file.id.to_string()),
MediaItemPayload::Animation { .. } => message.animation().map(|a| a.file.id.to_string()),
}
}
fn kind_of_item(item: &MediaItemPayload) -> CachedMediaKind {
match item {
MediaItemPayload::Photo { .. } => CachedMediaKind::Photo,
MediaItemPayload::Video { .. } => CachedMediaKind::Video,
MediaItemPayload::Animation { .. } => CachedMediaKind::Animation,
}
}
/// Collects the Telegram file ids of a sent media group, aligned to the
/// batch's items.
fn collect_file_ids(messages: &[Message], batch: &[MediaItemPayload], out: &mut Vec<CachedMedia>) {
for (message, item) in messages.iter().zip(batch.iter()) {
if let Some(file_id) = file_id_of_message(message, item) {
out.push(CachedMedia {
kind: kind_of_item(item),
file_id,
});
}
}
}
/// Persists a successful send under the post's cache key. Only runs for a
/// fresh (non-resumed) task that carried raw cache data with no file ids yet.
async fn cache_sent_task(task: &Task, media: Vec<CachedMedia>) {
let Some(cache_data) = task.cache_data() else {
return;
};
if !cache_data.media.is_empty() || media.is_empty() {
return;
}
let mut post = cache_data.clone();
post.media = media;
if let Some(key) = x_media::site::cache_key(&post.url) {
LINK_CACHE.put(&key, &post).await;
log::info!("cached send for {}", post.url);
}
}
/// Persists a lone animation send under the post's cache key.
async fn cache_animation_send(task: &Task, message: &Message) {
if let Some(file_id) = message.animation().map(|a| a.file.id.to_string()) {
cache_sent_task(
task,
vec![CachedMedia {
kind: CachedMediaKind::Animation,
file_id,
}],
)
.await;
}
}
/// A cached Telegram file id failed permanently (stale/expired); drop the
/// cache entry so the next request re-fetches instead of repeating it.
pub async fn invalidate_cache(task: &Task) {
if task.is_cached_send()
&& let Some(url) = task.source_url()
&& let Some(key) = x_media::site::cache_key(url)
{
log::info!("removing stale link cache entry for {url}");
LINK_CACHE.remove(&key).await;
}
}
pub const MAX_MEDIA_GROUP: usize = 9;
pub const MAX_UPLOAD_BYTES: u64 = 50 * 1024 * 1024; // Telegram Bot API upload cap
/// Upload cap (bytes): files above this are not uploaded; the bot falls back
/// to a smaller media URL instead.
pub const MAX_UPLOAD_BYTES: u64 = 10 * 1024 * 1024; // 10485760
/// Splits media into batches of at most [`MAX_MEDIA_GROUP`] items.
pub fn chunk_media_items<T: Clone>(items: Vec<T>) -> Vec<Vec<T>> {
@@ -102,6 +240,18 @@ pub fn is_media_fetch_failure(e: &ApiError) -> bool {
MARKERS.iter().any(|marker| description.contains(marker))
}
/// Telegram reported the media file as too large (HTTP 413 on multipart
/// upload, or a "too large" message for URL-fetched media). These errors are
/// handled by the size-check fallback (use a smaller media URL), NOT by a
/// queue retry.
pub fn is_size_error(e: &ApiError) -> bool {
if matches!(e, ApiError::RequestEntityTooLarge) {
return true;
}
let description = e.to_string().to_lowercase();
["too large", "too big"].iter().any(|marker| description.contains(marker))
}
/// Task-free classification of a Telegram request error. The callers attach
/// the (updated) task when building a [`SendError`].
pub enum Classification {
@@ -168,6 +318,32 @@ fn input_file_for(media: &str) -> Result<InputFile, String> {
}
}
impl MediaItemPayload {
/// The input for a send: a cached file id goes out as `InputFile::file_id`
/// (no fetch, no upload), URLs go to Telegram, anything else is a local
/// path (transient upload fallback).
fn input_file(&self) -> Result<InputFile, String> {
match self {
MediaItemPayload::Photo {
media,
file_id: true,
..
}
| MediaItemPayload::Video {
media,
file_id: true,
..
}
| MediaItemPayload::Animation {
media,
file_id: true,
..
} => Ok(InputFile::file_id(media.clone().into())),
_ => input_file_for(item_url(self)),
}
}
}
fn photo_media(file: InputFile, caption: Option<&str>, spoiler: bool) -> InputMedia {
let mut photo = InputMediaPhoto::new(file).parse_mode(ParseMode::Html);
if let Some(caption) = caption {
@@ -214,24 +390,22 @@ fn build_media_group(
let item_caption = if i == 0 { caption } else { None };
Ok(match item {
MediaItemPayload::Photo {
media,
has_spoiler,
} => photo_media(input_file_for(media)?, item_caption, *has_spoiler),
has_spoiler, ..
} => photo_media(item.input_file()?, item_caption, *has_spoiler),
MediaItemPayload::Video {
media,
has_spoiler,
thumbnail,
..
} => {
let mut video = video_media(input_file_for(media)?, item_caption, *has_spoiler);
let mut video = video_media(item.input_file()?, item_caption, *has_spoiler);
if let (Some(thumb), InputMedia::Video(v)) = (thumbnail, &mut video) {
*v = v.clone().thumbnail(input_file_for(thumb)?);
}
video
}
MediaItemPayload::Animation {
media,
has_spoiler,
} => animation_media(input_file_for(media)?, item_caption, *has_spoiler),
has_spoiler, ..
} => animation_media(item.input_file()?, item_caption, *has_spoiler),
})
})
.collect()
@@ -258,6 +432,9 @@ fn sniff_ext(bytes: &[u8]) -> &'static str {
enum FallbackError {
Retryable { delay_seconds: f64 },
Permanent { message: String },
/// The downloaded file exceeds the upload cap; the caller falls back to
/// the item's smaller URL.
MediaTooLarge,
}
/// Downloads one media item to a temp file (deleted on drop). Network errors
@@ -282,9 +459,7 @@ async fn download_to_temp(item: &MediaItemPayload) -> Result<NamedTempFile, Fall
}
};
if bytes.len() as u64 > MAX_UPLOAD_BYTES {
return Err(FallbackError::Permanent {
message: "media too large".into(),
});
return Err(FallbackError::MediaTooLarge);
}
let ext = sniff_ext(&bytes);
let mut file = tempfile::Builder::new()
@@ -302,7 +477,48 @@ async fn download_to_temp(item: &MediaItemPayload) -> Result<NamedTempFile, Fall
Ok(file)
}
/// Download-and-reupload fallback for one media batch.
/// Builds the media group item from an uploaded file.
fn media_from_file(
item: &MediaItemPayload,
path: std::path::PathBuf,
caption: Option<&str>,
) -> InputMedia {
match item {
MediaItemPayload::Photo { has_spoiler, .. } => {
photo_media(InputFile::file(path), caption, *has_spoiler)
}
MediaItemPayload::Video { has_spoiler, .. } => {
video_media(InputFile::file(path), caption, *has_spoiler)
}
MediaItemPayload::Animation { has_spoiler, .. } => {
animation_media(InputFile::file(path), caption, *has_spoiler)
}
}
}
/// Builds the media group item from a (smaller) URL.
fn media_from_url(
item: &MediaItemPayload,
url: &str,
caption: Option<&str>,
) -> Result<InputMedia, String> {
Ok(match item {
MediaItemPayload::Photo { has_spoiler, .. } => {
photo_media(input_file_for(url)?, caption, *has_spoiler)
}
MediaItemPayload::Video { has_spoiler, .. } => {
video_media(input_file_for(url)?, caption, *has_spoiler)
}
MediaItemPayload::Animation { has_spoiler, .. } => {
animation_media(input_file_for(url)?, caption, *has_spoiler)
}
})
}
/// Download-and-reupload fallback for one media batch. Files over the upload
/// cap are not downloaded/uploaded; the item falls back to its smaller URL
/// (which Telegram fetches itself). Returns the fallback-error without the
/// task attached; callers wrap it with the updated task state.
async fn send_batch_via_upload(
bot: &Bot,
chat_id: i64,
@@ -313,22 +529,51 @@ async fn send_batch_via_upload(
let mut files = Vec::new();
let mut items = Vec::new();
for (i, item) in batch.iter().enumerate() {
let file = download_to_temp(item).await?;
let path = file.path().to_path_buf();
let item_caption = if i == 0 { caption } else { None };
let media = match item {
MediaItemPayload::Photo { has_spoiler, .. } => {
photo_media(InputFile::file(path), item_caption, *has_spoiler)
// Size check before downloading/uploading: over the cap, use the
// smaller URL instead of the file.
let too_large = match x_media::site::media_size(item_url(item)).await {
Ok(Some(size)) => size > MAX_UPLOAD_BYTES,
_ => false,
};
let media = if too_large {
match item.fallback_url() {
Some(url) => match media_from_url(item, url, item_caption) {
Ok(media) => media,
Err(message) => {
return Err(FallbackError::Permanent { message });
}
},
None => {
return Err(FallbackError::Permanent {
message: "media too large".into(),
});
}
}
MediaItemPayload::Video { has_spoiler, .. } => {
video_media(InputFile::file(path), item_caption, *has_spoiler)
}
MediaItemPayload::Animation { has_spoiler, .. } => {
animation_media(InputFile::file(path), item_caption, *has_spoiler)
} else {
match download_to_temp(item).await {
Ok(file) => {
let path = file.path().to_path_buf();
files.push(file);
media_from_file(item, path, item_caption)
}
Err(FallbackError::MediaTooLarge) => match item.fallback_url() {
Some(url) => match media_from_url(item, url, item_caption) {
Ok(media) => media,
Err(message) => {
return Err(FallbackError::Permanent { message });
}
},
None => {
return Err(FallbackError::Permanent {
message: "media too large".into(),
});
}
},
Err(e) => return Err(e),
}
};
items.push(media);
files.push(file);
}
let result = bot
.send_media_group(ChatId(chat_id), items)
@@ -355,12 +600,14 @@ fn updated_sequence_task(task: &Task, batch_index: usize, sent_message_ids: Vec<
reply_to_message_id,
caption,
media_batches,
batch_index: _,
sent_message_ids: _,
source_url,
edit_before_forward,
forward_channel_id,
notify_chat_id,
notify_message_id,
..
cache_data,
} => Task::SendMediaSequence {
chat_id: *chat_id,
reply_to_message_id: *reply_to_message_id,
@@ -373,6 +620,7 @@ fn updated_sequence_task(task: &Task, batch_index: usize, sent_message_ids: Vec<
forward_channel_id: *forward_channel_id,
notify_chat_id: *notify_chat_id,
notify_message_id: *notify_message_id,
cache_data: cache_data.clone(),
},
_ => unreachable!("updated_sequence_task requires a SendMediaSequence task"),
}
@@ -397,6 +645,10 @@ pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result<Vec<i64>, Sen
let chat_id = *chat_id;
let reply_to = *reply_to_message_id;
let mut sent = sent_message_ids.clone();
// File ids accumulated across batches for the link cache. Only a fresh
// (non-resumed) full send populates the cache.
let mut cached_media: Vec<CachedMedia> = Vec::new();
let fresh_send = *batch_index == 0 && sent.is_empty();
for idx in *batch_index..media_batches.len() {
let batch = &media_batches[idx];
let caption = if idx == 0 { Some(caption.as_str()) } else { None };
@@ -420,18 +672,21 @@ pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result<Vec<i64>, Sen
media_batches.len(),
batch.len()
);
collect_file_ids(&messages, batch, &mut cached_media);
sent.extend(messages.into_iter().map(|m| m.id.0 as i64));
}
Err(RequestError::Api(api)) if is_media_fetch_failure(&api) => {
Err(RequestError::Api(api))
if is_media_fetch_failure(&api) || is_size_error(&api) =>
{
log::info!(
"Telegram could not fetch media for batch {idx} ({}), downloading and reuploading",
batch
.first()
.map(|item| item_url(item))
.unwrap_or("?")
batch.first().map(item_url).unwrap_or("?")
);
match send_batch_via_upload(bot, chat_id, reply_to, batch, caption).await {
Ok(messages) => sent.extend(messages.into_iter().map(|m| m.id.0 as i64)),
Ok(messages) => {
collect_file_ids(&messages, batch, &mut cached_media);
sent.extend(messages.into_iter().map(|m| m.id.0 as i64));
}
Err(FallbackError::Retryable { delay_seconds }) => {
return Err(SendError::Retryable {
delay_seconds,
@@ -444,6 +699,7 @@ pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result<Vec<i64>, Sen
task: updated_sequence_task(task, idx, sent),
});
}
Err(FallbackError::MediaTooLarge) => unreachable!("handled inside upload"),
}
}
Err(e) => {
@@ -454,6 +710,9 @@ pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result<Vec<i64>, Sen
}
}
}
if fresh_send {
cache_sent_task(task, cached_media).await;
}
Ok(sent)
}
@@ -494,6 +753,7 @@ pub async fn send_animation(bot: &Bot, task: &Task) -> Result<Vec<i64>, SendErro
MediaItemPayload::Animation {
media,
has_spoiler,
..
} => (media, *has_spoiler),
MediaItemPayload::Photo { .. } | MediaItemPayload::Video { .. } => {
unreachable!("SendAnimation carries an Animation payload")
@@ -506,34 +766,76 @@ pub async fn send_animation(bot: &Bot, task: &Task) -> Result<Vec<i64>, SendErro
match send_animation_inner(bot, chat_id, reply_to, caption, has_spoiler, url_file)
.await
{
Ok(message) => Ok(vec![message.id.0 as i64]),
Err(RequestError::Api(api)) if is_media_fetch_failure(&api) => {
Ok(message) => {
let id = message.id.0 as i64;
cache_animation_send(task, &message).await;
Ok(vec![id])
}
Err(RequestError::Api(api))
if is_media_fetch_failure(&api) || is_size_error(&api) =>
{
log::info!(
"Telegram could not fetch animation URL, downloading and reuploading: {}",
media_url
);
let file = match download_to_temp(animation).await {
Ok(file) => file,
match download_to_temp(animation).await {
Ok(file) => {
let path = file.path().to_path_buf();
match send_animation_inner(
bot,
chat_id,
reply_to,
caption,
has_spoiler,
InputFile::file(path),
)
.await
{
Ok(message) => {
let id = message.id.0 as i64;
cache_animation_send(task, &message).await;
Ok(vec![id])
}
Err(e) => Err(classify_to_send_error(&e, task.clone())),
}
}
// Over the upload cap: fall back to the smaller URL.
Err(FallbackError::MediaTooLarge) => match animation.fallback_url() {
Some(url) => match input_file_for(url) {
Ok(file) => {
match send_animation_inner(
bot,
chat_id,
reply_to,
caption,
has_spoiler,
file,
)
.await
{
Ok(message) => {
let id = message.id.0 as i64;
cache_animation_send(task, &message).await;
Ok(vec![id])
}
Err(e) => Err(classify_to_send_error(&e, task.clone())),
}
}
Err(message) => {
Err(SendError::Permanent { message, task: task.clone() })
}
},
None => Err(SendError::Permanent {
message: "media too large".into(),
task: task.clone(),
}),
},
Err(FallbackError::Retryable { delay_seconds }) => {
return Err(SendError::Retryable { delay_seconds, task: task.clone() });
Err(SendError::Retryable { delay_seconds, task: task.clone() })
}
Err(FallbackError::Permanent { message }) => {
return Err(SendError::Permanent { message, task: task.clone() });
Err(SendError::Permanent { message, task: task.clone() })
}
};
let path = file.path().to_path_buf();
match send_animation_inner(
bot,
chat_id,
reply_to,
caption,
has_spoiler,
InputFile::file(path),
)
.await
{
Ok(message) => Ok(vec![message.id.0 as i64]),
Err(e) => Err(classify_to_send_error(&e, task.clone())),
}
}
Err(e) => Err(classify_to_send_error(&e, task.clone())),
@@ -731,6 +1033,7 @@ pub async fn handle_task(payload: serde_json::Value) -> Result<(), QueueError> {
});
}
Err(SendError::Permanent { message, task }) => {
invalidate_cache(&task).await;
return Err(QueueError::Permanent {
message,
payload: serde_json::to_value(task).expect("task serializes"),
@@ -819,6 +1122,36 @@ mod tests {
}
}
#[test]
fn is_size_error_matches_known_errors() {
// 413 upload cap.
let e = ApiError::RequestEntityTooLarge;
assert!(is_size_error(&e), "{e:?}");
// Unknown descriptions with size wording.
for description in [
"Bad Request: file is too large",
"Bad Request: media is too big",
"Bad Request: url file size is too big",
] {
let api = ApiError::Unknown(description.to_string());
assert!(is_size_error(&api), "{description}");
}
// Unrelated errors must not match.
for description in ["Bad Request: WEBPAGE_MEDIA_EMPTY", "Bad Request: message is not modified"] {
let api = ApiError::Unknown(description.to_string());
assert!(!is_size_error(&api), "{description}");
}
}
#[test]
fn media_item_payload_fallback_url_serde_default() {
// Old queued payloads without the field deserialize with None.
let json = serde_json::json!({"kind": "photo", "media": "https://a/b.jpg", "has_spoiler": false});
let photo: MediaItemPayload = serde_json::from_value(json).unwrap();
assert!(matches!(photo, MediaItemPayload::Photo { fallback_url: None, .. }));
assert_eq!(photo.fallback_url(), None);
}
#[test]
fn classification_mapping() {
use teloxide::types::Seconds;
@@ -858,11 +1191,15 @@ mod tests {
vec![MediaItemPayload::Photo {
media: "https://a/b.jpg".into(),
has_spoiler: true,
fallback_url: Some("https://a/b_small.jpg".into()),
file_id: false,
}],
vec![MediaItemPayload::Video {
media: "https://a/v.mp4".into(),
has_spoiler: false,
thumbnail: Some("https://a/t.jpg".into()),
fallback_url: None,
file_id: false,
}],
],
batch_index: 1,
@@ -872,6 +1209,7 @@ mod tests {
forward_channel_id: Some(333),
notify_chat_id: Some(111),
notify_message_id: Some(222),
cache_data: None,
};
let json = serde_json::to_value(&task).unwrap();
assert_eq!(json["type"], "send_media_sequence");
@@ -900,6 +1238,8 @@ mod tests {
let photo = MediaItemPayload::Photo {
media: "https://a/b.jpg".into(),
has_spoiler: false,
fallback_url: None,
file_id: false,
};
let json = serde_json::to_value(&photo).unwrap();
assert_eq!(json["kind"], "photo");
+5
View File
@@ -74,6 +74,10 @@ impl ChatStore {
let db_path = self.db_path.clone();
let payload = tokio::task::spawn_blocking(move || -> rusqlite::Result<Option<String>> {
let conn = Connection::open(&db_path)?;
// Concurrent handler tasks (batch-forwards) may write chat_state
// while this read runs; without a busy timeout a write lock
// collision fails the query immediately.
conn.busy_timeout(std::time::Duration::from_secs(5))?;
let mut stmt = conn.prepare("SELECT payload FROM chat_state WHERE chat_id = ?1")?;
let mut rows = stmt.query(params![chat_id.to_string()])?;
match rows.next()? {
@@ -100,6 +104,7 @@ impl ChatStore {
let db_path = self.db_path.clone();
tokio::task::spawn_blocking(move || -> rusqlite::Result<()> {
let conn = Connection::open(&db_path)?;
conn.busy_timeout(std::time::Duration::from_secs(5))?;
conn.execute(
"INSERT OR REPLACE INTO chat_state (chat_id, payload) VALUES (?1, ?2)",
params![chat_id.to_string(), payload],
+54 -14
View File
@@ -1,30 +1,70 @@
services:
nginx-proxy:
image: nginxproxy/nginx-proxy:1.11.6-alpine
restart: always
ports:
- '80:80'
- '443:443'
environment:
# Bare-IP access only.
# DEFAULT_HOST: 'bot.example.com'
volumes:
- /var/run/docker.sock:/tmp/docker.sock:ro
- ./nginx-certs:/etc/nginx/certs:ro
- ./nginx-vhost.d:/etc/nginx/vhost.d:ro
- ./nginx-html:/usr/share/nginx/html:ro
networks: [proxy]
labels:
- 'com.github.jrcs.letsencrypt_nginx_proxy_companion.nginx_proxy=true'
container_name: nginx-proxy
acme-companion:
image: nginxproxy/acme-companion
restart: always
environment:
DEFAULT_EMAIL: 'admin@yoursfunny.top'
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
- ./nginx-certs:/etc/nginx/certs:rw
- ./nginx-vhost.d:/etc/nginx/vhost.d:rw
- ./nginx-html:/usr/share/nginx/html:rw
- ./nginx-acme:/etc/acme.sh
networks: [proxy]
container_name: acme-companion
depends_on:
- nginx-proxy
tgxmb:
image: yoursfunny/telegram-twitter-media-bot:latest
restart: always
# ports:
# - "8443:8443"
environment:
# docker-entrypoint.sh drops privileges to this uid.
LOCAL_USER_ID: '1000'
# Bot token (BotFather). Required.
TELOXIDE_TOKEN: ''
# Comma-separated admin chat ids; receives startup/shutdown notices.
BOT_ADMIN: ''
# Required for pixiv support; pixiv is disabled when unset.
PIXIV_REFRESH_TOKEN: ''
# Edit-before-forward records expire after this many seconds (default 86400 = 24h).
# Optional: x.com session cookie (auth_token) — fetches NSFW tweets
# that the public syndication endpoint withholds.
TWITTER_AUTH_TOKEN: ''
EDIT_MESSAGE_TTL_SECONDS: '86400'
# Link-result cache TTL (default 604800 = 7 days).
LINK_CACHE_TTL_SECONDS: '604800'
RUST_LOG: 'info'
# Webhook mode is off by default (polling). The listener binds inside the
# container, so use 0.0.0.0 and publish the port if you enable it.
WEBHOOK: 'false'
VIRTUAL_HOST: 'bot.example.com'
VIRTUAL_PORT: '8443'
# LETSENCRYPT_HOST: 'bot.example.com'
WEBHOOK: 'true'
WEBHOOK_LISTEN: '0.0.0.0'
WEBHOOK_PORT: '8443'
WEBHOOK_URL: 'https://example.com'
WEBHOOK_CERT: './cert/cert.pem'
WEBHOOK_SECRET_TOKEN: 'secret-token'
WEBHOOK_URL: 'https://bot.example.com/'
# WEBHOOK_CERT: './cert/cert.pem'
WEBHOOK_SECRET_TOKEN: ''
volumes:
- ./data:/app/data
# - ./cert:/app/cert
networks: [proxy]
depends_on:
- nginx-proxy
container_name: tgxmb
networks:
proxy:
name: proxy