# ============================================================================== # Author: Simon-Pierre Boucher # File: creaka/ingest.py # Desc: Pipeline d'ingestion — découverte puis enrichissement par plateforme, # synchronisation de la base (ajouts / fusions / mises à jour) # ============================================================================== """Pipeline Créa-Ka (calqué sur louka/ingest.py). Ordre d'un passage complet : 1. connecteurs DÉCOUVERTE (listes médias, hashtags, agences…) → nouvelles fiches ; 2. connecteurs ENRICHISSEMENT par plateforme (instagram, tiktok, youtube, twitch) → métriques fraîches + liens croisés, créateur par créateur ; 3. link-in-bio EN DERNIER : il consomme les `link_in_bio_url` découvertes à l'étape 2 et révèle d'un coup tous les comptes restants (§8-9). Logging structuré par source : créateurs, comptes, confiance moyenne, durée, erreurs (§16) ; alerte si une source ne retourne plus rien (§18 monitoring). """ from __future__ import annotations import sys import time import traceback from . import db from .connectors import CONNECTORS from .connectors.base import SkipSource from .dedup import dedupe # ordre d'exécution des connecteurs d'enrichissement — link-in-bio vers la # fin (il consomme les link_in_bio_url des étapes précédentes) ; bio-liens en # tout dernier (il relit les bios enrichies, sans aucune requête) ENRICH_ORDER = ["instagram-apify", "snowball-ig", "tiktok-apify", "x-apify", "facebook-apify", "threads-apify", "snapchat-apify", "youtube-apify", "twitch-apify", "kick-apify", "fansly-apify", "patreon-apify", "discord-apify", "instagram-profil", "tiktok-profil", "youtube", "twitch", "x-profil", "balados-rss", "podcastindex", "link-in-bio", "bio-liens", "region-bio"] def _finalize_batch(creators: list) -> tuple[list, int]: finalized, dropped = [], 0 for cr in creators: try: finalized.append(cr.finalize()) except Exception: # une fiche malformée ne bloque pas la source dropped += 1 return finalized, dropped def load_active(con) -> list: """Fiches actives (jamais les opted_out : exclues des ré-agrégations, §15).""" rows = con.execute( "SELECT * FROM creators WHERE status='active' ORDER BY total_reach DESC" ).fetchall() return [db.creator_from_row(r) for r in rows] def run(sources: list[str] | None = None) -> list[dict]: """Exécute l'ingestion pour toutes les sources (ou celles demandées).""" con = db.connect() results: list[dict] = [] discovery = [sid for sid, cls in CONNECTORS.items() if cls.kind == "discovery"] enrichment = [sid for sid in ENRICH_ORDER if sid in CONNECTORS] enrichment += [sid for sid, cls in CONNECTORS.items() if cls.kind == "enrichment" and sid not in enrichment] targets = sources or (discovery + enrichment) for sid in targets: cls = CONNECTORS.get(sid) if cls is None: print(f"[crea-ka] connecteur inconnu : {sid}", file=sys.stderr) continue t0 = time.time() print(f"[crea-ka] sync {sid} ({cls.kind}) ...") try: connector = cls() if cls.kind == "discovery": batch = connector.fetch() else: batch = connector.enrich(load_active(con)) finalized, dropped = _finalize_batch(batch) finalized = dedupe(finalized) # dédup intra-lot avant la base stats = db.sync_source(con, sid, finalized, started=t0, errors=getattr(connector, "errors", 0)) if dropped: stats["dropped"] = dropped if stats.get("alert") and cls.kind == "discovery": print(f"[crea-ka] ⚠ ALERTE {sid} : {stats['alert']}") print(f"[crea-ka] {stats}") results.append(stats) except SkipSource as exc: # passage sauté volontairement (ex. clés API absentes) : log clair, # PAS d'alerte de blocage ni de run compté à 0 dans sync_log print(f"[crea-ka] ↷ {sid} sauté : {exc}") results.append({"source": sid, "skipped": str(exc)}) except Exception as exc: # robustesse : une source ne bloque pas les autres db.log_failure(con, sid, f"{exc}") traceback.print_exc() results.append({"source": sid, "error": str(exc)}) if not sources: # passage COMPLET seulement : archivage des disparus (§13) db.archive_missing(con) return results def watch(interval_seconds: int = 86400) -> None: """Synchronisation en boucle (cadence §14 — défaut : quotidienne).""" while True: try: run() except Exception: traceback.print_exc() print(f"[crea-ka] prochain passage dans {interval_seconds}s") time.sleep(interval_seconds)