Compare commits

...

9 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
29 changed files with 1920 additions and 227 deletions
+72 -19
View File
@@ -61,7 +61,7 @@
│ ├── templates/ # HTML-шаблоны │ ├── templates/ # HTML-шаблоны
│ └── workers/ │ └── workers/
│ ├── parser.py # VK-парсер │ ├── parser.py # VK-парсер
│ ├── vk_storage_uploader.py# Telegram media storage uploader │ ├── media_uploader.py# Telegram media storage uploader
│ ├── ai_qualifier.py # AI-квалификатор │ ├── ai_qualifier.py # AI-квалификатор
│ └── ai_writer.py # AI-райтер │ └── ai_writer.py # AI-райтер
├── original_project/ # старый проект как справочник ├── original_project/ # старый проект как справочник
@@ -196,7 +196,8 @@
Источники для парсинга. Источники для парсинга.
Сейчас реально используется только `platform='vk'`, но схема заложена под другие площадки. Используются `platform='vk'` и `platform='site'`. Для сайта конфигурация
конкретного адаптера хранится вместе с источником в `settings_json`.
Важные поля: Важные поля:
@@ -215,9 +216,14 @@
- `parse_from` - `parse_from`
- `priority` - `priority`
- `archived_at` - `archived_at`
- `settings_json` для `platform='site'`
`tag` нужен для будущих хэштегов и финального оформления. `tag` нужен для будущих хэштегов и финального оформления.
Для сайта отсутствие или ошибка JSON-конфига блокирует запуск с видимой
ошибкой источника. Конфиг задаёт способ получения именно этого сайта; общее
расписание и секреты остаются в `app_settings`.
### 5.4. `raw_posts` ### 5.4. `raw_posts`
Главная таблица жизненного цикла поста. Главная таблица жизненного цикла поста.
@@ -338,7 +344,7 @@ AI-райтер:
Сейчас основной тип: Сейчас основной тип:
- `vk.storage.copy` - `media.storage.copy`
Название историческое: сначала планировался VK storage, потом медиа-сторедж вернулся в Telegram. Тип задачи пока не переименован. Название историческое: сначала планировался VK storage, потом медиа-сторедж вернулся в Telegram. Тип задачи пока не переименован.
@@ -357,7 +363,7 @@ AI-райтер:
Имена воркеров: Имена воркеров:
- `vk-parser` - `vk-parser`
- `vk-storage-uploader` - `media-uploader`
- `ai-qualifier` - `ai-qualifier`
- `ai-writer` - `ai-writer`
- `tg-poster` - `tg-poster`
@@ -556,7 +562,44 @@ Retention: записи старше 30 дней удаляются при ст
6. При включенном `parser_dedupe_content_hash` также убирает дубли по `content_hash`. 6. При включенном `parser_dedupe_content_hash` также убирает дубли по `content_hash`.
7. Применяет политику мусорных постов. 7. Применяет политику мусорных постов.
8. Сохраняет raw post и media. 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. Политика мусорных постов ### 7.2. Политика мусорных постов
@@ -615,13 +658,13 @@ Retention: записи старше 30 дней удаляются при ст
## 8. Media storage uploader ## 8. Media storage uploader
Файл: `src/vk_parser_app/workers/vk_storage_uploader.py`. Файл: `src/vk_parser_app/workers/media_uploader.py`.
Название файла историческое: сейчас фактическое хранилище медиа - Telegram, а не VK. Название файла историческое: сейчас фактическое хранилище медиа - Telegram, а не VK.
### 8.1. Что делает ### 8.1. Что делает
1. Забирает job типа `vk.storage.copy`. 1. Забирает job типа `media.storage.copy`.
2. Загружает `raw_post` и media. 2. Загружает `raw_post` и media.
3. Скачивает изображения/видео. 3. Скачивает изображения/видео.
4. Отправляет в Telegram storage channel. 4. Отправляет в Telegram storage channel.
@@ -1091,6 +1134,22 @@ Worker учитывает:
- active - active
- parse_from - 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 ### 12.3. Raw
Raw-страница показывает: Raw-страница показывает:
@@ -1247,7 +1306,7 @@ uvicorn vk_parser_app.admin:app --host 0.0.0.0 --port 8080
```powershell ```powershell
$env:PYTHONPATH="src" $env:PYTHONPATH="src"
python -m vk_parser_app.workers.parser 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_qualifier
python -m vk_parser_app.workers.ai_writer python -m vk_parser_app.workers.ai_writer
python -m vk_parser_app.workers.tg_poster python -m vk_parser_app.workers.tg_poster
@@ -1376,16 +1435,10 @@ systemctl start ai-writer.service
Они не влияют на текущую очередь, но могут путать в админке. Их стоит аккуратно пометить как failed/cancelled или скрывать служебные тесты. Они не влияют на текущую очередь, но могут путать в админке. Их стоит аккуратно пометить как failed/cancelled или скрывать служебные тесты.
### 20.2. Название `vk_storage_uploader.py` ### 20.2. Media uploader
Файл исторически называется VK storage uploader, но фактически загружает в Telegram. Можно позже переименовать: Uploader переименован в `media_uploader.py`, worker — в `media-uploader`, а jobs — в
`media.storage.copy`. Миграция сохраняет уже созданные задания.
```text
vk_storage_uploader.py -> tg_storage_uploader.py
JOB_TYPE_VK_STORAGE_COPY -> tg.storage.copy
```
Делать только миграцией и аккуратно, чтобы не потерять jobs.
### 20.3. Настройки Telegram uploader ### 20.3. Настройки Telegram uploader
@@ -1449,7 +1502,7 @@ JOB_TYPE_VK_STORAGE_COPY -> tg.storage.copy
sources sources
-> parser -> parser
-> raw_posts + raw_post_media -> raw_posts + raw_post_media
-> jobs(vk.storage.copy) -> jobs(media.storage.copy)
-> Telegram storage uploader -> Telegram storage uploader
-> status=storage_ready -> status=storage_ready
-> AI qualifier -> AI qualifier
@@ -1485,7 +1538,7 @@ admin_sessions
```text ```text
src/vk_parser_app/admin.py src/vk_parser_app/admin.py
src/vk_parser_app/workers/parser.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_qualifier.py
src/vk_parser_app/workers/ai_writer.py src/vk_parser_app/workers/ai_writer.py
src/vk_parser_app/workers/tg_poster.py src/vk_parser_app/workers/tg_poster.py
+25 -2
View File
@@ -1,6 +1,6 @@
# RAA deployment # 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 RAA is not a separate code project anymore. The source of truth for code and
documentation is: documentation is:
@@ -26,6 +26,7 @@ continue feature work there.
- Destination: `raa-fn8`, CTID `105`, IP `192.168.1.105` - 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` - 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` - Health check: `https://raa.panel.f-n8.ru/health`
- Current verified deploy: `f866a13a8a83598e26a7a6a3406b8bc1c01c9a80`
## Per-Project Settings ## Per-Project Settings
@@ -55,6 +56,10 @@ Database settings:
- `vk_poster_access_token` is the RAA community token - `vk_poster_access_token` is the RAA community token
- `site_poster_enabled=false` - `site_poster_enabled=false`
- `site_poster_provider=""` until a website adapter is configured - `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: - Optional social text headers/footers:
- `tg_poster_header_text` / `tg_poster_footer_text` - `tg_poster_header_text` / `tg_poster_footer_text`
- `vk_poster_header_text` / `vk_poster_footer_text` - `vk_poster_header_text` / `vk_poster_footer_text`
@@ -70,7 +75,7 @@ Database settings:
As last checked on 2026-08-03: As last checked on 2026-08-03:
- `vk-parser=true` - `vk-parser=true`
- `vk-storage-uploader=true` - `media-uploader=true`
- `ai-qualifier=false` - `ai-qualifier=false`
- `ai-writer=false` - `ai-writer=false`
- `tg-poster=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 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. 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
RAA categories live in the RAA `content_categories` table. The UI has only RAA categories live in the RAA `content_categories` table. The UI has only
+29 -13
View File
@@ -124,13 +124,29 @@ https://vk.com/wall-239548476_123
} }
``` ```
Если 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 воркера, его токен, ключ В глобальном разделе `Site Parser` задаются URL воркера, его токен, ключ
RuCaptcha и таймаут. Внешний модуль запускается командой: RuCaptcha и таймаут. Внешний модуль запускается командой:
```bash ```bash
docker build -t site-parser-worker site_parser_worker cd site_parser_worker
docker run -d --restart unless-stopped -p 8080:8080 \ printf 'WORKER_TOKEN=replace-me\n' > .env
-e WORKER_TOKEN=replace-me site-parser-worker chmod 600 .env
docker-compose up -d --build
docker-compose ps
``` ```
## Запуск локально ## Запуск локально
@@ -147,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 PYTHONPATH=src .venv/bin/python -m vk_parser_app.workers.parser
``` ```
VK storage uploader: Media uploader:
```bash ```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 ## systemd units
@@ -191,17 +207,17 @@ RestartSec=10
WantedBy=multi-user.target WantedBy=multi-user.target
``` ```
`vk-storage-uploader.service`: `media-uploader.service`:
```ini ```ini
[Unit] [Unit]
Description=VK Storage Uploader Worker Description=Media Uploader Worker
After=network.target postgresql.service After=network.target postgresql.service
[Service] [Service]
WorkingDirectory=/opt/vk-parser WorkingDirectory=/opt/vk-parser
Environment=PYTHONPATH=/opt/vk-parser/src 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 Restart=always
RestartSec=10 RestartSec=10
@@ -212,13 +228,13 @@ WantedBy=multi-user.target
## Текущий pipeline ## Текущий pipeline
```text ```text
sources(vk) sources(vk/site)
-> vk-parser -> vk-parser
-> raw_posts + raw_post_media -> raw_posts + raw_post_media
-> jobs(type='vk.storage.copy') -> jobs(type='media.storage.copy')
-> vk-storage-uploader -> media-uploader
-> wall.post в storage-группе -> локальный Telegram Bot API
-> raw_posts.storage_post_url -> raw_posts.storage_post_url
``` ```
Фото копируются в storage-группу. Видео в первой версии сохраняются как `link_only`/исходное VK-вложение, чтобы сначала стабилизировать основной контур. Фото и видео загружаются в Telegram media storage; исходные URL остаются в `raw_post_media`.
+57 -3
View File
@@ -1,6 +1,6 @@
# N8 Parser: current state # 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. 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` - Templates: `src/vk_parser_app/templates`
- Workers: - Workers:
- VK parser: `src/vk_parser_app/workers/parser.py` - 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 qualifier: `src/vk_parser_app/workers/ai_qualifier.py`
- AI writer: `src/vk_parser_app/workers/ai_writer.py` - AI writer: `src/vk_parser_app/workers/ai_writer.py`
- DB migrations: `db/migrations` - DB migrations: `db/migrations`
@@ -37,6 +37,8 @@ Current runtime is local infrastructure, not the old VPS.
1. Sources are configured in admin. 1. Sources are configured in admin.
2. VK parser reads enabled sources and stores raw posts. 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. 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. 4. AI qualifier scores raw posts and marks accepted/rejected/maybe.
5. AI writer rewrites accepted posts into editor-ready drafts. 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`. - 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. - `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`. - 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 and publication workers are intentionally disabled until explicitly started:
`ai-qualifier=false`, `ai-writer=false`, `tg-poster=false`, `vk-poster=false`, `ai-qualifier=false`, `ai-writer=false`, `tg-poster=false`, `vk-poster=false`,
`site-poster=false`, `tg_poster_enabled=false`, `vk_poster_enabled=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-sender`) remains a valid legacy VK reposting app and is not the shared
RAA parser/writer/poster app. 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 ## Important Text Rules
- AI writer output text must be clean: no physical hashtags at the end. - AI writer output text must be clean: no physical hashtags at the end.
@@ -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.constants import MEDIA_STATUS_FAILED, MEDIA_STATUS_LINK_ONLY
from vk_parser_app.db import get_pool 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: def jlog(event: str, **payload: Any) -> None:
print(json.dumps({"event": event, **payload}, ensure_ascii=False, default=str), flush=True) 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( rows = await worker.pool.fetch(
""" """
SELECT m.*, rp.original_url AS post_url 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] 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 = ( caption = (
f"Backfill video media #{int(media['id'])} for raw_post #{int(media['raw_post_id'])}\n" 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>' 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: async def main() -> None:
limit = int(os.getenv("BACKFILL_LIMIT", "200")) limit = int(os.getenv("BACKFILL_LIMIT", "200"))
worker = TelegramStorageUploader() worker = MediaUploader()
await worker.init() await worker.init()
videos = await load_pending_videos(worker, limit) videos = await load_pending_videos(worker, limit)
jlog("start", count=len(videos), limit=limit) jlog("start", count=len(videos), limit=limit)
+1 -1
View File
@@ -4,7 +4,7 @@ WORKDIR /app
COPY requirements.txt . COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt RUN pip install --no-cache-dir -r requirements.txt
RUN playwright install chrome RUN playwright install chrome
COPY app.py . COPY app.py extractor.py ./
ENV PYTHONUNBUFFERED=1 ENV PYTHONUNBUFFERED=1
EXPOSE 8080 EXPOSE 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.
+379 -83
View File
@@ -1,23 +1,30 @@
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import logging
import os import os
import secrets import secrets
from contextlib import asynccontextmanager
from datetime import datetime, timezone from datetime import datetime, timezone
from time import struct_time from time import struct_time
from typing import Annotated, Any from typing import Annotated, Any, AsyncIterator
from urllib.parse import urljoin
import feedparser import feedparser
import httpx import httpx
from bs4 import BeautifulSoup from bs4 import BeautifulSoup
from curl_cffi.requests import AsyncSession as CurlAsyncSession
from fastapi import FastAPI, Header, HTTPException from fastapi import FastAPI, Header, HTTPException
from playwright.async_api import async_playwright from playwright.async_api import BrowserContext, Page, TimeoutError as PlaywrightTimeoutError, async_playwright
from playwright_captcha import CaptchaType, FrameworkType, TwoCaptchaSolver from playwright_captcha import CaptchaType, FrameworkType, TwoCaptchaSolver
from pydantic import BaseModel, Field, HttpUrl, SecretStr from pydantic import BaseModel, Field, HttpUrl, SecretStr
from twocaptcha import AsyncTwoCaptcha from twocaptcha import AsyncTwoCaptcha
app = FastAPI(title="Site Parser Worker", version="0.1.0") 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() browser_lock = asyncio.Lock()
logger = logging.getLogger("site_parser")
class ParseRequest(BaseModel): class ParseRequest(BaseModel):
@@ -25,6 +32,7 @@ class ParseRequest(BaseModel):
config: dict[str, Any] config: dict[str, Any]
rucaptcha_token: SecretStr | None = None rucaptcha_token: SecretStr | None = None
browser_state: dict[str, Any] | None = None browser_state: dict[str, Any] | None = None
browser_user_agent: str | None = None
class ParsedItem(BaseModel): class ParsedItem(BaseModel):
@@ -38,6 +46,14 @@ class ParsedItem(BaseModel):
media: list[dict[str, str]] = Field(default_factory=list) 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: def require_token(worker_token: str | None) -> None:
expected = os.getenv("WORKER_TOKEN", "") expected = os.getenv("WORKER_TOKEN", "")
if not expected: if not expected:
@@ -59,74 +75,354 @@ def iso_date(value: struct_time | None) -> str | None:
return datetime(*value[:6], tzinfo=timezone.utc).isoformat() return datetime(*value[:6], tzinfo=timezone.utc).isoformat()
def clean_html(value: Any) -> tuple[str, str, list[dict[str, str]]]: def clean_feed_content(value: Any, base_url: str) -> tuple[str, str, list[dict[str, str]]]:
html = str(value or "").strip() html = str(value or "").strip()
soup = BeautifulSoup(html, "html.parser") soup = BeautifulSoup(html, "html.parser")
for element in soup.find_all(["script", "style", "noscript"]):
element.decompose()
text = soup.get_text("\n", strip=True) text = soup.get_text("\n", strip=True)
media = [{"type": "photo", "url": str(image["src"])} for image in soup.find_all("img", src=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 return html, text, media
def parse_rss(xml: str, max_items: int) -> list[ParsedItem]: def default_rss_fields() -> dict[str, Any]:
feed = feedparser.parse(xml) 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: if feed.bozo and not feed.entries:
raise ValueError(f"Invalid RSS: {feed.bozo_exception}") raise ValueError(f"Invalid RSS: {feed.bozo_exception}")
items = [] discovery = config["discovery"]
for entry in feed.entries[:max_items]: fields = discovery.get("fields") or default_rss_fields()
title = BeautifulSoup(str(entry.get("title") or ""), "html.parser").get_text(" ", strip=True) items: list[dict[str, Any]] = []
html, text, media = clean_html(entry.get("description") or entry.get("summary") or "") for entry in feed.entries[: discovery["limit"]]:
for enclosure in entry.get("enclosures") or []: raw = dict(entry)
url = str(enclosure.get("href") or enclosure.get("url") or "").strip() 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 "") media_type = str(enclosure.get("type") or "")
if url and (not media_type or media_type.startswith("image/")): if enclosure_url and media_type.startswith("video/"):
media.append({"type": "photo", "url": url}) media.append({"type": "video", "url": enclosure_url})
url = str(entry.get("link") or "").strip() elif enclosure_url:
media.append({"type": "photo", "url": enclosure_url})
items.append( items.append(
ParsedItem( {
external_id=str(entry.get("id") or entry.get("guid") or url).strip(), "external_id": external_id,
url=url, "url": url,
title=title, "title": str(values.get("title") or BeautifulSoup(str(raw.get("title") or ""), "html.parser").get_text(" ", strip=True)),
text=text, "text": text,
html=html, "html": html,
published_at=iso_date(entry.get("published_parsed") or entry.get("updated_parsed")), "published_at": values.get("published_at") or iso_date(raw.get("published_parsed") or raw.get("updated_parsed")),
author=str(entry.get("author") or "").strip() or None, "author": values.get("author") or str(raw.get("author") or "").strip() or None,
media=list({item["url"]: item for item in media}.values()), "media": list({entry["url"]: entry for entry in media}.values()),
) "source_context": {"rss": raw},
"discovery_errors": errors,
}
) )
return items return items
async def fetch_in_browser( def discover_html(html: str, source_url: str, config: dict[str, Any]) -> list[dict[str, Any]]:
url: str, soup = BeautifulSoup(html, "html.parser")
rucaptcha_token: str, discovery = config["discovery"]
browser_state: dict[str, Any] | None, fields = discovery.get("fields") or {}
) -> tuple[str, dict[str, Any]]: items: list[dict[str, Any]] = []
# ponytail: one browser at a time; use a queue only when parallel source parsing is needed. for element in soup.select(str(discovery["item_selector"]))[: discovery["limit"]]:
async with browser_lock: values, errors = extract_fields(element, fields, {"page": {"url": source_url}})
async with async_playwright() as playwright: url = urljoin(source_url, str(values.get("url") or "").strip())
browser = await playwright.chromium.launch(channel="chrome", headless=False) external_id = str(values.get("external_id") or url).strip()
context = await browser.new_context(storage_state=browser_state) if browser_state else await browser.new_context() if not url or not external_id:
page = await context.new_page() 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( async with TwoCaptchaSolver(
framework=FrameworkType.PLAYWRIGHT, framework=FrameworkType.PLAYWRIGHT,
page=page, page=page,
async_two_captcha_client=AsyncTwoCaptcha(rucaptcha_token), async_two_captcha_client=AsyncTwoCaptcha(token),
max_attempts=1, max_attempts=1,
) as solver: ) as solver:
response = await page.goto(url, wait_until="domcontentloaded", timeout=60_000) yield solver
if response and response.status in {403, 429, 503} and "Just a moment" in await page.title():
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( await solver.solve_captcha(
captcha_container=page, captcha_container=page,
captcha_type=CaptchaType.CLOUDFLARE_INTERSTITIAL, captcha_type=CaptchaType.CLOUDFLARE_INTERSTITIAL,
) )
if wait_for:
try: try:
await page.locator("pre").wait_for(state="visible", timeout=60_000) 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: except Exception as exc:
raise RuntimeError(f"RSS did not load after Cloudflare challenge: {await page.title()}") from exc logger.warning("Detail failed %s: %s", item.get("url"), exc)
xml = await page.locator("pre").inner_text() errors.append({"url": item.get("url"), "stage": "detail", "attempts": attempts, "error": str(exc)})
items = enriched
state = await context.storage_state() state = await context.storage_state()
user_agent = await page.evaluate("navigator.userAgent")
await browser.close() await browser.close()
return xml, state return items, errors, state, user_agent
@app.get("/health") @app.get("/health")
@@ -140,53 +436,53 @@ async def parse_source(
worker_token: Annotated[str | None, Header(alias="X-Worker-Token")] = None, worker_token: Annotated[str | None, Header(alias="X-Worker-Token")] = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
require_token(worker_token) require_token(worker_token)
config = request.config
if not config:
raise HTTPException(status_code=422, detail="config is required and cannot be empty")
if config.get("format") != "rss":
raise HTTPException(status_code=422, detail="Only config.format=rss is supported")
access = str(config.get("access") or "auto")
if access not in {"auto", "http", "cloudflare"}:
raise HTTPException(status_code=422, detail="config.access must be auto, http or cloudflare")
try: try:
max_items = int(config.get("max_items", 20)) config = normalize_config(request.config)
except (TypeError, ValueError) as exc: except ConfigError as exc:
raise HTTPException(status_code=422, detail="config.max_items must be an integer") from exc raise HTTPException(status_code=422, detail=str(exc)) from exc
if not 1 <= max_items <= 100:
raise HTTPException(status_code=422, detail="config.max_items must be between 1 and 100") source_url = str(request.url)
url = str(request.url) access = config["access"]["type"]
xml = ""
state = request.browser_state state = request.browser_state
browser_user_agent = request.browser_user_agent
fetched_via = "http" fetched_via = "http"
if access != "cloudflare": logger.info("Parse request source=%s discovery=%s detail=%s", source_url, config["discovery"]["type"], config["detail"]["enabled"])
try: try:
async with httpx.AsyncClient(follow_redirects=True, timeout=30) as client: if access in {"browser", "cloudflare"}:
response = await client.get(url, headers={"User-Agent": "Mozilla/5.0 SiteParser/0.1"}) items, errors, state, browser_user_agent = await parse_in_browser(
except httpx.HTTPError as exc: source_url,
raise HTTPException(status_code=502, detail=f"Source request failed: {exc}") from exc config,
if not is_cloudflare_challenge(response.status_code, response.text): request.rucaptcha_token.get_secret_value() if request.rucaptcha_token else None,
state,
browser_user_agent,
)
fetched_via = "browser"
else:
try: try:
response.raise_for_status() items, errors = await parse_over_http(source_url, config)
except httpx.HTTPStatusError as exc: except CloudflareRequired:
raise HTTPException(status_code=502, detail=f"Source returned HTTP {response.status_code}") from exc if access == "http":
xml = response.text raise
elif access == "http": items, errors, state, browser_user_agent = await parse_in_browser(
raise HTTPException(status_code=502, detail="Cloudflare challenge received in http-only mode") source_url,
if not xml: config,
if not request.rucaptcha_token: request.rucaptcha_token.get_secret_value() if request.rucaptcha_token else None,
raise HTTPException(status_code=422, detail="rucaptcha_token is required for Cloudflare") state,
fetched_via = "cloudflare" browser_user_agent,
try: )
xml, state = await fetch_in_browser(url, request.rucaptcha_token.get_secret_value(), state) fetched_via = "browser"
except Exception as exc: except Exception as exc:
raise HTTPException(status_code=502, detail=str(exc)) from exc raise HTTPException(status_code=502, detail=str(exc)) from exc
try:
items = parse_rss(xml, max_items) if not items and not errors:
except ValueError as exc: errors = [{"url": source_url, "stage": "discovery", "attempts": 1, "error": "no items found"}]
raise HTTPException(status_code=502, detail=str(exc)) from exc status = "partial" if errors and items else "failed" if errors else "ok"
return { return {
"source_url": url, "status": status,
"source_url": source_url,
"fetched_via": fetched_via, "fetched_via": fetched_via,
"items": [item.model_dump() for item in items], "items": [public_item(item).model_dump() for item in items],
"errors": errors,
"browser_state": state, "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
+1
View File
@@ -1,5 +1,6 @@
2captcha-python-async==1.5.1 2captcha-python-async==1.5.1
beautifulsoup4==4.15.0 beautifulsoup4==4.15.0
curl_cffi==0.16.0
fastapi==0.141.1 fastapi==0.141.1
feedparser==6.0.14 feedparser==6.0.14
httpx==0.28.1 httpx==0.28.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()
+3 -1
View File
@@ -4,12 +4,14 @@ from app import parse_rss
def test_parse_rss() -> None: def test_parse_rss() -> None:
xml = """<rss><channel><item><guid>1</guid><title>Title</title> xml = """<rss><channel><item><guid>1</guid><title>Title</title>
<link>https://example.test/1</link> <link>https://example.test/1</link>
<description>&lt;p&gt;Body&lt;/p&gt;&lt;img src="https://example.test/1.jpg"&gt;</description> <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>""" <pubDate>Mon, 10 Aug 2026 10:00:00 +0000</pubDate></item></channel></rss>"""
item = parse_rss(xml, 20)[0] item = parse_rss(xml, 20)[0]
assert item.external_id == "1" assert item.external_id == "1"
assert item.text == "Body" assert item.text == "Body"
assert item.media[0]["url"] == "https://example.test/1.jpg" 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__": if __name__ == "__main__":
+107 -9
View File
@@ -9,14 +9,18 @@ import re
import secrets import secrets
import time import time
from datetime import date, datetime, timedelta, timezone from datetime import date, datetime, timedelta, timezone
from io import BytesIO
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
from urllib.parse import urlencode from urllib.parse import urlencode
from zoneinfo import ZoneInfo from zoneinfo import ZoneInfo
import aiohttp import aiohttp
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 import FastAPI, File, Form, HTTPException, Request, UploadFile, status
from fastapi.responses import FileResponse, HTMLResponse, JSONResponse, RedirectResponse from fastapi.responses import FileResponse, HTMLResponse, JSONResponse, RedirectResponse, Response
from fastapi.staticfiles import StaticFiles from fastapi.staticfiles import StaticFiles
from fastapi.templating import Jinja2Templates from fastapi.templating import Jinja2Templates
from loguru import logger from loguru import logger
@@ -25,7 +29,7 @@ from .config import settings
from .constants import PLATFORM_SITE, PLATFORM_VK from .constants import PLATFORM_SITE, PLATFORM_VK
from .db import fetch_int_setting, fetch_setting, get_pool from .db import fetch_int_setting, fetch_setting, get_pool
from .security import hash_password, new_token, token_hash, verify_password from .security import hash_password, new_token, token_hash, verify_password
from .source_adapters import validate_source_config from .source_adapters import json_object, validate_source_config
from .text_utils import build_publication_text, normalize_hash_tag, parse_categories from .text_utils import build_publication_text, normalize_hash_tag, parse_categories
from .vk_api import VKAPIClient, normalize_vk_source from .vk_api import VKAPIClient, normalize_vk_source
from .workers.ai_qualifier import AIQualifierWorker, normalize_model, response_usage from .workers.ai_qualifier import AIQualifierWorker, normalize_model, response_usage
@@ -37,7 +41,7 @@ from .workers.site_poster import SitePoster
from .workers.tg_poster import TelegramPoster from .workers.tg_poster import TelegramPoster
from .workers.tg_reactor import TelegramReactor from .workers.tg_reactor import TelegramReactor
from .workers.vk_poster import VKPoster from .workers.vk_poster import VKPoster
from .workers.vk_storage_uploader import TelegramStorageUploader from .workers.media_uploader import MediaUploader
COOKIE_NAME = "vk_parser_admin" COOKIE_NAME = "vk_parser_admin"
VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier" VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier"
@@ -366,6 +370,14 @@ SETTING_ORDER = {
"parser_dedupe_content_hash", "parser_dedupe_content_hash",
"parser_source_pause_sec", "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": [
"vk_requests_per_second", "vk_requests_per_second",
"vk_wall_page_size", "vk_wall_page_size",
@@ -494,6 +506,40 @@ async def get_current_user(request: Request) -> dict | None:
return dict(row) if row else 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: def require_csrf(user: dict, csrf_token: str) -> None:
if not user or csrf_token != user["csrf_token"]: if not user or csrf_token != user["csrf_token"]:
raise PermissionError("bad csrf") raise PermissionError("bad csrf")
@@ -1813,7 +1859,7 @@ async def startup() -> None:
("tg-poster", TelegramPoster()), ("tg-poster", TelegramPoster()),
("tg-reactor", TelegramReactor()), ("tg-reactor", TelegramReactor()),
("vk-poster", VKPoster()), ("vk-poster", VKPoster()),
("vk-storage-uploader", TelegramStorageUploader()), ("media-uploader", MediaUploader()),
] ]
for name, worker in workers: for name, worker in workers:
asyncio.create_task(start_worker_task(worker, name)) asyncio.create_task(start_worker_task(worker, name))
@@ -2196,6 +2242,7 @@ async def sources_list(request: Request, q: str = "", status_filter: str = "", a
sources = [] sources = []
for row in rows: for row in rows:
item = dict(row) 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_parsed_at_fmt"] = format_dt(item.get("last_parsed_at"))
item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at")) item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at"))
item["created_at_fmt"] = format_dt(item.get("created_at")) item["created_at_fmt"] = format_dt(item.get("created_at"))
@@ -2274,6 +2321,7 @@ async def sources_preview(
sources = [] sources = []
for row in rows: for row in rows:
item = dict(row) 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_parsed_at_fmt"] = format_dt(item.get("last_parsed_at"))
item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at")) item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at"))
item["created_at_fmt"] = format_dt(item.get("created_at")) item["created_at_fmt"] = format_dt(item.get("created_at"))
@@ -2439,7 +2487,7 @@ async def source_edit(request: Request, source_id: int):
base_context( base_context(
request, request,
user, user,
source=dict(source), source={**dict(source), "settings_json": json_object(source["settings_json"])},
action=f"/sources/{source_id}/edit", action=f"/sources/{source_id}/edit",
title=f"Источник #{source_id}", title=f"Источник #{source_id}",
), ),
@@ -2532,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) source_row = await pool.fetchrow("SELECT * FROM sources WHERE id=$1", source_id)
if source_row: if source_row:
s = dict(source_row) s = dict(source_row)
s["settings_json"] = json_object(s.get("settings_json"))
# Fetch stats just for this source to render correctly # Fetch stats just for this source to render correctly
stats = await pool.fetchrow( stats = await pool.fetchrow(
""" """
@@ -2549,6 +2598,25 @@ async def source_toggle(request: Request, source_id: int, csrf_token: str = Form
return redirect("/sources") 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") @app.post("/sources/{source_id}/delete")
async def source_delete(request: Request, source_id: int, csrf_token: str = Form(...)): async def source_delete(request: Request, source_id: int, csrf_token: str = Form(...)):
user = await get_current_user(request) user = await get_current_user(request)
@@ -2650,7 +2718,11 @@ async def raw_posts(
jsonb_build_object( jsonb_build_object(
'id', rpm.id, 'id', rpm.id,
'type', rpm.media_type, '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, 'status', rpm.status,
'error', rpm.error, 'error', rpm.error,
'duration_sec', rpm.duration_sec 'duration_sec', rpm.duration_sec
@@ -2816,7 +2888,11 @@ async def fetch_single_editor_row(pool, post_id: int):
jsonb_build_object( jsonb_build_object(
'id', rpm.id, 'id', rpm.id,
'type', rpm.media_type, '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, 'status', rpm.status,
'error', rpm.error, 'error', rpm.error,
'duration_sec', rpm.duration_sec 'duration_sec', rpm.duration_sec
@@ -2957,7 +3033,11 @@ async def editor_feed(
jsonb_build_object( jsonb_build_object(
'id', rpm.id, 'id', rpm.id,
'type', rpm.media_type, '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, 'status', rpm.status,
'error', rpm.error, 'error', rpm.error,
'duration_sec', rpm.duration_sec 'duration_sec', rpm.duration_sec
@@ -3297,7 +3377,11 @@ async def raw_post_detail(request: Request, post_id: int):
jsonb_build_object( jsonb_build_object(
'id', rpm.id, 'id', rpm.id,
'type', rpm.media_type, '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, 'status', rpm.status,
'error', rpm.error, 'error', rpm.error,
'duration_sec', rpm.duration_sec 'duration_sec', rpm.duration_sec
@@ -3546,6 +3630,18 @@ async def workers(request: Request):
continue continue
settings.append(item) settings.append(item)
settings.sort(key=setting_sort_key) 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_provider = str(setting_value(settings, "ai_qualifier_provider", "openrouter") or "openrouter")
ai_model = str(setting_value(settings, "ai_qualifier_model", "") or "") ai_model = str(setting_value(settings, "ai_qualifier_model", "") or "")
ai_api_key = str(setting_value(settings, "ai_qualifier_api_key", "") or "") ai_api_key = str(setting_value(settings, "ai_qualifier_api_key", "") or "")
@@ -3580,6 +3676,8 @@ async def workers(request: Request):
for r in rows for r in rows
], ],
settings=settings, settings=settings,
site_parser_url=site_parser_url,
site_parser_health=site_parser_health,
provider_options=PROVIDER_OPTIONS, provider_options=PROVIDER_OPTIONS,
ai_provider=ai_provider, ai_provider=ai_provider,
ai_model=ai_model, ai_model=ai_model,
+2 -2
View File
@@ -27,10 +27,10 @@ JOB_STATUS_RETRY = "retry"
JOB_STATUS_DONE = "done" JOB_STATUS_DONE = "done"
JOB_STATUS_DEAD = "dead" 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_PARSER = "vk-parser"
WORKER_STORAGE_UPLOADER = "vk-storage-uploader" WORKER_MEDIA_UPLOADER = "media-uploader"
WORKER_AI_QUALIFIER = "ai-qualifier" WORKER_AI_QUALIFIER = "ai-qualifier"
WORKER_AI_WRITER = "ai-writer" WORKER_AI_WRITER = "ai-writer"
WORKER_TG_POSTER = "tg-poster" WORKER_TG_POSTER = "tg-poster"
+99 -5
View File
@@ -1,5 +1,6 @@
from __future__ import annotations from __future__ import annotations
import json
from dataclasses import dataclass, field from dataclasses import dataclass, field
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any from typing import Any
@@ -9,6 +10,18 @@ import aiohttp
from .constants import PLATFORM_SITE, PLATFORM_VK 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 @dataclass
class SourceMedia: class SourceMedia:
url: str url: str
@@ -24,6 +37,10 @@ class SourceItem:
media: list[SourceMedia] = field(default_factory=list) media: list[SourceMedia] = field(default_factory=list)
raw: dict[str, Any] = field(default_factory=dict) 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: def validate_source_config(platform: str, config: dict[str, Any]) -> None:
if platform == PLATFORM_VK: if platform == PLATFORM_VK:
@@ -32,6 +49,50 @@ def validate_source_config(platform: str, config: dict[str, Any]) -> None:
raise ValueError("Неподдерживаемая площадка") raise ValueError("Неподдерживаемая площадка")
if not config: if not config:
raise ValueError("Для сайта нужен конфиг JSON") 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": if config.get("format") != "rss":
raise ValueError('Сейчас поддерживается только "format": "rss"') raise ValueError('Сейчас поддерживается только "format": "rss"')
if str(config.get("access") or "auto") not in {"auto", "http", "cloudflare"}: if str(config.get("access") or "auto") not in {"auto", "http", "cloudflare"}:
@@ -42,6 +103,17 @@ def validate_source_config(platform: str, config: dict[str, Any]) -> None:
raise ValueError("max_items должен быть целым числом") from exc raise ValueError("max_items должен быть целым числом") from exc
if not 1 <= max_items <= 100: if not 1 <= max_items <= 100:
raise ValueError("max_items должен быть от 1 до 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: def _posted_at(value: Any) -> datetime:
@@ -69,17 +141,18 @@ class SiteParserClient:
self.rucaptcha_token = rucaptcha_token self.rucaptcha_token = rucaptcha_token
self.timeout = aiohttp.ClientTimeout(total=max(10, timeout_sec)) self.timeout = aiohttp.ClientTimeout(total=max(10, timeout_sec))
async def fetch(self, source: dict) -> tuple[list[SourceItem], dict[str, Any] | None]: async def fetch(self, source: dict) -> tuple[list[SourceItem], dict[str, Any] | None, str | None]:
if not self.base_url or not self.token: if not self.base_url or not self.token:
raise RuntimeError("Site Parser URL или токен не настроены") raise RuntimeError("Site Parser URL или токен не настроены")
config = dict(source.get("settings_json") or {}) config = json_object(source.get("settings_json"))
validate_source_config(PLATFORM_SITE, config) validate_source_config(PLATFORM_SITE, config)
runtime_state = dict(source.get("runtime_state_json") or {}) runtime_state = json_object(source.get("runtime_state_json"))
payload = { payload = {
"url": source["url"], "url": source["url"],
"config": config, "config": config,
"rucaptcha_token": self.rucaptcha_token or None, "rucaptcha_token": self.rucaptcha_token or None,
"browser_state": runtime_state.get("browser_state"), "browser_state": runtime_state.get("browser_state"),
"browser_user_agent": runtime_state.get("browser_user_agent"),
} }
try: try:
async with self.session.post( async with self.session.post(
@@ -96,11 +169,17 @@ class SiteParserClient:
except aiohttp.ClientError as exc: except aiohttp.ClientError as exc:
raise RuntimeError(f"Site Parser недоступен: {exc}") from exc raise RuntimeError(f"Site Parser недоступен: {exc}") from exc
result_status = str(data.get("status") or "ok")
result_errors = data.get("errors") or []
items = [] items = []
for raw in data.get("items") or []: for raw in data.get("items") or []:
title = str(raw.get("title") or "").strip() title = str(raw.get("title") or "").strip()
body = str(raw.get("text") or "").strip() body = str(raw.get("text") or "").strip()
text = "\n\n".join(part for part in (title, body) if part) 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() url = str(raw.get("url") or source["url"]).strip()
external_id = str(raw.get("external_id") or url).strip() external_id = str(raw.get("external_id") or url).strip()
if not external_id: if not external_id:
@@ -112,4 +191,19 @@ class SiteParserClient:
] ]
items.append(SourceItem(external_id, url, text, _posted_at(raw.get("published_at")), media, raw)) items.append(SourceItem(external_id, url, text, _posted_at(raw.get("published_at")), media, raw))
state = data.get("browser_state") state = data.get("browser_state")
return items, ({"browser_state": state} if isinstance(state, dict) else None) 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
+29 -6
View File
@@ -43,15 +43,32 @@
<div id="site-config" class="form-control md:col-span-2"> <div id="site-config" class="form-control md:col-span-2">
<label class="label"><span class="label-text font-bold">Конфигурация JSON</span></label> <label class="label"><span class="label-text font-bold">Конфигурация JSON</span></label>
<textarea name="settings_json" rows="8" class="textarea textarea-bordered w-full font-mono" placeholder='{"format":"rss","access":"auto","max_items":20}'>{{ source.settings_json | tojson(indent=2) if source else '{}' }}</textarea> <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"> <details class="mt-2 text-sm text-base-content/70">
<summary class="cursor-pointer">Пример и параметры</summary> <summary class="cursor-pointer">Пример и параметры</summary>
<pre class="mt-2 p-3 bg-base-200 overflow-x-auto">{ <pre class="mt-2 p-3 bg-base-200 overflow-x-auto">{
"format": "rss", "version": 1,
"access": "auto", "discovery": {
"max_items": 20 "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> }</pre>
<p class="mt-2"><code>access</code>: <code>auto</code> сначала пробует обычный запрос и при Cloudflare использует RuCaptcha; <code>http</code> запрещает браузер; <code>cloudflare</code> сразу запускает браузер.</p> <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> </details>
</div> </div>
</div> </div>
@@ -72,7 +89,13 @@
<script> <script>
const platform = document.querySelector('[name="platform"]'); const platform = document.querySelector('[name="platform"]');
const siteConfig = document.getElementById('site-config'); const siteConfig = document.getElementById('site-config');
const syncConfig = () => siteConfig.hidden = platform.value !== 'site'; 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); platform.addEventListener('change', syncConfig);
syncConfig(); syncConfig();
</script> </script>
@@ -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="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"> <td class="px-3 sm:px-4 py-3 text-center">
<div class="flex justify-center items-center gap-1"> <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="Править"> <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> <i data-lucide="edit-2" class="w-4 h-4"></i>
</a> </a>
+5 -2
View File
@@ -1,11 +1,14 @@
{% extends "base.html" %} {% extends "base.html" %}
{% block body %} {% block body %}
<div class="mb-8"> <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"> <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> <i data-lucide="database" class="text-app-primary w-8 h-8"></i>
Источники Источники
</h1> </h1>
<div class="text-app-textMuted text-sm">Источники для парсинга. Название идёт в prompt, тэг — в будущие хэштеги.</div> <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> </div>
<details class="card mb-8 group/details" {% if source_preview %}open{% endif %}> <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"> <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> <i data-lucide="plus" class="w-5 h-5"></i>
</div> </div>
<span class="text-lg font-bold text-white">Добавить источники</span> <span class="text-lg font-bold text-white">Добавить VK списком</span>
</div> </div>
<i data-lucide="chevron-down" class="w-5 h-5 text-app-textMuted transition-transform group-open/details:rotate-180"></i> <i data-lucide="chevron-down" class="w-5 h-5 text-app-textMuted transition-transform group-open/details:rotate-180"></i>
</summary> </summary>
+18 -2
View File
@@ -8,6 +8,18 @@
<div class="text-app-textMuted text-sm">Управление фоновыми процессами, категориями, расписанием и настройками AI.</div> <div class="text-app-textMuted text-sm">Управление фоновыми процессами, категориями, расписанием и настройками AI.</div>
</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 --> <!-- Workers Panel -->
<div class="card overflow-hidden mb-8"> <div class="card overflow-hidden mb-8">
<div class="overflow-x-auto"> <div class="overflow-x-auto">
@@ -25,7 +37,7 @@
<tbody class="divide-y divide-app-border text-xs"> <tbody class="divide-y divide-app-border text-xs">
{% for w in workers %} {% for w in workers %}
<tr class="hover:bg-app-surfaceHover transition-colors"> <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"> <td class="px-4 py-3">
{% if w.effective_enabled %} {% if w.effective_enabled %}
<span class="badge badge-success px-2 py-0.5 flex items-center gap-1.5 w-max"> <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"> <div class="flex flex-col gap-6">
{% for category, rows in settings|groupby("category") %} {% 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"> <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> <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> <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 scrollKey = "workers-scroll-y";
const detailsKey = "workers-open-details"; const detailsKey = "workers-open-details";
const detailItems = Array.from(document.querySelectorAll("details[data-workers-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); const savedDetails = sessionStorage.getItem(detailsKey);
if (savedDetails) { if (savedDetails) {
try { try {
@@ -7,6 +7,7 @@ import re
import sys import sys
import time import time
from pathlib import Path from pathlib import Path
from urllib.parse import urlparse
import aiohttp import aiohttp
from aiogram import Bot from aiogram import Bot
@@ -18,17 +19,18 @@ from loguru import logger
from ..config import settings from ..config import settings
from ..constants import ( from ..constants import (
JOB_TYPE_VK_STORAGE_COPY, JOB_TYPE_MEDIA_STORAGE_COPY,
MEDIA_STATUS_FAILED, MEDIA_STATUS_FAILED,
MEDIA_STATUS_LINK_ONLY, MEDIA_STATUS_LINK_ONLY,
MEDIA_STATUS_UPLOADED, MEDIA_STATUS_UPLOADED,
POST_STATUS_FAILED, POST_STATUS_FAILED,
POST_STATUS_STORAGE_READY, POST_STATUS_STORAGE_READY,
WORKER_STORAGE_UPLOADER, WORKER_MEDIA_UPLOADER,
) )
from ..db import fetch_float_setting, fetch_int_setting, fetch_setting, get_pool from ..db import fetch_float_setting, fetch_int_setting, fetch_setting, get_pool
from ..heartbeat import HeartbeatReporter from ..heartbeat import HeartbeatReporter
from ..jobs import ack_done, ack_retry, claim_job, is_worker_enabled, recover_stale_jobs 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_DIR = Path("/tmp")
TMP_PREFIX = "vkparser_tg_media_" 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)) 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]: def split_message_chunks(text: str, limit: int) -> list[str]:
text = str(text or "").strip() text = str(text or "").strip()
if not text: 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"])}' return f'<a href="{url}">#{int(post["id"])}</a>' if url else f'#{int(post["id"])}'
class TelegramStorageUploader: class MediaUploader:
def __init__(self) -> None: def __init__(self) -> None:
self.pool = None self.pool = None
self.bot: Bot | None = None self.bot: Bot | None = None
self.storage_chat_id: int | None = None self.storage_chat_id: int | None = None
self.storage_thread_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_group_max_items = MAX_MEDIA_GROUP
self.media_upload_delay_sec = 1.0 self.media_upload_delay_sec = 1.0
self.tg_retry_max_attempts = 4 self.tg_retry_max_attempts = 4
@@ -100,7 +116,7 @@ class TelegramStorageUploader:
async def init(self) -> None: async def init(self) -> None:
self.pool = await get_pool() 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: if recovered:
logger.warning("Recovered stale storage jobs: {}", 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]: async def load_media(self, raw_post_id: int) -> list[dict]:
rows = await self.pool.fetch( rows = await self.pool.fetch(
""" """
SELECT * SELECT m.*, s.url AS source_url, s.runtime_state_json AS source_runtime_state
FROM raw_post_media FROM raw_post_media m
WHERE raw_post_id=$1 JOIN raw_posts rp ON rp.id=m.raw_post_id
ORDER BY sort_order ASC, id ASC 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, raw_post_id,
) )
@@ -262,15 +280,70 @@ class TelegramStorageUploader:
return f"media not uploaded: {details}" return f"media not uploaded: {details}"
return None 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: 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: if response.status != 200:
return None return None
return await response.read() return await response.read()
except Exception: except Exception:
return None 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: async def mark_media_uploaded(self, media_id: int, file_id: str, unique_id: str | None) -> None:
await self.pool.execute( await self.pool.execute(
""" """
@@ -320,37 +393,43 @@ class TelegramStorageUploader:
error[:1000], error[:1000],
) )
async def download_video(self, vk_url: str, output_path: str) -> dict | None: async def download_video(self, video_url: str, output_path: str) -> dict | None:
parsed = parse_vk_video_url(vk_url) provider = video_provider(video_url)
if not parsed: if not provider:
return None return None
owner_id, video_id = parsed
max_size = self.video_max_size_mb * 1024 * 1024 max_size = self.video_max_size_mb * 1024 * 1024
netrc_path = f"{output_path}.netrc" target_url = video_url
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 = [ cmd = [
sys.executable, sys.executable,
"-m", "-m",
"yt_dlp", "yt_dlp",
"--netrc-location", ]
netrc_path, netrc_path = None
f"https://vk.com/video{owner_id}_{video_id}", if provider == "vk":
"-o", owner_id, video_id = parse_vk_video_url(video_url) or (0, 0)
output_path, 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", "--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}][filesize<{max_size}]"
f"/best[height<={self.video_max_height}]" 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}][filesize<{max_size}]+bestaudio/best"
f"/bestvideo[height<={self.video_max_height}]+bestaudio/best" f"/bestvideo[height<={self.video_max_height}]+bestaudio/best"
f"/best[filesize<{max_size}]" f"/best[filesize<{max_size}]"
), ),
"--quiet", "--quiet", "--no-warnings",
"--no-warnings", ])
] proc = None
stderr = b""
try: try:
proc = await asyncio.create_subprocess_exec( proc = await asyncio.create_subprocess_exec(
*cmd, *cmd,
@@ -359,20 +438,24 @@ class TelegramStorageUploader:
) )
_, stderr = await asyncio.wait_for(proc.communicate(), timeout=self.yt_dlp_timeout_sec) _, stderr = await asyncio.wait_for(proc.communicate(), timeout=self.yt_dlp_timeout_sec)
except asyncio.TimeoutError: except asyncio.TimeoutError:
if proc:
proc.kill() proc.kill()
await proc.communicate() await proc.communicate()
remove_download_files(output_path)
return {"error": "video download timeout", "permanent": False} return {"error": "video download timeout", "permanent": False}
finally: finally:
try: if netrc_path:
os.unlink(netrc_path) Path(netrc_path).unlink(missing_ok=True)
except FileNotFoundError:
pass
if proc.returncode != 0 or not os.path.exists(output_path): if proc.returncode != 0 or not os.path.exists(output_path):
err = (stderr or b"").decode("utf-8", errors="ignore").lower() 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} return {"error": "video unavailable or download failed", "permanent": permanent}
size = os.path.getsize(output_path) size = os.path.getsize(output_path)
if size > max_size: if size > max_size:
remove_download_files(output_path)
return {"error": "video too large", "permanent": True} return {"error": "video too large", "permanent": True}
return {"path": output_path, "size_bytes": size} return {"path": output_path, "size_bytes": size}
@@ -383,6 +466,7 @@ class TelegramStorageUploader:
media_id = int(item["id"]) media_id = int(item["id"])
media_type = str(item["media_type"]) media_type = str(item["media_type"])
url = str(item.get("original_url") or "") url = str(item.get("original_url") or "")
request_options = self.site_request_options(item)
if item.get("tg_file_id"): if item.get("tg_file_id"):
prepared.append({"media_id": media_id, "media_type": media_type, "media": item["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 continue
if media_type == "photo": if media_type == "photo":
data = await self.download_bytes(session, url) data = await self.download_bytes(session, url, request_options)
if not data: if not data:
await self.mark_media_failed_attempt(media_id, "photo download failed") await self.mark_media_failed_attempt(media_id, "photo download failed")
continue continue
@@ -410,7 +494,11 @@ class TelegramStorageUploader:
await self.mark_media_link_only(media_id, "video too long") await self.mark_media_link_only(media_id, "video too long")
continue continue
temp_path = str(TMP_DIR / f"{TMP_PREFIX}{raw_post_id}_{media_id}.mp4") 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: if not info:
await self.mark_media_failed_attempt(media_id, "video download failed") await self.mark_media_failed_attempt(media_id, "video download failed")
continue continue
@@ -581,8 +669,7 @@ class TelegramStorageUploader:
tmp_path = item.get("tmp_path") tmp_path = item.get("tmp_path")
if tmp_path: if tmp_path:
try: try:
if os.path.exists(tmp_path): remove_download_files(tmp_path)
os.remove(tmp_path)
except Exception: except Exception:
pass pass
if not blocking_error.startswith("media still pending:"): if not blocking_error.startswith("media still pending:"):
@@ -600,8 +687,7 @@ class TelegramStorageUploader:
tmp_path = item.get("tmp_path") tmp_path = item.get("tmp_path")
if tmp_path: if tmp_path:
try: try:
if os.path.exists(tmp_path): remove_download_files(tmp_path)
os.remove(tmp_path)
except Exception: except Exception:
pass pass
await self.mark_post_ready(raw_post_id, message_ids, meta_message_id) 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) logger.info("Telegram storage done: raw_post={} messages={}", raw_post_id, message_ids)
async def run_once(self, worker_id: str) -> bool: 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: if not enabled:
await self.heartbeat.beat(self.pool, status="disabled", force=True) await self.heartbeat.beat(self.pool, status="disabled", force=True)
return False 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: if not job:
await self.heartbeat.beat(self.pool, status="idle") await self.heartbeat.beat(self.pool, status="idle")
return False return False
@@ -637,7 +723,7 @@ class TelegramStorageUploader:
async def run_loop(self) -> None: async def run_loop(self) -> None:
await self.init() await self.init()
worker_id = f"{WORKER_STORAGE_UPLOADER}:{os.getpid()}" worker_id = f"{WORKER_MEDIA_UPLOADER}:{os.getpid()}"
logger.info("{} started", worker_id) logger.info("{} started", worker_id)
try: try:
while True: while True:
@@ -654,7 +740,7 @@ class TelegramStorageUploader:
async def main() -> None: async def main() -> None:
logger.remove() logger.remove()
logger.add(sys.stdout, level=settings.log_level) logger.add(sys.stdout, level=settings.log_level)
worker = TelegramStorageUploader() worker = MediaUploader()
await worker.run_loop() await worker.run_loop()
+44 -13
View File
@@ -10,7 +10,7 @@ from loguru import logger
from ..config import settings from ..config import settings
from ..constants import ( from ..constants import (
JOB_TYPE_VK_STORAGE_COPY, JOB_TYPE_MEDIA_STORAGE_COPY,
PLATFORM_SITE, PLATFORM_SITE,
PLATFORM_VK, PLATFORM_VK,
POST_STATUS_SKIPPED, POST_STATUS_SKIPPED,
@@ -22,7 +22,7 @@ from ..constants import (
from ..db import fetch_bool_setting, fetch_float_setting, fetch_int_setting, fetch_setting, get_pool from ..db import fetch_bool_setting, fetch_float_setting, fetch_int_setting, fetch_setting, get_pool
from ..heartbeat import HeartbeatReporter from ..heartbeat import HeartbeatReporter
from ..jobs import is_worker_enabled from ..jobs import is_worker_enabled
from ..source_adapters import SiteParserClient, SourceItem from ..source_adapters import SiteParserClient, SourceItem, json_object
from ..vk_api import ( from ..vk_api import (
VKAPIClient, VKAPIClient,
VKAPIError, VKAPIError,
@@ -55,7 +55,7 @@ class VKParserWorker:
async def init(self) -> None: async def init(self) -> None:
self.pool = await get_pool() 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( rows = await self.pool.fetch(
""" """
SELECT * SELECT *
@@ -63,9 +63,21 @@ class VKParserWorker:
WHERE platform=ANY($1::text[]) WHERE platform=ANY($1::text[])
AND active=TRUE AND active=TRUE
AND archived_at IS NULL 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 ORDER BY last_checked_at NULLS FIRST, priority ASC, id ASC
""", """,
[PLATFORM_VK, PLATFORM_SITE], [PLATFORM_VK, PLATFORM_SITE],
PLATFORM_VK,
vk_interval_sec,
PLATFORM_SITE,
site_interval_minutes,
) )
return [dict(r) for r in rows] return [dict(r) for r in rows]
@@ -211,7 +223,7 @@ class VKParserWorker:
VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending') VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending')
ON CONFLICT DO NOTHING ON CONFLICT DO NOTHING
""", """,
JOB_TYPE_VK_STORAGE_COPY, JOB_TYPE_MEDIA_STORAGE_COPY,
raw_post_id, raw_post_id,
) )
return int(raw_post_id) return int(raw_post_id)
@@ -230,11 +242,13 @@ class VKParserWorker:
if since_dt.tzinfo is None: if since_dt.tzinfo is None:
since_dt = since_dt.replace(tzinfo=timezone.utc) since_dt = since_dt.replace(tzinfo=timezone.utc)
fetched, runtime_state = await client.fetch(source) fetched, runtime_state, partial_warning = await client.fetch(source)
recent = [item for item in fetched if item.posted_at > since_dt] 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]) 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] candidates = [item for item in recent if item.external_id not in known]
min_text_length = max(0, await fetch_int_setting("parser_min_text_length", 0)) 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_empty_text = await fetch_bool_setting("parser_skip_empty_text", True)
skip_no_media = await fetch_bool_setting("parser_skip_no_media", 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) skip_short_text = await fetch_bool_setting("parser_skip_text_too_short", True)
@@ -248,11 +262,11 @@ class VKParserWorker:
if dedupe_content_hash and content_hash in known_hashes: if dedupe_content_hash and content_hash in known_hashes:
continue continue
skip_reason = None skip_reason = None
if skip_empty_text and not item.text: if skip_empty_text and not item.body_text:
skip_reason = "empty_text" skip_reason = "empty_text"
elif skip_no_media and not item.media: elif skip_no_media and not item.media:
skip_reason = "no_media" skip_reason = "no_media"
elif skip_short_text and len(item.text) < min_text_length: elif skip_short_text and len(item.body_text) < min_text_length:
skip_reason = "text_too_short" skip_reason = "text_too_short"
if skip_reason and not store_skipped: if skip_reason and not store_skipped:
continue continue
@@ -267,6 +281,22 @@ class VKParserWorker:
saved += 1 saved += 1
known_hashes.add(content_hash) known_hashes.add(content_hash)
max_seen = max((item.posted_at for item in fetched), default=last_parsed_at) 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) await self.mark_source_ok(source_id, max_seen, runtime_state)
logger.info( logger.info(
"Parsed site source {}: fetched={} recent={} known={} saved={}", "Parsed site source {}: fetched={} recent={} known={} saved={}",
@@ -377,7 +407,7 @@ class VKParserWorker:
VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending') VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending')
ON CONFLICT DO NOTHING ON CONFLICT DO NOTHING
""", """,
JOB_TYPE_VK_STORAGE_COPY, JOB_TYPE_MEDIA_STORAGE_COPY,
raw_post_id, raw_post_id,
) )
return int(raw_post_id) return int(raw_post_id)
@@ -512,6 +542,8 @@ class VKParserWorker:
await self.heartbeat.beat(self.pool, status="disabled", force=True) await self.heartbeat.beat(self.pool, status="disabled", force=True)
return 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) rps = await fetch_int_setting("vk_requests_per_second", 3)
timeout_total = await fetch_int_setting("vk_api_timeout_total_sec", 15) timeout_total = await fetch_int_setting("vk_api_timeout_total_sec", 15)
timeout_connect = await fetch_int_setting("vk_api_timeout_connect_sec", 5) timeout_connect = await fetch_int_setting("vk_api_timeout_connect_sec", 5)
@@ -520,10 +552,10 @@ class VKParserWorker:
retry_min_delay = await fetch_float_setting("vk_api_retry_min_delay_sec", 2.0) 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) 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)) 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)}) await self.heartbeat.beat(self.pool, meta={"sources": len(sources)})
if not sources: if not sources:
logger.info("No active sources") logger.debug("No sources due for parsing")
return return
logger.info( logger.info(
@@ -592,8 +624,7 @@ class VKParserWorker:
await self.run_once() await self.run_once()
except Exception as e: except Exception as e:
logger.exception("Parser loop error: {}", e) logger.exception("Parser loop error: {}", e)
interval = max(10, await fetch_int_setting("parser_interval_sec", 300)) await asyncio.sleep(10)
await asyncio.sleep(interval)
async def main() -> None: async def main() -> None:
+79 -3
View File
@@ -25,31 +25,107 @@ class FakeResponse:
"media": [{"type": "photo", "url": "https://example.test/1.jpg"}], "media": [{"type": "photo", "url": "https://example.test/1.jpg"}],
}], }],
"browser_state": {"cookies": [{"name": "cf_clearance"}]}, "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: class FakeSession:
def post(self, *_args, **_kwargs): def post(self, *_args, **_kwargs):
return FakeResponse() 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): class SourceAdapterTests(unittest.IsolatedAsyncioTestCase):
async def test_worker_response_is_normalized(self) -> None: async def test_worker_response_is_normalized(self) -> None:
client = SiteParserClient(FakeSession(), "http://worker", "token", "captcha", 30) client = SiteParserClient(FakeSession(), "http://worker", "token", "captcha", 30)
items, state = await client.fetch({ items, state, warning = await client.fetch({
"url": "https://example.test/rss.xml", "url": "https://example.test/rss.xml",
"settings_json": {"format": "rss", "access": "auto"}, "settings_json": '{"format":"rss","access":"auto"}',
"runtime_state_json": {}, "runtime_state_json": '{}',
}) })
self.assertEqual(items[0].text, "Title\n\nBody") self.assertEqual(items[0].text, "Title\n\nBody")
self.assertEqual(items[0].media[0].url, "https://example.test/1.jpg") 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_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: def test_site_config_is_required(self) -> None:
with self.assertRaisesRegex(ValueError, "нужен конфиг"): with self.assertRaisesRegex(ValueError, "нужен конфиг"):
validate_source_config("site", {}) 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__": if __name__ == "__main__":
unittest.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())