"""Per-origin freshness watchdog — alert gdy aktywne źródło przestało dawać nowe treści. Globalny monitor źródeł (ingest_runs per `Source`) tego NIE łapie, bo wszystkie tube scrapery dzielą jeden `Source` = "tube-scraper" — pojedynczy origin może zamarznąć (np. freshporno: scraper browsował z roota `/`, który KVS rotuje → cold-session dostawała stary zestaw → new=0/skipped=N przez 2 dni), a zagregowany run nadal raportuje success. Sygnał per-origin: `max(created_at)` na źródłach playbacku. Pokrywamy TRZY klasy, każda z własnym progiem: - **browse** (`ALL_BROWSE_SCRAPERS`) — crawlowane codziennie z listingu, próg 48h. - **search** (`ALL_DIRECT_SCRAPERS`) — performer-driven, nierówna kadencja (~30d refresh per performer), próg wyższy (domyślnie 7d). - **movies** (`_MOVIE_CONNECTORS`) — dodane 2026-08-03. Wcześniej filmy NIE BYŁY pokryte w ogóle: streamporn.vip stał 23 dni (czytał listing nieposortowany po dacie) i żaden automat nie miał jak tego zauważyć. Tag obecny w kilku listach liczymy wg najostrzejszego progu. **Eskalacja zamiast jednego cichego issue.** Do 2026-08-03 zdarzenie szło zawsze jako `warning` ze stabilnym fingerprintem per origin. Fingerprint jest tam po to, żeby nie tworzyć nowego issue co 6h — ale skutkiem ubocznym było JEDNO issue przy pierwszym wystąpieniu, potem tylko rosnący licznik. Sentry powiadamia o nowych i regresjach, nie o kolejnych wystąpieniach otwartego issue, a reguły alertów zwykle celują w `error`, nie `warning`. Efekt: vjav stał 12 dni, watchdog go poprawnie wypisywał co cykl, i nikt się nie dowiedział. Teraz fingerprint zawiera KUBEŁEK WIEKU, więc przekroczenie każdego kolejnego progu zakłada nowe issue (= nowe powiadomienie), a w obrębie kubełka dalej nie ma spamu. Od 7 dni ciszy poziom idzie na `error`, żeby wpaść w standardowe reguły alertów. Patrz [[reference_kvs_root_rotates_use_latest_updates]] dla klasy błędu, którą to łapie. """ from __future__ import annotations import logging from datetime import UTC, datetime from typing import Any from sqlalchemy import text from app.db import session_scope log = logging.getLogger(__name__) #: Kubełki wieku ciszy (godziny, malejąco) → etykieta + poziom Sentry. Wejście w #: kolejny kubełek zmienia fingerprint, czyli zakłada NOWE issue i wysyła alert. _BUCKETS: tuple[tuple[int, str, str], ...] = ( (720, "30d+", "error"), (168, "7d+", "error"), (48, "2d+", "warning"), (0, "swiezo", "warning"), ) def _bucket(age_h: float) -> tuple[str, str]: for floor, label, level in _BUCKETS: if age_h >= floor: return label, level return "swiezo", "warning" def run_ingest_freshness_watchdog( *, max_age_hours: int = 48, search_max_age_hours: int = 168, movie_max_age_hours: int = 72, min_history: int = 100, to_sentry: bool = True, to_slack: bool = False, ) -> dict[str, Any]: """Sprawdź każde aktywne źródło: czy dostało nową treść < próg dla swojej klasy. `min_history` odsiewa świeżo dodane źródła bez ustalonej kadencji (za mało pozycji, by wiedzieć, czy cisza to anomalia). `to_sentry` / `to_slack` sterują kanałami — job co 6h woła Sentry, a osobny job dobowy woła Slacka, żeby nie wysyłać tej samej listy cztery razy dziennie. """ from app.connectors import get_movie_connectors # noqa: PLC0415 from app.connectors.direct_scrapers import ( # noqa: PLC0415 ALL_BROWSE_SCRAPERS, ALL_DIRECT_SCRAPERS, ) browse_tags = {cls.sitetag for cls in ALL_BROWSE_SCRAPERS} search_tags = {cls.sitetag for cls in ALL_DIRECT_SCRAPERS} - browse_tags # (origin, tabela, próg_h, klasa) checks: list[tuple[str, str, int, str]] = ( [(f"tube:{t}", "playback_sources", max_age_hours, "browse") for t in sorted(browse_tags)] + [(f"tube:{t}", "playback_sources", search_max_age_hours, "search") for t in sorted(search_tags)] ) try: for name, _cls in get_movie_connectors(): checks.append((name, "movie_playback_sources", movie_max_age_hours, "movies")) except Exception as e: # pragma: no cover - rejestr filmów niedostępny log.warning("ingest-watchdog: nie udało się pobrać konektorów filmowych: %s", e) now = datetime.now(UTC) stale: list[dict[str, Any]] = [] with session_scope() as s: for origin, table, threshold, kind in checks: # Filmy trzymają origin jako `:`, sceny jako dokładne # `tube:` — stąd LIKE dla filmów, równość dla scen. if table == "movie_playback_sources": sql = ( f"SELECT max(created_at) AS newest, count(*) AS total FROM {table} " "WHERE origin = :o OR origin LIKE :p" ) params = {"o": origin, "p": f"{origin}:%"} else: sql = ( f"SELECT max(created_at) AS newest, count(*) AS total FROM {table} " "WHERE origin = :o" ) params = {"o": origin} row = s.execute(text(sql), params).one() newest, total = row.newest, row.total if total < min_history or newest is None: continue age_h = (now - newest).total_seconds() / 3600.0 if age_h >= threshold: label, level = _bucket(age_h) stale.append( { "origin": origin, "age_hours": round(age_h, 1), "total": total, "kind": kind, "max_age_hours": threshold, "bucket": label, "level": level, } ) if not stale: log.info("ingest-watchdog: wszystkie %d źródła świeże", len(checks)) return {"checked": len(checks), "stale": []} stale.sort(key=lambda x: -x["age_hours"]) log.warning( "ingest-watchdog: %d/%d źródeł bez nowych treści: %s", len(stale), len(checks), ", ".join(f"{x['origin']}({x['age_hours']}h,{x['kind']})" for x in stale), ) if to_sentry: try: import sentry_sdk for x in stale: with sentry_sdk.push_scope() as scope: scope.level = x["level"] scope.set_tag("ingest_origin", x["origin"]) scope.set_tag("ingest_scraper_kind", x["kind"]) scope.set_tag("ingest_stale_bucket", x["bucket"]) scope.set_extra("age_hours", x["age_hours"]) scope.set_extra("total_items", x["total"]) scope.set_extra("max_age_hours", x["max_age_hours"]) # Kubełek W fingerprincie: nowe issue (= nowy alert) przy każdym # kolejnym progu, ale bez spamu co 6h w obrębie tego samego progu. scope.fingerprint = ["ingest-stale-origin", x["origin"], x["bucket"]] sentry_sdk.capture_message( f"ingest-watchdog: {x['origin']} ({x['kind']}) bez nowych treści " f"({x['age_hours']:.0f}h, próg {x['max_age_hours']}h)" ) except Exception: # pragma: no cover - Sentry off / brak DSN log.exception("ingest-watchdog: Sentry capture failed") if to_slack: from app.notify.slack import send_slack # noqa: PLC0415 lines = [f"*Zamrożone źródła ingestu* ({len(stale)}/{len(checks)})"] for x in stale: days = x["age_hours"] / 24.0 lines.append( f"• `{x['origin']}` ({x['kind']}) — {days:.1f} dnia bez nowych treści " f"(próg {x['max_age_hours']}h, {x['total']} pozycji w bazie)" ) send_slack("\n".join(lines)) return {"checked": len(checks), "stale": stale}