From dd57d9501843340b18a611b8cf75b8959ec5ed89 Mon Sep 17 00:00:00 2001 From: Patrik Wikstrom Date: Mon, 21 Sep 2026 10:20:18 +1000 Subject: [PATCH] Fetch login-gated Instagram posts, and let corroborated verdicts drain a YouTube queue of dead videos Two unrelated scrapers had both stopped draining. Instagram cleared 10 of 75 queued posts per run; YouTube cleared none of 190, run after run. - Instagram scraped public posts only. The yt-dlp legs had run cookie-less since 2026-07, when attaching cookies made yt-dlp take Instagram's authenticated web API and that 404'd on every post. Under yt-dlp 2026.8.19 the path works again, and anonymous-only had become the binding constraint: 65 of 75 posts failed every attempt, 50 on Instagram's own ruling ("This content isn't available to everyone", which matched no rule and churned as a retryable unknown) and 15 on yt-dlp's "empty media response". Sampled live, 13/13 failed anonymously and 13/13 extracted with the cookies. The scraper now goes anonymous first and retries once with the session cookies for a post hidden from logged-out viewers -- immediately, rather than burning three anonymous attempts on a wall it cannot pass -- and the media leg follows the metadata leg's auth mode. Anonymous-first keeps the account off every public post and survives either path breaking again. With the cookies attached an empty media response means throttling, as it did before. - A YouTube queue that retries had distilled down to dead videos could never drain. The metadata leg's ignore_no_formats_error swallowed each refusal, so a dead video returned an info dict and was saved as an empty placeholder row (no author, -1 plays, created 2000-01-01); the media leg then answered a bare "Video unavailable", rightly distrusted since 2026-09-18; and fifteen of those in a row tripped the permanent-storm guard, whose abort charges no retry budget. The metadata leg now also asks the tv player client -- the only one that states why YouTube will not play a video -- and captures the reason the flag hides. A video the platform has no record of is a failure, never a placeholder row. - Corroborated verdicts. The storm guards read a homogeneous run as a broken session, but a queue of nothing but retries is homogeneous by construction, so retrying failures guaranteed the guard would trip. A permanent verdict resting on per-item evidence independent of the error text is now marked corroborated: it neither extends nor resets a storm run and is pruned even when the guard trips. YouTube corroborates no-record-plus-permanent-reason, and a record kept but refused here by region or rights claim (scraped metadata-only, media leg skipped). A bare "Video unavailable" with the record intact is never corroborated, so the 2026-09-18 protection stands. - Classification: "This video is unavailable" was a retryable unknown while "Video unavailable" was a permanent removal; all nine seen were gone for good. Content ID blocks get their own permanent "blocked" category, kept distinct from geo_blocked and removed so another vantage point can single them out; copyright takedowns read as removals; captcha reads as a bot check. - Retry strikes: a batch that pruned anything deleted the whole per-platform sidecar, so one item's success reset a never-succeeding tail's strikes. Progress now clears only the strikes of the ids it pruned, and ids that leave the queue drop their media strikes even after an aborted batch. - Instagram and YouTube are marked residential_ip_only: neither works from a datacenter IP, and a Cloud Run run would burn the queue and trip guards whose shared state then holds off the local install. The worker refuses and the enrichment supervisor skips the queue without charging the plan a stall. Replayed against the live queue, the next YouTube run takes it from 190 to 2. --- CHANGELOG.md | 65 ++++ README.md | 5 + docs/architecture.md | 7 +- docs/installation.md | 3 +- docs/pipeline.md | 62 ++++ fyp/scrape/instagram_dl.py | 118 ++++-- fyp/scrape/platform_scraper.py | 51 ++- fyp/scrape/scrape.py | 39 +- fyp/scrape/scrape_queues.py | 48 ++- fyp/scrape/youtube_dl.py | 211 ++++++++++- tests/unit/test_instagram_image_posts.py | 17 +- tests/unit/test_instagram_login_fallback.py | 160 ++++++++ tests/unit/test_scrape_retry_budget.py | 79 ++-- .../unit/test_scrape_verdict_corroboration.py | 349 ++++++++++++++++++ tests/unit/test_youtube_scraper.py | 25 ++ web_interface/run_enrichment_supervisor.py | 22 ++ web_interface/run_queue_scraper.py | 31 +- 17 files changed, 1181 insertions(+), 111 deletions(-) create mode 100644 tests/unit/test_instagram_login_fallback.py create mode 100644 tests/unit/test_scrape_verdict_corroboration.py diff --git a/CHANGELOG.md b/CHANGELOG.md index d0551159..549d4891 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -112,6 +112,71 @@ public version. Entries below describe the Hub as it stands at that release. ### Fixed +- **Instagram scraped public posts only, and said nothing about it.** Since + 2026-07 the scraper ran fully anonymously, because attaching cookies then + made yt-dlp take Instagram's authenticated web API, which 404'd on every + post. That path works again, and anonymous-only had quietly become the + binding constraint: a queue drained 10 of 75 posts per run while the other + 65 failed every attempt, 50 of them on Instagram's own ruling ("This + content isn't available to everyone: It can't be seen by certain + audiences", which matched no classifier rule and churned as a retryable + `unknown`) and 15 on yt-dlp's "empty media response". Sampled the day of + the fix, 13 of 13 such posts failed anonymously and 13 of 13 extracted with + the operator's cookies. The scraper now goes anonymous first and retries + once with the session cookies for a post hidden from logged-out viewers — + immediately, instead of burning three anonymous attempts on a wall it + cannot pass — and the media leg follows the metadata leg's auth mode. The + ruling classifies as a login wall; with cookies attached an empty media + response means throttling again, as it did before. Instagram's health check + reports the cookies' real state instead of "anonymous access". + +- **A YouTube queue of dead videos could never drain.** Every id in the queue + had already failed, so retries had distilled it down to videos that were + gone or blocked. Three things then interlocked: the metadata leg runs with + `ignore_no_formats_error`, so a refused video still returned an info dict + and was saved as an empty placeholder row (no author, -1 plays, created + 2000-01-01 — 224 such rows in the last 25 scrape files); the media leg then + answered a bare "Video unavailable", which reads as a removal but is also + what a throttled session returns, so the verdict is deliberately + distrusted; and 15 of those in a row tripped the permanent-storm guard, + which aborts the batch — and an aborted batch charges no retry budget. The + result was 190 items queued, 45 attempted, 0 drained, run after run, with + the queue file untouched for two days. Fixed by giving the metadata leg the + tv player client, the only one that states *why* YouTube will not play a + video, and capturing the reason the flag otherwise swallows. A video the + platform has no record of is now a failure rather than a placeholder row, + and both it and a video whose record is intact but which names a region + whitelist or a rights claim are treated as *corroborated*: verdicts backed + by per-item evidence, which neither feed the storm guards nor are demoted + by them. A bare "Video unavailable" with the record intact is still + distrusted, so the protection added after 2026-09-18 stands. Replayed + against the live queue, the next run prunes 187 of 190 and downloads the + one video that plays. + +- **YouTube read two spellings of the same failure oppositely.** "Video + unavailable" classified as a permanent removal while "This video is + unavailable" matched no rule and churned as a retryable `unknown` — every + one of the nine seen was gone for good. Takedown and block phrasings the tv + client reports are classified too: copyright claim blocks as their own + `blocked` category (kept distinct from `geo_blocked` and `removed` so a run + from another vantage point can single them out), copyright *takedowns* as + removals, and a captcha challenge as a bot check. + +- **One item's success reset every other item's retry strikes.** A batch that + pruned anything deleted the whole per-platform strike sidecar, so a queue + that trickled forward could carry a tail that failed every run + indefinitely — no strikes ever accumulated against it. Progress now clears + only the strikes of the ids it actually pruned, and every id that leaves + the queue drops its media-retry strikes even when the batch was aborted. + +- **Instagram and YouTube scraping is refused on Cloud Run.** Neither works + from a datacenter IP whatever cookies or PO tokens are attached, and a run + there would burn the queue and trip guards whose state, shared through the + bucket, then holds off the local install that can actually drain it. Both + scrapers are marked `residential_ip_only`: the queue worker declines to + start and the enrichment supervisor leaves the queue alone without charging + the plan a stall. TikTok is unaffected. + - **Viability floors counted rows that are not viewing.** A donation was admitted on 10+ rows of *any* kind, so an export could become a collection with no viewing activity in it at all. TikTok was the clearest case: diff --git a/README.md b/README.md index 80e38e37..90f66b27 100644 --- a/README.md +++ b/README.md @@ -107,6 +107,11 @@ python web_interface/run_queue_annotator.py python web_interface/run_queue_scraper.py --platform tiktok ``` +The Instagram and YouTube scrapers have to be run this way, from a +residential connection: both platforms wall off Cloud Run's datacenter IPs +whatever cookies are attached, so the deployed services decline to scrape +them and leave those queues to a local install. + ## Verification Every change should pass the gate before merging: diff --git a/docs/architecture.md b/docs/architecture.md index 74b575f1..5522d5a6 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -163,7 +163,12 @@ identical *permanent* classifications (a flagged session mis-reporting live items as removed) or identical *transient* ones (a bot wall failing every item retryably) abort the batch, stop self-chaining, and raise a persistent per-platform scraper alert; the failed-scrapes record stores each item's -failure category so storms are diagnosable after the fact. On the analysis +failure category so storms are diagnosable after the fact. Because those +guards read a homogeneous run as a broken session, and a queue of nothing but +retries is homogeneous by construction, a scraper can mark a verdict +**corroborated** by per-item evidence — it then neither extends nor resets a +storm run and is pruned even when the guard trips (see +[pipeline.md](pipeline.md)). On the analysis side, every study refresh writes a **methods/provenance note** (`{study}_methods.json`) summarising filters, counts, and the contract/model versions behind the data — see diff --git a/docs/installation.md b/docs/installation.md index 4d503dbd..0bef3155 100644 --- a/docs/installation.md +++ b/docs/installation.md @@ -30,7 +30,8 @@ Slack and email integrations silently no-op when unconfigured. Everything stores | **Python 3.12** | everything | Matches production (`python:3.12-slim` on Cloud Run); `ruff` and CI target 3.12 too. `brew install python@3.12` / `apt install python3.12` | | `ffmpeg` | YouTube HD media only | yt-dlp needs it to merge DASH video+audio. TikTok/Instagram downloads and photo-slideshow assembly work without it (bundled `imageio-ffmpeg`). `brew install ffmpeg` / `apt install ffmpeg` | | `node` *or* `deno` | YouTube media from datacenter IPs | Runs yt-dlp's JS challenge solver. Usually unnecessary on a home (residential) connection. | -| Google Chrome, logged in | authenticated scraping | Cookies are read from the local Chrome profile — **macOS only** (approve the Keychain prompt on first use). Instagram scraping effectively requires this; TikTok/YouTube degrade to public-content access. On Linux, provide a Netscape cookies file instead: `YTDLP_COOKIE_FILE_TIKTOK` / `_INSTAGRAM` / `_YOUTUBE` for one platform, or `YTDLP_COOKIE_FILE` for all of them (the per-platform form wins; either is ignored unless the file exists). | +| Google Chrome, logged in | authenticated scraping | Cookies are read from the local Chrome profile — **macOS only** (approve the Keychain prompt on first use). YouTube scrapes signed in; Instagram needs the cookies for posts it hides from logged-out viewers (public posts scrape anonymously); TikTok degrades to public-content access. On Linux, provide a Netscape cookies file instead: `YTDLP_COOKIE_FILE_TIKTOK` / `_INSTAGRAM` / `_YOUTUBE` for one platform, or `YTDLP_COOKIE_FILE` for all of them (the per-platform form wins; either is ignored unless the file exists). | +| A residential connection | Instagram and YouTube scraping | Both wall off datacenter IPs whatever cookies or PO tokens are attached, so their scrape queues are drained by a local install and the Cloud Run deployment declines to run them (`BaseScraper.residential_ip_only`). TikTok scrapes fine from Cloud Run. | ## Install diff --git a/docs/pipeline.md b/docs/pipeline.md index d2432a91..31890f38 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -114,6 +114,68 @@ plus optional overrides (throttle limits, health check, slideshow hooks). All three current scrapers (TikTok, Instagram, YouTube) are yt-dlp-based; cookies are managed per-platform by `scraper_cookies.py`. +**Where each scraper runs.** Instagram and YouTube do not scrape from +datacenter IPs in practice — both wall Cloud Run off whatever cookies or +PO tokens are attached — so their queues are drained by a local install on a +residential IP. That is a property of the scraper +(`BaseScraper.residential_ip_only`): on Cloud Run the queue worker refuses to +start and the enrichment supervisor leaves those queues alone, instead of +burning them against the wall and tripping guards whose state, shared through +the bucket, would then hold the local install off too. + +**Authentication.** TikTok and YouTube scrape signed in (locally, from the +operator's own Chrome profile). Instagram goes anonymous first and retries +with the session cookies only for a post it hides from logged-out viewers, +which keeps the account's footprint to the posts that need it and survives +either path breaking — both have (2026-07: attaching cookies broke every +extraction; 2026-09: anonymous alone left 65 of 75 queued posts unfetchable). + +**Permanent vs transient.** Each scraper classifies a failure into its own +taxonomy, and `classify_error` maps it to `permanent:` (pruned from +the queue, recorded in the failed-scrapes ledger) or `transient:` +(kept for a later run). One broken session can make every item read as +permanently gone, so three guards sit on top — all of them abort the batch, +and the storm guards also raise a scraper alert for a human: + +| Guard | Trips on | The items | +|---|---|---| +| circuit breaker | 15 consecutive throttle verdicts | stay queued | +| permanent-storm guard | 15 consecutive identical *permanent* verdicts | demoted to transient, stay queued | +| transient-storm guard | 25 consecutive identical *transient* verdicts | already transient; chaining stops | + +**Corroborated verdicts.** Those guards assume a healthy queue produces +heterogeneous outcomes — which a queue of nothing but retries never does, +since retrying only the failures distils it down to items that fail. A +scraper may therefore mark a permanent verdict *corroborated* +(`attrs['verdict_corroborated']`, see `BaseScraper.fetch`) when it rests on +per-item evidence independent of the error text: such a verdict neither +extends nor resets a storm run, and is pruned even when the guard trips. +YouTube corroborates two cases, both read from the metadata leg, which adds +the tv player client — the only one that states *why* a video will not play — +and captures the reason yt-dlp otherwise swallows under +`ignore_no_formats_error`: + +* the platform has no record of the video (no channel, no view count, no + duration) **and** the stated reason is itself a removal. A video with no + record is a failure, never the empty placeholder row it used to be saved as; +* the record is intact but YouTube refuses to play it here, naming a region + whitelist or a rights claim. Its metadata is scraped, the media leg is + skipped, and the id leaves the queue with its metadata-only row standing. + +A bare "Video unavailable" with the record intact is never corroborated: that +is exactly what a throttled session returns. + +**Retry budgets** bound everything that stays queued, in per-platform sidecars +(`scrape_queues.py`). An item transiently failing through +`MAX_ZERO_PROGRESS_STRIKES` runs in which the queue as a whole made no +progress is given up on and recorded as failed; an item whose metadata scraped +but whose media did not is retried for `MAX_MEDIA_RETRY_STRIKES` healthy runs +and then pruned with its metadata-only row standing. An aborted batch never +charges either budget — the verdicts implicate the session, not the items — +but ids that left the queue always drop their strikes, and a batch that makes +progress clears only the strikes of the ids it pruned (one item's success is +no evidence for another's). + The canonical cross-platform scrape schema lives in `config/scrape_contract.toml`: base fields every platform emits (including generic popularity counts `fave_count`/`comment_count`/... and per-K diff --git a/fyp/scrape/instagram_dl.py b/fyp/scrape/instagram_dl.py index eaa6920d..1c84a8fd 100644 --- a/fyp/scrape/instagram_dl.py +++ b/fyp/scrape/instagram_dl.py @@ -5,13 +5,23 @@ Fetches metadata + media for Instagram posts/reels identified by their URL shortcode (the ``item_id`` produced by :class:`fyp.ingest.InstagramDDPCollection`). -Extraction runs **anonymously** (no cookies): as of 2026-07 Instagram killed -its authenticated web API (``api/v1/media/{pk}/info/`` 404s for web sessions -and post pages render as an empty SPA shell), so attaching session cookies -makes every yt-dlp extraction fail — while the logged-out GraphQL path -yt-dlp ≥2026.7.4 uses works. Follow-gated/private content is therefore -permanently inaccessible (classified ``private``); the donated enrichment -seed still surfaces its caption/author. +Run it from a residential IP. Instagram (like YouTube) does not scrape from +Cloud Run's datacenter IPs in practice; the local install drains this queue. + +Extraction runs **anonymously first, then with the session cookies** for a +post Instagram hides from logged-out viewers ("This content isn't available to +everyone: It can't be seen by certain audiences", or yt-dlp's "Instagram sent +an empty media response"). History: in 2026-07 Instagram's authenticated web +API (``api/v1/media/{pk}/info/``) 404'd for web sessions, and yt-dlp takes that +path whenever cookies are attached — so every cookie-bearing extraction failed +and scraping went fully anonymous. By 2026-09 (yt-dlp 2026.8.19) that path +works again, and anonymous extraction had left 65 of 75 queued posts failing +every run: measured 2026-09-21, 13/13 sampled posts failed anonymously and +13/13 extracted with the cookies. Anonymous-first keeps the logged-in account's +footprint to the posts that need it, and survives either path breaking again. +Follow-gated content ("only available for registered users who follow this +account") stays permanently ``private``; the donated enrichment seed still +surfaces its caption/author. Image-only posts (single photos and carousels): extraction uses yt-dlp's ``ignore_no_formats_error`` so an image post returns a full info dict; the @@ -87,7 +97,8 @@ def _classify_error(exc: Exception) -> tuple[str, str]: (category, detail) where category is one of: - "rate_limited" — HTTP 429/403, empty media response, IG's ambiguous "rate-limit reached or login required" catch-all - - "login_required" — bare login wall (usually dead cookies, account-wide) + - "login_required" — login wall: the post is shown only to logged-in + viewers (retried with the session cookies) - "no_video" — image-only post (no video to extract) - "private" — private account/post - "removed" — post deleted or id nonexistent @@ -141,7 +152,13 @@ def _classify_error(exc: Exception) -> tuple[str, str]: if 'private' in msg_lower or 'only available for registered users' in msg_lower: return "private", msg - if 'login required' in msg_lower or 'log in' in msg_lower or 'logged-in' in msg_lower: + # Instagram's own ruling for a post it shows only to logged-in viewers + # ("This content isn't available to everyone: It can't be seen by certain + # audiences"). It fell through to "unknown" until 2026-09-21; as a login + # wall it now triggers the retry with the session cookies. + if ("isn't available to everyone" in msg_lower or 'certain audiences' in msg_lower + or 'login required' in msg_lower or 'log in' in msg_lower + or 'logged-in' in msg_lower): return "login_required", msg if any(kw in msg_lower for kw in ('unavailable', 'removed', 'deleted', 'not found', @@ -209,17 +226,51 @@ def _info_to_row(info: dict, item_id: str) -> pd.DataFrame: +def _login_gated(category: str, detail: str) -> bool: + """True when a failure means "Instagram shows this post only to logged-in viewers". + + "Instagram sent an empty media response" classifies ``rate_limited`` — with + the cookies attached it means throttling — but anonymously yt-dlp itself + says the post may need a login, and on 2026-09-21 every sampled one did. + """ + return category == "login_required" or 'empty media response' in (detail or '').lower() + + def _extract_metadata(url: str, item_id: str, verbose: bool = False): """yt-dlp metadata extraction with retry. Returns (info, None) or (None, fail_df). - Runs anonymously (see module docstring — session cookies make every - extraction fail since Instagram's 2026-07 web-API change). + Anonymous first (see module docstring). A post Instagram shows only to + logged-in viewers is retried once with the session cookies; its info dict + is then marked ``_fyp_authenticated`` so the media leg follows suit. ``ignore_no_formats_error`` lets image-only posts return their info dict (metadata + image thumbnails) instead of raising ``no_video``. """ + info, fail = _extract_metadata_as(url, item_id, {}, verbose=verbose) + if fail is None or not _login_gated(fail.attrs.get('error_type'), fail.attrs.get('error_detail')): + return info, fail + cookies = scraper_cookies.cookie_opts("instagram") + if not cookies: + return info, fail + logger.info("Scrape %s: hidden from logged-out viewers — retrying with the session cookies", + item_id) + info, fail = _extract_metadata_as(url, item_id, cookies, verbose=verbose) + if info is not None: + info['_fyp_authenticated'] = True + return info, fail + + +def _extract_metadata_as(url: str, item_id: str, cookies: dict, verbose: bool = False): + """One metadata pass, anonymous (``cookies={}``) or with the session cookies. + + A login wall ends the pass at once: repeating the same request cannot get + past it (the anonymous pass used to burn three attempts per gated post). + With the cookies attached, an empty media response is throttling again + and retries with backoff like any rate limit. + """ ydl_opts: dict = { 'quiet': True, 'no_warnings': not verbose, + **cookies, 'skip_download': True, 'no_color': True, 'ignore_no_formats_error': True, @@ -233,8 +284,11 @@ def _extract_metadata(url: str, item_id: str, verbose: bool = False): return ydl.extract_info(url, download=False), None except (yt_dlp.utils.DownloadError, ExtractorError) as e: category, detail = _classify_error(e) - logger.warning("Scrape %s metadata attempt %d/%d failed: [%s] %s", - item_id, attempt + 1, _META_MAX_RETRIES, category, detail) + logger.warning("Scrape %s metadata attempt %d/%d failed%s: [%s] %s", + item_id, attempt + 1, _META_MAX_RETRIES, + " (with cookies)" if cookies else "", category, detail) + if category == "login_required" or (not cookies and _login_gated(category, detail)): + return None, _empty_fail(category, detail) if category in _RETRYABLE and attempt < _META_MAX_RETRIES - 1: backoff = 3 * (2 ** attempt) logger.info("Retrying %s in %ds...", item_id, backoff) @@ -279,9 +333,15 @@ def _download_media( save_path: str, stream_to_bucket=None, verbose: bool = False, + authenticated: bool = False, ) -> tuple[bool, str | None, str, float | None]: """Download the post's video to temp and move/upload it. + Args: + authenticated: Attach the session cookies — set when the metadata leg + needed them (the download re-extracts the post, so a post hidden + from logged-out viewers fails anonymously here too). + Returns: ``(ok, error_category, error_detail, duration)`` — category/detail are ``None``/"" on success, otherwise the :func:`_classify_error` result of @@ -294,6 +354,7 @@ def _download_media( dl_opts: dict = { 'quiet': True, 'no_warnings': not verbose, + **(scraper_cookies.cookie_opts("instagram") if authenticated else {}), 'outtmpl': out_template, 'no_color': True, 'overwrites': True, @@ -824,20 +885,23 @@ def _download_images( class InstagramScraper(BaseScraper): - """Instagram platform scraper (yt-dlp, anonymous logged-out extraction). + """Instagram platform scraper (yt-dlp; anonymous first, cookies for gated posts). - Handles video posts and reels via yt-dlp's anonymous GraphQL path; image - posts (photos and carousels) extract to format-less info dicts whose + Handles video posts and reels via yt-dlp's anonymous GraphQL path, falling + back to the session cookies for posts hidden from logged-out viewers; + image posts (photos and carousels) extract to format-less info dicts whose thumbnails carry the source images, downloaded for the orchestrator's silent-slideshow assembly (see module docstring). Anonymous access is tightly rate-limited by Instagram, so concurrency stays capped hard — - Instagram is the most ban-happy of the supported platforms. + Instagram is the most ban-happy of the supported platforms. Residential + IP only: it does not scrape from Cloud Run in practice. """ platform = "instagram" # /p/ serves reel and tv shortcodes too (Instagram redirects). url_template = "https://www.instagram.com/p/{item_id}/" slideshow_image_column = "image_list" + residential_ip_only = True def item_url(self, item_id: str) -> str: @@ -905,7 +969,8 @@ def fetch( ok, media_category, media_detail, media_duration = _download_media( url, item_id, save_path, - stream_to_bucket=stream_to_bucket, verbose=verbose) + stream_to_bucket=stream_to_bucket, verbose=verbose, + authenticated=bool(info.get('_fyp_authenticated'))) if ok: data_row.loc[0, 'video_downloaded'] = True # Backfill the duration metadata extraction no longer returns @@ -1002,14 +1067,15 @@ def throttle_limits(self, max_workers: int) -> tuple[int, int, int]: def health_check(self) -> dict | None: - # Instagram scraping is anonymous (see module docstring) — there is no - # login session to monitor, and reporting cookie state here would send - # users chasing cookie renewals that have no effect. - return { - "present": False, - "status": "healthy", - "message": "Anonymous access — Instagram scraping does not use login cookies.", - } + # Public posts scrape anonymously, but posts Instagram hides from + # logged-out viewers need the session cookies (see module docstring), + # so their state is worth monitoring again. Without them those posts + # churn in the queue until the retry budget gives up on them. + health = scraper_cookies.cookie_health("instagram", session_cookie="sessionid") + health["message"] = (f"{health.get('message', '')} Public posts scrape anonymously; " + f"the cookies are needed for posts hidden from logged-out " + f"viewers.").strip() + return health def media_probe_url(self, item_id: str) -> dict | None: diff --git a/fyp/scrape/platform_scraper.py b/fyp/scrape/platform_scraper.py index b73ffbc2..c3f19e01 100644 --- a/fyp/scrape/platform_scraper.py +++ b/fyp/scrape/platform_scraper.py @@ -15,6 +15,7 @@ """ import logging +import os import threading from abc import ABC, abstractmethod from glob import glob @@ -30,7 +31,8 @@ logger = logging.getLogger(__name__) -def empty_fail(error_type: str = "unknown", error_detail: str = "") -> pd.DataFrame: +def empty_fail(error_type: str = "unknown", error_detail: str = "", *, + corroborated: bool = False) -> pd.DataFrame: """Return an empty DataFrame tagged with error classification metadata. Shared by every platform scraper (hoisted from the per-platform copies in @@ -40,13 +42,20 @@ def empty_fail(error_type: str = "unknown", error_detail: str = "") -> pd.DataFr Args: error_type: Scraper error category (e.g. ``"rate_limited"``). error_detail: Free-text detail for the failure row. + corroborated: The permanent verdict rests on per-item evidence beyond + the error message (see the corroboration clause of + :meth:`BaseScraper.fetch`). Stamped as + ``attrs['verdict_corroborated']``. Returns: - An empty DataFrame with ``error_type``/``error_detail`` in ``attrs``. + An empty DataFrame with ``error_type``/``error_detail`` (and + ``verdict_corroborated`` when set) in ``attrs``. """ df = pd.DataFrame() df.attrs['error_type'] = error_type df.attrs['error_detail'] = error_detail + if corroborated: + df.attrs['verdict_corroborated'] = True return df @@ -109,12 +118,19 @@ class BaseScraper(ABC): slideshow_image_column: raw column holding the ``" | "``-joined image URLs of a photo/carousel post, or ``None`` when the platform has no carousel concept. + residential_ip_only: the platform does not scrape from a datacenter IP + in practice, so its queue is drained by a local install on a + residential IP: on Cloud Run the queue worker refuses to run and + the enrichment supervisor leaves the queue alone. Instagram and + YouTube (operator experience, 2026-09): both wall off Cloud Run's + IPs whatever the cookies or PO tokens. base_columns: ``{column: pyarrow_dtype}`` for the canonical base fields. platform_columns: ``{column: pyarrow_dtype}`` for this platform's fields. """ platform: str | None = None slideshow_image_column: str | None = None + residential_ip_only: bool = False _registry: list[type] = [] def __init_subclass__(cls, **kwargs): @@ -170,6 +186,23 @@ def fetch( the media-retry budget in :func:`scrape_queues.charge_media_retry` bounds the retries instead. The category also feeds the throttle controller and the storm guards. + + Corroboration clause: a failure frame (via ``empty_fail(..., + corroborated=True)``), or a metadata-only row for its media verdict, + may carry ``attrs['verdict_corroborated'] = True`` when the permanent + verdict rests on per-item evidence independent of the error text — + e.g. YouTube answering with no record of the video at all AND a stated + reason that is itself a removal, or keeping the record but naming a + region whitelist or rights claim as the reason it will not play. The + storm guards exist because one broken session can make every item read + as removed; a verdict the item itself corroborates is explained by the + item, so it neither extends nor resets a storm run. A corroborated + failure is pruned even when the guard trips; a corroborated media + verdict prunes the id with its metadata row standing, instead of + queueing a media retry. Without this, a queue that retries have + distilled down to dead items trips the guard on every run and never + drains (2026-09-21, YouTube). Never set it on a verdict that rests on + the message alone. """ @@ -243,6 +276,20 @@ def health_check(self) -> dict | None: return None + def unavailable_here(self) -> str | None: + """Why this scraper must not run in the current environment, or ``None``. + + A ``residential_ip_only`` platform on Cloud Run (``K_SERVICE`` set) + would only burn its queue against the datacenter-IP wall and trip the + storm guards — whose tripped state, in the shared task status, would + then hold off the local install that can actually drain it. + """ + if self.residential_ip_only and os.environ.get("K_SERVICE"): + return (f"the {self.platform} scraper needs a residential IP and does not " + f"work from Cloud Run — drain this queue from a local install") + return None + + def media_probe_url(self, item_id: str) -> dict | None: """Resolve an item's direct media URL for a lightweight reachability probe. diff --git a/fyp/scrape/scrape.py b/fyp/scrape/scrape.py index 95b19cd8..7762df6f 100644 --- a/fyp/scrape/scrape.py +++ b/fyp/scrape/scrape.py @@ -1101,7 +1101,7 @@ def download_video_threads( mem_stop_event = threading.Event() inter_delay = scraper.inter_request_delay() - def _breaker_track(category) -> None: + def _breaker_track(category, corroborated: bool = False) -> None: with breaker_lock: if category in THROTTLE_CATEGORIES: breaker_state["consecutive"] += 1 @@ -1118,6 +1118,12 @@ def _breaker_track(category) -> None: # transient and would otherwise wipe the storm classification. return classification = scraper.classify_error(category) + if corroborated and classification.startswith("permanent"): + # Explained by the item, not the session (the corroboration + # clause of BaseScraper.fetch): no evidence either way, so it + # neither extends nor resets a storm run. A retry-only queue + # of dead videos otherwise trips the guard on every run. + return if classification.startswith("permanent"): if classification == storm_state["classification"]: storm_state["consecutive"] += 1 @@ -1195,8 +1201,9 @@ def worker(idx_video): error_cat = res.attrs.get('media_error_type') else: error_cat = None + corroborated = isinstance(res, pd.DataFrame) and bool(res.attrs.get('verdict_corroborated')) throttle.report_result(error_cat) - _breaker_track(error_cat) + _breaker_track(error_cat, corroborated) if inter_delay > 0: # Sleep while holding the throttle slot: paces the whole # session, not just this thread. @@ -1346,6 +1353,7 @@ def _mem_watch(): permanent_failed_ids: list[str] = [] transient_failed_ids: list[str] = [] media_retry_ids: list[str] = [] + media_unplayable = 0 storm_demoted = 0 storm_media_kept = 0 for idx in range(len(interesting_videos)): @@ -1357,7 +1365,12 @@ def _mem_watch(): # media is retried next run. attrs don't survive pd.concat, so # this is the last place they're visible. media_error = res.attrs.get('media_error_type') - if media_error is not None: + if media_error is not None and res.attrs.get('verdict_corroborated'): + # The platform named why this item will never play here (e.g. + # a region whitelist or a rights claim, record intact): the + # metadata row stands and the id is pruned like any success. + media_unplayable += 1 + elif media_error is not None: # Whatever the category. A permanent verdict on the media leg # is not trusted on its own: a throttled session's bare "Video # unavailable" reads exactly like a removal, and on 2026-09-18 @@ -1372,10 +1385,12 @@ def _mem_watch(): else: vid = interesting_videos[idx] error_type = res.attrs.get('error_type', 'unknown') if isinstance(res, pd.DataFrame) else 'unknown' + corroborated = isinstance(res, pd.DataFrame) and bool(res.attrs.get('verdict_corroborated')) # The scraper owns its platform's permanent-vs-transient taxonomy. classification = scraper.classify_error(error_type) if classification.startswith('permanent'): - if storm_state["tripped"] and classification == storm_state["classification"]: + if (storm_state["tripped"] and classification == storm_state["classification"] + and not corroborated): # Suspect storm verdict: keep the id queued and off the # failed record — a later healthy session re-scrapes it. transient_failed_ids.append(vid) @@ -1439,6 +1454,10 @@ def _mem_watch(): elif results: scraper_alerts.clear_alert(scraper.platform, reason="healthy batch") + if media_unplayable: + logger.info(f" Unplayable here: {media_unplayable} items scraped metadata-only " + f"(the platform named why — e.g. region or rights block) — removed from queue") + if media_retry_ids: logger.info(f" Media retries: {len(media_retry_ids)} items scraped metadata-only " f"(media download failed) — kept in queue for media retry") @@ -1678,7 +1697,7 @@ def _on_threads_change(n): f"({len(good_scrapes)} OK, {len(all_permanent_failed)} permanent fail). " f"{len(all_transient_failed)} transient failures remain for retry. " f"Queue length: {remaining}") - scrape_queues.clear_zero_progress(platform_resolved) + scrape_queues.clear_zero_progress(platform_resolved, items_to_remove) elif all_transient_failed and not aborted and not dry_run: # Zero-progress run: every item failed "transiently" and nothing was # pruned, so without intervention the queue would never drain (and the @@ -1699,12 +1718,12 @@ def _on_threads_change(n): # ---------------- # Media-retry budget: metadata-only rows (media failed) stay queued, but not - # forever. An aborted run never charges — the verdicts implicate the session. + # forever. An aborted run never charges — the verdicts implicate the session + # — but every id that left the queue drops its strikes either way. # ----------------- - if all_media_retry and not aborted and not dry_run: - retry_set = set(all_media_retry) - got_media = [v for v in good_scrapes if v not in retry_set] - exhausted = scrape_queues.charge_media_retry(platform_resolved, all_media_retry, got_media) + if (all_media_retry or items_to_remove) and not dry_run: + exhausted = scrape_queues.charge_media_retry( + platform_resolved, [] if aborted else all_media_retry, items_to_remove) if exhausted: _, remaining = scrape_queues.prune_scrape_queue(platform_resolved, set(exhausted)) logger.warning( diff --git a/fyp/scrape/scrape_queues.py b/fyp/scrape/scrape_queues.py index 94549612..cbcaa323 100644 --- a/fyp/scrape/scrape_queues.py +++ b/fyp/scrape/scrape_queues.py @@ -336,18 +336,39 @@ def _mutate(current): -def clear_zero_progress(platform: str) -> None: - """Drop one platform's retry strikes after a batch that made progress. +def clear_zero_progress(platform: str, resolved_ids) -> None: + """Drop the retry strikes of the items that left the queue this batch. + + Only the resolved items' strikes go. Wiping the whole sidecar on any + progress let a queue that trickled forward carry a never-succeeding tail + indefinitely: on 2026-09-21 Instagram drained 10 of 75 items in a run + while 65 gated posts, failing every attempt, had their strikes reset to + zero by those 10 successes. Another item's success says nothing in an + item's favour. An item that still carries a strike has not succeeded + since it was charged — anything that does succeed is pruned, and cleared, + here first — so keeping its strike cannot burn a fetchable item. - A draining queue means the transient failures are riding along with - successes — today's semantics (retry indefinitely) are right for those, - and keeping stale strikes would burn them spuriously if the queue later - stalls for an unrelated reason. + Args: + platform: Platform whose sidecar to update. + resolved_ids: Ids pruned this batch (scraped OK or permanently failed). """ data_io = _data_io() target = strikes_filename(platform) - if data_io.exists(storage_location=QUEUE_LOCATION, filename=target): - data_io.remove(storage_location=QUEUE_LOCATION, filename=target) + resolved = set(_dedup(list(resolved_ids or []))) + if not resolved or not data_io.exists(storage_location=QUEUE_LOCATION, filename=target): + return + + def _mutate(current): + counts = current if isinstance(current, dict) else {} + kept = {vid: n for vid, n in counts.items() if vid not in resolved} + return None if kept == counts else kept + + data_io.update_json( + storage_location=QUEUE_LOCATION, + filename=target, + mutate=_mutate, + default={}, + ) @@ -386,15 +407,18 @@ def charge_media_retry( ) -> list[str]: """Charge one media-retry strike per id and drop the ids that resolved. - Callers invoke this after a batch that was NOT aborted by a storm, + Callers pass no ``retry_ids`` after a batch that was aborted by a storm, circuit breaker or memory stop — an abort implicates the session rather - than the items, and must not burn retry budget. + than the items, and must not burn retry budget — but still pass the ids + the batch pruned, so the sidecar never keeps strikes for ids that are no + longer queued. Args: platform: Platform whose sidecar to update. retry_ids: Ids scraped metadata-only this batch (media failed). - resolved_ids: Ids whose media was downloaded this batch — their - strikes, if any, are cleared in the same write. + resolved_ids: Ids that left the queue this batch (media downloaded, + unplayable here, or permanently failed) — their strikes, if any, + are cleared in the same write. Returns: The ids whose strike count reached ``MAX_MEDIA_RETRY_STRIKES`` — diff --git a/fyp/scrape/youtube_dl.py b/fyp/scrape/youtube_dl.py index e4fe6a77..bb6ddf91 100644 --- a/fyp/scrape/youtube_dl.py +++ b/fyp/scrape/youtube_dl.py @@ -5,13 +5,20 @@ Fetches metadata + media for YouTube videos identified by their 11-character video id (the ``item_id`` produced by :class:`fyp.ingest.YouTubeDDPCollection`). -Datacenter IPs (Cloud Run) frequently hit YouTube's bot wall ("Sign in to -confirm you're not a bot"); research-account cookies partially mitigate it -(see :mod:`fyp.scraper_cookies`) and the distinct ``bot_check`` category is a -throttle signal so concurrency backs off. Media streams additionally require -proof-of-origin (PO) tokens from datacenter IPs: the bgutil provider -(pip plugin + script built in Dockerfile.base, wired via -:func:`_pot_extractor_args`) supplies them in production. +Run it from a residential IP. In practice YouTube (like Instagram) does not +scrape from Cloud Run's datacenter IPs — the bot wall ("Sign in to confirm +you're not a bot") holds even with research-account cookies and proof-of-origin +(PO) tokens — so the local install, signed in through the operator's own +Chrome, drains this queue. The Cloud Run pieces remain wired but are not a +working path: the bgutil PO-token provider (pip plugin + script built in +Dockerfile.base, via :func:`_pot_extractor_args`), and ``bot_check`` as a +throttle signal so concurrency backs off. + +Refused videos: the metadata leg adds the tv player client, the only one that +states WHY YouTube will not play a video, and captures that reason (see +:class:`_ReasonLog`). A video YouTube has no record of is a failure — never a +placeholder row — and one it keeps but will not play here (region, rights +claim) is scraped metadata-only and leaves the queue. Most watch-history items are long-form and exceed the media duration cap — they are deliberately scraped metadata-only; Shorts and clips get media. @@ -58,9 +65,21 @@ def _cf(): # "bot_check" is transient AND a throttle signal (_THROTTLE_CATEGORIES in # platform_scraper): the batch backs off instead of burning the whole queue # against the bot wall. HTTP 403 is typically YouTube throttling (unlike -# TikTok, where it means an IP block) — kept retryable. +# TikTok, where it means an IP block) — kept retryable. "blocked" is a +# copyright (Content ID) block — "It was blocked due to the claimed content +# by " — distinct from "removed" so a later run from another +# vantage point can single it out, like "geo_blocked". _RETRYABLE = {"bot_check", "rate_limited", "network", "server_error", "unknown"} -_PERMANENT = {"removed", "private", "age_restricted", "members_only", "geo_blocked"} +_PERMANENT = {"removed", "private", "age_restricted", "members_only", "geo_blocked", + "blocked"} + +# Permanent verdicts that stand even when YouTube still has the video's record +# (channel, views, duration): the player refuses it HERE, on grounds that name +# the video itself — a region whitelist or a rights holder's claim. No +# throttled session has ever produced these; it answers a bare "Video +# unavailable" (2026-09-18). Anything else with the record intact goes down +# the ordinary media leg, whose verdict is distrusted and budgeted. +_UNPLAYABLE_WITH_RECORD = {"geo_blocked", "blocked"} _META_MAX_RETRIES = 3 _DL_MAX_RETRIES = 2 @@ -76,6 +95,15 @@ def _cf(): # not used, so enabling both is safe everywhere. _JS_RUNTIMES = {'deno': {'path': None}, 'node': {'path': None}} +# Player clients for the metadata leg: yt-dlp's defaults plus "tv". When +# YouTube refuses a video, the default web clients all report a bare "Video +# unavailable" — the very text a throttled session returns — whereas the tv +# client states the reason: "removed by the uploader", "The uploader has not +# made this video available in your country", "It was blocked due to the +# claimed content by …". Measured 2026-09-21: ~1 s more per item, and a +# playable video resolves the same formats. The media leg keeps the defaults. +_METADATA_PLAYER_CLIENTS = ['default', 'tv'] + # The bgutil PO-token provider's script directory (Dockerfile.base builds it # and sets this env var). YouTube requires proof-of-origin tokens for media # streams from datacenter IPs — cookies alone don't pass the bot wall. @@ -103,13 +131,14 @@ def _classify_error(exc: Exception) -> tuple[str, str]: Returns: (category, detail) where category is one of: - - "bot_check" — "Sign in to confirm you're not a bot" wall + - "bot_check" — "Sign in to confirm you're not a bot" wall, captcha - "rate_limited" — HTTP 429/403, too many requests - "removed" — video deleted/unavailable, account terminated - "private" — private video - "age_restricted" — age gate (cookies already applied → permanent) - "members_only" — channel-membership gate - "geo_blocked" — GeoRestrictedError + - "blocked" — copyright (Content ID) claim block - "network" — timeout, connection refused, DNS failure, SSL - "server_error" — HTTP 5xx - "unknown" — unrecognised (kept retryable) @@ -130,11 +159,22 @@ def _classify_error(exc: Exception) -> tuple[str, str]: if isinstance(cause, TransportError): return "network", f"Transport error: {msg}" + return _classify_message(msg) + + +def _classify_message(msg: str) -> tuple[str, str]: + """Classify a yt-dlp error or warning text; see :func:`_classify_error`. + + Split out so the metadata leg can classify the playability reason it + captures from a warning (see :class:`_ReasonLog`) — there is no exception + object there, only the text. + """ # YouTube uses typographic apostrophes ("confirm you’re not a bot") — # normalize so ASCII keyword matching works. msg_lower = msg.lower().replace('’', "'") - if "confirm you're not a bot" in msg_lower or 'not a robot' in msg_lower: + if ("confirm you're not a bot" in msg_lower or 'not a robot' in msg_lower + or 'captcha' in msg_lower): return "bot_check", msg # Rate-limit detection must precede the "removed" keywords: YouTube's @@ -156,13 +196,30 @@ def _classify_error(exc: Exception) -> tuple[str, str]: return "members_only", msg # Geo restrictions sometimes surface as a flattened message instead of a - # GeoRestrictedError instance. + # GeoRestrictedError instance. A territorial copyright block ("…who has + # blocked it in your country on copyright grounds") lands here too — it + # is region-bound, which is what the category records. if 'in your country' in msg_lower or 'geo restriction' in msg_lower: return "geo_blocked", msg + # A Content ID block names the rights holder: "It was blocked due to the + # claimed content by Paramount Global (PMN)." / "…who has blocked it on + # copyright grounds." Only the tv player client states it (the web clients + # say a bare "Video unavailable"). A copyright TAKEDOWN ("no longer + # available due to a copyright claim") says nothing of blocking and falls + # through to "removed". + if 'claimed content' in msg_lower or ('copyright' in msg_lower and 'blocked' in msg_lower): + return "blocked", msg + # "This content isn't available, try again later" without the rate-limit # sentence is YouTube's soft-block/removal phrasing — kept as removed. - if any(kw in msg_lower for kw in ('video unavailable', 'has been removed', + # "This video is unavailable" is the phrasing for an id YouTube has no + # record of; it matched none of these until 2026-09-21 and churned as + # "unknown" for days (every one of the nine seen was gone for good). The + # tv client words takedowns as "It was removed following a copyright + # removal request by ". + if any(kw in msg_lower for kw in ('video unavailable', 'video is unavailable', + 'has been removed', 'was removed', 'removal request', 'no longer available', 'account associated', 'terminated', 'does not exist', 'not available')): return "removed", msg @@ -176,9 +233,10 @@ def _classify_error(exc: Exception) -> tuple[str, str]: -def _empty_fail(error_type: str = "unknown", error_detail: str = "") -> pd.DataFrame: +def _empty_fail(error_type: str = "unknown", error_detail: str = "", *, + corroborated: bool = False) -> pd.DataFrame: """Return an empty DataFrame tagged with error classification metadata.""" - return empty_fail(error_type, error_detail) + return empty_fail(error_type, error_detail, corroborated=corroborated) def _cleanup_temp_files(temp_dir: str, item_id: str) -> None: @@ -238,8 +296,94 @@ def _info_to_row(info: dict, item_id: str) -> pd.DataFrame: +def _metadata_extractor_args() -> dict: + """``extractor_args`` for the metadata leg: the PO-token wiring + tv client.""" + args = dict(_pot_extractor_args().get('extractor_args', {})) + args['youtube'] = {'player_client': list(_METADATA_PLAYER_CLIENTS)} + return {'extractor_args': args} + + +class _ReasonLog: + """yt-dlp logger that keeps the playability reasons of one extraction. + + The metadata leg runs with ``ignore_no_formats_error``, under which yt-dlp + downgrades a player's refusal ("This video has been removed by the + uploader") to a warning and returns an info dict anyway — the reason + exists nowhere else. Only ``[youtube]`` extractor warnings are kept; + plugin chatter (the PO-token provider) and the generic no-formats + follow-ups are not reasons. yt-dlp routes errors here too once a logger + is set; our own "attempt failed" line re-reports them, so they go to + debug. + """ + + _NOT_REASONS = ('[pot', 'no video formats found', 'requested format is not available', + 'n challenge', 'formats have been skipped', 'sabr') + + def __init__(self): + self.reasons: list[str] = [] + + def debug(self, msg): + pass + + def info(self, msg): + pass + + def warning(self, msg): + text = str(msg) + if not text.startswith('[youtube] '): + return + text = text[len('[youtube] '):] + lowered = text.lower() + if any(marker in lowered for marker in self._NOT_REASONS): + return + self.reasons.append(text) + + def error(self, msg): + logger.debug("yt-dlp: %s", msg) + + +def _playability_verdict(reasons: list[str]) -> tuple[str, str] | None: + """Classify the reasons one extraction reported; ``None`` when there were none. + + Safety first: if any reason reads as throttling (bot wall, captcha, rate + limit), that verdict wins — a permanent reason must never mask a session + problem. Otherwise the first reason that names a category wins over an + unrecognised one. + """ + if not reasons: + return None + verdicts = [_classify_message(r) for r in reasons] + for category, detail in verdicts: + if category in ("bot_check", "rate_limited"): + return category, detail + for category, detail in verdicts: + if category != "unknown": + return category, detail + return verdicts[0] + + +def _has_no_record(info: dict) -> bool: + """True when yt-dlp returned only a placeholder for the video. + + For a live video the player refuses here (geo- or copyright-blocked) the + info dict still carries the channel, view count and duration. For one that + no longer exists it carries none of them — only a synthesised title + ("youtube video #") — and until 2026-09-21 that shell was saved as a + scraped row: no author, -1 plays, created 2000-01-01. + """ + return all(info.get(k) is None for k in ('channel_id', 'uploader_id', 'view_count', 'duration')) + + def _extract_metadata(url: str, item_id: str, verbose: bool = False): - """yt-dlp metadata extraction with retry. Returns (info, None) or (None, fail_df).""" + """yt-dlp metadata extraction with retry. Returns (info, None) or (None, fail_df). + + A video YouTube has no record of is a failure, not a row: the verdict is + the stated reason, and when that reason is itself permanent the failure is + corroborated (see :meth:`BaseScraper.fetch`) — two independent signals + agree that the item, not the session, is the problem. When the record + exists but no format does, the classified reason travels to + :meth:`YouTubeScraper.fetch` as ``info['_fyp_unplayable']``. + """ ydl_opts: dict = { 'quiet': True, 'no_warnings': not verbose, @@ -250,18 +394,32 @@ def _extract_metadata(url: str, item_id: str, verbose: bool = False): 'extractor_retries': 3, 'socket_timeout': 30, 'js_runtimes': _JS_RUNTIMES, - **_pot_extractor_args(), + **_metadata_extractor_args(), # Metadata must never depend on the n-challenge solver: without a JS # runtime + yt-dlp-ejs, format extraction fails ("No video formats # found") even though all metadata fields are present. The media phase - # runs its own extraction and does need the solver. + # runs its own extraction and does need the solver. The flag also + # swallows a refused video's reason, hence the capturing logger. 'ignore_no_formats_error': True, } for attempt in range(_META_MAX_RETRIES): + reason_log = _ReasonLog() try: - with yt_dlp.YoutubeDL(ydl_opts) as ydl: - return ydl.extract_info(url, download=False), None + with yt_dlp.YoutubeDL({**ydl_opts, 'logger': reason_log}) as ydl: + info = ydl.extract_info(url, download=False) + if info and not info.get('formats'): + verdict = _playability_verdict(reason_log.reasons) + if _has_no_record(info): + category, detail = verdict or ( + "unknown", "no record of the video and no stated reason") + logger.warning("Scrape %s metadata: no record of the video — [%s] %s", + item_id, category, detail) + return None, _empty_fail(category, detail, + corroborated=category in _PERMANENT) + if verdict is not None: + info['_fyp_unplayable'] = verdict + return info, None except (yt_dlp.utils.DownloadError, ExtractorError) as e: category, detail = _classify_error(e) logger.warning("Scrape %s metadata attempt %d/%d failed: [%s] %s", @@ -394,6 +552,7 @@ class YouTubeScraper(BaseScraper): platform = "youtube" url_template = "https://www.youtube.com/watch?v={item_id}" slideshow_image_column = None + residential_ip_only = True def item_url(self, item_id: str) -> str: @@ -428,6 +587,18 @@ def fetch( item_id, duration, self.media_duration_cap()) return data_row + unplayable = info.get('_fyp_unplayable') + if unplayable is not None and unplayable[0] in _UNPLAYABLE_WITH_RECORD: + # YouTube keeps the video's record but will not play it here, and + # says why in terms of the video itself. The media leg could only + # repeat that, so the verdict is corroborated: the metadata row + # stands and the id leaves the queue instead of burning retries. + logger.info("Item '%s' is unplayable here — [%s] %s. Metadata only.", + item_id, *unplayable) + data_row.attrs['media_error_type'], data_row.attrs['media_error_detail'] = unplayable + data_row.attrs['verdict_corroborated'] = True + return data_row + ok, media_category, media_detail = _download_media( url, item_id, save_path, stream_to_bucket=stream_to_bucket, verbose=verbose) diff --git a/tests/unit/test_instagram_image_posts.py b/tests/unit/test_instagram_image_posts.py index af88b54b..1f13d4ea 100644 --- a/tests/unit/test_instagram_image_posts.py +++ b/tests/unit/test_instagram_image_posts.py @@ -385,11 +385,20 @@ def test_fetch_backfills_duration_from_downloaded_file(monkeypatch): -def test_health_check_reports_anonymous_no_cookies(): +def test_health_check_reports_the_session_cookies(monkeypatch): + """Posts hidden from logged-out viewers need the cookies, so their health is real again.""" + seen = {} + + def fake_health(platform, session_cookie="sessionid"): + seen.update(platform=platform, session_cookie=session_cookie) + return {"present": False, "status": "missing", "message": "Local dev: not logged in."} + + monkeypatch.setattr(instagram_dl.scraper_cookies, "cookie_health", fake_health) h = InstagramScraper().health_check() - assert h["status"] == "healthy" - assert "cookie" in h["message"].lower() - assert h["present"] is False + assert seen == {"platform": "instagram", "session_cookie": "sessionid"} + assert h["status"] == "missing" and h["present"] is False + assert h["message"].startswith("Local dev: not logged in.") + assert "hidden from logged-out viewers" in h["message"] diff --git a/tests/unit/test_instagram_login_fallback.py b/tests/unit/test_instagram_login_fallback.py new file mode 100644 index 00000000..fc9ca2db --- /dev/null +++ b/tests/unit/test_instagram_login_fallback.py @@ -0,0 +1,160 @@ +"""Instagram: anonymous first, the session cookies for posts hidden from logged-out viewers. + +Covers the 2026-09-21 field log. The scraper had run fully anonymously since +2026-07 (Instagram's authenticated web API 404'd then, and yt-dlp takes that +path whenever cookies are attached). By September 65 of 75 queued posts failed +every run: 50 with Instagram's own ruling "This content isn't available to +everyone: It can't be seen by certain audiences" (fell through to a retryable +``unknown``) and 15 with yt-dlp's "Instagram sent an empty media response" +(``rate_limited``). Measured that day, 13/13 sampled posts failed anonymously +and 13/13 extracted with the operator's Chrome cookies. + +Pinned here: the ruling classifies as a login wall; a login-gated anonymous +failure is retried once with the cookies — straight away, without burning the +anonymous retries — and the media leg follows the metadata leg's auth mode; +with the cookies attached an empty media response is throttling again. + +Run: pytest tests/unit/test_instagram_login_fallback.py +""" + +import pytest +from yt_dlp.utils import ExtractorError + +from fyp.scrape import instagram_dl +from fyp.scrape.instagram_dl import InstagramScraper + +AUDIENCE_RULING = ("ERROR: [Instagram] DUi1MEGieRX: This content isn't available to " + "everyone: It can't be seen by certain audiences.") +EMPTY_MEDIA = ("ERROR: [Instagram] DdXvQ_7HFSh: Instagram sent an empty media response. " + "Check if this post is accessible in your browser without being logged-in.") +COOKIES = {'cookiefile': '/tmp/instagram_cookies.txt'} +_POST = {'id': '1', 'description': 'caption', 'timestamp': 1750000000, + 'uploader_id': '99', 'channel': 'someuser', 'uploader': 'Some User', + 'view_count': 10, 'like_count': 2, 'comment_count': 1, 'duration': 12.0, + 'formats': [{'format_id': 'dash', 'url': 'https://x'}]} + + +class _FakeYDL: + """yt_dlp.YoutubeDL stand-in whose outcome depends on the auth mode. + + ``anon`` / ``authed`` are either an error message (raised as an + ExtractorError) or an info dict to return. + """ + + anon: object = None + authed: object = None + calls: list[str] = [] + + def __init__(self, opts): + self.opts = opts + + def __enter__(self): + return self + + def __exit__(self, *exc): + return False + + def extract_info(self, url, download=False): + authed = 'cookiefile' in self.opts + _FakeYDL.calls.append("cookies" if authed else "anonymous") + outcome = _FakeYDL.authed if authed else _FakeYDL.anon + if isinstance(outcome, str): + raise ExtractorError(outcome, expected=True) + return dict(outcome) + + +@pytest.fixture +def ig(monkeypatch): + _FakeYDL.calls = [] + monkeypatch.setattr(instagram_dl.yt_dlp, "YoutubeDL", _FakeYDL) + monkeypatch.setattr(instagram_dl, "sleep", lambda s: None) + monkeypatch.setattr(instagram_dl.scraper_cookies, "cookie_opts", lambda platform: dict(COOKIES)) + + def script(anon, authed=None): + _FakeYDL.anon, _FakeYDL.authed = anon, authed + return script + + +def test_the_audience_ruling_is_a_login_wall(): + category, _ = instagram_dl._classify_error(ExtractorError(AUDIENCE_RULING, expected=True)) + assert category == "login_required", "it fell through to 'unknown' until 2026-09-21" + assert category in instagram_dl._RETRYABLE + + +@pytest.mark.parametrize("anonymous_failure", [AUDIENCE_RULING, EMPTY_MEDIA]) +def test_a_login_gated_post_is_retried_once_with_the_cookies(ig, anonymous_failure): + ig(anon=anonymous_failure, authed=_POST) + info, fail = instagram_dl._extract_metadata("u", "x") + assert fail is None and info['_fyp_authenticated'] is True + assert _FakeYDL.calls == ["anonymous", "cookies"], \ + "one anonymous attempt — repeating it cannot get past the wall" + + +def test_a_public_post_never_touches_the_cookies(ig): + ig(anon=_POST, authed="must not be called") + info, fail = instagram_dl._extract_metadata("u", "x") + assert fail is None and "_fyp_authenticated" not in info + assert _FakeYDL.calls == ["anonymous"] + + +def test_other_failures_do_not_try_the_cookies(ig): + ig(anon="ERROR: [Instagram] x: This post is unavailable", authed=_POST) + _, fail = instagram_dl._extract_metadata("u", "x") + assert fail.attrs["error_type"] == "removed" + assert _FakeYDL.calls == ["anonymous"] + + +def test_without_cookies_the_anonymous_verdict_stands(ig, monkeypatch): + monkeypatch.setattr(instagram_dl.scraper_cookies, "cookie_opts", lambda platform: {}) + ig(anon=AUDIENCE_RULING) + _, fail = instagram_dl._extract_metadata("u", "x") + assert fail.attrs["error_type"] == "login_required" + assert _FakeYDL.calls == ["anonymous"] + + +def test_with_the_cookies_an_empty_media_response_is_throttling(ig): + ig(anon=EMPTY_MEDIA, authed=EMPTY_MEDIA) + _, fail = instagram_dl._extract_metadata("u", "x") + assert fail.attrs["error_type"] == "rate_limited" + assert _FakeYDL.calls == ["anonymous"] + ["cookies"] * instagram_dl._META_MAX_RETRIES, \ + "authenticated, it retries with backoff like any rate limit" + + +def test_hidden_even_from_the_session_ends_after_one_authenticated_attempt(ig): + ig(anon=AUDIENCE_RULING, authed=AUDIENCE_RULING) + _, fail = instagram_dl._extract_metadata("u", "x") + assert fail.attrs["error_type"] == "login_required" + assert _FakeYDL.calls == ["anonymous", "cookies"] + + +@pytest.mark.parametrize("anon,expect_authenticated", [(AUDIENCE_RULING, True), (_POST, False)]) +def test_the_media_leg_follows_the_metadata_legs_auth_mode(ig, monkeypatch, anon, + expect_authenticated): + ig(anon=anon, authed=_POST) + seen = {} + monkeypatch.setattr(instagram_dl, "_download_media", + lambda *a, **k: seen.update(k) or (True, None, "", 12.0)) + row = InstagramScraper().fetch("x", save_media=True, save_path="/tmp") + assert row.loc[0, "video_downloaded"] == True # noqa: E712 + assert seen["authenticated"] is expect_authenticated + + +def test_download_media_attaches_the_cookies_only_when_authenticated(monkeypatch, tmp_path): + opts_seen = [] + + class _RecordingYDL(_FakeYDL): + def __init__(self, opts): + opts_seen.append(opts) + super().__init__(opts) + + def download(self, urls): + raise ExtractorError("stop here", expected=True) + + monkeypatch.setattr(instagram_dl.yt_dlp, "YoutubeDL", _RecordingYDL) + monkeypatch.setattr(instagram_dl.scraper_cookies, "cookie_opts", lambda platform: dict(COOKIES)) + monkeypatch.setattr(instagram_dl, "sleep", lambda s: None) + monkeypatch.setattr(instagram_dl, "_cf", lambda: {"paths": {"temp": str(tmp_path)}}) + for authenticated in (False, True): + opts_seen.clear() + instagram_dl._download_media("u", "x", str(tmp_path), authenticated=authenticated) + assert all(("cookiefile" in o) is authenticated for o in opts_seen), opts_seen diff --git a/tests/unit/test_scrape_retry_budget.py b/tests/unit/test_scrape_retry_budget.py index d1166306..c3a744f5 100644 --- a/tests/unit/test_scrape_retry_budget.py +++ b/tests/unit/test_scrape_retry_budget.py @@ -94,18 +94,20 @@ def test_charge_zero_progress_accumulates_then_exhausts(): -def test_clear_zero_progress_resets_the_sidecar(): - """A progressing batch wipes all strikes (and tolerates a missing file).""" +def test_clear_zero_progress_drops_only_resolved_ids(): + """Progress clears the strikes of the items that left the queue — only those.""" with tempfile.TemporaryDirectory() as tmp: io = _fake_data_io(tmp) with patch.object(scrape_queues, "_data_io", return_value=io): - scrape_queues.clear_zero_progress("tiktok") # no file: no-op - scrape_queues.charge_zero_progress("tiktok", ["a"]) - scrape_queues.clear_zero_progress("tiktok") + scrape_queues.clear_zero_progress("tiktok", ["a"]) # no file: no-op assert not io.exists(filename=scrape_queues.strikes_filename("tiktok")) - assert scrape_queues.charge_zero_progress("tiktok", ["a"]) == [], \ - "after a clear the count restarts from zero" - print("PASS: clear_zero_progress resets the sidecar") + scrape_queues.charge_zero_progress("tiktok", ["a", "b"]) + scrape_queues.clear_zero_progress("tiktok", ["a"]) + sidecar = io.load_json(filename=scrape_queues.strikes_filename("tiktok")) + assert sidecar == {"b": 1}, f"b did not succeed, so it keeps its strike: {sidecar}" + assert scrape_queues.charge_zero_progress("tiktok", ["a", "b"]) == ["b"], \ + "a restarts from zero; b exhausts on its second zero-progress run" + print("PASS: clear_zero_progress drops only the resolved ids") @@ -141,6 +143,10 @@ def max_batch_size(): def health_check(): return None + @staticmethod + def unavailable_here(): + return None + def _all_transient_threads(**kwargs): """Every item fails transiently; no storm, no breaker, no memory stop.""" @@ -196,27 +202,51 @@ def test_cloud_zero_progress_runs_burn_the_stuck_tail(): -def test_cloud_progressing_batch_clears_strikes(): - """A batch that prunes something wipes earlier strikes.""" +def _mixed_batch(**kwargs): + """'good' scrapes, 'flaky' fails transiently; no storm, breaker or memory stop.""" + frame = pd.DataFrame({"item_id": ["good"]}) + for k in ("circuit_breaker_tripped", "permanent_storm_tripped", + "transient_storm_tripped", "memory_stop"): + frame.attrs[k] = False + return frame, [], ["flaky"] + + +def test_cloud_progressing_batch_clears_only_the_pruned_strikes(): + """A batch that prunes something clears the pruned ids' strikes, not the rest.""" with tempfile.TemporaryDirectory() as tmp: io = _fake_data_io(tmp) io.save_json(data=["good", "flaky"], filename=scrape_queues.queue_filename("tiktok")) + io.save_json(data={"good": 1, "flaky": 1}, + filename=scrape_queues.strikes_filename("tiktok")) recorded = [] - def mixed(**kwargs): - empty = pd.DataFrame({"item_id": ["good"]}) - for k in ("circuit_breaker_tripped", "permanent_storm_tripped", - "transient_storm_tripped", "memory_stop"): - empty.attrs[k] = False - return empty, [], ["flaky"] - - io.save_json(data={"flaky": 1}, filename=scrape_queues.strikes_filename("tiktok")) - _run_cloud_batch(io, mixed, recorded) - assert not io.exists(filename=scrape_queues.strikes_filename("tiktok")), \ - "queue progress must reset the strike counts" + _run_cloud_batch(io, _mixed_batch, recorded) + assert io.load_json(filename=scrape_queues.strikes_filename("tiktok")) == {"flaky": 1}, \ + "another item's success must not reset flaky's strike" assert recorded == [] assert io.load_json(filename=scrape_queues.queue_filename("tiktok")) == ["flaky"] - print("PASS: a progressing batch clears strikes") + print("PASS: a progressing batch clears only the pruned ids' strikes") + + + + +def test_cloud_trickling_queue_still_sheds_its_stuck_tail(): + """2026-09-21: a queue that drains a few items per run while a tail fails every + attempt must still give that tail up — successes elsewhere used to reset it.""" + with tempfile.TemporaryDirectory() as tmp: + io = _fake_data_io(tmp) + io.save_json(data=["flaky"], filename=scrape_queues.queue_filename("tiktok")) + recorded = [] + + _run_cloud_batch(io, _all_transient_threads, recorded) # stall: strike 1 + io.save_json(data=["good", "flaky"], filename=scrape_queues.queue_filename("tiktok")) + _run_cloud_batch(io, _mixed_batch, recorded) # progress elsewhere + assert io.load_json(filename=scrape_queues.strikes_filename("tiktok")) == {"flaky": 1} + _run_cloud_batch(io, _all_transient_threads, recorded) # stall: strike 2 + + assert len(recorded) == 1 and [r["item_id"] for r in recorded[0]] == ["flaky"] + assert io.load_json(filename=scrape_queues.queue_filename("tiktok")) == [] + print("PASS: a trickling queue still sheds its stuck tail") @@ -267,9 +297,10 @@ def test_no_video_formats_found_is_permanent(): if __name__ == "__main__": test_charge_zero_progress_accumulates_then_exhausts() - test_clear_zero_progress_resets_the_sidecar() + test_clear_zero_progress_drops_only_resolved_ids() test_cloud_zero_progress_runs_burn_the_stuck_tail() - test_cloud_progressing_batch_clears_strikes() + test_cloud_progressing_batch_clears_only_the_pruned_strikes() + test_cloud_trickling_queue_still_sheds_its_stuck_tail() test_cloud_storm_abort_does_not_charge() test_no_video_formats_found_is_permanent() print("All scrape retry-budget tests passed.") diff --git a/tests/unit/test_scrape_verdict_corroboration.py b/tests/unit/test_scrape_verdict_corroboration.py new file mode 100644 index 00000000..77308843 --- /dev/null +++ b/tests/unit/test_scrape_verdict_corroboration.py @@ -0,0 +1,349 @@ +"""Corroborated permanent verdicts, YouTube's refusal reasons, and the residential-IP guard. + +Covers the 2026-09-21 YouTube deadlock. Every id in the queue had already +failed, so retries had distilled it down to dead and blocked videos. The +metadata leg runs with ``ignore_no_formats_error``, which swallowed each +refusal and saved an empty placeholder row (no author, -1 plays, created +2000-01-01); the media leg then answered a bare "Video unavailable" → +permanent:removed; fifteen in a row tripped the permanent-storm guard; and an +aborted batch charges no retry budget — so 190 ids stayed queued run after run +with zero drained. A queue-wide sweep with the tv player client found every +one of them genuinely gone or blocked (137 with no record left at all, 23 kept +but region- or rights-blocked), so the guard was right about the session but +wrong about the items. + +Pinned here: + * the metadata leg captures the swallowed reason; a video with no record is a + failure (never a placeholder row), corroborated when its reason is itself + permanent, and never when the reason reads as throttling; + * a video YouTube keeps but will not play here (region, rights claim) is + scraped metadata-only without a pointless media attempt, and leaves the + queue; + * a corroborated verdict neither extends nor resets a storm run, and is + pruned even when the guard trips; uncorroborated verdicts keep the old + protection; + * Instagram and YouTube refuse to run on Cloud Run, and the enrichment + supervisor leaves their queues to the local install. + +Run: pytest tests/unit/test_scrape_verdict_corroboration.py +""" + +import sys +from pathlib import Path +from unittest.mock import patch + +sys.path.insert(0, str(Path(__file__).resolve().parents[2])) + +import pandas as pd +import pytest + +from fyp.scrape import scrape, youtube_dl +from fyp.scrape.youtube_dl import YouTubeScraper + +STORM_THRESHOLD = 5 + +# Real reasons and shapes, from the 2026-09-21 sweep of the live queue. +_GONE = {'id': 'x', 'title': 'youtube video #x', 'formats': []} +_KEPT = {'id': 'x', 'title': 'ICE-COLD FROM HAALAND', 'formats': [], + 'channel_id': 'UC1', 'uploader_id': '@c', 'channel': 'C', + 'view_count': 1265231, 'duration': 11, 'description': ''} +_PLAYABLE = {**_KEPT, 'formats': [{'format_id': '18', 'url': 'https://x'}]} + + +class _FakeYDL: + """Stands in for yt_dlp.YoutubeDL: replays one scripted extraction. + + ``script`` is ``(warnings, info)``: each warning goes to the logger the + metadata leg installs, exactly as yt-dlp reports a swallowed refusal. + """ + + script: tuple[list[str], dict] = ([], {}) + opts_seen: list[dict] = [] + + def __init__(self, opts): + self.opts = opts + _FakeYDL.opts_seen.append(opts) + + def __enter__(self): + return self + + def __exit__(self, *exc): + return False + + def extract_info(self, url, download=False): + warnings, info = _FakeYDL.script + for w in warnings: + self.opts['logger'].warning(w) + return dict(info) + + +@pytest.fixture +def fake_ydl(monkeypatch): + _FakeYDL.opts_seen = [] + monkeypatch.setattr(youtube_dl.yt_dlp, "YoutubeDL", _FakeYDL) + monkeypatch.setattr(youtube_dl.scraper_cookies, "cookie_opts", lambda platform: {}) + + def play(warnings, info): + _FakeYDL.script = (warnings, info) + return play + + +# --------------------------------------------------------------------------- # +# The metadata leg +# --------------------------------------------------------------------------- # + +def test_no_record_with_a_removal_reason_is_a_corroborated_failure(fake_ydl): + fake_ydl(["[youtube] This video has been removed by the uploader"], _GONE) + info, fail = youtube_dl._extract_metadata("u", "x") + assert info is None + assert fail.empty, "a video with no record must not become a row" + assert fail.attrs["error_type"] == "removed" + assert fail.attrs["verdict_corroborated"] is True + + +@pytest.mark.parametrize("reason,category", [ + ("[youtube] Sign in to confirm you're not a bot. Use --cookies", "bot_check"), + ("[youtube] Video unavailable. This content isn't available, try again later. " + "The current session has been rate-limited by YouTube for up to an hour.", "rate_limited"), +]) +def test_no_record_with_a_throttle_reason_is_never_corroborated(fake_ydl, reason, category): + """A walled session may return no record for a live video — that must stay transient.""" + fake_ydl([reason], _GONE) + _, fail = youtube_dl._extract_metadata("u", "x") + assert fail.attrs["error_type"] == category + assert "verdict_corroborated" not in fail.attrs + assert not YouTubeScraper().classify_error(category).startswith("permanent") + + +def test_no_record_and_no_reason_is_unknown(fake_ydl): + fake_ydl([], _GONE) + _, fail = youtube_dl._extract_metadata("u", "x") + assert fail.attrs["error_type"] == "unknown" + assert "verdict_corroborated" not in fail.attrs + + +def test_a_playable_video_ignores_stray_warnings(fake_ydl): + fake_ydl(["[youtube] Video unavailable"], _PLAYABLE) + info, fail = youtube_dl._extract_metadata("u", "x") + assert fail is None + assert "_fyp_unplayable" not in info + + +def test_the_metadata_leg_asks_the_tv_client_and_keeps_the_po_token_wiring(fake_ydl, monkeypatch): + pot = {'youtubepot-bgutilscript': {'server_home': ['/srv/pot']}} + monkeypatch.setattr(youtube_dl, "_pot_extractor_args", lambda: {'extractor_args': dict(pot)}) + fake_ydl([], _PLAYABLE) + youtube_dl._extract_metadata("u", "x") + args = _FakeYDL.opts_seen[-1]['extractor_args'] + assert args['youtube'] == {'player_client': ['default', 'tv']} + assert args['youtubepot-bgutilscript'] == pot['youtubepot-bgutilscript'] + assert _FakeYDL.opts_seen[-1]['ignore_no_formats_error'] is True + + +def test_reason_log_keeps_only_playability_reasons(): + log = youtube_dl._ReasonLog() + for msg in ("[youtube] [pot:bgutil:http] Error reaching GET http://127.0.0.1:4416/ping", + "No video formats found!", + "[youtube] No video formats found!", + "[youtube] Requested format is not available", + "[youtube] x: n challenge solving failed: Some formats may be missing", + "[generic] something else", + "[youtube] The uploader has not made this video available in your country"): + log.warning(msg) + assert log.reasons == ["The uploader has not made this video available in your country"] + + +def test_playability_verdict_lets_a_throttle_signal_win(): + """A permanent reason must never mask a session problem.""" + verdict = youtube_dl._playability_verdict( + ["This video has been removed by the uploader", "Sign in to confirm you're not a bot"]) + assert verdict[0] == "bot_check" + assert youtube_dl._playability_verdict([]) is None + assert youtube_dl._playability_verdict(["weird", "This video is private"])[0] == "private" + + +# --------------------------------------------------------------------------- # +# fetch(): kept but unplayable here +# --------------------------------------------------------------------------- # + +@pytest.mark.parametrize("reason,category", [ + ("[youtube] The uploader has not made this video available in your country", "geo_blocked"), + ("[youtube] It was blocked due to the claimed content by UFC.", "blocked"), +]) +def test_kept_but_blocked_here_is_metadata_only_and_skips_the_media_leg( + fake_ydl, monkeypatch, reason, category): + fake_ydl([reason], _KEPT) + monkeypatch.setattr(youtube_dl, "_download_media", + lambda *a, **k: pytest.fail("the media leg could only repeat the refusal")) + row = YouTubeScraper().fetch("x", save_media=True, save_path="/tmp") + assert not row.empty and row.loc[0, "author_name_raw"] == "C" + assert row.loc[0, "video_downloaded"] == False # noqa: E712 + assert row.attrs["media_error_type"] == category + assert row.attrs["verdict_corroborated"] is True + + +def test_kept_with_a_bare_unavailable_still_takes_the_distrusted_media_leg(fake_ydl, monkeypatch): + """The 2026-09-18 signature: record intact, bare reason — never corroborated.""" + fake_ydl(["[youtube] This video is not available"], _KEPT) + calls = [] + monkeypatch.setattr(youtube_dl, "_download_media", + lambda *a, **k: calls.append(1) or (False, "removed", "Video unavailable")) + row = YouTubeScraper().fetch("x", save_media=True, save_path="/tmp") + assert calls, "the media leg must still be tried" + assert row.attrs["media_error_type"] == "removed" + assert "verdict_corroborated" not in row.attrs + + +# --------------------------------------------------------------------------- # +# The orchestrator +# --------------------------------------------------------------------------- # + +def _failure(category: str, corroborated: bool = False) -> pd.DataFrame: + return youtube_dl._empty_fail(category, "simulated", corroborated=corroborated) + + +def _metadata_row(item_id: str, media_error: str | None = None, + corroborated: bool = False) -> pd.DataFrame: + row = pd.DataFrame([{ + "item_id": item_id, "desc": "x", "create_time_raw": pd.Timestamp("2026-01-01"), + "duration_raw": 30, "author_id": "a", "yt_author_handle": "@a", + "author_name_raw": "A", "play_count_raw": 1, "yt_like_count": 0, + "yt_comment_count": 0, "yt_channel_follower_count": 0, + "yt_categories": "", "video_downloaded": media_error is None, + }]) + if media_error is not None: + row.attrs["media_error_type"] = media_error + row.attrs["media_error_detail"] = "simulated" + if corroborated: + row.attrs["verdict_corroborated"] = True + return row + + +def _run_batch(ids, fake_dl, max_workers=1): + with patch.object(scrape, "download_single_video", side_effect=fake_dl), \ + patch.object(scrape, "_permanent_storm_threshold", return_value=STORM_THRESHOLD), \ + patch.object(scrape.scrape_versioning, "ensure_active_version_registered", + lambda: None), \ + patch.object(YouTubeScraper, "inter_request_delay", return_value=0.0): + return scrape.download_video_threads( + interesting_videos=ids, max_workers=max_workers, + dry_run=True, platform="youtube") + + +def test_a_queue_of_dead_videos_drains_instead_of_storming(): + """The 2026-09-21 queue: nothing but corroborated removals, far past the threshold.""" + ids = [f"v{i}" for i in range(STORM_THRESHOLD * 4)] + + results, perm, trans = _run_batch(ids, lambda video_id=None, **k: _failure("removed", True)) + + assert results.attrs["permanent_storm_tripped"] is False + assert set(perm) == set(ids), "corroborated removals are pruned as permanent" + assert trans == [] + + +def test_corroborated_verdicts_neither_extend_nor_reset_a_storm_run(): + """3 bare removals, 5 corroborated, 2 bare: the bare run reaches 5 on the last call.""" + # By call order, not id: the pool's threads need not take the single + # throttle slot in submission order (same trick as the storm-guard tests). + plan = iter([False] * 3 + [True] * 5 + [False] * 2) + ids = [f"v{i}" for i in range(10)] + proven = set() + + def fake_dl(video_id=None, **kwargs): + corroborated = next(plan) + if corroborated: + proven.add(video_id) + return _failure("removed", corroborated=corroborated) + + results, perm, trans = _run_batch(ids, fake_dl) + + assert results.attrs["permanent_storm_tripped"] is True + assert len(proven) == 5 + assert set(perm) == proven, \ + "corroborated ids are pruned even though the guard tripped on their category" + assert set(trans) == set(ids) - proven, \ + "the uncorroborated storm ids are demoted and stay queued, as before" + + +def test_a_corroborated_media_verdict_prunes_with_its_row(): + ids = ["kept_blocked", "media_flaky"] + + def fake_dl(video_id=None, **kwargs): + if video_id == "kept_blocked": + return _metadata_row(video_id, media_error="geo_blocked", corroborated=True) + return _metadata_row(video_id, media_error="removed") + + results, perm, trans = _run_batch(ids, fake_dl) + + assert set(results["item_id"]) == set(ids), "both metadata rows are saved" + assert results.attrs["media_retry_ids"] == ["media_flaky"] + assert trans == ["media_flaky"], "only the distrusted media failure stays queued" + assert perm == [] + + +# --------------------------------------------------------------------------- # +# Residential IP only +# --------------------------------------------------------------------------- # + +def test_instagram_and_youtube_refuse_cloud_run_and_tiktok_does_not(monkeypatch): + from fyp.scrape.platform_scraper import get_scraper + + monkeypatch.delenv("K_SERVICE", raising=False) + assert all(get_scraper(p).unavailable_here() is None + for p in ("tiktok", "instagram", "youtube")) + + monkeypatch.setenv("K_SERVICE", "fyp-data-hub") + assert get_scraper("tiktok").unavailable_here() is None + for platform in ("instagram", "youtube"): + why = get_scraper(platform).unavailable_here() + assert why and "residential IP" in why and platform in why + + +class _Reporter: + def __init__(self): + self.lines = [] + + def log(self, msg): + self.lines.append(str(msg)) + + def update_progress(self, *a, **k): + pass + + def emit_data(self, payload): + pass + + def check_cancelled(self): + return False + + +def test_the_cloud_run_worker_leaves_a_residential_queue_untouched(monkeypatch): + import fyp.scrape as fyp_scrape + from fyp.scrape import scrape_queues + from web_interface.run_queue_scraper import run_queue_scraper + + monkeypatch.setenv("K_SERVICE", "fyp-data-hub") + monkeypatch.setattr(scrape_queues, "load_scrape_queue", + lambda platform: pytest.fail("the queue must not be read")) + monkeypatch.setattr(fyp_scrape, "download_video_threads", + lambda **k: pytest.fail("nothing may be scraped")) + reporter = _Reporter() + + assert run_queue_scraper(reporter, {"platform": "youtube"}) is None + assert any("residential IP" in line for line in reporter.lines), reporter.lines + + +def test_the_cloud_run_supervisor_leaves_a_residential_queue_to_the_local_install(monkeypatch): + from fyp.scrape import scrape_queues + from web_interface import run_enrichment_supervisor as sup + + monkeypatch.setenv("K_SERVICE", "fyp-data-hub") + monkeypatch.setattr(scrape_queues, "queue_lengths", lambda: {"youtube": 190}) + monkeypatch.setattr(sup, "_scrape_lane_busy", lambda platform: False) + monkeypatch.setattr(sup, "_scraper_blocked", lambda platform: None) + monkeypatch.setattr(sup, "_queue_stalled", + lambda *a, **k: pytest.fail("no stall may be charged for a skipped queue")) + monkeypatch.setattr(sup, "_start", lambda *a, **k: pytest.fail("no worker may be started")) + reporter = _Reporter() + + assert sup._drain(reporter, {"c1": {"platform": "youtube"}}) is None + assert any("residential IP" in line for line in reporter.lines), reporter.lines diff --git a/tests/unit/test_youtube_scraper.py b/tests/unit/test_youtube_scraper.py index f7b5295f..d8869695 100644 --- a/tests/unit/test_youtube_scraper.py +++ b/tests/unit/test_youtube_scraper.py @@ -128,6 +128,31 @@ def test_classify_error_truth_table(): "HTTP Error 429: Too Many Requests": ("rate_limited", "transient"), "Connection timed out": ("network", "transient"), "brand new failure mode": ("unknown", "transient"), + # Every reason the tv player client gave across the live queue on + # 2026-09-21 (see test_scrape_verdict_corroboration.py). + "This video is unavailable": ("removed", "permanent"), + "This video is private": ("private", "permanent"), + "This video is not available": ("removed", "permanent"), + "This video is no longer available because the YouTube account associated " + "with this video has been terminated.": ("removed", "permanent"), + "This video is no longer available because the uploader has closed their " + "YouTube account.": ("removed", "permanent"), + "This video is no longer available due to a privacy claim by a third party.": + ("removed", "permanent"), + "This video has been removed for violating YouTube's Terms of Service": + ("removed", "permanent"), + "This video has been removed for violating YouTube's policy on violent or " + "graphic content": ("removed", "permanent"), + "It was removed following a copyright removal request by NBC Universal.": + ("removed", "permanent"), + "It was blocked due to the claimed content by Paramount Global (PMN).": + ("blocked", "permanent"), + "This video contains content from UFC, who has blocked it on copyright grounds.": + ("blocked", "permanent"), + "This video contains content from UFC, who has blocked it in your country on " + "copyright grounds.": ("geo_blocked", "permanent"), + "Video unavailable. YouTube is requiring a captcha challenge before playback": + ("bot_check", "transient"), } for msg, (category, bucket) in cases.items(): got_cat, _ = _classify_error(Exception(msg)) diff --git a/web_interface/run_enrichment_supervisor.py b/web_interface/run_enrichment_supervisor.py index 70853c83..985d802e 100644 --- a/web_interface/run_enrichment_supervisor.py +++ b/web_interface/run_enrichment_supervisor.py @@ -246,6 +246,20 @@ def _start(name: str, task_args: dict | None = None) -> tuple[bool, str]: started_by="enrichment_supervisor") +def _unavailable_here(platform: str) -> str | None: + """Why this platform's scraper must not run here, if it must not. + + See :meth:`fyp.scrape.platform_scraper.BaseScraper.unavailable_here` — + Instagram and YouTube on Cloud Run. Never raises: an unknown platform is + left to the checks that follow. + """ + try: + from fyp.scrape.platform_scraper import get_scraper + return get_scraper(platform).unavailable_here() + except Exception: + return None + + def _scraper_blocked(platform: str) -> str | None: """A storm/circuit-breaker abort the operator has to clear, if any. @@ -470,6 +484,14 @@ def _drain(reporter, plans: dict) -> dict | None: continue if _scrape_lane_busy(platform): continue # already being drained + unavailable = _unavailable_here(platform) + if unavailable: + # Left for a local install on a residential IP. Skipping here — + # not starting a worker that would refuse — keeps the refusal + # from re-ticking the supervisor into a dispatch loop, and it + # charges no stall, so the plan waits rather than parks. + reporter.log(f"Leaving the '{platform}' queue ({count} item(s)): {unavailable}.") + continue tripped = _scraper_blocked(platform) if tripped: reporter.log(f"Scraper for '{platform}' is held off: {tripped}. " diff --git a/web_interface/run_queue_scraper.py b/web_interface/run_queue_scraper.py index 14d5bf48..b510e968 100644 --- a/web_interface/run_queue_scraper.py +++ b/web_interface/run_queue_scraper.py @@ -81,6 +81,13 @@ def run_queue_scraper(reporter: TaskStatusReporter, task_args: dict | None = Non task_args = {} platform: str = str(task_args.get("platform") or "") or scrape_queues.default_platform() + # Instagram and YouTube only scrape from a residential IP: on Cloud Run a + # run would burn the queue against the wall and trip storm guards that then + # hold off the local install (BaseScraper.residential_ip_only). + unavailable = get_scraper(platform).unavailable_here() + if unavailable: + reporter.log(f"Not scraping — {unavailable}. The queue is untouched.") + return None batch_size: int = min(int(task_args.get("batch_size", 500)), MAX_BATCH_SIZE) # A platform may cap the batch below that (YouTube: one signed-in session). platform_cap = get_scraper(platform).max_batch_size() @@ -236,17 +243,19 @@ def _on_threads_change(n: int) -> None: # supervisor restarts the worker, gets the same verdict, and its no-drain # guard then parks every armed plan. A misclassified permanent failure # (e.g. yt-dlp's "No video formats found") stays "transient" forever. - # Items that produce MAX_ZERO_PROGRESS_STRIKES such batches in a row are - # given up on: recorded in the failed-scrapes ledger (consolidation then - # marks them scrape_fail like any other permanent failure) and pruned so - # the queue drains. Storm / circuit-breaker / memory aborts never charge - # strikes — those verdicts implicate the scraper, not the items. + # Items charged in MAX_ZERO_PROGRESS_STRIKES such batches without + # succeeding in between are given up on: recorded in the failed-scrapes + # ledger (consolidation then marks them scrape_fail like any other + # permanent failure) and pruned so the queue drains. A batch that makes + # progress clears only the strikes of the items it pruned. Storm / + # circuit-breaker / memory aborts never charge strikes — those verdicts + # implicate the scraper, not the items. batch_aborted = any(results_df.attrs.get(k) for k in ( 'circuit_breaker_tripped', 'permanent_storm_tripped', 'transient_storm_tripped', 'memory_stop')) given_up: list[str] = [] if pruned_this_batch > 0: - scrape_queues.clear_zero_progress(platform) + scrape_queues.clear_zero_progress(platform, items_to_remove) elif transient_failed and not batch_aborted: exhausted = scrape_queues.charge_zero_progress(platform, transient_failed) if exhausted: @@ -270,12 +279,12 @@ def _on_threads_change(n: int) -> None: # Metadata-only rows (media failed, any category) stay queued for a media # retry, but not forever: after MAX_MEDIA_RETRY_STRIKES healthy runs the # item is pruned and its metadata-only row stands (no ledger entry — the - # metadata did scrape). An aborted batch never charges. + # metadata did scrape). An aborted batch never charges, but every id that + # left the queue drops its strikes either way. media_retry = list(results_df.attrs.get('media_retry_ids') or []) - if media_retry and not batch_aborted: - retry_set = set(media_retry) - got_media = [v for v in good_ids if v not in retry_set] - media_exhausted = scrape_queues.charge_media_retry(platform, media_retry, got_media) + if media_retry or items_to_remove: + media_exhausted = scrape_queues.charge_media_retry( + platform, [] if batch_aborted else media_retry, items_to_remove) if media_exhausted: gave_up, queue_remaining = scrape_queues.prune_scrape_queue( platform, set(media_exhausted))