# ============================================================================== # Author: Simon-Pierre Boucher # File: creaka/connectors/apify_social.py # Desc: Connecteurs Apify MAISON — acteurs KA (~/Desktop/ka-apify-actors, # compte gorgeous_thistle) exécutés sur la plateforme Apify avec proxy # RÉSIDENTIEL : enrichissement ULTRA DÉTAILLÉ (abonnés, bio, badges, # contenu récent avec métriques, engagement) pour Instagram, TikTok, # X, Facebook, Threads et Snapchat + découverte Instagram (topsearch). # Un passage = UN run d'acteur par plateforme (lot de handles), au lieu # d'une requête Scrapfly par profil — moins cher, plus fiable (IP # résidentielles + empreinte TLS Chrome côté acteur). # ============================================================================== """Enrichissement multi-plateforme via les acteurs Apify KA. Chaque connecteur d'enrichissement : 1. sélectionne les comptes de SA plateforme dont `last_checked` est plus vieux que `revisit_days` (rotation douce = maîtrise du coût proxy) ; 2. lance l'acteur `gorgeous_thistle/ka-` avec le lot de handles (proxy résidentiel Apify forcé) et attend le dataset ; 3. applique aux fiches : abonnés/badge (§13), métriques étendues + CONTENU récent (posts/vidéos/tweets avec likes, commentaires, vues) dans `account.metrics`, avatar/bio manquants, liens de bio → `cross_link` (§12.1). Nécessite APIFY_TOKEN (.env). Sans clé : SkipSource propre (§16/§18). Uniquement des profils PUBLICS (§15) ; un profil privé/introuvable est ignoré. """ from __future__ import annotations import os import random import time from concurrent.futures import ThreadPoolExecutor, as_completed from datetime import datetime, timedelta, timezone import requests from ..identity import account, merge_accounts from ..normalize import parse_count, platform_from_url from ..schema import Creator, now_iso from .base import BaseConnector, SkipSource from .linkinbio import is_supported as is_linkinbio APIFY_API = "https://api.apify.com/v2" ACTOR_OWNER = "gorgeous_thistle" RESIDENTIAL = {"useApifyProxy": True, "apifyProxyGroups": ["RESIDENTIAL"]} # erreurs DÉFINITIVES des acteurs (handle mort, profil privé, mur de login…) : # on tamponne le compte quand même pour qu'il ne soit re-sondé qu'à la rotation # `revisit_days`, au lieu d'être re-payé chaque jour en pure perte. Les erreurs # transitoires (429, 5xx, shell, proxy) ne sont PAS tamponnées → re-tentées. MISS_PREFIXES = ("user_not_found", "private", "not_found", "no_profile", "no_channel", "login_wall", "invite_invalide", "http_404", "wall_") # --- parallélisme : on dispose de 256 Go de RAM Apify → au lieu d'UN run à # 1 Go, on éclate le lot en tranches lancées EN PARALLÈLE, chacune avec plus de # mémoire (donc plus de CPU/vCPU côté Apify) et une concurrence interne accrue. # profil « max parallélisme » : ces acteurs sont I/O-bound (requêtes proxy), # donc BEAUCOUP de petits runs > peu de gros runs. 32 × 2 Go = 64 Go de pointe # (large sous les 256 Go / 128 runs concurrents du plan). SHARD_SIZE = int(os.environ.get("APIFY_SHARD_SIZE", "60")) # handles / run MAX_SHARDS = int(os.environ.get("APIFY_MAX_SHARDS", "32")) # runs simultanés RUN_MEMORY_MB = int(os.environ.get("APIFY_RUN_MEMORY_MB", "2048")) # /run RUN_CONCURRENCY = int(os.environ.get("APIFY_RUN_CONCURRENCY", "10")) # interne (max 10) # rafraîchissement : un compte est re-sondé s'il n'a pas été enrichi depuis # REVISIT_DAYS jours. Le watch quotidien couvre ainsi TOUT le bassin en ≤ N j. REVISIT_DAYS = int(os.environ.get("APIFY_REVISIT_DAYS", "7")) # passage FORCÉ (APIFY_FORCE=1) : ignore fraîcheur ET caps — couvre TOUT le # bassin de chaque plateforme en un seul passage (rattrapage de couverture) FORCE_PASS = os.environ.get("APIFY_FORCE", "") == "1" def _token() -> str: token = os.environ.get("APIFY_TOKEN") if not token: raise SkipSource("APIFY_TOKEN manquant (voir .env)") return token def _launch(name: str, run_input: dict, memory_mb: int, retries: int = 6) -> str: """Démarre un run. Retry sur 400/429 (limite mémoire concurrente atteinte : elle se libère quand d'autres tranches terminent) avec backoff.""" headers = {"Authorization": f"Bearer {_token()}"} for attempt in range(retries): resp = requests.post( f"{APIFY_API}/acts/{ACTOR_OWNER}~{name}/runs", json=run_input, headers=headers, params={"memory": memory_mb}, timeout=60) if resp.status_code in (400, 429, 402) or resp.status_code >= 500: if attempt < retries - 1: # backoff + jitter : désynchronise les tranches (évite le # lockstep où toutes redemandent la RAM au même instant) time.sleep(min(60, 8 * (attempt + 1)) + random.uniform(0, 6)) continue resp.raise_for_status() return resp.json()["data"]["id"] resp.raise_for_status() return resp.json()["data"]["id"] def _await_items(run_id: str, name: str, timeout_s: int) -> list[dict]: headers = {"Authorization": f"Bearer {_token()}"} t0 = time.time() while True: time.sleep(15) r = requests.get(f"{APIFY_API}/actor-runs/{run_id}", headers=headers, timeout=60) r.raise_for_status() status = r.json()["data"] if status["status"] in ("SUCCEEDED", "FAILED", "ABORTED", "TIMED-OUT"): break if time.time() - t0 > timeout_s: requests.post(f"{APIFY_API}/actor-runs/{run_id}/abort", headers=headers, timeout=60) raise RuntimeError(f"acteur {name} : délai {timeout_s}s dépassé") if status["status"] != "SUCCEEDED": raise RuntimeError(f"acteur {name} : run {status['status']}") dataset = status["defaultDatasetId"] items: list[dict] = [] offset = 0 while True: r = requests.get(f"{APIFY_API}/datasets/{dataset}/items", params={"offset": offset, "limit": 1000, "clean": "true"}, headers=headers, timeout=120) r.raise_for_status() chunk = r.json() items.extend(chunk) if len(chunk) < 1000: return items offset += 1000 def run_actor(name: str, run_input: dict, *, timeout_s: int = 2400, memory_mb: int = RUN_MEMORY_MB) -> list[dict]: """Un seul run (recherche, petit lot).""" return _await_items(_launch(name, run_input, memory_mb), name, timeout_s) def run_actor_sharded(name: str, base_input: dict, usernames: list[str], *, timeout_s: int = 2400, shard_size: int = SHARD_SIZE, max_shards: int = MAX_SHARDS, memory_mb: int = RUN_MEMORY_MB) -> list[dict]: """Éclate `usernames` en tranches lancées EN PARALLÈLE sur Apify. Chaque tranche = un run indépendant (mémoire dédiée → plus de CPU). On plafonne à `max_shards` runs simultanés (256 Go de RAM = large marge : 12 × 4 Go = 48 Go). Les datasets sont fusionnés. Une tranche qui échoue n'annule pas les autres. """ if not usernames: return [] n = max(1, min(max_shards, (len(usernames) + shard_size - 1) // shard_size)) shards = [usernames[i::n] for i in range(n)] # répartition équilibrée items: list[dict] = [] with ThreadPoolExecutor(max_workers=n) as pool: futs = {pool.submit(run_actor, name, {**base_input, "usernames": shard}, timeout_s=timeout_s, memory_mb=memory_mb): i for i, shard in enumerate(shards) if shard} for fut in as_completed(futs): try: items.extend(fut.result()) except Exception as exc: # une tranche morte ne bloque pas le reste print(f"[crea-ka] {name} tranche {futs[fut]} échouée : {exc}") return items class _ApifySocialEnrich(BaseConnector): """Base commune des enrichissements Apify (un acteur par plateforme).""" kind = "enrichment" platform = "" # plateforme canonique (§6.2) actor = "" # nom court de l'acteur (ka-instagram…) # concurrence interne du run — DOIT respecter le `maximum` du schéma # d'input de l'acteur, sinon Apify rejette le run en 400 (cause des # passages facebook/threads à zéro, diagnostiqué le 2026-08-21) run_concurrency = RUN_CONCURRENCY cap = 3000 # comptes max par passage (shardé en parallèle) revisit_days = REVISIT_DAYS # re-sonde hebdo par défaut (env APIFY_REVISIT_DAYS) content_key = "" # clé du contenu récent dans l'item acteur # clés d'item copiées telles quelles dans account.metrics si non nulles metric_map: dict[str, str] = {} @property def _stamp_key(self) -> str: # tampon PROPRE au connecteur : quand CE connecteur a enrichi ce compte. # On ne se base PAS sur acc.last_checked, re-tamponné par la découverte # (kick-decouverte etc.) même sans enrichissement → fausse fraîcheur. return f"{self.source_id}_at" def _stale(self, acc) -> bool: if FORCE_PASS: return True stamp = (acc.metrics or {}).get(self._stamp_key) if not stamp: return True try: seen = datetime.fromisoformat(str(stamp).replace("Z", "+00:00")) except ValueError: return True return (datetime.now(timezone.utc) - seen > timedelta(days=self.revisit_days)) def extra_input(self) -> dict: return {} def target_of(self, acc) -> str: """Valeur envoyée à l'acteur pour ce compte (défaut : le handle). Certaines plateformes exigent autre chose que le handle minusculé — ex. Discord, dont le code d'invitation est SENSIBLE À LA CASSE et doit être repris depuis l'URL. Le rapprochement se fait ensuite en minuscules des deux côtés. """ return acc.handle def enrich(self, creators: list[Creator]) -> list[Creator]: targets: dict[str, list] = {} for cr in creators: acc = next((a for a in cr.platforms if a.platform == self.platform and a.handle), None) if acc is None or not self._stale(acc): continue targets.setdefault(self.target_of(acc), []).append((cr, acc)) if not FORCE_PASS and len(targets) >= self.cap: break if not targets: return [] items = run_actor_sharded( self.actor, {"proxyConfiguration": RESIDENTIAL, "concurrency": self.run_concurrency, **self.extra_input()}, sorted(targets), # passage forcé : tranches plus grosses (bassin entier / 32 runs) # → délai par run élargi en proportion timeout_s=7200 if FORCE_PASS else 2400) by_handle = {str(it.get("username") or "").lower(): it for it in items if it.get("kind") == "profile"} enriched: list[Creator] = [] self.errors = 0 for handle, pairs in targets.items(): it = by_handle.get(handle.lower()) if it is None: continue if not it.get("found"): err = str(it.get("error", "")) if err.startswith(("http_5", "shell")): self.errors += 1 elif err.startswith(MISS_PREFIXES): for cr, acc in pairs: # miss définitif → rotation douce acc.metrics[self._stamp_key] = now_iso() acc.metrics[f"{self.source_id}_miss"] = err enriched.append(cr) continue for cr, acc in pairs: self.apply(cr, acc, it) enriched.append(cr) return enriched # -- application aux fiches ------------------------------------------------ def apply(self, cr: Creator, acc, it: dict) -> None: acc.followers = parse_count(it.get("followers")) or acc.followers if it.get("is_verified") is not None: acc.verified = bool(it["is_verified"]) acc.last_checked = now_iso() metrics = {dst: it.get(src) for src, dst in self.metric_map.items() if it.get(src) is not None} if self.content_key and it.get(self.content_key): metrics[self.content_key] = it[self.content_key][:12] # images du profil : avatar/bannière conservés PAR COMPTE (affichage # par plateforme + chaîne de repli côté frontend) for img_key in ("avatar", "banner"): if it.get(img_key): metrics[img_key] = it[img_key] metrics[self._stamp_key] = now_iso() # tampon d'enrichissement propre acc.metrics.pop(f"{self.source_id}_miss", None) # ressuscité acc.metrics.update(metrics) # avatar/bannière de la FICHE : rafraîchis à chaque passage depuis la # plateforme principale (les URLs CDN signées expirent — ex. Instagram) ; # sinon on remplit seulement les manquants if it.get("avatar") and (not cr.avatar_url or acc.platform == cr.primary_platform): cr.avatar_url = it["avatar"] if it.get("banner") and (not getattr(cr, "banner_url", None) or acc.platform == cr.primary_platform): cr.banner_url = it["banner"] if not cr.bio and it.get("biography"): cr.bio = it["biography"] self.cross_links(cr, it) def cross_links(self, cr: Creator, it: dict) -> None: """Liens de bio → link-in-bio ou compte `cross_link` (§12.1).""" for url in self.bio_urls(it): url = (url or "").strip() if url and not url.startswith("http"): url = "https://" + url if not url: continue if is_linkinbio(url) and not cr.link_in_bio_url: cr.link_in_bio_url = url continue hit = platform_from_url(url) if hit and hit[0] != self.platform: cr.platforms = merge_accounts(cr.platforms, [ account(hit[0], hit[1], "cross_link", url=url).finalize()]) def bio_urls(self, it: dict) -> list[str]: return [] class ApifyInstagram(_ApifySocialEnrich): source_id = "instagram-apify" platform = "instagram" actor = "ka-instagram" cap = 1800 content_key = "recent_posts" metric_map = {"following": "following", "posts_count": "posts", "category": "category", "is_business": "is_business", "business_category": "business_category", "highlight_reels": "highlight_reels", "pronouns": "pronouns", "bio_links": "bio_links", "avg_likes": "avg_likes", "avg_comments": "avg_comments", "avg_video_views": "avg_video_views", "engagement_rate_pct": "engagement_rate_pct", "posts_per_week": "posts_per_week", "last_post_at": "last_post_at", "video_share_pct": "video_share_pct", "top_post": "top_post", "bio_mentions": "bio_mentions", "bio_hashtags": "bio_hashtags", "igtv_videos_count": "igtv_videos", "business_email": "business_email"} def bio_urls(self, it: dict) -> list[str]: return [it.get("external_url") or ""] + (it.get("bio_links") or []) def apply(self, cr, acc, it): super().apply(cr, acc, it) related = [r.get("username") for r in (it.get("related_profiles") or []) if r.get("username")] if related: # matière première de découvertes futures (§12.1 mention) acc.metrics["related_profiles"] = related[:10] # flag Meta « compte Threads relié » : même handle → cross_link fort if it.get("has_threads") and it.get("username"): u = it["username"] cr.platforms = merge_accounts(cr.platforms, [ account("threads", u, "cross_link", url=f"https://www.threads.com/@{u}").finalize()]) class ApifyTikTok(_ApifySocialEnrich): source_id = "tiktok-apify" platform = "tiktok" actor = "ka-tiktok" cap = 1800 content_key = "recent_videos" metric_map = {"following": "following", "total_likes": "likes", "videos_count": "videos", "region": "region", "friends": "friends", "is_seller": "is_seller", "is_live_now": "is_live_now", "commerce_category": "commerce_category", "account_created_at": "account_created_at", "avg_views": "avg_views", "avg_likes": "avg_likes", "avg_comments": "avg_comments", "avg_shares": "avg_shares", "engagement_rate_pct": "engagement_rate_pct", "videos_per_week": "videos_per_week", "last_video_at": "last_video_at", "top_hashtags": "top_hashtags", "language": "language", "top_video": "top_video", "likes_given": "likes_given", "avg_saves": "avg_saves", "pinned_videos_count": "pinned_videos"} def bio_urls(self, it: dict) -> list[str]: return [it.get("bio_link") or ""] def apply(self, cr, acc, it): super().apply(cr, acc, it) # drapeau mineur DÉCLARÉ par TikTok → régime restreint §15 if it.get("is_under_18") is True: cr.is_minor = True class ApifyX(_ApifySocialEnrich): source_id = "x-apify" platform = "x" actor = "ka-x" # l'endpoint syndication de X rend des 429 en rafale au-delà de ~4 # requêtes simultanées par run (constat 2026-08-21) run_concurrency = 4 cap = 1800 content_key = "recent_tweets" metric_map = {"following": "following", "tweets_count": "tweets", "location": "location", "created_at": "created_at", "listed_count": "listed", "account_age_days": "account_age_days", "avg_likes": "avg_likes", "avg_retweets": "avg_retweets", "avg_replies": "avg_replies", "engagement_rate_pct": "engagement_rate_pct", "tweets_per_week": "tweets_per_week", "last_tweet_at": "last_tweet_at", "media_share_pct": "media_share_pct", "reply_share_pct": "reply_share_pct", "is_blue_verified": "is_blue_verified", "likes_given": "likes_given", "website": "website", "top_hashtags": "top_hashtags", "top_tweet": "top_tweet", "avg_quotes": "avg_quotes", "verified_type": "verified_type", "link_share_pct": "link_share_pct", "retweet_share_pct": "retweet_share_pct", "top_mentions": "top_mentions"} def extra_input(self) -> dict: return {"maxTweets": 20} def bio_urls(self, it: dict) -> list[str]: return [it.get("website") or ""] class ApifyFacebook(_ApifySocialEnrich): source_id = "facebook-apify" platform = "facebook" actor = "ka-facebook" run_concurrency = 6 # maximum du schéma d'input de ka-facebook cap = 400 revisit_days = 7 metric_map = {"category": "category", "page_id": "page_id", "website": "website", "rating": "rating"} def bio_urls(self, it: dict) -> list[str]: # site web auto-déclaré de la page → cross_link (§12.1) return [it.get("website") or ""] class ApifyThreads(_ApifySocialEnrich): source_id = "threads-apify" platform = "threads" actor = "ka-threads" run_concurrency = 6 # maximum du schéma d'input de ka-threads cap = 60 revisit_days = 7 content_key = "recent_posts" metric_map = {"avg_likes": "avg_likes", "top_post": "top_post", "bio_links": "bio_links"} def bio_urls(self, it: dict) -> list[str]: # liens de bio AUTO-DÉCLARÉS Threads → cross_link fort (§12.1) return list(it.get("bio_links") or []) class ApifySnapchat(_ApifySocialEnrich): source_id = "snapchat-apify" platform = "snapchat" actor = "ka-snapchat" cap = 60 revisit_days = 7 metric_map = {"category": "category", "has_story": "has_story", "has_spotlight": "has_spotlight", "address": "address", "subcategory": "subcategory", "publisher_type": "publisher_type", "story_snaps_count": "story_snaps", "lenses_count": "lenses", "story_previews": "story_previews", "spotlight_previews": "spotlight_previews", "website": "website", "snapcode": "snapcode", "spotlight_highlights_count": "spotlight_highlights"} def bio_urls(self, it: dict) -> list[str]: return [it.get("website") or ""] class ApifyYouTube(_ApifySocialEnrich): source_id = "youtube-apify" platform = "youtube" actor = "ka-youtube" cap = 1800 content_key = "recent_videos" metric_map = {"videos_count": "videos", "country": "country", "keywords": "keywords", "channel_id": "channel_id", "avg_views": "avg_views", "top_video": "top_video", "total_views": "total_views", "joined_date": "joined_date", "external_links": "external_links", "last_video_published": "last_video_published", "videos_per_month": "videos_per_month", "has_shorts": "has_shorts", "is_live_now": "is_live_now"} def bio_urls(self, it: dict) -> list[str]: # liens externes AUTO-DÉCLARÉS de la page À propos → cross_link (§12.1) return [ln.get("url") or "" for ln in (it.get("external_links") or []) if isinstance(ln, dict)] class ApifyTwitch(_ApifySocialEnrich): source_id = "twitch-apify" platform = "twitch" actor = "ka-twitch" cap = 1200 content_key = "recent_videos" metric_map = {"is_partner": "is_partner", "is_affiliate": "is_affiliate", "team": "team", "created_at": "created_at", "is_live_now": "is_live_now", "live_viewers": "live_viewers", "live_game": "live_game", "last_broadcast_at": "last_broadcast_at", "last_broadcast_game": "last_broadcast_game", "last_broadcast_title": "last_broadcast_title", "live_title": "live_title", "live_thumbnail": "live_thumbnail", "avg_video_views": "avg_video_views", "social_links": "social_links", "recent_games": "recent_games", "videos_count": "videos", "live_type": "live_type", "live_started_at": "live_started_at"} def bio_urls(self, it: dict) -> list[str]: # panneau « À propos » Twitch : liens sociaux AUTO-DÉCLARÉS (§12.1) return [ln.get("url") or "" for ln in (it.get("social_links") or []) if isinstance(ln, dict)] class ApifyKick(_ApifySocialEnrich): source_id = "kick-apify" platform = "kick" actor = "ka-kick" cap = 600 metric_map = {"is_live_now": "is_live_now", "live_viewers": "live_viewers", "subscription_enabled": "subscription_enabled", "live_title": "live_title", "live_thumbnail": "live_thumbnail", "live_category": "live_category", "vod_enabled": "vod_enabled", "recent_categories": "recent_categories", "live_started_at": "live_started_at", "live_language": "live_language"} def cross_links(self, cr: Creator, it: dict) -> None: """Liens sociaux AUTO-DÉCLARÉS de la fiche Kick → cross_link (§12.1).""" socials = it.get("social_links") or {} tmpl = {"instagram": "https://instagram.com/{}", "twitter": "https://x.com/{}", "youtube": "https://youtube.com/{}", "tiktok": "https://www.tiktok.com/@{}", "facebook": "https://facebook.com/{}", "discord": "https://discord.gg/{}"} for key, val in socials.items(): val = (val or "").strip() if not val: continue url = val if val.startswith("http") else \ tmpl.get(key, "{}").format(val.lstrip("@")) hit = platform_from_url(url) if hit and hit[0] != self.platform: cr.platforms = merge_accounts(cr.platforms, [ account(hit[0], hit[1], "cross_link", url=url).finalize()]) class ApifyOnlyFans(_ApifySocialEnrich): source_id = "onlyfans-apify" platform = "onlyfans" actor = "ka-onlyfans" cap = 100 revisit_days = 7 metric_map = {"posts_count": "posts", "photos_count": "photos", "videos_count": "videos", "likes": "likes", "streams_count": "streams", "subscribe_price_usd": "subscribe_price_usd", "is_free": "is_free", "location": "location", "join_date": "join_date", "audios_count": "audios"} def bio_urls(self, it: dict) -> list[str]: return [it.get("website") or ""] class ApifyFansly(_ApifySocialEnrich): source_id = "fansly-apify" platform = "fansly" actor = "ka-fansly" cap = 600 revisit_days = 7 metric_map = {"posts_count": "posts", "images_count": "images", "videos_count": "videos", "likes": "likes", "following": "following", "location": "location", "subscription_price": "subscription_price", "account_created_at": "account_created_at"} class ApifyPatreon(_ApifySocialEnrich): source_id = "patreon-apify" platform = "patreon" actor = "ka-patreon" cap = 400 revisit_days = 7 metric_map = {"patrons": "patrons", "posts_count": "posts", "is_monthly": "is_monthly", "is_nsfw": "is_nsfw", "creation_name": "creation_name", "social_links": "social_links", "published_at": "published_at", "pay_per_name": "pay_per_name"} def bio_urls(self, it: dict) -> list[str]: # liens sociaux AUTO-DÉCLARÉS de la campagne Patreon → cross_link return list((it.get("social_links") or {}).values()) class ApifyDiscord(_ApifySocialEnrich): source_id = "discord-apify" platform = "discord" actor = "ka-discord" cap = 400 revisit_days = 7 metric_map = {"members": "members", "online": "online", "boosts": "boosts", "partnered": "partnered", "guild_id": "guild_id", "channel": "channel", "splash": "splash", "vanity_url_code": "vanity_url_code", "premium_tier": "premium_tier", "nsfw_level": "nsfw_level"} def target_of(self, acc) -> str: """Code d'invitation SENSIBLE À LA CASSE, repris de l'URL (pas du handle minusculé) : discord.gg/ ou discord.com/invite/.""" url = acc.url or "" for sep in ("discord.gg/", "discord.com/invite/", "/invite/"): if sep in url: return url.split(sep, 1)[1].split("/")[0].split("?")[0] return acc.handle # --- découverte Instagram (topsearch, mots-clés QC) --------------------------- # thèmes × marqueurs QC — mêmes familles que youtube-recherche, adaptées à la # recherche d'UTILISATEURS Instagram (courtes, telles qu'on tape dans l'app) _THEMES = ["humoriste", "gamer", "gaming", "maquillage", "beauté", "fitness", "cuisine", "chef", "mode", "voyage", "musique", "chanteuse", "chanteur", "danse", "maman", "papa", "plein air", "chasse", "pêche", "photographe", "artiste", "peintre", "tatoueur", "coiffure", "esthétique", "auto", "moto", "déco", "immobilier", "entrepreneure", "podcast", "comédien", "comédienne", "drag", "twitch", "youtubeur", "youtubeuse", "tiktokeuse", "influenceuse", "créateur de contenu", "créatrice de contenu"] _MARKERS = ["québec", "qc", "montréal", "mtl"] class InstagramDecouverteApify(BaseConnector): """Découverte Instagram par recherche de mots-clés QC (acteur ka-instagram). L'endpoint topsearch du web public retourne les meilleurs comptes pour un mot-clé (nom, handle, abonnés approximatifs, badge). Requêtes tournantes (~40/passage sur ~160 combinaisons) pour accumuler sans re-payer les mêmes résultats chaque jour. Signal : profil_source (0.98) — le compte découvert EST la source de la fiche. """ source_id = "instagram-decouverte" kind = "discovery" queries_per_pass = 40 def build_queries(self) -> list[str]: combos = [f"{t} {m}" for t in _THEMES for m in _MARKERS] day = datetime.now(timezone.utc).timetuple().tm_yday start = (day * self.queries_per_pass) % len(combos) rotated = combos[start:] + combos[:start] return rotated[:self.queries_per_pass] def fetch(self) -> list[Creator]: # NOTE 2026-08-18 : l'endpoint topsearch du web public exige désormais # une session authentifiée (401 « Server Error » depuis toute IP non # connectée). La découverte par mots-clés est donc suspendue proprement # jusqu'à ce qu'une voie publique existe ; l'enrichissement Instagram # (instagram-apify) alimente `related_profiles` comme graine future. raise SkipSource("topsearch Instagram exige une session connectée " "(voie publique fermée le 2026-08-18)") items = run_actor("ka-instagram", { # noqa: unreachable (repli futur) "usernames": [], "queries": self.build_queries(), "proxyConfiguration": RESIDENTIAL, }) creators: list[Creator] = [] seen: set[str] = set() for it in items: if it.get("kind") != "search_user" or not it.get("username"): continue handle = str(it["username"]).lower() if handle in seen: continue seen.add(handle) acc = account("instagram", handle, "profil_source", followers=parse_count(it.get("followers")), verified=it.get("is_verified"), url=f"https://www.instagram.com/{handle}/") creators.append(Creator( source=self.source_id, external_id=handle, display_name=it.get("full_name") or handle, avatar_url=it.get("avatar") or None, platforms=[acc], notes=f"découvert via recherche Instagram « {it.get('query')} »", )) return creators