"""Native Discord ingester — the Discord ADAPTER over `ingest_core.Ingester`. Thin counterpart to subscribestar_ingester (milestone 428). The walk's modes, both ledgers, cursor checkpointing and the post-first capture live in the core; this wires in the Discord client, downloader, ledger models and key. Two things differ from the cookie platforms: - Discord authenticates with a user TOKEN, so `auth_token` is the credential here rather than an argument accepted and ignored. - The body canary is off. It fails a walk whose first 30+ captured posts all came back without text, on the theory that a creator nearly always writes something; a Discord drop is routinely files and nothing else, so on Discord that is an ordinary backfill, not a broken parser. And two things only a Discord source has: `discord_authors` in its config_overrides limits it to the creator's own messages (a server's other members post pictures too), and it learns its own name — the server and channel it walks — into `source.display_name`. `campaign_id` is the source URL (a server, channel, thread or category link). FC runs on a plain-HTTP homelab; nothing here uses a secure-context Web API. """ from __future__ import annotations import asyncio import logging from collections.abc import Callable from pathlib import Path from sqlalchemy import select, update from ..models import DiscordFailedMedia, DiscordSeenMedia, Source from .discord_client import DiscordAPIError, DiscordClient, MediaItem from .discord_downloader import DiscordDownloader from .ingest_core import Ingester log = logging.getLogger(__name__) _LEDGER_KEY_MAX = 128 # `config_overrides` key: the people whose messages this source takes (user # ids, usernames or display names). Absent or empty takes everyone's. AUTHORS_KEY = "discord_authors" def _ledger_key(media: MediaItem) -> str: """`:` — stable across edits (see MediaItem).""" return f"{media.post_id}:{media.media_id}"[:_LEDGER_KEY_MAX] class DiscordIngester(Ingester): """Walk a Discord source's channels, download unseen files, return a `DownloadResult`. `client` / `downloader` are injectable for tests.""" def __init__( self, images_root: Path, cookies_path: str | None, session_factory: Callable[[], object], *, validate: bool = True, rate_limit: float = 0.0, request_sleep: float = 0.0, auth_token: str | None = None, client: DiscordClient | None = None, downloader: DiscordDownloader | None = None, ): del cookies_path # Discord authenticates by token (uniform signature) self.images_root = Path(images_root) super().__init__( client=client if client is not None else DiscordClient( auth_token, request_sleep=request_sleep, ), downloader=downloader if downloader is not None else DiscordDownloader( self.images_root, validate=validate, rate_limit=rate_limit, ), session_factory=session_factory, seen_model=DiscordSeenMedia, failed_model=DiscordFailedMedia, seen_constraint="uq_discord_seen_media_source_id", failed_constraint="uq_discord_failed_media_source_id", ledger_key=_ledger_key, platform="discord", error_base=DiscordAPIError, drift_label="Discord API", body_canary=False, ) def run(self, **kwargs): """The core walk, bracketed by the two things only a Discord source has: whose messages it takes, read before the walk, and the name of what it walks, written after it (the walk is what loads that name).""" source_id = kwargs["source_id"] self.client.only_from(self._source_authors(source_id)) try: return super().run(**kwargs) finally: self._record_display_name(source_id, kwargs["campaign_id"]) def _source_authors(self, source_id: int) -> list: if self.session_factory is None: return [] with self.session_factory() as session: overrides = session.execute( select(Source.config_overrides).where(Source.id == source_id) ).scalar_one_or_none() or {} authors = overrides.get(AUTHORS_KEY) or [] return authors if isinstance(authors, list) else [] def _record_display_name(self, source_id: int, url: str) -> None: """Refreshed on every walk, never written once: a renamed channel should read as its new name. A name the walk couldn't read leaves the stored one alone, and a failure here never fails the walk.""" label = self.client.source_label(url) if not label or self.session_factory is None: return try: with self.session_factory() as session: session.execute( update(Source) .where(Source.id == source_id) .where(Source.display_name.is_distinct_from(label)) .values(display_name=label) ) session.commit() except Exception as exc: # a name is decoration — never fail the walk log.warning("Discord: couldn't record source %s's name: %s", source_id, exc) async def verify_discord_credential(url: str, auth_token: str | None) -> tuple[bool | None, str]: """The uniform `(ok, message)` probe: is the token valid, and can its account see the channel or server the source names?""" if not auth_token: return False, "No Discord token is saved — add one under Credentials." client = DiscordClient(auth_token) loop = asyncio.get_running_loop() return await loop.run_in_executor(None, client.verify_auth, url)