fix: reuse Cloudflare session for site media
This commit is contained in:
@@ -142,9 +142,11 @@ https://vk.com/wall-239548476_123
|
|||||||
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
|
||||||
```
|
```
|
||||||
|
|
||||||
## Запуск локально
|
## Запуск локально
|
||||||
@@ -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
|
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.media_uploader
|
PYTHONPATH=src .venv/bin/python -m vk_parser_app.workers.media_uploader
|
||||||
|
|||||||
@@ -27,6 +27,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):
|
||||||
@@ -133,12 +134,16 @@ async def fetch_in_browser(
|
|||||||
url: str,
|
url: str,
|
||||||
rucaptcha_token: str,
|
rucaptcha_token: str,
|
||||||
browser_state: dict[str, Any] | None,
|
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.
|
# ponytail: one browser at a time; use a queue only when parallel source parsing is needed.
|
||||||
async with browser_lock:
|
async with browser_lock:
|
||||||
async with async_playwright() as playwright:
|
async with async_playwright() as playwright:
|
||||||
browser = await playwright.chromium.launch(channel="chrome", headless=False)
|
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()
|
page = await context.new_page()
|
||||||
async with TwoCaptchaSolver(
|
async with TwoCaptchaSolver(
|
||||||
framework=FrameworkType.PLAYWRIGHT,
|
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
|
raise RuntimeError(f"RSS did not load after Cloudflare challenge: {await page.title()}") from exc
|
||||||
xml = await page.locator("pre").inner_text()
|
xml = await page.locator("pre").inner_text()
|
||||||
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 xml, state, user_agent
|
||||||
|
|
||||||
|
|
||||||
async def enrich_items_in_browser(
|
async def enrich_items_in_browser(
|
||||||
items: list[ParsedItem],
|
items: list[ParsedItem],
|
||||||
rucaptcha_token: str,
|
rucaptcha_token: str,
|
||||||
browser_state: dict[str, Any] | None,
|
browser_state: dict[str, Any] | None,
|
||||||
|
browser_user_agent: str | None,
|
||||||
content_selector: str,
|
content_selector: str,
|
||||||
text_selector: str,
|
text_selector: str,
|
||||||
) -> tuple[list[ParsedItem], dict[str, Any]]:
|
) -> tuple[list[ParsedItem], dict[str, Any], str]:
|
||||||
async with browser_lock:
|
async with browser_lock:
|
||||||
async with async_playwright() as playwright:
|
async with async_playwright() as playwright:
|
||||||
browser = await playwright.chromium.launch(channel="chrome", headless=False)
|
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()
|
page = await context.new_page()
|
||||||
async with TwoCaptchaSolver(
|
async with TwoCaptchaSolver(
|
||||||
framework=FrameworkType.PLAYWRIGHT,
|
framework=FrameworkType.PLAYWRIGHT,
|
||||||
@@ -202,8 +212,9 @@ async def enrich_items_in_browser(
|
|||||||
item.text = text
|
item.text = text
|
||||||
item.media = list({entry["url"]: entry for entry in [*item.media, *media]}.values())
|
item.media = list({entry["url"]: entry for entry in [*item.media, *media]}.values())
|
||||||
state = await context.storage_state()
|
state = await context.storage_state()
|
||||||
|
user_agent = await page.evaluate("navigator.userAgent")
|
||||||
await browser.close()
|
await browser.close()
|
||||||
return items, state
|
return items, state, user_agent
|
||||||
|
|
||||||
|
|
||||||
@app.get("/health")
|
@app.get("/health")
|
||||||
@@ -241,6 +252,7 @@ async def parse_source(
|
|||||||
url = str(request.url)
|
url = str(request.url)
|
||||||
xml = ""
|
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":
|
if access != "cloudflare":
|
||||||
try:
|
try:
|
||||||
@@ -261,7 +273,12 @@ async def parse_source(
|
|||||||
raise HTTPException(status_code=422, detail="rucaptcha_token is required for Cloudflare")
|
raise HTTPException(status_code=422, detail="rucaptcha_token is required for Cloudflare")
|
||||||
fetched_via = "cloudflare"
|
fetched_via = "cloudflare"
|
||||||
try:
|
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:
|
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:
|
try:
|
||||||
@@ -272,10 +289,11 @@ async def parse_source(
|
|||||||
if not request.rucaptcha_token:
|
if not request.rucaptcha_token:
|
||||||
raise HTTPException(status_code=422, detail="rucaptcha_token is required when follow_links=true")
|
raise HTTPException(status_code=422, detail="rucaptcha_token is required when follow_links=true")
|
||||||
try:
|
try:
|
||||||
items, state = await enrich_items_in_browser(
|
items, state, browser_user_agent = await enrich_items_in_browser(
|
||||||
items,
|
items,
|
||||||
request.rucaptcha_token.get_secret_value(),
|
request.rucaptcha_token.get_secret_value(),
|
||||||
state,
|
state,
|
||||||
|
browser_user_agent,
|
||||||
content_selector,
|
content_selector,
|
||||||
text_selector,
|
text_selector,
|
||||||
)
|
)
|
||||||
@@ -287,4 +305,5 @@ async def parse_source(
|
|||||||
"fetched_via": fetched_via,
|
"fetched_via": fetched_via,
|
||||||
"items": [item.model_dump() for item in items],
|
"items": [item.model_dump() for item in items],
|
||||||
"browser_state": state,
|
"browser_state": state,
|
||||||
|
"browser_user_agent": browser_user_agent,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -108,6 +108,7 @@ class SiteParserClient:
|
|||||||
"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(
|
||||||
@@ -144,4 +145,11 @@ 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)
|
return items, (
|
||||||
|
{
|
||||||
|
"browser_state": state,
|
||||||
|
"browser_user_agent": str(data.get("browser_user_agent") or "") or None,
|
||||||
|
}
|
||||||
|
if isinstance(state, dict)
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ from ..constants import (
|
|||||||
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_"
|
||||||
@@ -210,10 +211,12 @@ class MediaUploader:
|
|||||||
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,
|
||||||
)
|
)
|
||||||
@@ -277,9 +280,35 @@ class MediaUploader:
|
|||||||
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()
|
||||||
@@ -291,10 +320,15 @@ class MediaUploader:
|
|||||||
session: aiohttp.ClientSession,
|
session: aiohttp.ClientSession,
|
||||||
url: str,
|
url: str,
|
||||||
output_path: str,
|
output_path: str,
|
||||||
|
request_options: dict | None = None,
|
||||||
) -> dict:
|
) -> dict:
|
||||||
max_size = self.video_max_size_mb * 1024 * 1024
|
max_size = self.video_max_size_mb * 1024 * 1024
|
||||||
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 {"error": f"video returned HTTP {response.status}", "permanent": False}
|
return {"error": f"video returned HTTP {response.status}", "permanent": False}
|
||||||
size = 0
|
size = 0
|
||||||
@@ -432,6 +466,7 @@ class MediaUploader:
|
|||||||
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"]})
|
||||||
@@ -440,7 +475,7 @@ class MediaUploader:
|
|||||||
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
|
||||||
@@ -462,7 +497,7 @@ class MediaUploader:
|
|||||||
info = (
|
info = (
|
||||||
await self.download_video(url, temp_path)
|
await self.download_video(url, temp_path)
|
||||||
if video_provider(url)
|
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:
|
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")
|
||||||
|
|||||||
@@ -25,6 +25,7 @@ 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",
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -69,6 +70,7 @@ class SourceAdapterTests(unittest.IsolatedAsyncioTestCase):
|
|||||||
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")
|
||||||
|
|
||||||
async def test_followed_page_does_not_repeat_title(self) -> None:
|
async def test_followed_page_does_not_repeat_title(self) -> None:
|
||||||
client = SiteParserClient(FollowedPageSession(), "http://worker", "token", "captcha", 30)
|
client = SiteParserClient(FollowedPageSession(), "http://worker", "token", "captcha", 30)
|
||||||
|
|||||||
@@ -33,6 +33,22 @@ def test_video_provider() -> None:
|
|||||||
assert video_provider("http://127.0.0.1/video.mp4") 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:
|
async def test_http_video_download() -> None:
|
||||||
path = "/tmp/vkparser_tg_media_http_test.mp4"
|
path = "/tmp/vkparser_tg_media_http_test.mp4"
|
||||||
result = await MediaUploader().download_http_video(FakeSession(), "https://example.test/a.mp4", path)
|
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__":
|
if __name__ == "__main__":
|
||||||
test_video_provider()
|
test_video_provider()
|
||||||
|
test_site_request_options()
|
||||||
asyncio.run(test_http_video_download())
|
asyncio.run(test_http_video_download())
|
||||||
|
|||||||
Reference in New Issue
Block a user