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))