SPB Git forge

spb/crea-ka

Public

Créa·Ka — annuaire public cross-plateforme des créateurs de contenu québécois (crea-ka.com)

52commits 1branches 0releases
11.3 MBsize
maindefault branch
19 days agolast push
Python 73.6% HTML 13.2% TypeScript 6% JavaScript 4.5% CSS 1.7% Dockerfile 0.6%
5.0 KB · 116 lines python
Raw Blame History
1# ==============================================================================2# Author: Simon-Pierre Boucher <contact@spboucher.ai>3# File:   creaka/ingest.py4# Desc:   Pipeline d'ingestion — découverte puis enrichissement par plateforme,5#         synchronisation de la base (ajouts / fusions / mises à jour)6# ==============================================================================7"""Pipeline Créa-Ka (calqué sur louka/ingest.py).89Ordre d'un passage complet :101. connecteurs DÉCOUVERTE (listes médias, hashtags, agences…) → nouvelles fiches ;112. connecteurs ENRICHISSEMENT par plateforme (instagram, tiktok, youtube,12   twitch) → métriques fraîches + liens croisés, créateur par créateur ;133. link-in-bio EN DERNIER : il consomme les `link_in_bio_url` découvertes à14   l'étape 2 et révèle d'un coup tous les comptes restants (§8-9).1516Logging structuré par source : créateurs, comptes, confiance moyenne, durée,17erreurs (§16) ; alerte si une source ne retourne plus rien (§18 monitoring).18"""19from __future__ import annotations2021import sys22import time23import traceback2425from . import db26from .connectors import CONNECTORS27from .connectors.base import SkipSource28from .dedup import dedupe2930# ordre d'exécution des connecteurs d'enrichissement — link-in-bio vers la31# fin (il consomme les link_in_bio_url des étapes précédentes) ; bio-liens en32# tout dernier (il relit les bios enrichies, sans aucune requête)33ENRICH_ORDER = ["instagram-apify", "snowball-ig", "tiktok-apify", "x-apify",34                "facebook-apify", "threads-apify", "snapchat-apify",35                "youtube-apify", "twitch-apify", "kick-apify",36                "fansly-apify", "patreon-apify", "discord-apify",37                "instagram-profil", "tiktok-profil", "youtube", "twitch",38                "x-profil", "balados-rss", "podcastindex",39                "link-in-bio", "bio-liens", "region-bio"]404142def _finalize_batch(creators: list) -> tuple[list, int]:43    finalized, dropped = [], 044    for cr in creators:45        try:46            finalized.append(cr.finalize())47        except Exception:  # une fiche malformée ne bloque pas la source48            dropped += 149    return finalized, dropped505152def load_active(con) -> list:53    """Fiches actives (jamais les opted_out : exclues des ré-agrégations, §15)."""54    rows = con.execute(55        "SELECT * FROM creators WHERE status='active' ORDER BY total_reach DESC"56    ).fetchall()57    return [db.creator_from_row(r) for r in rows]585960def run(sources: list[str] | None = None) -> list[dict]:61    """Exécute l'ingestion pour toutes les sources (ou celles demandées)."""62    con = db.connect()63    results: list[dict] = []64    discovery = [sid for sid, cls in CONNECTORS.items() if cls.kind == "discovery"]65    enrichment = [sid for sid in ENRICH_ORDER if sid in CONNECTORS]66    enrichment += [sid for sid, cls in CONNECTORS.items()67                   if cls.kind == "enrichment" and sid not in enrichment]68    targets = sources or (discovery + enrichment)6970    for sid in targets:71        cls = CONNECTORS.get(sid)72        if cls is None:73            print(f"[crea-ka] connecteur inconnu : {sid}", file=sys.stderr)74            continue75        t0 = time.time()76        print(f"[crea-ka] sync {sid} ({cls.kind}) ...")77        try:78            connector = cls()79            if cls.kind == "discovery":80                batch = connector.fetch()81            else:82                batch = connector.enrich(load_active(con))83            finalized, dropped = _finalize_batch(batch)84            finalized = dedupe(finalized)  # dédup intra-lot avant la base85            stats = db.sync_source(con, sid, finalized, started=t0,86                                   errors=getattr(connector, "errors", 0))87            if dropped:88                stats["dropped"] = dropped89            if stats.get("alert") and cls.kind == "discovery":90                print(f"[crea-ka]   ⚠ ALERTE {sid} : {stats['alert']}")91            print(f"[crea-ka]   {stats}")92            results.append(stats)93        except SkipSource as exc:94            # passage sauté volontairement (ex. clés API absentes) : log clair,95            # PAS d'alerte de blocage ni de run compté à 0 dans sync_log96            print(f"[crea-ka]   ↷ {sid} sauté : {exc}")97            results.append({"source": sid, "skipped": str(exc)})98        except Exception as exc:  # robustesse : une source ne bloque pas les autres99            db.log_failure(con, sid, f"{exc}")100            traceback.print_exc()101            results.append({"source": sid, "error": str(exc)})102    if not sources:  # passage COMPLET seulement : archivage des disparus (§13)103        db.archive_missing(con)104    return results105106107def watch(interval_seconds: int = 86400) -> None:108    """Synchronisation en boucle (cadence §14 — défaut : quotidienne)."""109    while True:110        try:111            run()112        except Exception:113            traceback.print_exc()114        print(f"[crea-ka] prochain passage dans {interval_seconds}s")115        time.sleep(interval_seconds)116