Compare commits

...

11 Commits

Author SHA1 Message Date
Your Name 2d6fa3c99a feat: add declarative site parser engine 2026-08-11 00:46:11 +05:00
Your Name 01eea99c88 docs: document external site parser 2026-08-10 23:28:47 +05:00
Your Name f866a13a8a fix: serve admin media previews from telegram 2026-08-10 23:22:49 +05:00
Your Name 55f3bb51b6 fix: reuse Cloudflare session for site media 2026-08-10 23:11:09 +05:00
Your Name aad9ee46eb feat: filter site bodies before media upload 2026-08-10 22:52:23 +05:00
Your Name e5ef3fb9e8 fix: stream direct site videos 2026-08-10 22:37:24 +05:00
Your Name faded0132b feat: ingest site videos through raw pipeline 2026-08-10 22:35:33 +05:00
Your Name 7889d4d376 fix: keep site parsing interval global 2026-08-10 21:32:46 +05:00
Your Name 5417737c08 fix: complete site parser source workflow 2026-08-10 21:11:43 +05:00
Your Name 3dea6f43f3 feat: add external site parser worker 2026-08-10 20:48:23 +05:00
Your Name fbc1e3fd1b chore: add Coolify Docker cleanup script 2026-08-10 18:48:47 +05:00
32 changed files with 2537 additions and 151 deletions
+72 -19
View File
@@ -61,7 +61,7 @@
│ ├── templates/ # HTML-шаблоны
│ └── workers/
│ ├── parser.py # VK-парсер
│ ├── vk_storage_uploader.py# Telegram media storage uploader
│ ├── media_uploader.py# Telegram media storage uploader
│ ├── ai_qualifier.py # AI-квалификатор
│ └── ai_writer.py # AI-райтер
├── original_project/ # старый проект как справочник
@@ -196,7 +196,8 @@
Источники для парсинга.
Сейчас реально используется только `platform='vk'`, но схема заложена под другие площадки.
Используются `platform='vk'` и `platform='site'`. Для сайта конфигурация
конкретного адаптера хранится вместе с источником в `settings_json`.
Важные поля:
@@ -215,9 +216,14 @@
- `parse_from`
- `priority`
- `archived_at`
- `settings_json` для `platform='site'`
`tag` нужен для будущих хэштегов и финального оформления.
Для сайта отсутствие или ошибка JSON-конфига блокирует запуск с видимой
ошибкой источника. Конфиг задаёт способ получения именно этого сайта; общее
расписание и секреты остаются в `app_settings`.
### 5.4. `raw_posts`
Главная таблица жизненного цикла поста.
@@ -338,7 +344,7 @@ AI-райтер:
Сейчас основной тип:
- `vk.storage.copy`
- `media.storage.copy`
Название историческое: сначала планировался VK storage, потом медиа-сторедж вернулся в Telegram. Тип задачи пока не переименован.
@@ -357,7 +363,7 @@ AI-райтер:
Имена воркеров:
- `vk-parser`
- `vk-storage-uploader`
- `media-uploader`
- `ai-qualifier`
- `ai-writer`
- `tg-poster`
@@ -556,7 +562,44 @@ Retention: записи старше 30 дней удаляются при ст
6. При включенном `parser_dedupe_content_hash` также убирает дубли по `content_hash`.
7. Применяет политику мусорных постов.
8. Сохраняет raw post и media.
9. Создаёт job `vk.storage.copy`.
9. Создаёт job `media.storage.copy`.
### 7.1.1. Источники сайтов
Тот же scheduling loop обрабатывает `platform='site'`, но сетевую работу
делегирует внешнему stateless-модулю Site Parser. Основная коробка хранит БД,
источники, конфиги, расписание, статусы и очередь. Внешний модуль получает один
запрос и возвращает унифицированный JSON с текстом, исходной ссылкой, датой и
медиа. Готовые источники запускаются строго последовательно.
Настройки основной коробки:
- `site_parser_url`
- `site_parser_token`
- `site_parser_rucaptcha_token`
- `site_parser_timeout_sec`
- `site_parser_interval_minutes`
- `site_parser_min_text_length`
Персональный `min_text_length` в JSON источника переопределяет общий порог.
Проверяется только тело материала: заголовок не помогает пройти фильтр.
Записи без достаточного текста, включая video-only, не сохраняются и не
попадают в media uploader.
Конфиг v1 описывает `discovery` (`rss` или `html`), точные `detail.fields`,
`detail.media`, транспорт detail-страниц и ограниченные ретраи. Каждый field
может иметь ordered `candidates`: CSS/attribute/JSON/date fallback. Ошибка
одного detail возвращает `partial`, сохраняет успешных соседей и показывает
диагностику в статусе источника.
`access.type='cloudflare'` включает Playwright и RuCaptcha для конкретного
источника; обычные RSS и сайты идут без браузера. Полученные `cf_clearance` и
точный User-Agent браузера сохраняются в runtime state источника и повторно
используются media uploader при скачивании защищённых картинок.
Внешний модуль работает в LXC `108` (`192.168.1.113:8080`) через Docker
Compose `/opt/site-parser/docker-compose.yml`. Репозиторный compose находится
в `site_parser_worker/docker-compose.yml`.
### 7.2. Политика мусорных постов
@@ -615,13 +658,13 @@ Retention: записи старше 30 дней удаляются при ст
## 8. Media storage uploader
Файл: `src/vk_parser_app/workers/vk_storage_uploader.py`.
Файл: `src/vk_parser_app/workers/media_uploader.py`.
Название файла историческое: сейчас фактическое хранилище медиа - Telegram, а не VK.
### 8.1. Что делает
1. Забирает job типа `vk.storage.copy`.
1. Забирает job типа `media.storage.copy`.
2. Загружает `raw_post` и media.
3. Скачивает изображения/видео.
4. Отправляет в Telegram storage channel.
@@ -1091,6 +1134,22 @@ Worker учитывает:
- active
- parse_from
Для сайта важны:
- platform `site`
- name
- url
- active
- обязательный JSON-конфиг конкретного источника
Полный проверенный конфиг PopularAirsoft находится в
`site_parser_worker/README.md`. Для него используется HTML discovery главной
страницы, потому что RSS смешивает короткие video pages и полные новости.
Ошибки конфигурации, внешнего worker или сайта показываются в статусе
источника. После загрузки защищённые изображения отображаются в raw/editor по
авторизованному маршруту `/raw/media/{media_id}` из Telegram storage.
### 12.3. Raw
Raw-страница показывает:
@@ -1247,7 +1306,7 @@ uvicorn vk_parser_app.admin:app --host 0.0.0.0 --port 8080
```powershell
$env:PYTHONPATH="src"
python -m vk_parser_app.workers.parser
python -m vk_parser_app.workers.vk_storage_uploader
python -m vk_parser_app.workers.media_uploader
python -m vk_parser_app.workers.ai_qualifier
python -m vk_parser_app.workers.ai_writer
python -m vk_parser_app.workers.tg_poster
@@ -1376,16 +1435,10 @@ systemctl start ai-writer.service
Они не влияют на текущую очередь, но могут путать в админке. Их стоит аккуратно пометить как failed/cancelled или скрывать служебные тесты.
### 20.2. Название `vk_storage_uploader.py`
### 20.2. Media uploader
Файл исторически называется VK storage uploader, но фактически загружает в Telegram. Можно позже переименовать:
```text
vk_storage_uploader.py -> tg_storage_uploader.py
JOB_TYPE_VK_STORAGE_COPY -> tg.storage.copy
```
Делать только миграцией и аккуратно, чтобы не потерять jobs.
Uploader переименован в `media_uploader.py`, worker — в `media-uploader`, а jobs — в
`media.storage.copy`. Миграция сохраняет уже созданные задания.
### 20.3. Настройки Telegram uploader
@@ -1449,7 +1502,7 @@ JOB_TYPE_VK_STORAGE_COPY -> tg.storage.copy
sources
-> parser
-> raw_posts + raw_post_media
-> jobs(vk.storage.copy)
-> jobs(media.storage.copy)
-> Telegram storage uploader
-> status=storage_ready
-> AI qualifier
@@ -1485,7 +1538,7 @@ admin_sessions
```text
src/vk_parser_app/admin.py
src/vk_parser_app/workers/parser.py
src/vk_parser_app/workers/vk_storage_uploader.py
src/vk_parser_app/workers/media_uploader.py
src/vk_parser_app/workers/ai_qualifier.py
src/vk_parser_app/workers/ai_writer.py
src/vk_parser_app/workers/tg_poster.py
+25 -2
View File
@@ -1,6 +1,6 @@
# RAA deployment
Last updated: 2026-08-04.
Last updated: 2026-08-11.
RAA is not a separate code project anymore. The source of truth for code and
documentation is:
@@ -26,6 +26,7 @@ continue feature work there.
- Destination: `raa-fn8`, CTID `105`, IP `192.168.1.105`
- Edge proxy: Caddy in CTID `110` proxies `raa.panel.f-n8.ru` to Coolify Traefik at `192.168.1.105:80`
- Health check: `https://raa.panel.f-n8.ru/health`
- Current verified deploy: `f866a13a8a83598e26a7a6a3406b8bc1c01c9a80`
## Per-Project Settings
@@ -55,6 +56,10 @@ Database settings:
- `vk_poster_access_token` is the RAA community token
- `site_poster_enabled=false`
- `site_poster_provider=""` until a website adapter is configured
- Site parser settings are stored here as well:
`site_parser_url`, `site_parser_token`, `site_parser_rucaptcha_token`,
`site_parser_timeout_sec`, `site_parser_interval_minutes`, and
`site_parser_min_text_length`.
- Optional social text headers/footers:
- `tg_poster_header_text` / `tg_poster_footer_text`
- `vk_poster_header_text` / `vk_poster_footer_text`
@@ -70,7 +75,7 @@ Database settings:
As last checked on 2026-08-03:
- `vk-parser=true`
- `vk-storage-uploader=true`
- `media-uploader=true`
- `ai-qualifier=false`
- `ai-writer=false`
- `tg-poster=false`
@@ -83,6 +88,24 @@ AI qualifier/writer are configured with user-provided token, models, and prompts
but both the worker controls and double-safety settings were still off. Turning
them on will spend LLM tokens and process the pending RAA raw posts.
## Website Sources
Website sources use the existing `vk-parser` scheduling loop but are sent one
at a time to a stateless external worker in LXC `108` (`192.168.1.113:8080`).
All source records, JSON configs, schedules, errors, raw posts, and media state
remain in the RAA database. The worker is managed by Docker Compose from
`/opt/site-parser` and should be healthy as `site-parser-worker`.
PopularAirsoft source id `72` uses `https://popularairsoft.com/` and the full
v1 config from `site_parser_worker/README.md`. Its RSS is not used for
discovery because it mixes short video pages with full articles. HTML discovery
selects the `The Latest News` cards and the worker then opens every `/news/...`
detail. Verified live 2026-08-11: nine details parsed, zero errors.
Cloudflare session state is returned to the main app and is also reused by
`media-uploader` for protected images. Admin previews remain authenticated
Telegram-backed `/raw/media/{media_id}` URLs.
## RAA Categories
RAA categories live in the RAA `content_categories` table. The UI has only
+51 -10
View File
@@ -108,6 +108,47 @@ https://vk.com/wall-239548476_123
- админ/редактор группы может открыть и отредактировать пост;
- вложенное фото видно в браузере.
## Site Parser
Сайты обрабатывает отдельный stateless-воркер из `site_parser_worker`. Основная
коробка хранит источники, расписание, RuCaptcha-ключ и браузерное состояние, а
воркер возвращает нормализованные материалы в существующий пайплайн.
Сейчас поддерживается RSS/Atom, включая Cloudflare:
```json
{
"format": "rss",
"access": "auto",
"max_items": 20
}
```
Если RSS отдаёт только заголовки и ссылки, можно явно дочитывать страницы записей:
```json
{
"format": "rss",
"access": "cloudflare",
"max_items": 20,
"follow_links": true,
"content_selector": "article.full",
"text_selector": ".field--name-body",
"min_text_length": 50
}
```
В глобальном разделе `Site Parser` задаются URL воркера, его токен, ключ
RuCaptcha и таймаут. Внешний модуль запускается командой:
```bash
cd site_parser_worker
printf 'WORKER_TOKEN=replace-me\n' > .env
chmod 600 .env
docker-compose up -d --build
docker-compose ps
```
## Запуск локально
Админка:
@@ -122,10 +163,10 @@ PYTHONPATH=src .venv/bin/uvicorn vk_parser_app.admin:app --host 0.0.0.0 --port 8
PYTHONPATH=src .venv/bin/python -m vk_parser_app.workers.parser
```
VK storage uploader:
Media uploader:
```bash
PYTHONPATH=src .venv/bin/python -m vk_parser_app.workers.vk_storage_uploader
PYTHONPATH=src .venv/bin/python -m vk_parser_app.workers.media_uploader
```
## systemd units
@@ -166,17 +207,17 @@ RestartSec=10
WantedBy=multi-user.target
```
`vk-storage-uploader.service`:
`media-uploader.service`:
```ini
[Unit]
Description=VK Storage Uploader Worker
Description=Media Uploader Worker
After=network.target postgresql.service
[Service]
WorkingDirectory=/opt/vk-parser
Environment=PYTHONPATH=/opt/vk-parser/src
ExecStart=/opt/vk-parser/.venv/bin/python -m vk_parser_app.workers.vk_storage_uploader
ExecStart=/opt/vk-parser/.venv/bin/python -m vk_parser_app.workers.media_uploader
Restart=always
RestartSec=10
@@ -187,13 +228,13 @@ WantedBy=multi-user.target
## Текущий pipeline
```text
sources(vk)
sources(vk/site)
-> vk-parser
-> raw_posts + raw_post_media
-> jobs(type='vk.storage.copy')
-> vk-storage-uploader
-> wall.post в storage-группе
-> jobs(type='media.storage.copy')
-> media-uploader
-> локальный Telegram Bot API
-> raw_posts.storage_post_url
```
Фото копируются в storage-группу. Видео в первой версии сохраняются как `link_only`/исходное VK-вложение, чтобы сначала стабилизировать основной контур.
Фото и видео загружаются в Telegram media storage; исходные URL остаются в `raw_post_media`.
+60 -3
View File
@@ -1,6 +1,6 @@
# N8 Parser: current state
Last updated: 2026-08-03
Last updated: 2026-08-11
This file is the short handoff state for future Codex threads. Read this first before touching the project, so the whole chat history does not need to be carried forward.
@@ -27,7 +27,7 @@ Current runtime is local infrastructure, not the old VPS.
- Templates: `src/vk_parser_app/templates`
- Workers:
- VK parser: `src/vk_parser_app/workers/parser.py`
- Telegram media/storage uploader: `src/vk_parser_app/workers/vk_storage_uploader.py`
- Telegram media/storage uploader: `src/vk_parser_app/workers/media_uploader.py`
- AI qualifier: `src/vk_parser_app/workers/ai_qualifier.py`
- AI writer: `src/vk_parser_app/workers/ai_writer.py`
- DB migrations: `db/migrations`
@@ -37,6 +37,8 @@ Current runtime is local infrastructure, not the old VPS.
1. Sources are configured in admin.
2. VK parser reads enabled sources and stores raw posts.
The same worker also dispatches `platform='site'` sources to the external
site-parser module one at a time.
3. Media/storage uploader sends raw post copies to Telegram storage and stores media metadata.
4. AI qualifier scores raw posts and marks accepted/rejected/maybe.
5. AI writer rewrites accepted posts into editor-ready drafts.
@@ -66,7 +68,7 @@ canonical git repository as FN-8. Code work for both projects must happen in
- VK publication owner: `vk_poster_owner_id=-36860851`, `vk_poster_from_group=true`.
- `TG_BOT_TOKEN`, `TELEGRAM_API_ID`, and `TELEGRAM_API_HASH` were copied from FN-8 Coolify env.
- Local Bot API app setting: `local_bot_api_url=http://127.0.0.1:8081`.
- As last checked, RAA has `vk-parser=true`, `vk-storage-uploader=true`.
- As last checked, RAA has `vk-parser=true`, `media-uploader=true`.
AI and publication workers are intentionally disabled until explicitly started:
`ai-qualifier=false`, `ai-writer=false`, `tg-poster=false`, `vk-poster=false`,
`site-poster=false`, `tg_poster_enabled=false`, `vk_poster_enabled=false`.
@@ -75,6 +77,58 @@ canonical git repository as FN-8. Code work for both projects must happen in
(`raa-sender`) remains a valid legacy VK reposting app and is not the shared
RAA parser/writer/poster app.
## External Site Parser
As of 2026-08-10, website parsing is integrated into the normal raw-post
pipeline. The main application owns sources, per-source JSON configs,
schedules, settings, database state, deduplication, media jobs, and errors.
The external module is stateless: it receives one source request, fetches it,
and returns normalized JSON.
- Runtime: LXC `108`, IP `192.168.1.113`, port `8080`.
- Management: Docker Compose in `/opt/site-parser`.
- Repository compose file: `site_parser_worker/docker-compose.yml`.
- The module is protected by `WORKER_TOKEN`; the value is stored only in
deployment settings/secrets.
- RuCaptcha is stored in the main app setting
`site_parser_rucaptcha_token`; never copy the key into git documentation.
- Main settings are under `/workers` -> `Site Parser`: worker URL/token,
RuCaptcha token, request timeout, global interval, and default minimum body
length.
- Website sources are added on `/sources` with `platform='site'` and their own
JSON config. A missing or invalid config produces a visible source error and
does not run silently.
- Due website sources are processed sequentially, not launched together.
- `access='cloudflare'` enables the Playwright/RuCaptcha path for that source.
Plain RSS/sites do not use a browser.
- Browser cookies plus the exact browser User-Agent are returned to the main
app and reused by `media-uploader` when downloading protected images.
- Uploaded photos are previewed in the admin through authenticated
`/raw/media/{media_id}` URLs backed by Telegram storage, not by the original
Cloudflare-protected URL.
Site configs now use the versioned declarative v1 contract. The worker supports
RSS or HTML discovery, exact per-field selectors with ordered fallbacks,
detail-page extraction, media normalization, bounded retries, and per-item
errors. A failed detail does not discard successful siblings; the source is
marked with a visible partial diagnostic.
PopularAirsoft source id `72` should use `https://popularairsoft.com/`, not its
mixed-content RSS feed. The complete config is documented in
`site_parser_worker/README.md`. It discovers the `The Latest News` cards and
then extracts each `article.news.full` page. Verified live on 2026-08-11: nine
items, nine full details, zero errors. Twitch embeds return as
`external_video`; YouTube embeds are normalized for the existing downloader.
LXC 108 operations:
```bash
cd /opt/site-parser
docker-compose up -d --build
docker-compose ps
docker-compose logs -f
```
## Important Text Rules
- AI writer output text must be clean: no physical hashtags at the end.
@@ -185,6 +239,9 @@ As of 2026-07-29:
- Prefer small targeted file reads with `rg` and narrow ranges.
- Deployment status polling should be sparse. Coolify rebuilds can take 7-11 minutes
because the Docker build currently runs without cache and reinstalls system deps.
- LXC 105 Docker cleanup: old app images can be removed safely because rollback
is by git SHA + rebuild. Use `scripts/coolify_docker_cleanup.sh`; it removes
only unused Coolify app images for FN-8/RAA and never prunes Docker volumes.
- Avoid dumping long post texts from DB unless explicitly needed.
- For runtime checks, use the current Coolify/Gitea/LXC 105 path from `D:\DEVELOPMENT\.infra.md`.
- Do not `git reset --hard` or revert user/server hotfixes.
+10
View File
@@ -0,0 +1,10 @@
ALTER TABLE sources
ADD COLUMN IF NOT EXISTS runtime_state_json JSONB NOT NULL DEFAULT '{}'::jsonb;
INSERT INTO app_settings(key, value_json, value_type, title, description, category)
VALUES
('site_parser_url', '"http://192.168.1.113:8080"'::jsonb, 'str', 'URL воркера', 'Адрес внешнего Site Parser.', 'Site Parser'),
('site_parser_token', '""'::jsonb, 'secret', 'Токен воркера', 'Должен совпадать с WORKER_TOKEN внешнего модуля.', 'Site Parser'),
('site_parser_rucaptcha_token', '""'::jsonb, 'secret', 'Ключ RuCaptcha', 'Используется воркером только при встрече Cloudflare.', 'Site Parser'),
('site_parser_timeout_sec', '180'::jsonb, 'int', 'Таймаут, сек', 'Максимальное время одного запуска источника.', 'Site Parser')
ON CONFLICT (key) DO NOTHING;
@@ -0,0 +1,10 @@
INSERT INTO app_settings(key, value_json, value_type, title, description, category)
VALUES (
'site_parser_default_interval_minutes',
'30'::jsonb,
'int',
'Интервал парсинга сайтов, мин',
'Как часто Site Parser проверяет все активные сайты.',
'Site Parser'
)
ON CONFLICT (key) DO NOTHING;
@@ -0,0 +1,14 @@
INSERT INTO app_settings(key, value_json, value_type, title, description, category, updated_by)
SELECT
'site_parser_interval_minutes',
value_json,
'int',
'Интервал парсинга сайтов, мин',
'Как часто Site Parser проверяет все активные сайты.',
'Site Parser',
updated_by
FROM app_settings
WHERE key='site_parser_default_interval_minutes'
ON CONFLICT (key) DO NOTHING;
DELETE FROM app_settings WHERE key='site_parser_default_interval_minutes';
@@ -0,0 +1,9 @@
UPDATE app_settings
SET description='Максимальная высота видео VK и YouTube перед загрузкой в Telegram.',
updated_at=NOW()
WHERE key='video_max_height';
UPDATE app_settings
SET description='Таймаут скачивания видео VK и YouTube.',
updated_at=NOW()
WHERE key='uploader_yt_dlp_timeout_sec';
@@ -0,0 +1,10 @@
INSERT INTO app_settings(key, value_json, value_type, title, description, category)
VALUES (
'site_parser_min_text_length',
'50'::jsonb,
'int',
'Минимальная длина текста сайта',
'Материалы сайта с более коротким телом не попадают в raw и media uploader.',
'Site Parser'
)
ON CONFLICT (key) DO NOTHING;
@@ -0,0 +1,14 @@
UPDATE jobs
SET type='media.storage.copy',
updated_at=NOW()
WHERE type='vk.storage.copy';
UPDATE worker_controls
SET name='media-uploader',
updated_at=NOW()
WHERE name='vk-storage-uploader';
UPDATE worker_heartbeats
SET name='media-uploader',
updated_at=NOW()
WHERE name='vk-storage-uploader';
+4 -4
View File
@@ -11,14 +11,14 @@ from aiogram.types import FSInputFile
from vk_parser_app.constants import MEDIA_STATUS_FAILED, MEDIA_STATUS_LINK_ONLY
from vk_parser_app.db import get_pool
from vk_parser_app.workers.vk_storage_uploader import TMP_DIR, TMP_PREFIX, TelegramStorageUploader
from vk_parser_app.workers.media_uploader import TMP_DIR, TMP_PREFIX, MediaUploader
def jlog(event: str, **payload: Any) -> None:
print(json.dumps({"event": event, **payload}, ensure_ascii=False, default=str), flush=True)
async def load_pending_videos(worker: TelegramStorageUploader, limit: int) -> list[dict[str, Any]]:
async def load_pending_videos(worker: MediaUploader, limit: int) -> list[dict[str, Any]]:
rows = await worker.pool.fetch(
"""
SELECT m.*, rp.original_url AS post_url
@@ -38,7 +38,7 @@ async def load_pending_videos(worker: TelegramStorageUploader, limit: int) -> li
return [dict(row) for row in rows]
async def send_video_to_storage(worker: TelegramStorageUploader, media: dict[str, Any], path: str) -> None:
async def send_video_to_storage(worker: MediaUploader, media: dict[str, Any], path: str) -> None:
caption = (
f"Backfill video media #{int(media['id'])} for raw_post #{int(media['raw_post_id'])}\n"
f'<a href="{html.escape(str(media.get("post_url") or ""), quote=True)}">source post</a>'
@@ -58,7 +58,7 @@ async def send_video_to_storage(worker: TelegramStorageUploader, media: dict[str
async def main() -> None:
limit = int(os.getenv("BACKFILL_LIMIT", "200"))
worker = TelegramStorageUploader()
worker = MediaUploader()
await worker.init()
videos = await load_pending_videos(worker, limit)
jlog("start", count=len(videos), limit=limit)
+44
View File
@@ -0,0 +1,44 @@
#!/bin/sh
set -eu
# Safe cleanup for the Coolify Docker destination.
# Removes old app images that are not used by running containers.
# Volumes are intentionally never pruned: they hold Postgres/app data.
KEEP_PER_REPO="${KEEP_PER_REPO:-1}"
CONTAINER_UNTIL="${CONTAINER_UNTIL:-24h}"
BUILDER_UNTIL="${BUILDER_UNTIL:-24h}"
if [ "$#" -gt 0 ]; then
REPOS="$*"
else
REPOS="n6cnr60anmruiwkiey0pukye korokrhpoyqxf0nw2y8vk7lp"
fi
used_images="$(docker ps --format '{{.Image}}' | sort -u)"
for repo in $REPOS; do
count=0
docker image ls "$repo" --format '{{.Repository}}:{{.Tag}} {{.ID}}' |
while read -r image image_id; do
[ -n "$image" ] || continue
if printf '%s\n' "$used_images" | grep -qx "$image"; then
count=$((count + 1))
echo "keep running image: $image"
continue
fi
count=$((count + 1))
if [ "$count" -le "$KEEP_PER_REPO" ]; then
echo "keep recent image: $image"
continue
fi
echo "remove old image: $image"
docker rmi "$image_id" || true
done
done
docker container prune -f --filter "until=$CONTAINER_UNTIL"
docker image prune -f
docker builder prune -f --filter "until=$BUILDER_UNTIL"
+11
View File
@@ -0,0 +1,11 @@
FROM mcr.microsoft.com/playwright/python:v1.62.0-noble
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
RUN playwright install chrome
COPY app.py extractor.py ./
ENV PYTHONUNBUFFERED=1
EXPOSE 8080
CMD ["bash", "-lc", "Xvfb :99 -screen 0 1280x1024x24 -nolisten tcp & export DISPLAY=:99; sleep 1; exec uvicorn app:app --host 0.0.0.0 --port 8080"]
+146
View File
@@ -0,0 +1,146 @@
# Site Parser Worker
The worker is a stateless executor. The main application owns source configs,
schedules, database state, and queues. Each request gives the worker a source
URL, a versioned JSON instruction, and optional Cloudflare session state. The
worker returns normalized items plus per-item diagnostics.
## Config v1
- `discovery.type`: `rss` or `html`.
- `discovery.fields`: where list/RSS values are located.
- `detail.enabled`: open every discovered item when enabled.
- `detail.root_selector`: limits extraction to the actual article.
- `detail.fields`: exact field rules. A rule may contain ordered `candidates`.
- `detail.media`: exact selectors and fallback attributes for media.
- `detail.transport`: `http`, `fetch`, or `browser`.
- `access.type`: `auto`, `http`, `browser`, or `cloudflare`.
- `retry`: bounded attempts, delay, and per-attempt timeout.
Field extraction modes are `text`, `html`, `attr`, and `json`. Date rules can
provide `formats` and `timezone`. Relative links are resolved against the page
URL. Unknown iframe/video providers are returned as `external_video`; they do
not fail the item.
Errors are isolated to one item. Successful items are returned with
`status=partial`; the source receives a visible diagnostic in the main app.
Only discovery/network failure prevents the whole request from returning a
batch.
## PopularAirsoft
Use `https://popularairsoft.com/` as the source URL. Its RSS mixes full news
articles with short video pages, so the reliable discovery surface is the
`The Latest News` HTML block. The worker opens the main page through
Cloudflare, collects `/news/...` links, and downloads each article using the
same clearance and an impersonated browser TLS fingerprint.
```json
{
"version": 1,
"discovery": {
"type": "html",
"item_selector": "#block-views-block-latest-news-list-block-2 .feature-contents, #block-views-block-latest-news-list-block-1 .lt-teasure",
"limit": 10,
"fields": {
"url": {
"selector": "a.link-title",
"extract": "attr",
"attribute": "href",
"required": true
},
"external_id": {
"selector": "a.link-title",
"extract": "attr",
"attribute": "href",
"required": true
},
"title": {
"selector": "a.link-title",
"extract": "text",
"required": true
},
"published_at": {
"selector": "time[datetime]",
"extract": "attr",
"attribute": "datetime",
"required": true
}
},
"media": [
{
"type": "photo",
"selector": ".site-image img",
"attributes": ["src", "data-src", "srcset"]
}
]
},
"detail": {
"enabled": true,
"always": true,
"transport": "http",
"root_selector": "article.news.full",
"fields": {
"title": {
"candidates": [
{
"selector": ".feature-contents > .news-story-texts:first-child h2",
"extract": "text"
},
{"source": "list.title"}
],
"required": true
},
"published_at": {
"candidates": [
{
"selector": ".feature-contents > .news-story-texts:first-child .news-story-date",
"extract": "text",
"formats": ["%d %b %Y"],
"timezone": "UTC"
},
{"source": "list.published_at"}
],
"required": true
},
"author": {
"selector": ".feature-contents > .news-story-texts:first-child h4",
"extract": "text"
},
"text": {
"selector": ".feature-contents > .news-story-texts:last-child .field--name-body",
"extract": "text",
"required": true,
"min_length": 50
}
},
"media": [
{
"type": "photo",
"selector": ".feature-contents > .site-image .field--name-field-image img, .field--name-body img",
"attributes": ["src", "data-src", "srcset"]
},
{
"type": "video",
"selector": ".field--name-body iframe, .field--name-body video, .field--name-body source",
"attributes": ["src", "data-src"]
}
]
},
"access": {
"type": "cloudflare",
"wait_for": "#block-views-block-latest-news-list-block-1"
},
"retry": {
"attempts": 2,
"delay_seconds": 2,
"timeout_seconds": 90
},
"min_text_length": 50
}
```
Verified against the live site on 2026-08-11: nine items, nine successful
details, zero errors. The VFC sample returned 5,041 text characters, one photo,
and two Twitch links. The Double Bell sample returned 1,293 characters, one
photo, and one normalized YouTube URL.
+488
View File
@@ -0,0 +1,488 @@
from __future__ import annotations
import asyncio
import logging
import os
import secrets
from contextlib import asynccontextmanager
from datetime import datetime, timezone
from time import struct_time
from typing import Annotated, Any, AsyncIterator
from urllib.parse import urljoin
import feedparser
import httpx
from bs4 import BeautifulSoup
from curl_cffi.requests import AsyncSession as CurlAsyncSession
from fastapi import FastAPI, Header, HTTPException
from playwright.async_api import BrowserContext, Page, TimeoutError as PlaywrightTimeoutError, async_playwright
from playwright_captcha import CaptchaType, FrameworkType, TwoCaptchaSolver
from pydantic import BaseModel, Field, HttpUrl, SecretStr
from twocaptcha import AsyncTwoCaptcha
from extractor import ConfigError, detail_from_html, extract_fields, extract_media, normalize_config
app = FastAPI(title="Site Parser Worker", version="1.0.0")
browser_lock = asyncio.Lock()
logger = logging.getLogger("site_parser")
class ParseRequest(BaseModel):
url: HttpUrl
config: dict[str, Any]
rucaptcha_token: SecretStr | None = None
browser_state: dict[str, Any] | None = None
browser_user_agent: str | None = None
class ParsedItem(BaseModel):
external_id: str
url: str
title: str
text: str
html: str
published_at: str | None
author: str | None
media: list[dict[str, str]] = Field(default_factory=list)
class CloudflareRequired(RuntimeError):
pass
class PermanentPageError(RuntimeError):
pass
def require_token(worker_token: str | None) -> None:
expected = os.getenv("WORKER_TOKEN", "")
if not expected:
raise HTTPException(status_code=503, detail="WORKER_TOKEN is not configured")
if not worker_token or not secrets.compare_digest(worker_token, expected):
raise HTTPException(status_code=401, detail="Invalid worker token")
def is_cloudflare_challenge(status: int, body: str) -> bool:
sample = body[:20_000].lower()
return status in {403, 429, 503} and (
"just a moment" in sample or "cf-chl" in sample or "challenge-platform" in sample
)
def iso_date(value: struct_time | None) -> str | None:
if not value:
return None
return datetime(*value[:6], tzinfo=timezone.utc).isoformat()
def clean_feed_content(value: Any, base_url: str) -> tuple[str, str, list[dict[str, str]]]:
html = str(value or "").strip()
soup = BeautifulSoup(html, "html.parser")
for element in soup.find_all(["script", "style", "noscript"]):
element.decompose()
text = soup.get_text("\n", strip=True)
media = extract_media(
soup,
[
{"type": "photo", "selector": "img", "attributes": ["src", "data-src", "srcset"]},
{"type": "video", "selector": "iframe, video, source", "attributes": ["src", "data-src"]},
],
base_url,
)
return html, text, media
def default_rss_fields() -> dict[str, Any]:
return {
"url": {"candidates": [{"source": "rss.link"}], "required": True},
"external_id": {
"candidates": [{"source": "rss.id"}, {"source": "rss.guid"}, {"source": "rss.link"}],
"required": True,
},
"title": {"candidates": [{"source": "rss.title"}]},
"published_at": {
"candidates": [
{"source": "rss.published"},
{"source": "rss.updated"},
]
},
"author": {"candidates": [{"source": "rss.author"}]},
}
def discover_rss(xml: str, config: dict[str, Any]) -> list[dict[str, Any]]:
feed = feedparser.parse(xml, sanitize_html=False)
if feed.bozo and not feed.entries:
raise ValueError(f"Invalid RSS: {feed.bozo_exception}")
discovery = config["discovery"]
fields = discovery.get("fields") or default_rss_fields()
items: list[dict[str, Any]] = []
for entry in feed.entries[: discovery["limit"]]:
raw = dict(entry)
values, errors = extract_fields(None, fields, {"rss": raw})
url = str(values.get("url") or raw.get("link") or "").strip()
external_id = str(values.get("external_id") or raw.get("id") or raw.get("guid") or url).strip()
if not url or not external_id:
items.append({"url": url, "discovery_errors": errors or [{"field": "url", "error": "URL not found"}]})
continue
content = (raw.get("content") or [{}])[0].get("value") or raw.get("description") or raw.get("summary") or ""
html, text, media = clean_feed_content(content, url)
for enclosure in raw.get("enclosures") or []:
enclosure_url = urljoin(url, str(enclosure.get("href") or enclosure.get("url") or "").strip())
media_type = str(enclosure.get("type") or "")
if enclosure_url and media_type.startswith("video/"):
media.append({"type": "video", "url": enclosure_url})
elif enclosure_url:
media.append({"type": "photo", "url": enclosure_url})
items.append(
{
"external_id": external_id,
"url": url,
"title": str(values.get("title") or BeautifulSoup(str(raw.get("title") or ""), "html.parser").get_text(" ", strip=True)),
"text": text,
"html": html,
"published_at": values.get("published_at") or iso_date(raw.get("published_parsed") or raw.get("updated_parsed")),
"author": values.get("author") or str(raw.get("author") or "").strip() or None,
"media": list({entry["url"]: entry for entry in media}.values()),
"source_context": {"rss": raw},
"discovery_errors": errors,
}
)
return items
def discover_html(html: str, source_url: str, config: dict[str, Any]) -> list[dict[str, Any]]:
soup = BeautifulSoup(html, "html.parser")
discovery = config["discovery"]
fields = discovery.get("fields") or {}
items: list[dict[str, Any]] = []
for element in soup.select(str(discovery["item_selector"]))[: discovery["limit"]]:
values, errors = extract_fields(element, fields, {"page": {"url": source_url}})
url = urljoin(source_url, str(values.get("url") or "").strip())
external_id = str(values.get("external_id") or url).strip()
if not url or not external_id:
items.append({"url": url, "discovery_errors": errors or [{"field": "url", "error": "URL not found"}]})
continue
items.append(
{
"external_id": external_id,
"url": url,
"title": str(values.get("title") or ""),
"text": str(values.get("text") or ""),
"html": "",
"published_at": values.get("published_at"),
"author": values.get("author"),
"media": extract_media(element, discovery.get("media") or [], source_url),
"source_context": {"list": values},
"discovery_errors": errors,
}
)
return items
def apply_detail(item: dict[str, Any], html: str, config: dict[str, Any]) -> tuple[dict[str, Any], list[dict[str, str]]]:
values, errors = detail_from_html(html, item["url"], config, item.get("source_context"))
if errors:
return item, errors
for key in ("title", "text", "published_at", "author", "html"):
if values.get(key) not in {None, ""}:
item[key] = values[key]
item["media"] = list(
{(entry["type"], entry["url"]): entry for entry in [*item.get("media", []), *values.get("media", [])]}.values()
)
return item, []
def public_item(item: dict[str, Any]) -> ParsedItem:
return ParsedItem(
external_id=str(item["external_id"]),
url=str(item["url"]),
title=str(item.get("title") or ""),
text=str(item.get("text") or ""),
html=str(item.get("html") or ""),
published_at=str(item["published_at"]) if item.get("published_at") else None,
author=str(item["author"]) if item.get("author") else None,
media=item.get("media") or [],
)
async def retry(operation, attempts: int, delay: float, timeout: float):
last_error: Exception | None = None
for attempt in range(1, attempts + 1):
try:
return await asyncio.wait_for(operation(), timeout=timeout), attempt
except PermanentPageError:
raise
except Exception as exc:
last_error = exc
if attempt < attempts and delay:
await asyncio.sleep(delay)
raise RuntimeError(str(last_error or "operation failed")) from last_error
async def http_document(client: httpx.AsyncClient, url: str) -> str:
response = await client.get(url)
if is_cloudflare_challenge(response.status_code, response.text):
raise CloudflareRequired("Cloudflare challenge received")
response.raise_for_status()
return response.text
async def curl_document(client: CurlAsyncSession, url: str) -> str:
response = await client.get(url, allow_redirects=True, timeout=30)
body = response.text
if is_cloudflare_challenge(response.status_code, body) or "Attention Required" in body[:10_000]:
raise CloudflareRequired("Cloudflare blocked impersonated request")
if response.status_code < 200 or response.status_code >= 300:
raise RuntimeError(f"detail returned HTTP {response.status_code}")
return body
async def enrich_items(
items: list[dict[str, Any]],
config: dict[str, Any],
document_loader,
) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]:
attempts = config["retry"]["attempts"]
delay = config["retry"]["delay_seconds"]
timeout = config["retry"]["timeout_seconds"]
errors: list[dict[str, Any]] = []
enriched = []
for item in items:
if item.get("discovery_errors"):
errors.append({"url": item.get("url"), "stage": "discovery", "attempts": 1, "errors": item["discovery_errors"]})
continue
if not config["detail"]["always"] and item.get("text") and item.get("media"):
enriched.append(item)
continue
try:
html, used_attempts = await retry(lambda item=item: document_loader(item["url"]), attempts, delay, timeout)
item, field_errors = apply_detail(item, html, config)
if field_errors:
errors.append({"url": item["url"], "stage": "detail", "attempts": used_attempts, "errors": field_errors})
continue
enriched.append(item)
logger.info("Parsed detail %s", item["url"])
except CloudflareRequired:
raise
except Exception as exc:
errors.append({"url": item.get("url"), "stage": "detail", "attempts": attempts, "error": str(exc)})
logger.warning("Detail failed %s: %s", item.get("url"), exc)
return enriched, errors
async def parse_over_http(source_url: str, config: dict[str, Any]) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]:
attempts = config["retry"]["attempts"]
delay = config["retry"]["delay_seconds"]
timeout = config["retry"]["timeout_seconds"]
async with httpx.AsyncClient(
follow_redirects=True,
timeout=30,
headers={"User-Agent": "Mozilla/5.0 SiteParser/1.0"},
) as client:
document, _ = await retry(lambda: http_document(client, source_url), attempts, delay, timeout)
items = discover_rss(document, config) if config["discovery"]["type"] == "rss" else discover_html(document, source_url, config)
if config["detail"]["enabled"]:
return await enrich_items(items, config, lambda url: http_document(client, url))
return items, []
@asynccontextmanager
async def captcha_solver(page: Page, token: str | None) -> AsyncIterator[TwoCaptchaSolver | None]:
if not token:
yield None
return
async with TwoCaptchaSolver(
framework=FrameworkType.PLAYWRIGHT,
page=page,
async_two_captcha_client=AsyncTwoCaptcha(token),
max_attempts=1,
) as solver:
yield solver
async def browser_document(
page: Page,
url: str,
solver: TwoCaptchaSolver | None,
wait_for: str | None = None,
) -> str:
await page.goto(url, wait_until="domcontentloaded", timeout=60_000)
title = (await page.title()).lower()
if "just a moment" in title or "attention required" in title:
if solver is None:
raise RuntimeError("Cloudflare challenge received but RuCaptcha is not configured")
await solver.solve_captcha(
captcha_container=page,
captcha_type=CaptchaType.CLOUDFLARE_INTERSTITIAL,
)
if wait_for:
try:
await page.locator(wait_for).first.wait_for(state="attached", timeout=20_000)
except PlaywrightTimeoutError as exc:
raise PermanentPageError(
f"selector not found: {wait_for}; url={page.url}; title={await page.title()}"
) from exc
return await page.content()
async def browser_discovery_document(page: Page, source_url: str, config: dict[str, Any], solver) -> str:
wait_for = config["access"].get("wait_for") if config["discovery"]["type"] == "html" else None
html = await browser_document(page, source_url, solver, wait_for)
if config["discovery"]["type"] != "rss":
return html
pre = page.locator("pre").first
await pre.wait_for(state="attached", timeout=60_000)
return await pre.inner_text()
async def browser_fetch_document(page: Page, url: str) -> str:
result = await page.evaluate(
"""
async (url) => {
const response = await fetch(url, {credentials: 'include'});
return {status: response.status, text: await response.text()};
}
""",
url,
)
status = int(result.get("status") or 0)
body = str(result.get("text") or "")
if is_cloudflare_challenge(status, body) or "Attention Required" in body[:10_000]:
raise CloudflareRequired("Cloudflare blocked browser fetch")
if status < 200 or status >= 300:
raise RuntimeError(f"detail returned HTTP {status}")
return body
async def parse_in_browser(
source_url: str,
config: dict[str, Any],
rucaptcha_token: str | None,
browser_state: dict[str, Any] | None,
browser_user_agent: str | None,
) -> tuple[list[dict[str, Any]], list[dict[str, Any]], dict[str, Any], str]:
attempts = config["retry"]["attempts"]
delay = config["retry"]["delay_seconds"]
timeout = config["retry"]["timeout_seconds"]
errors: list[dict[str, Any]] = []
async with browser_lock:
async with async_playwright() as playwright:
browser = await playwright.chromium.launch(channel="chrome", headless=False)
context: BrowserContext = await browser.new_context(storage_state=browser_state, user_agent=browser_user_agent)
page = await context.new_page()
async with captcha_solver(page, rucaptcha_token) as solver:
document, _ = await retry(
lambda: browser_discovery_document(page, source_url, config, solver), attempts, delay, timeout
)
items = discover_rss(document, config) if config["discovery"]["type"] == "rss" else discover_html(document, source_url, config)
logger.info("Discovered %s items from %s", len(items), source_url)
if config["detail"]["enabled"]:
if config["detail"]["transport"] == "http":
original_state = browser_state if isinstance(browser_state, dict) else {}
original_cookies = original_state.get("cookies") or []
cookies = {
str(cookie["name"]): str(cookie["value"])
for cookie in (original_cookies or await context.cookies())
if cookie.get("name") and cookie.get("value")
}
user_agent = browser_user_agent or await page.evaluate("navigator.userAgent")
async with CurlAsyncSession(
impersonate="chrome",
cookies=cookies,
headers={"User-Agent": str(user_agent), "Referer": source_url},
) as client:
items, errors = await enrich_items(items, config, lambda url: curl_document(client, url))
elif config["detail"]["transport"] == "fetch":
items, errors = await enrich_items(items, config, lambda url: browser_fetch_document(page, url))
else:
enriched = []
wait_for = str(config["detail"].get("wait_for") or config["detail"]["root_selector"])
for item in items:
if item.get("discovery_errors"):
errors.append({"url": item.get("url"), "stage": "discovery", "attempts": 1, "errors": item["discovery_errors"]})
continue
if not config["detail"]["always"] and item.get("text") and item.get("media"):
enriched.append(item)
continue
try:
html, used_attempts = await retry(
lambda item=item: browser_document(page, item["url"], solver, wait_for), attempts, delay, timeout
)
item, field_errors = apply_detail(item, html, config)
if field_errors:
errors.append({"url": item["url"], "stage": "detail", "attempts": used_attempts, "errors": field_errors})
continue
enriched.append(item)
logger.info("Parsed detail %s", item["url"])
except Exception as exc:
logger.warning("Detail failed %s: %s", item.get("url"), exc)
errors.append({"url": item.get("url"), "stage": "detail", "attempts": attempts, "error": str(exc)})
items = enriched
state = await context.storage_state()
user_agent = await page.evaluate("navigator.userAgent")
await browser.close()
return items, errors, state, user_agent
@app.get("/health")
async def health() -> dict[str, bool]:
return {"ok": True}
@app.post("/v1/parse")
async def parse_source(
request: ParseRequest,
worker_token: Annotated[str | None, Header(alias="X-Worker-Token")] = None,
) -> dict[str, Any]:
require_token(worker_token)
try:
config = normalize_config(request.config)
except ConfigError as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
source_url = str(request.url)
access = config["access"]["type"]
state = request.browser_state
browser_user_agent = request.browser_user_agent
fetched_via = "http"
logger.info("Parse request source=%s discovery=%s detail=%s", source_url, config["discovery"]["type"], config["detail"]["enabled"])
try:
if access in {"browser", "cloudflare"}:
items, errors, state, browser_user_agent = await parse_in_browser(
source_url,
config,
request.rucaptcha_token.get_secret_value() if request.rucaptcha_token else None,
state,
browser_user_agent,
)
fetched_via = "browser"
else:
try:
items, errors = await parse_over_http(source_url, config)
except CloudflareRequired:
if access == "http":
raise
items, errors, state, browser_user_agent = await parse_in_browser(
source_url,
config,
request.rucaptcha_token.get_secret_value() if request.rucaptcha_token else None,
state,
browser_user_agent,
)
fetched_via = "browser"
except Exception as exc:
raise HTTPException(status_code=502, detail=str(exc)) from exc
if not items and not errors:
errors = [{"url": source_url, "stage": "discovery", "attempts": 1, "error": "no items found"}]
status = "partial" if errors and items else "failed" if errors else "ok"
return {
"status": status,
"source_url": source_url,
"fetched_via": fetched_via,
"items": [public_item(item).model_dump() for item in items],
"errors": errors,
"browser_state": state,
"browser_user_agent": browser_user_agent,
}
+21
View File
@@ -0,0 +1,21 @@
version: "3.8"
services:
site-parser:
build: .
container_name: site-parser-worker
restart: unless-stopped
env_file: .env
ports:
- "8080:8080"
shm_size: 512mb
healthcheck:
test:
- CMD
- python
- -c
- import urllib.request; urllib.request.urlopen('http://127.0.0.1:8080/health', timeout=5)
interval: 10s
timeout: 5s
retries: 6
start_period: 20s
+357
View File
@@ -0,0 +1,357 @@
from __future__ import annotations
import json
import re
from datetime import datetime, timezone
from email.utils import parsedate_to_datetime
from typing import Any
from urllib.parse import parse_qs, urljoin, urlparse
from zoneinfo import ZoneInfo
from bs4 import BeautifulSoup, Tag
class ConfigError(ValueError):
pass
def normalize_config(value: dict[str, Any]) -> dict[str, Any]:
if not isinstance(value, dict) or not value:
raise ConfigError("config is required and cannot be empty")
if "discovery" not in value:
follow_links = value.get("follow_links", False)
config = {
"version": 1,
"discovery": {"type": "rss", "limit": value.get("max_items", 20)},
"detail": {
"enabled": follow_links,
"always": follow_links,
"root_selector": value.get("content_selector", ""),
"fields": {
"text": {
"selector": value.get("text_selector") or ":root",
"extract": "text",
"required": True,
}
},
"media": [
{"type": "photo", "selector": "img", "attributes": ["src", "data-src"]},
{"type": "video", "selector": "iframe, video, source", "attributes": ["src", "data-src"]},
],
},
"access": {"type": value.get("access", "auto")},
"retry": {"attempts": 1, "delay_seconds": 0},
"min_text_length": value.get("min_text_length", 0),
}
else:
config = dict(value)
if config.get("version", 1) != 1:
raise ConfigError("only config.version=1 is supported")
discovery = config.get("discovery")
if not isinstance(discovery, dict) or discovery.get("type") not in {"rss", "html"}:
raise ConfigError("discovery.type must be rss or html")
try:
limit = int(discovery.get("limit", 20))
except (TypeError, ValueError) as exc:
raise ConfigError("discovery.limit must be an integer") from exc
if not 1 <= limit <= 100:
raise ConfigError("discovery.limit must be between 1 and 100")
discovery["limit"] = limit
if discovery["type"] == "html" and not str(discovery.get("item_selector") or "").strip():
raise ConfigError("discovery.item_selector is required for html discovery")
access = config.get("access", {"type": "auto"})
if isinstance(access, str):
access = {"type": access}
if not isinstance(access, dict) or access.get("type", "auto") not in {"auto", "http", "browser", "cloudflare"}:
raise ConfigError("access.type must be auto, http, browser or cloudflare")
config["access"] = access
detail = config.get("detail") or {"enabled": False}
if not isinstance(detail, dict):
raise ConfigError("detail must be an object")
detail["enabled"] = bool(detail.get("enabled", False))
detail["always"] = bool(detail.get("always", detail["enabled"]))
if detail["enabled"]:
if not str(detail.get("root_selector") or "").strip():
raise ConfigError("detail.root_selector is required when detail is enabled")
if not isinstance(detail.get("fields") or {}, dict):
raise ConfigError("detail.fields must be an object")
if not isinstance(detail.get("media") or [], list):
raise ConfigError("detail.media must be an array")
if detail.get("transport", "browser") not in {"http", "fetch", "browser"}:
raise ConfigError("detail.transport must be http, fetch or browser")
detail["transport"] = detail.get("transport", "browser")
config["detail"] = detail
retry = config.get("retry") or {}
try:
attempts = int(retry.get("attempts", 2))
delay = float(retry.get("delay_seconds", 2))
timeout = float(retry.get("timeout_seconds", 90))
except (TypeError, ValueError) as exc:
raise ConfigError("retry values must be numeric") from exc
if not 1 <= attempts <= 5 or not 0 <= delay <= 60 or not 10 <= timeout <= 300:
raise ConfigError("retry attempts must be 1..5, delay_seconds 0..60 and timeout_seconds 10..300")
config["retry"] = {"attempts": attempts, "delay_seconds": delay, "timeout_seconds": timeout}
try:
min_length = int(config.get("min_text_length", 0))
except (TypeError, ValueError) as exc:
raise ConfigError("min_text_length must be an integer") from exc
if not 0 <= min_length <= 100_000:
raise ConfigError("min_text_length must be between 0 and 100000")
config["min_text_length"] = min_length
return config
def nested_value(value: Any, path: str) -> Any:
current = value
for part in path.split("."):
if isinstance(current, dict):
current = current.get(part)
else:
return None
return current
def json_path_values(value: Any, path: str) -> list[Any]:
parts = path.split(".")
found: list[Any] = []
def visit(node: Any, index: int) -> None:
if index == len(parts):
found.append(node)
return
if isinstance(node, list):
for entry in node:
visit(entry, index)
elif isinstance(node, dict):
if parts[index] in node:
visit(node[parts[index]], index + 1)
for entry in node.values():
if isinstance(entry, (dict, list)):
visit(entry, index)
visit(value, 0)
return found
def candidate_elements(root: Tag | BeautifulSoup, candidate: dict[str, Any]) -> list[Tag]:
selectors = candidate.get("selectors") or candidate.get("selector") or []
if isinstance(selectors, str):
selectors = [selectors]
elements: list[Tag] = []
for selector in selectors:
if selector == ":root":
elements.append(root)
elif str(selector).strip():
elements.extend(root.select(str(selector)))
return elements
def element_value(element: Tag, candidate: dict[str, Any]) -> Any:
mode = str(candidate.get("extract") or "text")
if mode == "text":
return element.get_text("\n", strip=True)
if mode == "html":
return element.decode_contents()
if mode == "attr":
attributes = candidate.get("attributes") or candidate.get("attribute") or []
if isinstance(attributes, str):
attributes = [attributes]
for attribute in attributes:
value = element.get(str(attribute))
if value:
return value
return None
if mode == "json":
try:
payload = json.loads(element.string or element.get_text("", strip=True))
except (TypeError, json.JSONDecodeError):
return None
values = json_path_values(payload, str(candidate.get("path") or ""))
return next((entry for entry in values if entry is not None and entry != ""), None)
raise ConfigError(f"unsupported extract mode: {mode}")
def parse_date(value: Any, candidate: dict[str, Any]) -> str | None:
if not value:
return None
raw = str(value).strip()
try:
parsed = datetime.fromisoformat(raw.replace("Z", "+00:00"))
except ValueError:
try:
parsed = parsedate_to_datetime(raw)
except (TypeError, ValueError, OverflowError):
parsed = None
if parsed is None:
formats = candidate.get("formats") or candidate.get("date_format") or []
if isinstance(formats, str):
formats = [formats]
for date_format in formats:
try:
parsed = datetime.strptime(raw, str(date_format))
break
except ValueError:
continue
if parsed is None:
return None
if parsed.tzinfo is None:
try:
parsed = parsed.replace(tzinfo=ZoneInfo(str(candidate.get("timezone") or "UTC")))
except Exception:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc).isoformat()
def apply_regex(value: Any, candidate: dict[str, Any]) -> Any:
pattern = candidate.get("regex")
if not pattern or value is None:
return value
match = re.search(str(pattern), str(value), flags=re.DOTALL)
if not match:
return None
group = candidate.get("group", 1 if match.lastindex else 0)
try:
return match.group(group)
except (IndexError, KeyError):
return None
def extract_field(
root: Tag | BeautifulSoup | None,
rule: dict[str, Any],
source: dict[str, Any] | None = None,
*,
field_name: str = "",
) -> Any:
candidates = rule.get("candidates") or [rule]
for candidate in candidates:
if not isinstance(candidate, dict):
continue
values: list[Any] = []
if candidate.get("source"):
values.append(nested_value(source or {}, str(candidate["source"])))
elif root is not None:
values.extend(element_value(element, candidate) for element in candidate_elements(root, candidate))
for value in values:
value = apply_regex(value, candidate)
if value is None or (isinstance(value, str) and not value.strip()):
continue
if field_name == "published_at" or candidate.get("type") == "date":
value = parse_date(value, candidate)
if not value:
continue
if isinstance(value, str):
value = value.strip()
if len(str(value)) < int(rule.get("min_length", 0)):
continue
return value
return None
def extract_fields(
root: Tag | BeautifulSoup | None,
fields: dict[str, Any],
source: dict[str, Any] | None = None,
) -> tuple[dict[str, Any], list[dict[str, str]]]:
values: dict[str, Any] = {}
errors: list[dict[str, str]] = []
for name, rule in fields.items():
if not isinstance(rule, dict):
errors.append({"field": name, "error": "field rule must be an object"})
continue
try:
value = extract_field(root, rule, source, field_name=name)
except Exception as exc:
errors.append({"field": name, "error": str(exc)})
continue
if value is None and rule.get("required"):
errors.append({"field": name, "error": "required field not found"})
elif value is not None:
values[name] = value
return values, errors
def youtube_url(value: str, base_url: str = "") -> str | None:
url = urljoin(base_url, str(value or "").strip())
parsed = urlparse(url)
host = (parsed.hostname or "").lower().removeprefix("www.").removeprefix("m.")
video_id = ""
if host == "youtu.be":
video_id = parsed.path.strip("/").split("/", 1)[0]
elif host in {"youtube.com", "youtube-nocookie.com"}:
if parsed.path == "/watch":
video_id = (parse_qs(parsed.query).get("v") or [""])[0]
elif parsed.path.startswith(("/embed/", "/shorts/", "/live/")):
video_id = parsed.path.strip("/").split("/", 1)[1]
if not re.fullmatch(r"[A-Za-z0-9_-]{6,20}", video_id):
return None
return f"https://www.youtube.com/watch?v={video_id}"
def srcset_urls(value: str) -> list[str]:
return [entry.strip().split()[0] for entry in value.split(",") if entry.strip()]
def extract_media(root: Tag | BeautifulSoup, specs: list[dict[str, Any]], base_url: str) -> list[dict[str, str]]:
media: list[dict[str, str]] = []
for spec in specs:
if not isinstance(spec, dict):
continue
media_type = str(spec.get("type") or "photo")
attributes = spec.get("attributes") or spec.get("attribute") or ["src"]
if isinstance(attributes, str):
attributes = [attributes]
for element in candidate_elements(root, spec):
raw_values: list[str] = []
for attribute in attributes:
raw = element.get(str(attribute))
if not raw:
continue
raw_values.extend(srcset_urls(str(raw)) if attribute == "srcset" else [str(raw)])
break
for raw in raw_values:
url = urljoin(base_url, raw.strip())
if not url.startswith(("http://", "https://")):
continue
item_type = media_type
provider = ""
if media_type == "video":
normalized = youtube_url(url)
if normalized:
url, provider = normalized, "youtube"
else:
host = (urlparse(url).hostname or "").lower()
provider = "twitch" if "twitch.tv" in host else host
if element.name in {"iframe", "a"} or provider == "twitch":
item_type = "external_video"
item = {"type": item_type, "url": url}
if provider:
item["provider"] = provider
media.append(item)
return list({(item["type"], item["url"]): item for item in media}.values())
def detail_from_html(
html: str,
url: str,
config: dict[str, Any],
discovery_values: dict[str, Any] | None = None,
) -> tuple[dict[str, Any], list[dict[str, str]]]:
soup = BeautifulSoup(html, "html.parser")
detail = config["detail"]
root = soup.select_one(str(detail["root_selector"]))
if root is None:
return {}, [{"field": "detail", "error": f"root selector not found: {detail['root_selector']}"}]
for selector in detail.get("remove_selectors") or []:
for element in root.select(str(selector)):
element.decompose()
values, errors = extract_fields(root, detail.get("fields") or {}, discovery_values)
values["media"] = extract_media(root, detail.get("media") or [], url)
values["html"] = str(root)
return values, errors
+9
View File
@@ -0,0 +1,9 @@
2captcha-python-async==1.5.1
beautifulsoup4==4.15.0
curl_cffi==0.16.0
fastapi==0.141.1
feedparser==6.0.14
httpx==0.28.1
playwright==1.62.0
playwright-captcha==0.1.5
uvicorn==0.52.1
+171
View File
@@ -0,0 +1,171 @@
from __future__ import annotations
import unittest
from bs4 import BeautifulSoup
from extractor import detail_from_html, extract_fields, normalize_config
POPULARAIRSOFT_CONFIG = {
"version": 1,
"discovery": {
"type": "html",
"item_selector": "#block-views-block-latest-news-list-block-2 .feature-contents, #block-views-block-latest-news-list-block-1 .lt-teasure",
"limit": 10,
"fields": {
"url": {"selector": "a.link-title", "extract": "attr", "attribute": "href", "required": True},
"external_id": {"selector": "a.link-title", "extract": "attr", "attribute": "href", "required": True},
"title": {"selector": "a.link-title", "extract": "text", "required": True},
"published_at": {
"selector": "time[datetime]",
"extract": "attr",
"attribute": "datetime",
"required": True,
},
},
"media": [
{"type": "photo", "selector": ".site-image img", "attributes": ["src", "data-src", "srcset"]}
],
},
"detail": {
"enabled": True,
"always": True,
"transport": "http",
"root_selector": "article.news.full",
"fields": {
"title": {
"selector": ".feature-contents > .news-story-texts:first-child h2",
"extract": "text",
"required": True,
},
"published_at": {
"candidates": [
{
"selector": ".feature-contents > .news-story-texts:first-child .news-story-date",
"extract": "text",
"formats": ["%d %b %Y"],
"timezone": "UTC",
},
{"source": "list.published_at"},
],
"required": True,
},
"text": {
"selector": ".feature-contents > .news-story-texts:last-child .field--name-body",
"extract": "text",
"required": True,
"min_length": 50,
},
},
"media": [
{
"type": "photo",
"selector": ".feature-contents > .site-image .field--name-field-image img",
"attributes": ["src", "data-src", "srcset"],
},
{
"type": "video",
"selector": ".field--name-body iframe, .field--name-body video, .field--name-body source",
"attributes": ["src", "data-src"],
},
],
},
"access": {"type": "cloudflare", "wait_for": "#block-views-block-latest-news-list-block-1"},
"retry": {"attempts": 3, "delay_seconds": 5, "timeout_seconds": 90},
"min_text_length": 50,
}
PAGE = """
<html><body>
<img src="/logo.png">
<article class="news full">
<div class="feature-contents">
<div class="news-story-texts">
<h2>Double Bell M16A2</h2>
<h4>OptimusPrime</h4>
<p class="news-story-date">10 Aug 2026</p>
</div>
<div class="site-image">
<div class="field--name-field-image"><img src="/cover.jpg"></div>
</div>
<div class="news-story-texts">
<div class="field--name-body">
<p>This is the complete article body with enough useful text to pass validation safely.</p>
<iframe src="https://www.youtube-nocookie.com/embed/t6mvlySpXNk?si=test"></iframe>
</div>
</div>
</div>
</article>
<div class="related"><img src="/garbage.jpg"></div>
</body></html>
"""
LIST = """
<section id="block-views-block-latest-news-list-block-2">
<div class="feature-contents">
<a class="link-title" href="/news/vfc-vityaz">VFC Vityaz</a>
<time datetime="2026-08-10T06:06:31+00:00">10 Aug 2026</time>
</div>
</section>
<section id="block-views-block-latest-news-list-block-1">
<div class="lt-teasure">
<a class="link-title" href="/news/double-bell">Double Bell</a>
<time datetime="2026-08-10T06:05:49+00:00">10 Aug 2026</time>
</div>
</section>
"""
class ExtractorTests(unittest.TestCase):
def test_popularairsoft_discovery_selects_news_links(self) -> None:
config = normalize_config(POPULARAIRSOFT_CONFIG)
soup = BeautifulSoup(LIST, "html.parser")
cards = soup.select(config["discovery"]["item_selector"])
values = [extract_fields(card, config["discovery"]["fields"])[0] for card in cards]
self.assertEqual([item["url"] for item in values], ["/news/vfc-vityaz", "/news/double-bell"])
self.assertEqual(values[1]["published_at"], "2026-08-10T06:05:49+00:00")
def test_popularairsoft_extracts_only_article_fields(self) -> None:
config = normalize_config(POPULARAIRSOFT_CONFIG)
item, errors = detail_from_html(
PAGE,
"https://popularairsoft.com/news/example",
config,
{"list": {"published_at": "2026-08-10T06:05:49+00:00"}},
)
self.assertEqual(errors, [])
self.assertEqual(item["title"], "Double Bell M16A2")
self.assertEqual(item["published_at"], "2026-08-10T00:00:00+00:00")
self.assertNotIn("garbage", item["text"])
self.assertEqual(item["media"][0]["url"], "https://popularairsoft.com/cover.jpg")
self.assertEqual(item["media"][1], {
"type": "video",
"url": "https://www.youtube.com/watch?v=t6mvlySpXNk",
"provider": "youtube",
})
def test_date_falls_back_to_rss(self) -> None:
fields = {
"published_at": {
"candidates": [
{"selector": "time", "extract": "attr", "attribute": "datetime"},
{"source": "rss.published"},
],
"required": True,
}
}
values, errors = extract_fields(None, fields, {"rss": {"published": "Mon, 10 Aug 2026 06:05:49 +0000"}})
self.assertEqual(errors, [])
self.assertEqual(values["published_at"], "2026-08-10T06:05:49+00:00")
def test_missing_required_field_is_an_item_error(self) -> None:
config = normalize_config(POPULARAIRSOFT_CONFIG)
item, errors = detail_from_html("<article class='news full'></article>", "https://example.test/1", config)
self.assertEqual(item["media"], [])
self.assertTrue(any(error["field"] == "text" for error in errors))
if __name__ == "__main__":
unittest.main()
+18
View File
@@ -0,0 +1,18 @@
from app import parse_rss
def test_parse_rss() -> None:
xml = """<rss><channel><item><guid>1</guid><title>Title</title>
<link>https://example.test/1</link>
<description>&lt;p&gt;Body&lt;/p&gt;&lt;img src="/1.jpg"&gt;
&lt;iframe src="https://www.youtube.com/embed/dQw4w9WgXcQ?feature=oembed"&gt;&lt;/iframe&gt;</description>
<pubDate>Mon, 10 Aug 2026 10:00:00 +0000</pubDate></item></channel></rss>"""
item = parse_rss(xml, 20)[0]
assert item.external_id == "1"
assert item.text == "Body"
assert item.media[0]["url"] == "https://example.test/1.jpg"
assert item.media[1] == {"type": "video", "url": "https://www.youtube.com/watch?v=dQw4w9WgXcQ"}
if __name__ == "__main__":
test_parse_rss()
+141 -12
View File
@@ -9,22 +9,27 @@ import re
import secrets
import time
from datetime import date, datetime, timedelta, timezone
from io import BytesIO
from pathlib import Path
from typing import Any
from urllib.parse import urlencode
from zoneinfo import ZoneInfo
import aiohttp
from fastapi import FastAPI, File, Form, Request, UploadFile, status
from fastapi.responses import FileResponse, HTMLResponse, JSONResponse, RedirectResponse
from aiogram import Bot
from aiogram.client.session.aiohttp import AiohttpSession
from aiogram.client.telegram import TelegramAPIServer
from fastapi import FastAPI, File, Form, HTTPException, Request, UploadFile, status
from fastapi.responses import FileResponse, HTMLResponse, JSONResponse, RedirectResponse, Response
from fastapi.staticfiles import StaticFiles
from fastapi.templating import Jinja2Templates
from loguru import logger
from .config import settings
from .constants import PLATFORM_VK
from .constants import PLATFORM_SITE, PLATFORM_VK
from .db import fetch_int_setting, fetch_setting, get_pool
from .security import hash_password, new_token, token_hash, verify_password
from .source_adapters import json_object, validate_source_config
from .text_utils import build_publication_text, normalize_hash_tag, parse_categories
from .vk_api import VKAPIClient, normalize_vk_source
from .workers.ai_qualifier import AIQualifierWorker, normalize_model, response_usage
@@ -36,7 +41,7 @@ from .workers.site_poster import SitePoster
from .workers.tg_poster import TelegramPoster
from .workers.tg_reactor import TelegramReactor
from .workers.vk_poster import VKPoster
from .workers.vk_storage_uploader import TelegramStorageUploader
from .workers.media_uploader import MediaUploader
COOKIE_NAME = "vk_parser_admin"
VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier"
@@ -196,6 +201,7 @@ CATEGORY_TITLES = {
"Publishing": "Публикации",
"Daily Report": "Ежедневный отчет",
"Parser": "Парсер",
"Site Parser": "Site Parser",
"Uploader": "Аплоадер",
"VK": "VK API",
"General": "Общие",
@@ -212,6 +218,7 @@ CATEGORY_ORDER = {
"Publishing": 53,
"Daily Report": 55,
"Parser": 60,
"Site Parser": 65,
"VK": 70,
"Uploader": 80,
"General": 100,
@@ -363,6 +370,14 @@ SETTING_ORDER = {
"parser_dedupe_content_hash",
"parser_source_pause_sec",
],
"Site Parser": [
"site_parser_url",
"site_parser_token",
"site_parser_rucaptcha_token",
"site_parser_timeout_sec",
"site_parser_interval_minutes",
"site_parser_min_text_length",
],
"VK": [
"vk_requests_per_second",
"vk_wall_page_size",
@@ -491,6 +506,40 @@ async def get_current_user(request: Request) -> dict | None:
return dict(row) if row else None
@app.get("/raw/media/{media_id}")
async def raw_media_preview(request: Request, media_id: int) -> Response:
if not await get_current_user(request):
raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED)
pool = await get_pool()
media = await pool.fetchrow(
"SELECT media_type, tg_file_id FROM raw_post_media WHERE id=$1",
media_id,
)
if not media or media["media_type"] != "photo" or not media["tg_file_id"]:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND)
local_bot_api_url = str(await fetch_setting("local_bot_api_url", settings.local_bot_api_url) or "").strip()
session = (
AiohttpSession(api=TelegramAPIServer.from_base(local_bot_api_url.rstrip("/"), is_local=True))
if local_bot_api_url
else AiohttpSession()
)
bot = Bot(token=settings.tg_bot_token, session=session)
output = BytesIO()
try:
await bot.download(str(media["tg_file_id"]), destination=output)
except Exception as exc:
logger.warning("Telegram media preview failed for media {}: {}", media_id, exc)
raise HTTPException(status_code=status.HTTP_502_BAD_GATEWAY, detail="Media preview unavailable") from exc
finally:
await bot.session.close()
return Response(
content=output.getvalue(),
media_type="image/jpeg",
headers={"Cache-Control": "private, max-age=3600"},
)
def require_csrf(user: dict, csrf_token: str) -> None:
if not user or csrf_token != user["csrf_token"]:
raise PermissionError("bad csrf")
@@ -1810,7 +1859,7 @@ async def startup() -> None:
("tg-poster", TelegramPoster()),
("tg-reactor", TelegramReactor()),
("vk-poster", VKPoster()),
("vk-storage-uploader", TelegramStorageUploader()),
("media-uploader", MediaUploader()),
]
for name, worker in workers:
asyncio.create_task(start_worker_task(worker, name))
@@ -2193,6 +2242,7 @@ async def sources_list(request: Request, q: str = "", status_filter: str = "", a
sources = []
for row in rows:
item = dict(row)
item["settings_json"] = json_object(item.get("settings_json"))
item["last_parsed_at_fmt"] = format_dt(item.get("last_parsed_at"))
item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at"))
item["created_at_fmt"] = format_dt(item.get("created_at"))
@@ -2271,6 +2321,7 @@ async def sources_preview(
sources = []
for row in rows:
item = dict(row)
item["settings_json"] = json_object(item.get("settings_json"))
item["last_parsed_at_fmt"] = format_dt(item.get("last_parsed_at"))
item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at"))
item["created_at_fmt"] = format_dt(item.get("created_at"))
@@ -2361,12 +2412,22 @@ async def source_create(
url: str = Form(...),
active: str = Form("off"),
priority: int = Form(100),
settings_json: str = Form("{}"),
):
user = await get_current_user(request)
if not user:
return redirect("/login")
require_csrf(user, csrf_token)
platform = platform.strip().lower() or PLATFORM_VK
if platform not in {PLATFORM_VK, PLATFORM_SITE}:
raise HTTPException(status_code=422, detail="Неподдерживаемая площадка")
try:
source_settings = json.loads(settings_json or "{}")
if not isinstance(source_settings, dict):
raise ValueError("Настройки должны быть JSON-объектом")
validate_source_config(platform, source_settings)
except (json.JSONDecodeError, ValueError) as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
external_id = normalize_vk_source(url) if platform == PLATFORM_VK else ""
external_owner_id = None
resolved_name = name.strip()
@@ -2389,8 +2450,8 @@ async def source_create(
pool = await get_pool()
row = await pool.fetchrow(
"""
INSERT INTO sources(platform, name, tag, url, external_id, external_owner_id, active, priority, status, status_msg, created_by)
VALUES($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
INSERT INTO sources(platform, name, tag, url, external_id, external_owner_id, active, priority, status, status_msg, settings_json, created_by)
VALUES($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11::jsonb, $12)
ON CONFLICT DO NOTHING
RETURNING id
""",
@@ -2404,6 +2465,7 @@ async def source_create(
priority,
status_value,
status_msg,
json.dumps(source_settings, ensure_ascii=False),
user["id"],
)
if row:
@@ -2425,7 +2487,7 @@ async def source_edit(request: Request, source_id: int):
base_context(
request,
user,
source=dict(source),
source={**dict(source), "settings_json": json_object(source["settings_json"])},
action=f"/sources/{source_id}/edit",
title=f"Источник #{source_id}",
),
@@ -2443,12 +2505,22 @@ async def source_update(
url: str = Form(...),
active: str = Form("off"),
priority: int = Form(100),
settings_json: str = Form("{}"),
):
user = await get_current_user(request)
if not user:
return redirect("/login")
require_csrf(user, csrf_token)
platform = platform.strip().lower() or PLATFORM_VK
if platform not in {PLATFORM_VK, PLATFORM_SITE}:
raise HTTPException(status_code=422, detail="Неподдерживаемая площадка")
try:
source_settings = json.loads(settings_json or "{}")
if not isinstance(source_settings, dict):
raise ValueError("Настройки должны быть JSON-объектом")
validate_source_config(platform, source_settings)
except (json.JSONDecodeError, ValueError) as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
external_id = normalize_vk_source(url) if platform == PLATFORM_VK else ""
pool = await get_pool()
await pool.execute(
@@ -2459,8 +2531,14 @@ async def source_update(
tag=$4,
url=$5,
external_id=$6,
external_owner_id=CASE WHEN platform=$2 AND url=$5 THEN external_owner_id ELSE NULL END,
active=$7,
priority=$8,
settings_json=$9::jsonb,
runtime_state_json=CASE
WHEN platform=$2 AND url=$5 AND settings_json=$9::jsonb THEN runtime_state_json
ELSE '{}'::jsonb
END,
updated_at=NOW()
WHERE id=$1
""",
@@ -2472,6 +2550,7 @@ async def source_update(
external_id,
active == "on",
priority,
json.dumps(source_settings, ensure_ascii=False),
)
await audit(user["id"], "source.update", "source", source_id, {"url": url, "platform": platform})
return redirect("/sources")
@@ -2501,6 +2580,7 @@ async def source_toggle(request: Request, source_id: int, csrf_token: str = Form
source_row = await pool.fetchrow("SELECT * FROM sources WHERE id=$1", source_id)
if source_row:
s = dict(source_row)
s["settings_json"] = json_object(s.get("settings_json"))
# Fetch stats just for this source to render correctly
stats = await pool.fetchrow(
"""
@@ -2518,6 +2598,25 @@ async def source_toggle(request: Request, source_id: int, csrf_token: str = Form
return redirect("/sources")
@app.post("/sources/{source_id}/parse-now")
async def source_parse_now(request: Request, source_id: int, csrf_token: str = Form(...)):
user = await get_current_user(request)
if not user:
return redirect("/login")
require_csrf(user, csrf_token)
pool = await get_pool()
await pool.execute(
"""
UPDATE sources
SET last_checked_at=NULL, status='new', status_msg=NULL, updated_at=NOW()
WHERE id=$1 AND archived_at IS NULL
""",
source_id,
)
await audit(user["id"], "source.parse_now", "source", source_id)
return redirect("/sources")
@app.post("/sources/{source_id}/delete")
async def source_delete(request: Request, source_id: int, csrf_token: str = Form(...)):
user = await get_current_user(request)
@@ -2619,7 +2718,11 @@ async def raw_posts(
jsonb_build_object(
'id', rpm.id,
'type', rpm.media_type,
'url', rpm.original_url,
'url', CASE
WHEN rpm.media_type='photo' AND rpm.tg_file_id IS NOT NULL
THEN '/raw/media/' || rpm.id::text
ELSE rpm.original_url
END,
'status', rpm.status,
'error', rpm.error,
'duration_sec', rpm.duration_sec
@@ -2785,7 +2888,11 @@ async def fetch_single_editor_row(pool, post_id: int):
jsonb_build_object(
'id', rpm.id,
'type', rpm.media_type,
'url', rpm.original_url,
'url', CASE
WHEN rpm.media_type='photo' AND rpm.tg_file_id IS NOT NULL
THEN '/raw/media/' || rpm.id::text
ELSE rpm.original_url
END,
'status', rpm.status,
'error', rpm.error,
'duration_sec', rpm.duration_sec
@@ -2926,7 +3033,11 @@ async def editor_feed(
jsonb_build_object(
'id', rpm.id,
'type', rpm.media_type,
'url', rpm.original_url,
'url', CASE
WHEN rpm.media_type='photo' AND rpm.tg_file_id IS NOT NULL
THEN '/raw/media/' || rpm.id::text
ELSE rpm.original_url
END,
'status', rpm.status,
'error', rpm.error,
'duration_sec', rpm.duration_sec
@@ -3266,7 +3377,11 @@ async def raw_post_detail(request: Request, post_id: int):
jsonb_build_object(
'id', rpm.id,
'type', rpm.media_type,
'url', rpm.original_url,
'url', CASE
WHEN rpm.media_type='photo' AND rpm.tg_file_id IS NOT NULL
THEN '/raw/media/' || rpm.id::text
ELSE rpm.original_url
END,
'status', rpm.status,
'error', rpm.error,
'duration_sec', rpm.duration_sec
@@ -3515,6 +3630,18 @@ async def workers(request: Request):
continue
settings.append(item)
settings.sort(key=setting_sort_key)
site_parser_url = str(setting_values.get("site_parser_url") or "").strip().rstrip("/")
site_parser_health = {"ok": False, "message": "URL не настроен"}
if site_parser_url:
try:
async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=3)) as session:
async with session.get(f"{site_parser_url}/health") as response:
site_parser_health = {
"ok": response.status == 200,
"message": "Доступен" if response.status == 200 else f"HTTP {response.status}",
}
except Exception as exc:
site_parser_health = {"ok": False, "message": str(exc)[:200] or "Недоступен"}
ai_provider = str(setting_value(settings, "ai_qualifier_provider", "openrouter") or "openrouter")
ai_model = str(setting_value(settings, "ai_qualifier_model", "") or "")
ai_api_key = str(setting_value(settings, "ai_qualifier_api_key", "") or "")
@@ -3549,6 +3676,8 @@ async def workers(request: Request):
for r in rows
],
settings=settings,
site_parser_url=site_parser_url,
site_parser_health=site_parser_health,
provider_options=PROVIDER_OPTIONS,
ai_provider=ai_provider,
ai_model=ai_model,
+3 -2
View File
@@ -1,4 +1,5 @@
PLATFORM_VK = "vk"
PLATFORM_SITE = "site"
SOURCE_STATUS_NEW = "new"
SOURCE_STATUS_OK = "ok"
@@ -26,10 +27,10 @@ JOB_STATUS_RETRY = "retry"
JOB_STATUS_DONE = "done"
JOB_STATUS_DEAD = "dead"
JOB_TYPE_VK_STORAGE_COPY = "vk.storage.copy"
JOB_TYPE_MEDIA_STORAGE_COPY = "media.storage.copy"
WORKER_PARSER = "vk-parser"
WORKER_STORAGE_UPLOADER = "vk-storage-uploader"
WORKER_MEDIA_UPLOADER = "media-uploader"
WORKER_AI_QUALIFIER = "ai-qualifier"
WORKER_AI_WRITER = "ai-writer"
WORKER_TG_POSTER = "tg-poster"
+209
View File
@@ -0,0 +1,209 @@
from __future__ import annotations
import json
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any
import aiohttp
from .constants import PLATFORM_SITE, PLATFORM_VK
def json_object(value: Any) -> dict[str, Any]:
if isinstance(value, dict):
return dict(value)
if isinstance(value, str):
try:
parsed = json.loads(value)
except json.JSONDecodeError:
return {}
return parsed if isinstance(parsed, dict) else {}
return {}
@dataclass
class SourceMedia:
url: str
media_type: str = "photo"
@dataclass
class SourceItem:
external_id: str
url: str
text: str
posted_at: datetime
media: list[SourceMedia] = field(default_factory=list)
raw: dict[str, Any] = field(default_factory=dict)
@property
def body_text(self) -> str:
return str(self.raw.get("text") or "").strip()
def validate_source_config(platform: str, config: dict[str, Any]) -> None:
if platform == PLATFORM_VK:
return
if platform != PLATFORM_SITE:
raise ValueError("Неподдерживаемая площадка")
if not config:
raise ValueError("Для сайта нужен конфиг JSON")
if "discovery" in config:
if config.get("version", 1) != 1:
raise ValueError("Поддерживается только version=1")
discovery = config.get("discovery")
if not isinstance(discovery, dict) or discovery.get("type") not in {"rss", "html"}:
raise ValueError('discovery.type должен быть "rss" или "html"')
try:
limit = int(discovery.get("limit", 20))
except (TypeError, ValueError) as exc:
raise ValueError("discovery.limit должен быть целым числом") from exc
if not 1 <= limit <= 100:
raise ValueError("discovery.limit должен быть от 1 до 100")
if discovery.get("type") == "html" and not str(discovery.get("item_selector") or "").strip():
raise ValueError("Для HTML discovery нужен item_selector")
detail = config.get("detail") or {"enabled": False}
if not isinstance(detail, dict):
raise ValueError("detail должен быть JSON-объектом")
if detail.get("enabled"):
if not str(detail.get("root_selector") or "").strip():
raise ValueError("Для detail нужен root_selector")
if not isinstance(detail.get("fields") or {}, dict):
raise ValueError("detail.fields должен быть JSON-объектом")
if not isinstance(detail.get("media") or [], list):
raise ValueError("detail.media должен быть массивом")
if detail.get("transport", "browser") not in {"http", "fetch", "browser"}:
raise ValueError('detail.transport должен быть "http", "fetch" или "browser"')
access = config.get("access") or {"type": "auto"}
if isinstance(access, str):
access = {"type": access}
if not isinstance(access, dict) or access.get("type", "auto") not in {"auto", "http", "browser", "cloudflare"}:
raise ValueError('access.type должен быть "auto", "http", "browser" или "cloudflare"')
retry = config.get("retry") or {}
try:
attempts = int(retry.get("attempts", 2))
delay = float(retry.get("delay_seconds", 2))
timeout = float(retry.get("timeout_seconds", 90))
min_text_length = int(config.get("min_text_length", 0))
except (TypeError, ValueError) as exc:
raise ValueError("retry и min_text_length должны быть числовыми") from exc
if not 1 <= attempts <= 5 or not 0 <= delay <= 60 or not 10 <= timeout <= 300:
raise ValueError("retry: attempts 1..5, delay_seconds 0..60, timeout_seconds 10..300")
if not 0 <= min_text_length <= 100_000:
raise ValueError("min_text_length должен быть от 0 до 100000")
return
if config.get("format") != "rss":
raise ValueError('Сейчас поддерживается только "format": "rss"')
if str(config.get("access") or "auto") not in {"auto", "http", "cloudflare"}:
raise ValueError('access должен быть "auto", "http" или "cloudflare"')
try:
max_items = int(config.get("max_items", 20))
except (TypeError, ValueError) as exc:
raise ValueError("max_items должен быть целым числом") from exc
if not 1 <= max_items <= 100:
raise ValueError("max_items должен быть от 1 до 100")
follow_links = config.get("follow_links", False)
if not isinstance(follow_links, bool):
raise ValueError("follow_links должен быть true или false")
if follow_links and not str(config.get("content_selector") or "").strip():
raise ValueError("При follow_links=true нужен content_selector")
try:
min_text_length = int(config.get("min_text_length", 0))
except (TypeError, ValueError) as exc:
raise ValueError("min_text_length должен быть целым числом") from exc
if not 0 <= min_text_length <= 100_000:
raise ValueError("min_text_length должен быть от 0 до 100000")
def _posted_at(value: Any) -> datetime:
try:
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
except (TypeError, ValueError):
return datetime.now(timezone.utc)
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
class SiteParserClient:
def __init__(
self,
session: aiohttp.ClientSession,
base_url: str,
token: str,
rucaptcha_token: str,
timeout_sec: int,
) -> None:
self.session = session
self.base_url = base_url.rstrip("/")
self.token = token
self.rucaptcha_token = rucaptcha_token
self.timeout = aiohttp.ClientTimeout(total=max(10, timeout_sec))
async def fetch(self, source: dict) -> tuple[list[SourceItem], dict[str, Any] | None, str | None]:
if not self.base_url or not self.token:
raise RuntimeError("Site Parser URL или токен не настроены")
config = json_object(source.get("settings_json"))
validate_source_config(PLATFORM_SITE, config)
runtime_state = json_object(source.get("runtime_state_json"))
payload = {
"url": source["url"],
"config": config,
"rucaptcha_token": self.rucaptcha_token or None,
"browser_state": runtime_state.get("browser_state"),
"browser_user_agent": runtime_state.get("browser_user_agent"),
}
try:
async with self.session.post(
f"{self.base_url}/v1/parse",
json=payload,
headers={"X-Worker-Token": self.token},
timeout=self.timeout,
) as response:
data = await response.json(content_type=None)
if response.status >= 400:
raise RuntimeError(f"Site Parser HTTP {response.status}: {data.get('detail', data)}")
except TimeoutError as exc:
raise RuntimeError(f"Site Parser превысил таймаут {int(self.timeout.total)} сек") from exc
except aiohttp.ClientError as exc:
raise RuntimeError(f"Site Parser недоступен: {exc}") from exc
result_status = str(data.get("status") or "ok")
result_errors = data.get("errors") or []
items = []
for raw in data.get("items") or []:
title = str(raw.get("title") or "").strip()
body = str(raw.get("text") or "").strip()
body_starts_with_title = title and (
body.casefold() == title.casefold()
or body.casefold().startswith(f"{title}\n".casefold())
)
text = body if body_starts_with_title else "\n\n".join(part for part in (title, body) if part)
url = str(raw.get("url") or source["url"]).strip()
external_id = str(raw.get("external_id") or url).strip()
if not external_id:
continue
media = [
SourceMedia(str(item["url"]), str(item.get("type") or "photo"))
for item in raw.get("media") or []
if isinstance(item, dict) and item.get("url")
]
items.append(SourceItem(external_id, url, text, _posted_at(raw.get("published_at")), media, raw))
state = data.get("browser_state")
runtime = (
{
"browser_state": state,
"browser_user_agent": str(data.get("browser_user_agent") or "") or None,
"last_result_status": result_status,
"last_result_errors": result_errors[:20],
}
if isinstance(state, dict)
else None
)
warning = (
f"Site Parser {result_status}: {len(result_errors)} ошибок; {str(result_errors)[:700]}"
if result_status in {"partial", "failed"}
else None
)
return items, runtime, warning
@@ -17,6 +17,7 @@
<label class="label"><span class="label-text font-bold">Площадка</span></label>
<select name="platform" class="select select-bordered w-full">
<option value="vk" {% if not source or source.platform == "vk" %}selected{% endif %}>VK</option>
<option value="site" {% if source and source.platform == "site" %}selected{% endif %}>Сайт</option>
</select>
</div>
@@ -39,6 +40,37 @@
<label class="label"><span class="label-text font-bold">Ссылка</span></label>
<input name="url" value="{{ source.url if source else '' }}" required class="input input-bordered w-full text-primary">
</div>
<div id="site-config" class="form-control md:col-span-2">
<label class="label"><span class="label-text font-bold">Конфигурация JSON</span></label>
<textarea name="settings_json" rows="12" class="textarea textarea-bordered w-full font-mono" placeholder='{"version":1,"discovery":{"type":"rss","limit":20}}'>{{ source.settings_json | tojson(indent=2) if source else '{}' }}</textarea>
<details class="mt-2 text-sm text-base-content/70">
<summary class="cursor-pointer">Пример и параметры</summary>
<pre class="mt-2 p-3 bg-base-200 overflow-x-auto">{
"version": 1,
"discovery": {
"type": "rss",
"limit": 20
},
"detail": {
"enabled": true,
"always": true,
"root_selector": "article",
"fields": {
"title": {"selector": "h1", "extract": "text", "required": true},
"text": {"selector": ".article-body", "extract": "text", "required": true}
}
},
"access": {"type": "auto"},
"retry": {"attempts": 2, "delay_seconds": 2, "timeout_seconds": 90},
"min_text_length": 50
}</pre>
<p class="mt-2"><code>discovery.type</code>: <code>rss</code> читает ссылки из ленты; <code>html</code> собирает карточки по <code>item_selector</code>. Секция <code>detail</code> описывает точные селекторы полей на странице материала.</p>
<p class="mt-2"><code>access.type</code>: <code>auto</code> сначала пробует обычный запрос; <code>http</code> запрещает браузер; <code>browser</code> выполняет страницу в браузере; <code>cloudflare</code> разрешает RuCaptcha.</p>
<p class="mt-2">Для fallback используйте <code>candidates</code>. Извлечение поддерживает <code>text</code>, <code>html</code>, <code>attr</code> и <code>json</code>. Ошибка одной detail-страницы не прерывает готовую пачку.</p>
<p class="mt-2"><code>min_text_length</code> переопределяет общий порог Site Parser только для этого источника.</p>
</details>
</div>
</div>
<div class="form-control mt-4">
@@ -54,4 +86,17 @@
</form>
</div>
</div>
<script>
const platform = document.querySelector('[name="platform"]');
const siteConfig = document.getElementById('site-config');
const configInput = document.querySelector('[name="settings_json"]');
const syncConfig = () => {
siteConfig.hidden = platform.value !== 'site';
if (platform.value === 'site' && configInput.value.trim() === '{}') {
configInput.value = '{\n "version": 1,\n "discovery": {\n "type": "rss",\n "limit": 20\n },\n "detail": {\n "enabled": false\n },\n "access": {\n "type": "auto"\n },\n "retry": {\n "attempts": 2,\n "delay_seconds": 2,\n "timeout_seconds": 90\n },\n "min_text_length": 50\n}';
}
};
platform.addEventListener('change', syncConfig);
syncConfig();
</script>
{% endblock %}
@@ -39,6 +39,12 @@
<td class="hidden xl:table-cell px-4 py-3 text-right text-[10px] text-app-textMuted">{{ s.last_parsed_at_fmt or "Никогда" }}</td>
<td class="px-3 sm:px-4 py-3 text-center">
<div class="flex justify-center items-center gap-1">
<form method="post" action="/sources/{{ s.id }}/parse-now" class="m-0">
<input type="hidden" name="csrf_token" value="{{ user.csrf_token }}">
<button class="btn btn-ghost btn-icon w-8 h-8 rounded-md hover:bg-app-surface text-app-textMuted hover:text-white" type="submit" title="Парсить сейчас">
<i data-lucide="refresh-cw" class="w-4 h-4"></i>
</button>
</form>
<a href="/sources/{{ s.id }}/edit" class="btn btn-ghost btn-icon w-8 h-8 rounded-md hover:bg-app-surface text-app-textMuted hover:text-white" title="Править">
<i data-lucide="edit-2" class="w-4 h-4"></i>
</a>
+10 -7
View File
@@ -1,11 +1,14 @@
{% extends "base.html" %}
{% block body %}
<div class="mb-8">
<h1 class="text-3xl font-bold tracking-tight text-white flex items-center gap-3 mb-2">
<i data-lucide="database" class="text-app-primary w-8 h-8"></i>
Источники
</h1>
<div class="text-app-textMuted text-sm">VK-источники для парсинга. Название идёт в prompt, тэг — в будущие хэштеги.</div>
<div class="mb-8 flex flex-wrap items-start justify-between gap-4">
<div>
<h1 class="text-3xl font-bold tracking-tight text-white flex items-center gap-3 mb-2">
<i data-lucide="database" class="text-app-primary w-8 h-8"></i>
Источники
</h1>
<div class="text-app-textMuted text-sm">Источники для парсинга. Название идёт в prompt, тэг — в будущие хэштеги.</div>
</div>
<a class="btn btn-primary" href="/sources/new"><i data-lucide="plus" class="w-4 h-4"></i> Добавить источник</a>
</div>
<details class="card mb-8 group/details" {% if source_preview %}open{% endif %}>
@@ -14,7 +17,7 @@
<div class="w-8 h-8 rounded-lg bg-app-primary/20 text-app-primary flex items-center justify-center">
<i data-lucide="plus" class="w-5 h-5"></i>
</div>
<span class="text-lg font-bold text-white">Добавить источники</span>
<span class="text-lg font-bold text-white">Добавить VK списком</span>
</div>
<i data-lucide="chevron-down" class="w-5 h-5 text-app-textMuted transition-transform group-open/details:rotate-180"></i>
</summary>
+18 -2
View File
@@ -8,6 +8,18 @@
<div class="text-app-textMuted text-sm">Управление фоновыми процессами, категориями, расписанием и настройками AI.</div>
</div>
<div class="card mb-8 p-4 flex flex-wrap items-center justify-between gap-4">
<div class="flex items-center gap-3 min-w-0">
<i data-lucide="rss" class="w-5 h-5 text-app-primary flex-shrink-0"></i>
<div class="min-w-0">
<div class="font-bold text-white">Site Parser</div>
<div class="text-xs text-app-textMuted truncate">{{ site_parser_url or "URL не настроен" }}</div>
</div>
<span class="badge {% if site_parser_health.ok %}badge-success{% else %}badge-error{% endif %}">{{ site_parser_health.message }}</span>
</div>
<a href="#site-parser-settings" class="btn btn-primary-outline btn-sm"><i data-lucide="settings" class="w-4 h-4"></i> Настройки</a>
</div>
<!-- Workers Panel -->
<div class="card overflow-hidden mb-8">
<div class="overflow-x-auto">
@@ -25,7 +37,7 @@
<tbody class="divide-y divide-app-border text-xs">
{% for w in workers %}
<tr class="hover:bg-app-surfaceHover transition-colors">
<td class="px-4 py-3 font-mono font-bold text-white">{{ w.name }}</td>
<td class="px-4 py-3 font-mono font-bold text-white">{{ "Парсер источников (VK + сайты)" if w.name == "vk-parser" else w.name }}</td>
<td class="px-4 py-3">
{% if w.effective_enabled %}
<span class="badge badge-success px-2 py-0.5 flex items-center gap-1.5 w-max">
@@ -294,7 +306,7 @@
<div class="flex flex-col gap-6">
{% for category, rows in settings|groupby("category") %}
<details class="card group/details" data-workers-details="settings:{{ category }}">
<details id="{% if category == 'Site Parser' %}site-parser-settings{% else %}settings-{{ category|lower|replace(' ', '-') }}{% endif %}" class="card group/details" data-workers-details="settings:{{ category }}" {% if category == "Site Parser" %}open{% endif %}>
<summary class="p-4 flex items-center justify-between cursor-pointer select-none hover:bg-app-surfaceHover transition-colors border-b border-app-border list-none">
<div class="text-lg font-bold text-white">{{ category_titles.get(category, category) }}</div>
<i data-lucide="chevron-down" class="w-5 h-5 text-app-textMuted transition-transform group-open/details:rotate-180"></i>
@@ -404,6 +416,10 @@
const scrollKey = "workers-scroll-y";
const detailsKey = "workers-open-details";
const detailItems = Array.from(document.querySelectorAll("details[data-workers-details]"));
const siteParserDetails = document.getElementById("site-parser-settings");
document.querySelector('a[href="#site-parser-settings"]')?.addEventListener("click", () => {
if (siteParserDetails) siteParserDetails.open = true;
});
const savedDetails = sessionStorage.getItem(detailsKey);
if (savedDetails) {
try {
+13 -11
View File
@@ -39,7 +39,7 @@ def parse_recipients(value: Any) -> list[int]:
return recipients
async def send_ai_worker_error_alert(worker_name: str, model: str, post_ids: list[int], error: str) -> None:
async def send_system_error_alert(text: str) -> None:
token = (
str(await fetch_setting("daily_report_bot_token", "") or "").strip()
or str(await fetch_setting("tg_poster_bot_token", "") or "").strip()
@@ -47,17 +47,8 @@ async def send_ai_worker_error_alert(worker_name: str, model: str, post_ids: lis
)
recipients = parse_recipients(await fetch_setting("daily_report_recipient_ids", [442509142]))
if not token or not recipients:
logger.warning("AI worker alert skipped: token or recipients are empty")
logger.warning("System alert skipped: token or recipients are empty")
return
post_part = ", ".join(str(post_id) for post_id in post_ids) if post_ids else "-"
text = (
"AI worker batch failed\n"
f"worker: {worker_name}\n"
f"model: {model or '-'}\n"
f"posts: {post_part}\n"
f"error: {error[:1000]}"
)
bot = Bot(token=token)
try:
for recipient_id in recipients:
@@ -69,3 +60,14 @@ async def send_ai_worker_error_alert(worker_name: str, model: str, post_ids: lis
await asyncio.sleep(float(exc.retry_after) + 1)
finally:
await bot.session.close()
async def send_ai_worker_error_alert(worker_name: str, model: str, post_ids: list[int], error: str) -> None:
post_part = ", ".join(str(post_id) for post_id in post_ids) if post_ids else "-"
await send_system_error_alert(
"AI worker batch failed\n"
f"worker: {worker_name}\n"
f"model: {model or '-'}\n"
f"posts: {post_part}\n"
f"error: {error[:1000]}"
)
@@ -7,6 +7,7 @@ import re
import sys
import time
from pathlib import Path
from urllib.parse import urlparse
import aiohttp
from aiogram import Bot
@@ -18,17 +19,18 @@ from loguru import logger
from ..config import settings
from ..constants import (
JOB_TYPE_VK_STORAGE_COPY,
JOB_TYPE_MEDIA_STORAGE_COPY,
MEDIA_STATUS_FAILED,
MEDIA_STATUS_LINK_ONLY,
MEDIA_STATUS_UPLOADED,
POST_STATUS_FAILED,
POST_STATUS_STORAGE_READY,
WORKER_STORAGE_UPLOADER,
WORKER_MEDIA_UPLOADER,
)
from ..db import fetch_float_setting, fetch_int_setting, fetch_setting, get_pool
from ..heartbeat import HeartbeatReporter
from ..jobs import ack_done, ack_retry, claim_job, is_worker_enabled, recover_stale_jobs
from ..source_adapters import json_object
TMP_DIR = Path("/tmp")
TMP_PREFIX = "vkparser_tg_media_"
@@ -52,6 +54,20 @@ def parse_vk_video_url(vk_url: str) -> tuple[int, int] | None:
return int(match.group(1)), int(match.group(2))
def video_provider(url: str) -> str | None:
host = (urlparse(str(url or "")).hostname or "").lower().removeprefix("www.").removeprefix("m.")
if host == "vk.com" and parse_vk_video_url(url):
return "vk"
if host in {"youtube.com", "youtube-nocookie.com", "youtu.be"}:
return "youtube"
return None
def remove_download_files(output_path: str) -> None:
for path in TMP_DIR.glob(f"{Path(output_path).name}*"):
path.unlink(missing_ok=True)
def split_message_chunks(text: str, limit: int) -> list[str]:
text = str(text or "").strip()
if not text:
@@ -75,13 +91,13 @@ def original_link(post: dict) -> str:
return f'<a href="{url}">#{int(post["id"])}</a>' if url else f'#{int(post["id"])}'
class TelegramStorageUploader:
class MediaUploader:
def __init__(self) -> None:
self.pool = None
self.bot: Bot | None = None
self.storage_chat_id: int | None = None
self.storage_thread_id: int | None = None
self.heartbeat = HeartbeatReporter(WORKER_STORAGE_UPLOADER, 30)
self.heartbeat = HeartbeatReporter(WORKER_MEDIA_UPLOADER, 30)
self.media_group_max_items = MAX_MEDIA_GROUP
self.media_upload_delay_sec = 1.0
self.tg_retry_max_attempts = 4
@@ -100,7 +116,7 @@ class TelegramStorageUploader:
async def init(self) -> None:
self.pool = await get_pool()
recovered = await recover_stale_jobs(self.pool, JOB_TYPE_VK_STORAGE_COPY, stale_minutes=20)
recovered = await recover_stale_jobs(self.pool, JOB_TYPE_MEDIA_STORAGE_COPY, stale_minutes=20)
if recovered:
logger.warning("Recovered stale storage jobs: {}", recovered)
@@ -195,10 +211,12 @@ class TelegramStorageUploader:
async def load_media(self, raw_post_id: int) -> list[dict]:
rows = await self.pool.fetch(
"""
SELECT *
FROM raw_post_media
WHERE raw_post_id=$1
ORDER BY sort_order ASC, id ASC
SELECT m.*, s.url AS source_url, s.runtime_state_json AS source_runtime_state
FROM raw_post_media m
JOIN raw_posts rp ON rp.id=m.raw_post_id
JOIN sources s ON s.id=rp.source_id
WHERE m.raw_post_id=$1
ORDER BY m.sort_order ASC, m.id ASC
""",
raw_post_id,
)
@@ -262,15 +280,70 @@ class TelegramStorageUploader:
return f"media not uploaded: {details}"
return None
async def download_bytes(self, session: aiohttp.ClientSession, url: str) -> bytes | None:
def site_request_options(self, media: dict) -> dict:
if str(media.get("platform") or "") != "site":
return {}
state = json_object(media.get("source_runtime_state"))
browser_state = json_object(state.get("browser_state"))
cookies = {
str(cookie["name"]): str(cookie["value"])
for cookie in browser_state.get("cookies") or []
if isinstance(cookie, dict) and cookie.get("name") and cookie.get("value")
}
headers = {}
if state.get("browser_user_agent"):
headers["User-Agent"] = str(state["browser_user_agent"])
if media.get("source_url"):
headers["Referer"] = str(media["source_url"])
return {"headers": headers, "cookies": cookies}
async def download_bytes(
self,
session: aiohttp.ClientSession,
url: str,
request_options: dict | None = None,
) -> bytes | None:
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=self.download_timeout_sec)) as response:
async with session.get(
url,
timeout=aiohttp.ClientTimeout(total=self.download_timeout_sec),
**(request_options or {}),
) as response:
if response.status != 200:
return None
return await response.read()
except Exception:
return None
async def download_http_video(
self,
session: aiohttp.ClientSession,
url: str,
output_path: str,
request_options: dict | None = None,
) -> dict:
max_size = self.video_max_size_mb * 1024 * 1024
try:
async with session.get(
url,
timeout=aiohttp.ClientTimeout(total=self.download_timeout_sec),
**(request_options or {}),
) as response:
if response.status != 200:
return {"error": f"video returned HTTP {response.status}", "permanent": False}
size = 0
with open(output_path, "wb") as fh:
async for chunk in response.content.iter_chunked(1024 * 1024):
size += len(chunk)
if size > max_size:
remove_download_files(output_path)
return {"error": "video too large", "permanent": True}
fh.write(chunk)
except Exception:
remove_download_files(output_path)
return {"error": "video download failed", "permanent": False}
return {"path": output_path, "size_bytes": size}
async def mark_media_uploaded(self, media_id: int, file_id: str, unique_id: str | None) -> None:
await self.pool.execute(
"""
@@ -320,37 +393,43 @@ class TelegramStorageUploader:
error[:1000],
)
async def download_video(self, vk_url: str, output_path: str) -> dict | None:
parsed = parse_vk_video_url(vk_url)
if not parsed:
async def download_video(self, video_url: str, output_path: str) -> dict | None:
provider = video_provider(video_url)
if not provider:
return None
owner_id, video_id = parsed
max_size = self.video_max_size_mb * 1024 * 1024
netrc_path = f"{output_path}.netrc"
fd = os.open(netrc_path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, "w", encoding="utf-8") as fh:
fh.write(f"machine vk.com login vk_token password {settings.vk_access_token}\n")
target_url = video_url
cmd = [
sys.executable,
"-m",
"yt_dlp",
"--netrc-location",
netrc_path,
f"https://vk.com/video{owner_id}_{video_id}",
"-o",
output_path,
]
netrc_path = None
if provider == "vk":
owner_id, video_id = parse_vk_video_url(video_url) or (0, 0)
target_url = f"https://vk.com/video{owner_id}_{video_id}"
netrc_path = f"{output_path}.netrc"
fd = os.open(netrc_path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, "w", encoding="utf-8") as fh:
fh.write(f"machine vk.com login vk_token password {settings.vk_access_token}\n")
cmd.extend(["--netrc-location", netrc_path])
cmd.extend([
target_url,
"-o", output_path,
"--no-playlist",
"-f",
(
"--match-filter", f"duration <= {self.video_max_duration_sec}",
"--merge-output-format", "mp4",
"-f", (
f"best[height<={self.video_max_height}][filesize<{max_size}]"
f"/best[height<={self.video_max_height}]"
f"/bestvideo[height<={self.video_max_height}][filesize<{max_size}]+bestaudio/best"
f"/bestvideo[height<={self.video_max_height}]+bestaudio/best"
f"/best[filesize<{max_size}]"
),
"--quiet",
"--no-warnings",
]
"--quiet", "--no-warnings",
])
proc = None
stderr = b""
try:
proc = await asyncio.create_subprocess_exec(
*cmd,
@@ -359,20 +438,24 @@ class TelegramStorageUploader:
)
_, stderr = await asyncio.wait_for(proc.communicate(), timeout=self.yt_dlp_timeout_sec)
except asyncio.TimeoutError:
proc.kill()
await proc.communicate()
if proc:
proc.kill()
await proc.communicate()
remove_download_files(output_path)
return {"error": "video download timeout", "permanent": False}
finally:
try:
os.unlink(netrc_path)
except FileNotFoundError:
pass
if netrc_path:
Path(netrc_path).unlink(missing_ok=True)
if proc.returncode != 0 or not os.path.exists(output_path):
err = (stderr or b"").decode("utf-8", errors="ignore").lower()
permanent = any(marker in err for marker in ("removed", "unavailable", "private", "access denied"))
permanent = any(marker in err for marker in (
"removed", "unavailable", "private", "access denied", "does not pass filter", "sign in",
))
remove_download_files(output_path)
return {"error": "video unavailable or download failed", "permanent": permanent}
size = os.path.getsize(output_path)
if size > max_size:
remove_download_files(output_path)
return {"error": "video too large", "permanent": True}
return {"path": output_path, "size_bytes": size}
@@ -383,6 +466,7 @@ class TelegramStorageUploader:
media_id = int(item["id"])
media_type = str(item["media_type"])
url = str(item.get("original_url") or "")
request_options = self.site_request_options(item)
if item.get("tg_file_id"):
prepared.append({"media_id": media_id, "media_type": media_type, "media": item["tg_file_id"]})
@@ -391,7 +475,7 @@ class TelegramStorageUploader:
continue
if media_type == "photo":
data = await self.download_bytes(session, url)
data = await self.download_bytes(session, url, request_options)
if not data:
await self.mark_media_failed_attempt(media_id, "photo download failed")
continue
@@ -410,7 +494,11 @@ class TelegramStorageUploader:
await self.mark_media_link_only(media_id, "video too long")
continue
temp_path = str(TMP_DIR / f"{TMP_PREFIX}{raw_post_id}_{media_id}.mp4")
info = await self.download_video(url, temp_path)
info = (
await self.download_video(url, temp_path)
if video_provider(url)
else await self.download_http_video(session, url, temp_path, request_options)
)
if not info:
await self.mark_media_failed_attempt(media_id, "video download failed")
continue
@@ -581,8 +669,7 @@ class TelegramStorageUploader:
tmp_path = item.get("tmp_path")
if tmp_path:
try:
if os.path.exists(tmp_path):
os.remove(tmp_path)
remove_download_files(tmp_path)
except Exception:
pass
if not blocking_error.startswith("media still pending:"):
@@ -600,8 +687,7 @@ class TelegramStorageUploader:
tmp_path = item.get("tmp_path")
if tmp_path:
try:
if os.path.exists(tmp_path):
os.remove(tmp_path)
remove_download_files(tmp_path)
except Exception:
pass
await self.mark_post_ready(raw_post_id, message_ids, meta_message_id)
@@ -609,12 +695,12 @@ class TelegramStorageUploader:
logger.info("Telegram storage done: raw_post={} messages={}", raw_post_id, message_ids)
async def run_once(self, worker_id: str) -> bool:
enabled = await is_worker_enabled(self.pool, WORKER_STORAGE_UPLOADER)
enabled = await is_worker_enabled(self.pool, WORKER_MEDIA_UPLOADER)
if not enabled:
await self.heartbeat.beat(self.pool, status="disabled", force=True)
return False
job = await claim_job(self.pool, JOB_TYPE_VK_STORAGE_COPY, worker_id)
job = await claim_job(self.pool, JOB_TYPE_MEDIA_STORAGE_COPY, worker_id)
if not job:
await self.heartbeat.beat(self.pool, status="idle")
return False
@@ -637,7 +723,7 @@ class TelegramStorageUploader:
async def run_loop(self) -> None:
await self.init()
worker_id = f"{WORKER_STORAGE_UPLOADER}:{os.getpid()}"
worker_id = f"{WORKER_MEDIA_UPLOADER}:{os.getpid()}"
logger.info("{} started", worker_id)
try:
while True:
@@ -654,7 +740,7 @@ class TelegramStorageUploader:
async def main() -> None:
logger.remove()
logger.add(sys.stdout, level=settings.log_level)
worker = TelegramStorageUploader()
worker = MediaUploader()
await worker.run_loop()
+222 -33
View File
@@ -5,11 +5,13 @@ import hashlib
import json
from datetime import datetime, timedelta, timezone
import aiohttp
from loguru import logger
from ..config import settings
from ..constants import (
JOB_TYPE_VK_STORAGE_COPY,
JOB_TYPE_MEDIA_STORAGE_COPY,
PLATFORM_SITE,
PLATFORM_VK,
POST_STATUS_SKIPPED,
POST_STATUS_STORAGE_PENDING,
@@ -17,9 +19,10 @@ from ..constants import (
SOURCE_STATUS_OK,
WORKER_PARSER,
)
from ..db import fetch_bool_setting, fetch_float_setting, fetch_int_setting, get_pool
from ..db import fetch_bool_setting, fetch_float_setting, fetch_int_setting, fetch_setting, get_pool
from ..heartbeat import HeartbeatReporter
from ..jobs import is_worker_enabled
from ..source_adapters import SiteParserClient, SourceItem, json_object
from ..vk_api import (
VKAPIClient,
VKAPIError,
@@ -29,6 +32,7 @@ from ..vk_api import (
is_repost,
post_vk_url,
)
from .ai_alerts import send_system_error_alert
def utc_from_ts(value: int) -> datetime:
@@ -51,17 +55,29 @@ class VKParserWorker:
async def init(self) -> None:
self.pool = await get_pool()
async def active_sources(self) -> list[dict]:
async def active_sources(self, vk_interval_sec: int, site_interval_minutes: int) -> list[dict]:
rows = await self.pool.fetch(
"""
SELECT *
FROM sources
WHERE platform=$1
WHERE platform=ANY($1::text[])
AND active=TRUE
AND archived_at IS NULL
AND (
last_checked_at IS NULL
OR (platform=$2 AND last_checked_at <= NOW() - $3::double precision * INTERVAL '1 second')
OR (
platform=$4
AND last_checked_at <= NOW() - $5::double precision * INTERVAL '1 minute'
)
)
ORDER BY last_checked_at NULLS FIRST, priority ASC, id ASC
""",
[PLATFORM_VK, PLATFORM_SITE],
PLATFORM_VK,
vk_interval_sec,
PLATFORM_SITE,
site_interval_minutes,
)
return [dict(r) for r in rows]
@@ -96,7 +112,12 @@ class VKParserWorker:
message[:1000],
)
async def mark_source_ok(self, source_id: int, last_parsed_at: datetime | None) -> None:
async def mark_source_ok(
self,
source_id: int,
last_parsed_at: datetime | None,
runtime_state: dict | None = None,
) -> None:
await self.pool.execute(
"""
UPDATE sources
@@ -104,12 +125,14 @@ class VKParserWorker:
status_msg=NULL,
last_checked_at=NOW(),
last_parsed_at=COALESCE($3, last_parsed_at),
runtime_state_json=COALESCE($4::jsonb, runtime_state_json),
updated_at=NOW()
WHERE id=$1
""",
source_id,
SOURCE_STATUS_OK,
last_parsed_at,
json.dumps(runtime_state, ensure_ascii=False) if runtime_state is not None else None,
)
async def resolve_source_if_needed(self, client: VKAPIClient, source: dict) -> dict:
@@ -136,6 +159,151 @@ class VKParserWorker:
source["name"] = resolved_name
return source
async def save_source_item(
self,
source: dict,
item: SourceItem,
*,
status: str,
skip_reason: str | None = None,
create_storage_job: bool = True,
) -> int | None:
source_id = int(source["id"])
platform = str(source["platform"])
media_urls = ",".join(media.url for media in item.media)
text_hash = make_hash(item.text)
content_hash = make_hash(item.text, media_urls)
raw = {**item.raw, "url": item.url, "media": [media.url for media in item.media]}
async with self.pool.acquire() as conn:
async with conn.transaction():
raw_post_id = await conn.fetchval(
"""
INSERT INTO raw_posts(
source_id, platform, external_post_id, original_url,
raw_text, raw_json, text_hash, content_hash,
posted_at, status, skip_reason
)
VALUES($1,$2,$3,$4,$5,$6::jsonb,$7,$8,$9,$10,$11)
ON CONFLICT (source_id, external_post_id) DO NOTHING
RETURNING id
""",
source_id,
platform,
item.external_id,
item.url,
item.text,
json.dumps(raw, ensure_ascii=False),
text_hash,
content_hash,
item.posted_at,
status,
skip_reason,
)
if raw_post_id is None:
return None
for order, media in enumerate(item.media):
await conn.execute(
"""
INSERT INTO raw_post_media(
raw_post_id, platform, media_type, original_url, sort_order
)
VALUES($1,$2,$3,$4,$5)
""",
raw_post_id,
platform,
media.media_type,
media.url,
order,
)
if create_storage_job:
await conn.execute(
"""
INSERT INTO jobs(type, entity_type, entity_id, payload_json, status)
VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending')
ON CONFLICT DO NOTHING
""",
JOB_TYPE_MEDIA_STORAGE_COPY,
raw_post_id,
)
return int(raw_post_id)
async def parse_site_source(self, client: SiteParserClient, source: dict) -> int:
source_id = int(source["id"])
lookback_days = max(1, await fetch_int_setting("parser_new_source_lookback_days", 14))
overlap_minutes = max(0, await fetch_int_setting("parser_reparse_overlap_minutes", 120))
last_parsed_at = source.get("last_parsed_at")
parse_from = source.get("parse_from")
since_dt = (
last_parsed_at - timedelta(minutes=overlap_minutes)
if last_parsed_at
else parse_from or datetime.now(timezone.utc) - timedelta(days=lookback_days)
)
if since_dt.tzinfo is None:
since_dt = since_dt.replace(tzinfo=timezone.utc)
fetched, runtime_state, partial_warning = await client.fetch(source)
recent = [item for item in fetched if item.posted_at > since_dt]
known = await self.known_post_ids(source_id, [item.external_id for item in recent])
candidates = [item for item in recent if item.external_id not in known]
default_min_text_length = max(0, await fetch_int_setting("site_parser_min_text_length", 50))
source_config = json_object(source.get("settings_json"))
min_text_length = max(0, int(source_config.get("min_text_length", default_min_text_length)))
skip_empty_text = await fetch_bool_setting("parser_skip_empty_text", True)
skip_no_media = await fetch_bool_setting("parser_skip_no_media", True)
skip_short_text = await fetch_bool_setting("parser_skip_text_too_short", True)
store_skipped = await fetch_bool_setting("parser_store_skipped_posts", False)
dedupe_content_hash = await fetch_bool_setting("parser_dedupe_content_hash", True)
hashes = {item.external_id: make_hash(item.text, ",".join(m.url for m in item.media)) for item in candidates}
known_hashes = await self.known_content_hashes(list(hashes.values())) if dedupe_content_hash else set()
saved = 0
for item in candidates:
content_hash = hashes[item.external_id]
if dedupe_content_hash and content_hash in known_hashes:
continue
skip_reason = None
if skip_empty_text and not item.body_text:
skip_reason = "empty_text"
elif skip_no_media and not item.media:
skip_reason = "no_media"
elif skip_short_text and len(item.body_text) < min_text_length:
skip_reason = "text_too_short"
if skip_reason and not store_skipped:
continue
raw_id = await self.save_source_item(
source,
item,
status=POST_STATUS_SKIPPED if skip_reason else POST_STATUS_STORAGE_PENDING,
skip_reason=skip_reason,
create_storage_job=not bool(skip_reason),
)
if raw_id:
saved += 1
known_hashes.add(content_hash)
max_seen = max((item.posted_at for item in fetched), default=last_parsed_at)
if partial_warning:
await self.pool.execute(
"""
UPDATE sources
SET status=$2, status_msg=$3, last_checked_at=NOW(),
last_parsed_at=COALESCE($4, last_parsed_at),
runtime_state_json=COALESCE($5::jsonb, runtime_state_json), updated_at=NOW()
WHERE id=$1
""",
source_id,
SOURCE_STATUS_ERROR,
partial_warning[:1000],
max_seen,
json.dumps(runtime_state, ensure_ascii=False) if runtime_state is not None else None,
)
else:
await self.mark_source_ok(source_id, max_seen, runtime_state)
logger.info(
"Parsed site source {}: fetched={} recent={} known={} saved={}",
source.get("name"), len(fetched), len(recent), len(known), saved,
)
return saved
async def known_post_ids(self, source_id: int, external_post_ids: list[str]) -> set[str]:
if not external_post_ids:
return set()
@@ -239,7 +407,7 @@ class VKParserWorker:
VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending')
ON CONFLICT DO NOTHING
""",
JOB_TYPE_VK_STORAGE_COPY,
JOB_TYPE_MEDIA_STORAGE_COPY,
raw_post_id,
)
return int(raw_post_id)
@@ -374,6 +542,8 @@ class VKParserWorker:
await self.heartbeat.beat(self.pool, status="disabled", force=True)
return
parser_interval_sec = max(10, await fetch_int_setting("parser_interval_sec", 300))
site_interval = max(1, await fetch_int_setting("site_parser_interval_minutes", 30))
rps = await fetch_int_setting("vk_requests_per_second", 3)
timeout_total = await fetch_int_setting("vk_api_timeout_total_sec", 15)
timeout_connect = await fetch_int_setting("vk_api_timeout_connect_sec", 5)
@@ -382,10 +552,10 @@ class VKParserWorker:
retry_min_delay = await fetch_float_setting("vk_api_retry_min_delay_sec", 2.0)
retry_max_delay = await fetch_float_setting("vk_api_retry_max_delay_sec", 10.0)
source_pause = max(0.0, await fetch_float_setting("parser_source_pause_sec", 0.0))
sources = await self.active_sources()
sources = await self.active_sources(parser_interval_sec, site_interval)
await self.heartbeat.beat(self.pool, meta={"sources": len(sources)})
if not sources:
logger.info("No active VK sources")
logger.debug("No sources due for parsing")
return
logger.info(
@@ -401,30 +571,50 @@ class VKParserWorker:
source_pause,
)
async with VKAPIClient(
rps=rps,
timeout_total_sec=timeout_total,
timeout_connect_sec=timeout_connect,
rate_limit_sleep_sec=rate_limit_sleep,
retry_attempts=retry_attempts,
retry_min_delay_sec=retry_min_delay,
retry_max_delay_sec=retry_max_delay,
) as client:
for source in sources:
try:
await self.parse_source(client, source)
except VKAPIError as e:
if is_fatal_source_error(e):
await self.deactivate_source(int(source["id"]), str(e))
logger.warning("VK source deactivated {}: {}", source.get("name"), e)
else:
async with aiohttp.ClientSession() as web_session:
site_client = SiteParserClient(
web_session,
str(await fetch_setting("site_parser_url", "") or "").strip(),
str(await fetch_setting("site_parser_token", "") or "").strip(),
str(await fetch_setting("site_parser_rucaptcha_token", "") or "").strip(),
await fetch_int_setting("site_parser_timeout_sec", 180),
)
async with VKAPIClient(
rps=rps,
timeout_total_sec=timeout_total,
timeout_connect_sec=timeout_connect,
rate_limit_sleep_sec=rate_limit_sleep,
retry_attempts=retry_attempts,
retry_min_delay_sec=retry_min_delay,
retry_max_delay_sec=retry_max_delay,
) as client:
for source in sources:
try:
if source.get("platform") == PLATFORM_VK:
await self.parse_source(client, source)
else:
await self.parse_site_source(site_client, source)
except VKAPIError as e:
if is_fatal_source_error(e):
await self.deactivate_source(int(source["id"]), str(e))
logger.warning("VK source deactivated {}: {}", source.get("name"), e)
else:
await self.mark_source_error(int(source["id"]), str(e))
logger.warning("VK source temporary error {}: {}", source.get("name"), e)
except Exception as e:
await self.mark_source_error(int(source["id"]), str(e))
logger.warning("VK source temporary error {}: {}", source.get("name"), e)
except Exception as e:
await self.mark_source_error(int(source["id"]), str(e))
logger.exception("Unexpected source error {}: {}", source.get("name"), e)
if source_pause:
await asyncio.sleep(source_pause)
logger.exception("Unexpected source error {}: {}", source.get("name"), e)
if source.get("platform") == PLATFORM_SITE and source.get("status") != SOURCE_STATUS_ERROR:
try:
await send_system_error_alert(
"Site Parser source failed\n"
f"source: {source.get('name') or source.get('url')}\n"
f"error: {str(e)[:1000]}"
)
except Exception as alert_exc:
logger.warning("Site Parser alert failed: {}", alert_exc)
if source_pause:
await asyncio.sleep(source_pause)
async def run_loop(self) -> None:
await self.init()
@@ -434,8 +624,7 @@ class VKParserWorker:
await self.run_once()
except Exception as e:
logger.exception("Parser loop error: {}", e)
interval = max(10, await fetch_int_setting("parser_interval_sec", 300))
await asyncio.sleep(interval)
await asyncio.sleep(10)
async def main() -> None:
+131
View File
@@ -0,0 +1,131 @@
from __future__ import annotations
import unittest
from vk_parser_app.source_adapters import SiteParserClient, validate_source_config
class FakeResponse:
status = 200
async def __aenter__(self):
return self
async def __aexit__(self, *_args):
return None
async def json(self, **_kwargs):
return {
"items": [{
"external_id": "post-1",
"url": "https://example.test/1",
"title": "Title",
"text": "Body",
"published_at": "2026-08-10T10:00:00+00:00",
"media": [{"type": "photo", "url": "https://example.test/1.jpg"}],
}],
"browser_state": {"cookies": [{"name": "cf_clearance"}]},
"browser_user_agent": "Test Browser",
}
class FollowedPageResponse(FakeResponse):
async def json(self, **_kwargs):
data = await super().json(**_kwargs)
data["items"][0]["text"] = "Title\nAuthor"
return data
class EmptyBodyResponse(FakeResponse):
async def json(self, **_kwargs):
data = await super().json(**_kwargs)
data["items"][0]["text"] = ""
return data
class FakeSession:
def post(self, *_args, **_kwargs):
return FakeResponse()
class FollowedPageSession(FakeSession):
def post(self, *_args, **_kwargs):
return FollowedPageResponse()
class EmptyBodySession(FakeSession):
def post(self, *_args, **_kwargs):
return EmptyBodyResponse()
class SourceAdapterTests(unittest.IsolatedAsyncioTestCase):
async def test_worker_response_is_normalized(self) -> None:
client = SiteParserClient(FakeSession(), "http://worker", "token", "captcha", 30)
items, state, warning = await client.fetch({
"url": "https://example.test/rss.xml",
"settings_json": '{"format":"rss","access":"auto"}',
"runtime_state_json": '{}',
})
self.assertEqual(items[0].text, "Title\n\nBody")
self.assertEqual(items[0].media[0].url, "https://example.test/1.jpg")
self.assertEqual(state["browser_state"]["cookies"][0]["name"], "cf_clearance")
self.assertEqual(state["browser_user_agent"], "Test Browser")
self.assertIsNone(warning)
async def test_followed_page_does_not_repeat_title(self) -> None:
client = SiteParserClient(FollowedPageSession(), "http://worker", "token", "captcha", 30)
items, _, _ = await client.fetch({
"url": "https://example.test/rss.xml",
"settings_json": '{"format":"rss","follow_links":true,"content_selector":"article.full"}',
"runtime_state_json": "{}",
})
self.assertEqual(items[0].text, "Title\nAuthor")
self.assertEqual(items[0].body_text, "Title\nAuthor")
async def test_title_is_not_counted_as_site_body(self) -> None:
client = SiteParserClient(EmptyBodySession(), "http://worker", "token", "captcha", 30)
items, _, _ = await client.fetch({
"url": "https://example.test/rss.xml",
"settings_json": '{"format":"rss"}',
"runtime_state_json": "{}",
})
self.assertEqual(items[0].text, "Title")
self.assertEqual(items[0].body_text, "")
def test_site_config_is_required(self) -> None:
with self.assertRaisesRegex(ValueError, "нужен конфиг"):
validate_source_config("site", {})
def test_follow_links_requires_selector(self) -> None:
with self.assertRaisesRegex(ValueError, "content_selector"):
validate_source_config("site", {"format": "rss", "follow_links": True})
def test_site_min_text_length_is_validated(self) -> None:
with self.assertRaisesRegex(ValueError, "min_text_length"):
validate_source_config("site", {"format": "rss", "min_text_length": "many"})
def test_versioned_site_config_is_validated(self) -> None:
validate_source_config("site", {
"version": 1,
"discovery": {"type": "rss", "limit": 10},
"detail": {
"enabled": True,
"root_selector": "article",
"fields": {"text": {"selector": ".body", "extract": "text"}},
},
"access": {"type": "cloudflare"},
"retry": {"attempts": 3, "delay_seconds": 5, "timeout_seconds": 90},
"min_text_length": 50,
})
def test_html_discovery_requires_item_selector(self) -> None:
with self.assertRaisesRegex(ValueError, "item_selector"):
validate_source_config("site", {
"version": 1,
"discovery": {"type": "html"},
})
if __name__ == "__main__":
unittest.main()
+63
View File
@@ -0,0 +1,63 @@
import asyncio
from pathlib import Path
from vk_parser_app.workers.media_uploader import MediaUploader, video_provider
class FakeContent:
async def iter_chunked(self, _size):
yield b"video"
class FakeResponse:
status = 200
content = FakeContent()
async def __aenter__(self):
return self
async def __aexit__(self, *_args):
return None
class FakeSession:
def get(self, *_args, **_kwargs):
return FakeResponse()
def test_video_provider() -> None:
assert video_provider("https://vk.com/video-1_2") == "vk"
assert video_provider("https://youtu.be/dQw4w9WgXcQ") == "youtube"
assert video_provider("https://www.youtube.com/watch?v=dQw4w9WgXcQ") == "youtube"
assert video_provider("file:///etc/passwd") is None
assert video_provider("http://127.0.0.1/video.mp4") is None
def test_site_request_options() -> None:
options = MediaUploader().site_request_options({
"platform": "site",
"source_url": "https://example.test/rss.xml",
"source_runtime_state": {
"browser_user_agent": "Test Browser",
"browser_state": {"cookies": [{"name": "cf_clearance", "value": "secret"}]},
},
})
assert options["headers"] == {
"User-Agent": "Test Browser",
"Referer": "https://example.test/rss.xml",
}
assert options["cookies"] == {"cf_clearance": "secret"}
async def test_http_video_download() -> None:
path = "/tmp/vkparser_tg_media_http_test.mp4"
result = await MediaUploader().download_http_video(FakeSession(), "https://example.test/a.mp4", path)
assert result["size_bytes"] == 5
assert Path(path).read_bytes() == b"video"
Path(path).unlink()
if __name__ == "__main__":
test_video_provider()
test_site_request_options()
asyncio.run(test_http_video_download())