diff --git a/README.md b/README.md index 6559c96..3d8e533 100644 --- a/README.md +++ b/README.md @@ -142,9 +142,11 @@ https://vk.com/wall-239548476_123 RuCaptcha и таймаут. Внешний модуль запускается командой: ```bash -docker build -t site-parser-worker site_parser_worker -docker run -d --restart unless-stopped -p 8080:8080 \ - -e WORKER_TOKEN=replace-me site-parser-worker +cd site_parser_worker +printf 'WORKER_TOKEN=replace-me\n' > .env +chmod 600 .env +docker-compose up -d --build +docker-compose ps ``` ## Запуск локально @@ -161,7 +163,7 @@ PYTHONPATH=src .venv/bin/uvicorn vk_parser_app.admin:app --host 0.0.0.0 --port 8 PYTHONPATH=src .venv/bin/python -m vk_parser_app.workers.parser ``` -VK storage uploader: +Media uploader: ```bash PYTHONPATH=src .venv/bin/python -m vk_parser_app.workers.media_uploader diff --git a/site_parser_worker/app.py b/site_parser_worker/app.py index 902659a..7ccf6b0 100644 --- a/site_parser_worker/app.py +++ b/site_parser_worker/app.py @@ -27,6 +27,7 @@ class ParseRequest(BaseModel): config: dict[str, Any] rucaptcha_token: SecretStr | None = None browser_state: dict[str, Any] | None = None + browser_user_agent: str | None = None class ParsedItem(BaseModel): @@ -133,12 +134,16 @@ async def fetch_in_browser( url: str, rucaptcha_token: str, browser_state: dict[str, Any] | None, -) -> tuple[str, dict[str, Any]]: + browser_user_agent: str | None, +) -> tuple[str, dict[str, Any], str]: # ponytail: one browser at a time; use a queue only when parallel source parsing is needed. async with browser_lock: async with async_playwright() as playwright: browser = await playwright.chromium.launch(channel="chrome", headless=False) - context = await browser.new_context(storage_state=browser_state) if browser_state else await browser.new_context() + context = await browser.new_context( + storage_state=browser_state, + user_agent=browser_user_agent, + ) page = await context.new_page() async with TwoCaptchaSolver( framework=FrameworkType.PLAYWRIGHT, @@ -158,21 +163,26 @@ async def fetch_in_browser( raise RuntimeError(f"RSS did not load after Cloudflare challenge: {await page.title()}") from exc xml = await page.locator("pre").inner_text() state = await context.storage_state() + user_agent = await page.evaluate("navigator.userAgent") await browser.close() - return xml, state + return xml, state, user_agent async def enrich_items_in_browser( items: list[ParsedItem], rucaptcha_token: str, browser_state: dict[str, Any] | None, + browser_user_agent: str | None, content_selector: str, text_selector: str, -) -> tuple[list[ParsedItem], dict[str, Any]]: +) -> tuple[list[ParsedItem], dict[str, Any], str]: async with browser_lock: async with async_playwright() as playwright: browser = await playwright.chromium.launch(channel="chrome", headless=False) - context = await browser.new_context(storage_state=browser_state) if browser_state else await browser.new_context() + context = await browser.new_context( + storage_state=browser_state, + user_agent=browser_user_agent, + ) page = await context.new_page() async with TwoCaptchaSolver( framework=FrameworkType.PLAYWRIGHT, @@ -202,8 +212,9 @@ async def enrich_items_in_browser( item.text = text item.media = list({entry["url"]: entry for entry in [*item.media, *media]}.values()) state = await context.storage_state() + user_agent = await page.evaluate("navigator.userAgent") await browser.close() - return items, state + return items, state, user_agent @app.get("/health") @@ -241,6 +252,7 @@ async def parse_source( url = str(request.url) xml = "" state = request.browser_state + browser_user_agent = request.browser_user_agent fetched_via = "http" if access != "cloudflare": try: @@ -261,7 +273,12 @@ async def parse_source( raise HTTPException(status_code=422, detail="rucaptcha_token is required for Cloudflare") fetched_via = "cloudflare" try: - xml, state = await fetch_in_browser(url, request.rucaptcha_token.get_secret_value(), state) + xml, state, browser_user_agent = await fetch_in_browser( + url, + request.rucaptcha_token.get_secret_value(), + state, + browser_user_agent, + ) except Exception as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc try: @@ -272,10 +289,11 @@ async def parse_source( if not request.rucaptcha_token: raise HTTPException(status_code=422, detail="rucaptcha_token is required when follow_links=true") try: - items, state = await enrich_items_in_browser( + items, state, browser_user_agent = await enrich_items_in_browser( items, request.rucaptcha_token.get_secret_value(), state, + browser_user_agent, content_selector, text_selector, ) @@ -287,4 +305,5 @@ async def parse_source( "fetched_via": fetched_via, "items": [item.model_dump() for item in items], "browser_state": state, + "browser_user_agent": browser_user_agent, } diff --git a/site_parser_worker/docker-compose.yml b/site_parser_worker/docker-compose.yml new file mode 100644 index 0000000..81986c1 --- /dev/null +++ b/site_parser_worker/docker-compose.yml @@ -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 diff --git a/src/vk_parser_app/source_adapters.py b/src/vk_parser_app/source_adapters.py index afabe72..705e505 100644 --- a/src/vk_parser_app/source_adapters.py +++ b/src/vk_parser_app/source_adapters.py @@ -108,6 +108,7 @@ class SiteParserClient: "config": config, "rucaptcha_token": self.rucaptcha_token or None, "browser_state": runtime_state.get("browser_state"), + "browser_user_agent": runtime_state.get("browser_user_agent"), } try: async with self.session.post( @@ -144,4 +145,11 @@ class SiteParserClient: ] items.append(SourceItem(external_id, url, text, _posted_at(raw.get("published_at")), media, raw)) state = data.get("browser_state") - return items, ({"browser_state": state} if isinstance(state, dict) else None) + return items, ( + { + "browser_state": state, + "browser_user_agent": str(data.get("browser_user_agent") or "") or None, + } + if isinstance(state, dict) + else None + ) diff --git a/src/vk_parser_app/workers/media_uploader.py b/src/vk_parser_app/workers/media_uploader.py index d3c56bf..a8bde82 100644 --- a/src/vk_parser_app/workers/media_uploader.py +++ b/src/vk_parser_app/workers/media_uploader.py @@ -30,6 +30,7 @@ from ..constants import ( from ..db import fetch_float_setting, fetch_int_setting, fetch_setting, get_pool from ..heartbeat import HeartbeatReporter from ..jobs import ack_done, ack_retry, claim_job, is_worker_enabled, recover_stale_jobs +from ..source_adapters import json_object TMP_DIR = Path("/tmp") TMP_PREFIX = "vkparser_tg_media_" @@ -210,10 +211,12 @@ class MediaUploader: async def load_media(self, raw_post_id: int) -> list[dict]: rows = await self.pool.fetch( """ - SELECT * - FROM raw_post_media - WHERE raw_post_id=$1 - ORDER BY sort_order ASC, id ASC + SELECT m.*, s.url AS source_url, s.runtime_state_json AS source_runtime_state + FROM raw_post_media m + JOIN raw_posts rp ON rp.id=m.raw_post_id + JOIN sources s ON s.id=rp.source_id + WHERE m.raw_post_id=$1 + ORDER BY m.sort_order ASC, m.id ASC """, raw_post_id, ) @@ -277,9 +280,35 @@ class MediaUploader: return f"media not uploaded: {details}" return None - async def download_bytes(self, session: aiohttp.ClientSession, url: str) -> bytes | None: + def site_request_options(self, media: dict) -> dict: + if str(media.get("platform") or "") != "site": + return {} + state = json_object(media.get("source_runtime_state")) + browser_state = json_object(state.get("browser_state")) + cookies = { + str(cookie["name"]): str(cookie["value"]) + for cookie in browser_state.get("cookies") or [] + if isinstance(cookie, dict) and cookie.get("name") and cookie.get("value") + } + headers = {} + if state.get("browser_user_agent"): + headers["User-Agent"] = str(state["browser_user_agent"]) + if media.get("source_url"): + headers["Referer"] = str(media["source_url"]) + return {"headers": headers, "cookies": cookies} + + async def download_bytes( + self, + session: aiohttp.ClientSession, + url: str, + request_options: dict | None = None, + ) -> bytes | None: try: - async with session.get(url, timeout=aiohttp.ClientTimeout(total=self.download_timeout_sec)) as response: + async with session.get( + url, + timeout=aiohttp.ClientTimeout(total=self.download_timeout_sec), + **(request_options or {}), + ) as response: if response.status != 200: return None return await response.read() @@ -291,10 +320,15 @@ class MediaUploader: 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)) as response: + async with session.get( + url, + timeout=aiohttp.ClientTimeout(total=self.download_timeout_sec), + **(request_options or {}), + ) as response: if response.status != 200: return {"error": f"video returned HTTP {response.status}", "permanent": False} size = 0 @@ -432,6 +466,7 @@ class MediaUploader: media_id = int(item["id"]) media_type = str(item["media_type"]) url = str(item.get("original_url") or "") + request_options = self.site_request_options(item) if item.get("tg_file_id"): prepared.append({"media_id": media_id, "media_type": media_type, "media": item["tg_file_id"]}) @@ -440,7 +475,7 @@ class MediaUploader: continue if media_type == "photo": - data = await self.download_bytes(session, url) + data = await self.download_bytes(session, url, request_options) if not data: await self.mark_media_failed_attempt(media_id, "photo download failed") continue @@ -462,7 +497,7 @@ class MediaUploader: info = ( await self.download_video(url, temp_path) if video_provider(url) - else await self.download_http_video(session, url, temp_path) + else await self.download_http_video(session, url, temp_path, request_options) ) if not info: await self.mark_media_failed_attempt(media_id, "video download failed") diff --git a/tests/test_source_adapters.py b/tests/test_source_adapters.py index dc253c7..5de33a2 100644 --- a/tests/test_source_adapters.py +++ b/tests/test_source_adapters.py @@ -25,6 +25,7 @@ class FakeResponse: "media": [{"type": "photo", "url": "https://example.test/1.jpg"}], }], "browser_state": {"cookies": [{"name": "cf_clearance"}]}, + "browser_user_agent": "Test Browser", } @@ -69,6 +70,7 @@ class SourceAdapterTests(unittest.IsolatedAsyncioTestCase): self.assertEqual(items[0].text, "Title\n\nBody") self.assertEqual(items[0].media[0].url, "https://example.test/1.jpg") self.assertEqual(state["browser_state"]["cookies"][0]["name"], "cf_clearance") + self.assertEqual(state["browser_user_agent"], "Test Browser") async def test_followed_page_does_not_repeat_title(self) -> None: client = SiteParserClient(FollowedPageSession(), "http://worker", "token", "captcha", 30) diff --git a/tests/test_storage_uploader.py b/tests/test_storage_uploader.py index a7ddae4..125dd13 100644 --- a/tests/test_storage_uploader.py +++ b/tests/test_storage_uploader.py @@ -33,6 +33,22 @@ def test_video_provider() -> 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) @@ -43,4 +59,5 @@ async def test_http_video_download() -> None: if __name__ == "__main__": test_video_provider() + test_site_request_options() asyncio.run(test_http_video_download())