feat: filter site bodies before media upload
This commit is contained in:
@@ -37,7 +37,7 @@ from .workers.site_poster import SitePoster
|
||||
from .workers.tg_poster import TelegramPoster
|
||||
from .workers.tg_reactor import TelegramReactor
|
||||
from .workers.vk_poster import VKPoster
|
||||
from .workers.vk_storage_uploader import TelegramStorageUploader
|
||||
from .workers.media_uploader import MediaUploader
|
||||
|
||||
COOKIE_NAME = "vk_parser_admin"
|
||||
VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier"
|
||||
@@ -372,6 +372,7 @@ SETTING_ORDER = {
|
||||
"site_parser_rucaptcha_token",
|
||||
"site_parser_timeout_sec",
|
||||
"site_parser_interval_minutes",
|
||||
"site_parser_min_text_length",
|
||||
],
|
||||
"VK": [
|
||||
"vk_requests_per_second",
|
||||
@@ -1820,7 +1821,7 @@ async def startup() -> None:
|
||||
("tg-poster", TelegramPoster()),
|
||||
("tg-reactor", TelegramReactor()),
|
||||
("vk-poster", VKPoster()),
|
||||
("vk-storage-uploader", TelegramStorageUploader()),
|
||||
("media-uploader", MediaUploader()),
|
||||
]
|
||||
for name, worker in workers:
|
||||
asyncio.create_task(start_worker_task(worker, name))
|
||||
|
||||
@@ -27,10 +27,10 @@ JOB_STATUS_RETRY = "retry"
|
||||
JOB_STATUS_DONE = "done"
|
||||
JOB_STATUS_DEAD = "dead"
|
||||
|
||||
JOB_TYPE_VK_STORAGE_COPY = "vk.storage.copy"
|
||||
JOB_TYPE_MEDIA_STORAGE_COPY = "media.storage.copy"
|
||||
|
||||
WORKER_PARSER = "vk-parser"
|
||||
WORKER_STORAGE_UPLOADER = "vk-storage-uploader"
|
||||
WORKER_MEDIA_UPLOADER = "media-uploader"
|
||||
WORKER_AI_QUALIFIER = "ai-qualifier"
|
||||
WORKER_AI_WRITER = "ai-writer"
|
||||
WORKER_TG_POSTER = "tg-poster"
|
||||
|
||||
@@ -37,6 +37,10 @@ class SourceItem:
|
||||
media: list[SourceMedia] = field(default_factory=list)
|
||||
raw: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
@property
|
||||
def body_text(self) -> str:
|
||||
return str(self.raw.get("text") or "").strip()
|
||||
|
||||
|
||||
def validate_source_config(platform: str, config: dict[str, Any]) -> None:
|
||||
if platform == PLATFORM_VK:
|
||||
@@ -60,6 +64,12 @@ def validate_source_config(platform: str, config: dict[str, Any]) -> None:
|
||||
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:
|
||||
|
||||
@@ -52,7 +52,8 @@
|
||||
"max_items": 20
|
||||
}</pre>
|
||||
<p class="mt-2"><code>access</code>: <code>auto</code> сначала пробует обычный запрос и при Cloudflare использует RuCaptcha; <code>http</code> запрещает браузер; <code>cloudflare</code> сразу запускает браузер.</p>
|
||||
<p class="mt-2">Если RSS не содержит текст или медиа, добавьте <code>"follow_links": true</code> и CSS-селектор содержимого страницы, например <code>"content_selector": "article.full"</code>.</p>
|
||||
<p class="mt-2">Если RSS не содержит текст или медиа, добавьте <code>"follow_links": true</code>, селектор материала <code>"content_selector": "article.full"</code> и при необходимости отдельный селектор текста <code>"text_selector": ".field--name-body"</code>.</p>
|
||||
<p class="mt-2"><code>min_text_length</code> переопределяет общий порог Site Parser только для этого источника.</p>
|
||||
</details>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
+9
-9
@@ -19,13 +19,13 @@ from loguru import logger
|
||||
|
||||
from ..config import settings
|
||||
from ..constants import (
|
||||
JOB_TYPE_VK_STORAGE_COPY,
|
||||
JOB_TYPE_MEDIA_STORAGE_COPY,
|
||||
MEDIA_STATUS_FAILED,
|
||||
MEDIA_STATUS_LINK_ONLY,
|
||||
MEDIA_STATUS_UPLOADED,
|
||||
POST_STATUS_FAILED,
|
||||
POST_STATUS_STORAGE_READY,
|
||||
WORKER_STORAGE_UPLOADER,
|
||||
WORKER_MEDIA_UPLOADER,
|
||||
)
|
||||
from ..db import fetch_float_setting, fetch_int_setting, fetch_setting, get_pool
|
||||
from ..heartbeat import HeartbeatReporter
|
||||
@@ -90,13 +90,13 @@ def original_link(post: dict) -> str:
|
||||
return f'<a href="{url}">#{int(post["id"])}</a>' if url else f'#{int(post["id"])}'
|
||||
|
||||
|
||||
class TelegramStorageUploader:
|
||||
class MediaUploader:
|
||||
def __init__(self) -> None:
|
||||
self.pool = None
|
||||
self.bot: Bot | None = None
|
||||
self.storage_chat_id: int | None = None
|
||||
self.storage_thread_id: int | None = None
|
||||
self.heartbeat = HeartbeatReporter(WORKER_STORAGE_UPLOADER, 30)
|
||||
self.heartbeat = HeartbeatReporter(WORKER_MEDIA_UPLOADER, 30)
|
||||
self.media_group_max_items = MAX_MEDIA_GROUP
|
||||
self.media_upload_delay_sec = 1.0
|
||||
self.tg_retry_max_attempts = 4
|
||||
@@ -115,7 +115,7 @@ class TelegramStorageUploader:
|
||||
|
||||
async def init(self) -> None:
|
||||
self.pool = await get_pool()
|
||||
recovered = await recover_stale_jobs(self.pool, JOB_TYPE_VK_STORAGE_COPY, stale_minutes=20)
|
||||
recovered = await recover_stale_jobs(self.pool, JOB_TYPE_MEDIA_STORAGE_COPY, stale_minutes=20)
|
||||
if recovered:
|
||||
logger.warning("Recovered stale storage jobs: {}", recovered)
|
||||
|
||||
@@ -660,12 +660,12 @@ class TelegramStorageUploader:
|
||||
logger.info("Telegram storage done: raw_post={} messages={}", raw_post_id, message_ids)
|
||||
|
||||
async def run_once(self, worker_id: str) -> bool:
|
||||
enabled = await is_worker_enabled(self.pool, WORKER_STORAGE_UPLOADER)
|
||||
enabled = await is_worker_enabled(self.pool, WORKER_MEDIA_UPLOADER)
|
||||
if not enabled:
|
||||
await self.heartbeat.beat(self.pool, status="disabled", force=True)
|
||||
return False
|
||||
|
||||
job = await claim_job(self.pool, JOB_TYPE_VK_STORAGE_COPY, worker_id)
|
||||
job = await claim_job(self.pool, JOB_TYPE_MEDIA_STORAGE_COPY, worker_id)
|
||||
if not job:
|
||||
await self.heartbeat.beat(self.pool, status="idle")
|
||||
return False
|
||||
@@ -688,7 +688,7 @@ class TelegramStorageUploader:
|
||||
|
||||
async def run_loop(self) -> None:
|
||||
await self.init()
|
||||
worker_id = f"{WORKER_STORAGE_UPLOADER}:{os.getpid()}"
|
||||
worker_id = f"{WORKER_MEDIA_UPLOADER}:{os.getpid()}"
|
||||
logger.info("{} started", worker_id)
|
||||
try:
|
||||
while True:
|
||||
@@ -705,7 +705,7 @@ class TelegramStorageUploader:
|
||||
async def main() -> None:
|
||||
logger.remove()
|
||||
logger.add(sys.stdout, level=settings.log_level)
|
||||
worker = TelegramStorageUploader()
|
||||
worker = MediaUploader()
|
||||
await worker.run_loop()
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ from loguru import logger
|
||||
|
||||
from ..config import settings
|
||||
from ..constants import (
|
||||
JOB_TYPE_VK_STORAGE_COPY,
|
||||
JOB_TYPE_MEDIA_STORAGE_COPY,
|
||||
PLATFORM_SITE,
|
||||
PLATFORM_VK,
|
||||
POST_STATUS_SKIPPED,
|
||||
@@ -22,7 +22,7 @@ from ..constants import (
|
||||
from ..db import fetch_bool_setting, fetch_float_setting, fetch_int_setting, fetch_setting, get_pool
|
||||
from ..heartbeat import HeartbeatReporter
|
||||
from ..jobs import is_worker_enabled
|
||||
from ..source_adapters import SiteParserClient, SourceItem
|
||||
from ..source_adapters import SiteParserClient, SourceItem, json_object
|
||||
from ..vk_api import (
|
||||
VKAPIClient,
|
||||
VKAPIError,
|
||||
@@ -223,7 +223,7 @@ class VKParserWorker:
|
||||
VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending')
|
||||
ON CONFLICT DO NOTHING
|
||||
""",
|
||||
JOB_TYPE_VK_STORAGE_COPY,
|
||||
JOB_TYPE_MEDIA_STORAGE_COPY,
|
||||
raw_post_id,
|
||||
)
|
||||
return int(raw_post_id)
|
||||
@@ -246,7 +246,9 @@ class VKParserWorker:
|
||||
recent = [item for item in fetched if item.posted_at > since_dt]
|
||||
known = await self.known_post_ids(source_id, [item.external_id for item in recent])
|
||||
candidates = [item for item in recent if item.external_id not in known]
|
||||
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_no_media = await fetch_bool_setting("parser_skip_no_media", True)
|
||||
skip_short_text = await fetch_bool_setting("parser_skip_text_too_short", True)
|
||||
@@ -260,11 +262,11 @@ class VKParserWorker:
|
||||
if dedupe_content_hash and content_hash in known_hashes:
|
||||
continue
|
||||
skip_reason = None
|
||||
if skip_empty_text and not item.text:
|
||||
if skip_empty_text and not item.body_text:
|
||||
skip_reason = "empty_text"
|
||||
elif skip_no_media and not item.media:
|
||||
skip_reason = "no_media"
|
||||
elif skip_short_text and len(item.text) < min_text_length:
|
||||
elif skip_short_text and len(item.body_text) < min_text_length:
|
||||
skip_reason = "text_too_short"
|
||||
if skip_reason and not store_skipped:
|
||||
continue
|
||||
@@ -389,7 +391,7 @@ class VKParserWorker:
|
||||
VALUES($1, 'raw_post', $2, '{}'::jsonb, 'pending')
|
||||
ON CONFLICT DO NOTHING
|
||||
""",
|
||||
JOB_TYPE_VK_STORAGE_COPY,
|
||||
JOB_TYPE_MEDIA_STORAGE_COPY,
|
||||
raw_post_id,
|
||||
)
|
||||
return int(raw_post_id)
|
||||
|
||||
Reference in New Issue
Block a user