Skip to main content

Public Stats Crawler (Twitch)

The public stats crawler records public, non-personal aggregate stream stats (viewer curves, category timelines, title/metadata, and non-identifying chat aggregates) for foreign channels - channels that are not connected Lumio accounts. Those rows land in the same channel_history* tables as first-party data, tagged source = 'external_crawler', so the existing /stats read model surfaces them with no reader change.

The Twitch arm is the app apps/crawler-twitch: a long-running Rust container (precedent: InnerTube Proxy) - not a Cloudflare Worker, not one of the *-bot apps, and not part of lumio-api.

Regimes

This crawler is Regime B (mass, foreign-channel ingest). First-party session recording for connected accounts (Regime A) is unrelated and writes directly to TimescaleDB. See Stream History Platforms.

Architecture​

Helix discovery ──▶ crawl-job bus (JobQueue) ──▶ sharded crawler replicas
(Get Streams, RabbitMQ in prod; ├─ anonymous IRC read
paged, app token) in-memory single-node ├─ Helix stats poll (60s)
pre-broker └─ Twitch-GQL panels (fail-open)
│ typed CrawlEvents
▼
result log (ResultBus)
Kafka in prod; in-memory single-node
│
▼
shared writer ──▶ batch upsert into
channel_history{,_streams,_stats}
(source='external_crawler')

Transport seam (broker-agnostic)​

Two traits keep the producer code identical on a single-node staging box and on a multi-node cluster:

  • JobQueue - the crawl-job bus. Production impl RabbitMQ (per-msg ack, per-platform routing, DLQ). A job is routed to the shard that owns its channel.
  • ResultBus - the durable, replayable result log. Production impl Kafka, partitioned by channel id, consumed by the shared writer.

Redis is cache only (never a bus or system-of-record). The crate ships an in-memory single-node implementation of both traits so the full discovery → shard → writer pipeline is runnable before the broker/cluster land. End-to-end validation on staging single-node is gated on broker provisioning; live global ingestion is a separate, capacity-gated step.

Components​

ComponentRole
DiscoveryResolves the live crawl scope and enumerates the channels it covers, publishing one crawl job per live channel. global pages Helix Get Streams (100/page, cursor, app token) over the whole live set; the bounded modes query a resolved channel set via Get Streams?user_id= (≤100 ids/request).
Shard workerOwns hash(channel_id) % num_shards. Opens one anonymous (justinfan) IRC client, JOINs its channels under a JOIN token-bucket, and folds chat into per-broadcast aggregates. No single-leader election.
Stats poller60s viewer sweep over the live crawl scope → StreamObservation events (viewer curve + title/category/tags). Branches on the mode exactly as Discovery does. Fires the panels fetch once per broadcast. The writer diffs each observation's title/category/tags into the per-broadcast stream-metadata change timeline (public_stream_metadata_history, ZAF-890 - old→new rows, idempotent on unchanged re-observe).
Channel GQLOne undocumented Twitch GQL user(login:) fetch per broadcast at session start, carrying the ad-marked panels, the public follower total (ZAF-796), the About box (description / social links / team, ZAF-799), and the channel type (roles { isPartner isAffiliate } → partner/affiliate, ZAF-841). Isolated and fail-open - a missing/renamed field simply reads as absent and never drops the session; each derived event is emitted independently. The panels event is emitted on any valid observation (a user node with an id) even when the ad-marked set is empty, so the writer can reconcile panel removals (ZAF-892); a fail-open null fetch emits nothing (no spurious removals).
WriterKafka consumer group. Batch-upserts into channel_history* and the crawled channel-level tables public_channel_{panels,metrics,profile}, all in TSDB (relocated from main Postgres, ZAF-843), with write-amplification guards. Panel lifecycle reconciliation (ZAF-892): after upserting the present ad panels, it marks removed_at on any stored panelId absent from the channel's current set and clears it on re-observation - coherent per panel, guarded so an unchanged state writes nothing.
Opt-out enforcerRefreshes the opt-out registry from public_channel_settings (main DB) on a fixed cadence; every producer consults it before crawling, and it erases newly opted-out channels' crawled rows from TSDB.

Crawl scope​

The crawler is not unconditionally global. A scope selector bounds which live channels the crawler touches - both the jobs Discovery publishes and the observations the Stats poller records - so /stats can launch on today's single-host capacity and later ramp to a full global sweep with one admin flip (founder ruling ZAF-594). Four modes:

ModeChannels crawled
lumioFirst-party roster only - channels in channel_connections.
listCurated crawl_watchlist only.
lumio_plus_listDefault - the union of the roster and the watchlist.
globalThe unfiltered Helix sweep (kill-switch / ramp target).
  • Fail-closed default. The default is lumio_plus_list (bounded), never global. global is an explicit, opt-in widening. An unreadable settings row or a corrupt value resolves to the configured fallback, never silently to global.
  • DB-backed & live-switchable. The active mode is the singleton crawler_settings.scope row in the main Postgres DB, re-read once per sweep by each producer - an operator change lands on the next tick, no redeploy. The TOML [crawler] scope key is only the fallback default used when the row is unreadable.
  • Every Helix producer branches on the mode. global keeps the paged whole-live-set sweep. The bounded modes resolve a channel-id set from channel_connections (lumio) and/or crawl_watchlist (list), deduplicate it, and query it directly via Helix Get Streams?user_id= in ≤100-id batches - skipping the global sweep. A bounded scope that resolves to zero channels does nothing (the fail-closed outcome), never a global fallback.
  • Both producers, not just Discovery. The Stats poller is the module that actually feeds channel_history_streams / _stats, so it must honour the scope too - while it paged the whole live set, a lumio_plus_list deployment still wrote global history rows regardless of what Discovery published (ZAF-790). A source-scan guard in apps/crawler-twitch/src/stats.rs fails the build if any crawler module reaches the unfiltered sweep without first resolving the scope.

Admin control surface​

The mode and the watchlist are steerable from the Admin app on both protocols (GraphQL primary + REST parity - same fields, validation, and error messages):

ActionGraphQLRESTPermission
Read scopecrawlerScopeGET /v1/crawler/scopecrawler:scope-read
Set scopesetCrawlerScope(scope)PUT /v1/crawler/scopecrawler:scope-edit
List watchlistcrawlerWatchlist(platform)GET /v1/crawler/watchlistcrawler:watchlist-read
Add channeladdCrawlerWatchlistChannel(input)POST /v1/crawler/watchlistcrawler:watchlist-create
Remove channelremoveCrawlerWatchlistChannel(platform, platformChannelId)DELETE /v1/crawler/watchlist/{platform}/{platform_channel_id}crawler:watchlist-delete

The crawler:* permissions are granular (the coarse :manage is disallowed by the ZAF-594 ruling; see RBAC Permissions) and admin-scope - steering the single global crawler is a platform-operator action, so they are enforced via require_admin_permission / AdminPermissionGuard (never account-scope grants) and seeded to the system_admin admin role (ZAF-725; migration 20260825000001_crawler_perms_to_admin_scope). Every scope change and watchlist mutation emits a system-scope audit event (crawler:scope_changed, crawler:watchlist_added, crawler:watchlist_removed). Adding a channel is an idempotent upsert on the (platform, platform_channel_id) identity.

Identity & write path​

Each crawled broadcast maps onto the frozen identity columns on channel_history_streams:

Helix fieldColumn
user_idplatform_channel_id (join key, frozen)
user_loginchannel_login (URL spelling, frozen)
user_namechannel_display_name (frozen)
id (stream)broadcast_id - public key (platform, broadcast_id)

Identity and source are captured once on INSERT and never rewritten (a later handle rename must not rewrite history), mirroring the first-party contract. Viewer samples upsert with ON CONFLICT … DO UPDATE … WHERE viewer_count IS DISTINCT FROM … so an unchanged re-poll is not a fresh row-version.

Notes on the data model (no migration owed)​

  • The data model already shipped; the crawler adds no migration.
  • channel_history.account_id / channel_history_stats.account_id are NOT NULL, so crawler rows (which have no owning account) use a nil-UUID sentinel account id. Whether to make account_id nullable instead is a write-path decision for the API owner.
  • There is no channel_history_panels table; panels land in channel_history_streams.platform_metrics -> 'twitch' -> 'panels'.
  • The panel lifecycle marker public_channel_panels.removed_at (ZAF-892) is the one exception to "no migration owed": that ALTER TABLE ... ADD COLUMN lives in apps/api/tsdb_migrations (API-owned), not the crawler. The crawler only writes it (the per-observation reconciliation above); image versions need no new column - they are already the content_hash-keyed version rows.

Regime-B capture scope (aggregates only)​

For channels crawled without consent, capture is restricted to genuinely public, non-personal aggregates:

  • Kept: unique-chatter count, message count, top words, top emotes, viewer curve, category timeline, title/metadata.

The viewer curve lands in channel_history_stats (raw 60s samples) and is rolled up into the channel_history_stats_5m / _1h continuous aggregates that long / cross-session reads use - see Stream History: storage, rollups & retention for which layer answers which read and how retention is bound by the privacy policy.

  • Excluded until a channel is verified-claimed / opted in: per-username "top chatters" and pinned-message author/text. These identity-bearing fields write NULL for source='external_crawler' rows.

GIF content data - Top GIFs (ZAF-965 decision)​

Top GIFs are the one Regime-B capture that is a content datum, not a pure distribution/count, so the ZAF-965 founder decision (2026-09-03) opened them as their own ingest / retention / legal case (capture + retention in ZAF-968; privacy, consent-category and retention decision in ZAF-969; public display in ZAF-970, which is blocked on the legal texts landing first).

  • What is stored: per broadcast, which GIF appeared and how often - the GIPHY id and the GIPHY image URL plus a frequency count. No sender reference is stored - no chatter username, no message text, no user/message id. The stored record identifies no natural person ("which GIF, how often"), the same non-personal class as top words / top emotes.
  • Retention: 90 days. Unlike the indefinitely-kept public-profile / stream-metadata change histories, GIF content data is bounded to 90 days (aligned with the raw retention of the other crawled broadcast data - the channel_history_stats raw-sample window, ZAF-689), then auto-deleted. Opt-out erases a channel's GIF content data like every other crawled row. This figure is fixed by the founder decision and is stated identically in ZAF-968; do not let the two drift.
  • Display is consent-gated (Giphy embed). The public Stats page renders the Top GIFs by loading the image directly from media*.giphy.com, which transmits the visitor's IP + user agent to GIPHY (US). Because stats.lumio.vision is public - the visitor is not a customer, so the Multichat Art. 6(1)(b) contract basis does not carry over - the embed loads only after the visitor grants the Functional consent category, honouring the cookies.mdx "loads only after consent, listed beforehand" promise. Consent category = Functional (the same class as the existing YouTube/Twitch embeds; see shared/cookies-consent cookieDetector.ts). Legal basis for the display is Art. 6(1)(a) GDPR consent; the US transfer rides GIPHY/Shutterstock's EU-US DPF certification + SCCs. Storing the id/URL does not itself hit GIPHY - only the display does.
  • Where it is documented: the published treatment lives in apps/web/content/{de,en}/legal/privacy.mdx (§"Public statistics (Lumio Stats)", content class (D)) and …/legal/cookies.mdx (§"Lumio Stats" embed registry).
  • Ingest (ZAF-968). Twitch delivers GIFs as a gifs IRC tag on the chat message - comma-separated <start>-<end>|<gifID>|<gifURL> segments. The shard worker reads it from the raw source.tags (the typed PrivmsgMessage drops it, same as source-room-id), parses it defensively (broken segments dropped, never panics), and apps/crawler-twitch/src/chat.rs folds it into per-broadcast gif_count (every occurrence) + top_gifs (top-N by frequency). Shared-chat GIFs from a foreign room are excluded from the host aggregate, exactly like emotes (ZAF-819), and no sender reference is ever kept (§12). writer.rs rolls both onto the channel_history session row - the same columns the first-party finalizer fills (ZAF-972) - in the canonical {id, url, provider, alt, count} shape (provider/alt are null on this IRC path, which carries only id + URL). They surface unchanged through the public read model (publicStreamStats / GET /v1/public-stats/streams/{platform}/{broadcast_id}).
  • Retention mechanism (ZAF-968). channel_history is a plain table with no TimescaleDB retention policy, so the 90-day cap is an app-side sweep: apps/api/src/workers/gif_retention.rs (gated by [public_stats] gif_retention_enabled, on in production) periodically erases gif_count → 0 / top_gifs → [] on crawler rows (account_id = nil) older than 90 days. Only the GIF content is cleared - the session row, viewer curve and emote/word aggregates stay; first-party (owner) history sharing the same columns is never touched (kept on its Art. 6(1)(b) contract basis); opt-out separately deletes the whole crawled row (GIF content included) via the opt-out enforcer.

Twitch Shared Chat attribution (ZAF-819)​

Twitch Shared Chat delivers messages from other channels into a host channel's room. The shard worker reads each PRIVMSG's source-room-id IRC tag (from the raw source.tags, which the typed PrivmsgMessage drops): when it is present and differs from the joined room, the message originated in a different channel and is attributed to that source channel instead of being folded into the host's aggregate. Concretely, a shared-chat message is kept out of the host broadcast's message_count / unique_chatters / top-words / top-emotes, and counted in a per-source-channel breakdown emitted on the ChatAggregate and written to channel_history_streams.shared_chat_sources as [{"platform_channel_id": …, "message_count": …}] (NULL for a normal broadcast). The public read model surfaces it on every stream summary (shared_chat_sources) so the stream page can show "who these stats are coming from". At read time each source is enriched with the source channel's display_name + avatar_url (ZAF-868), resolved from public_channel_profile keyed by (platform, platform_channel_id) - the crawler captures these from the same user(login:) GQL node (displayName / profileImageURL) it already fetches for the About box, refresh-and-overwrite, and a source the crawler never profiled resolves to null identity (fail-open, the client falls back to the id). This stays aggregates-only - the breakdown is by channel (public broadcaster identity), never by viewer.

Opt-out enforcement (two-sided)​

A channel opts out by setting is_public = false in public_channel_settings, via one of two writers: the dashboard toggle (a first-party account managing its own connected channels) or a verified !optout chat command typed by the broadcaster in their own channel (see below). The registry is public-by-default: the absence of a row means public. Enforcement is two-sided, and both sides must hold before global ingest is armed:

  • Read side (already shipped, apps/api): the public-stats read model serves nothing for an opted-out channel - single-channel reads return empty, the cross-channel browse excludes it, and the WS bootstrap is gated.
  • Capture side (crawler): the crawler stops producing data for the channel and erases what it already captured:
    • Roster-drop. All three producers consult an in-memory opt-out registry before acting - Discovery does not publish a crawl job, the Stats poller does not emit a viewer sample, and a Shard worker PARTs a channel it had already joined (dropping its in-flight aggregate). So an opted-out channel is neither joined, polled, nor written.
    • Serviceable deletion. When a channel newly opts out, its already-captured source='external_crawler' rows are deleted from channel_history, channel_history_streams, and channel_history_stats (scoped to the crawler account sentinel / source, so first-party history is never touched), and its channel-level rows are deleted from public_channel_panels (every ad-panel version), public_channel_metrics, every follower-growth snapshot row (public_channel_follower_history, ZAF-891), public_channel_profile, every About-box change-history row (public_channel_profile_history, ZAF-836), and every stream-metadata change-history row for its broadcasts (public_stream_metadata_history, ZAF-890 - keyed by broadcast for reads but carrying platform_channel_id so opt-out erases by channel identity like the siblings). All these crawled tables live in TSDB - the public_channel_* data tables were relocated there (_panels/_metrics/_profile, ZAF-843) or authored there directly (_profile_history ZAF-836, _follower_history ZAF-891, public_stream_metadata_history ZAF-890) per the founder ruling ZAF-791 - so every erasure targets the TSDB pool; only the public_channel_settings opt-out registry is read from the main Postgres db. A crawler restart reconciles the full opted-out set once, so rows that landed while it was down are still erased. Each per-table delete is independently fail-open with no wrapping transaction: if any chunk fails (a transient TSDB outage), that channel is held back from the reconciled set and re-attempted on every subsequent tick until all its rows are gone - it is never left partially erased until a restart (ZAF-844). The deletes are idempotent, so re-running them against an already-clean table is a harmless no-op.

The registry (public_channel_settings) is the one crawled surface that stays in the main Postgres DB, so the crawler opens a small pool onto it ([database]) in addition to its TSDB write pool (the target for channel_history*, the public_channel_* data tables, and every erasure). The main-DB pool is read-mostly: the enforcer reads the registry, and the only write the crawler ever makes to the main DB is the broadcaster !optout upsert below. The refresh + erasure loop is fail-open: a main-DB outage degrades to "enforce with the last-known snapshot", never a crash, and is retried on the next tick.

crawler_optout_* metrics expose the registry size, roster-drop skips, PARTs, rows erased, refresh/erasure errors, and honoured !optout commands (crawler_optout_commands_total).

Broadcaster !optout chat command​

A broadcaster can opt their own channel out from chat: they type !optout in their own Twitch channel and the crawler upserts public_channel_settings(is_public = false, opted_out_via = 'crawler_command') for that channel identity. Enforcement is then automatic - the enforcer loop above picks the new row up on its next refresh and PARTs + erases exactly as it does for a dashboard opt-out. This is the write counterpart to the capture-side enforcement; there is no separate command-side enforcer.

  • Verification - broadcaster badge alone (design §6 ruling). The command is honoured only when the PRIVMSG carries the server-authoritative broadcaster/1 badge, which Twitch stamps into the IRC tags for the channel owner (and only in their own room). It cannot be forged by a viewer, so it is strong self-service proof of this channel's ownership. Design §6 also names an out-of-band hash-verified claim; that stronger proof exists to bind a channel to a Lumio account for management (an access-granting link) and is not required for opt-out. !optout only ever removes a channel and triggers erasure - a privacy-protective, reversible action whose safe default is to honour it immediately (requiring an extra step would mean crawling a channel whose owner has already asked us to stop).
  • Non-clobbering, idempotent write. The upsert only flips a currently-public row to opted-out; it never overwrites an existing opt-out (a prior dashboard opt-out keeps its provenance), never touches account_id or show_top_chatters, and is a no-op once the channel is already opted out (... ON CONFLICT DO UPDATE ... WHERE is_public = true), so a spammed command does not amplify writes. It runs off the IRC hot path and is fail-open - a failed write is logged/counted and simply retried on the next !optout.
  • Audit parity (deliberate difference). The dashboard path emits a public_stats:channel_opted_out audit event (account scope, actor = the Lumio user). The command path does not emit that event: the actor is a Twitch platform user proven by a badge (no Lumio user_id) and a crawled channel has no account_id, so an account-scoped event has no valid keys; and db::audit::emit lives inside the lumio-api binary crate, unreachable from the crawler. Provenance is instead recorded on the row itself (opted_out_via = 'crawler_command', opted_out_by = <broadcaster Twitch user id>, updated_at). A system-scoped mirror event would require extracting audit into a shared crate and is tracked as a coordinated follow-up with the Rust Backend Engineer rather than a one-sided raw insert.

Rate limits & IP reputation​

  • Runs on a trusted-IP container host (never Cloudflare Edge / CI ranges, which platforms bot-gate - the InnerTube Proxy lesson).
  • Helix is a points bucket (~800/min per app token). Prefer paged sweeps over per-channel calls; on 429, honour Ratelimit-Reset.
  • Anonymous IRC has a per-connection JOIN limit (~20 JOIN/10s); a token-bucket paces JOINs. Reconnects are handled automatically by the IRC pool.
  • Crawl data may be dropped under writer lag (Kafka retention / 1-of-N sampling for hot channels). Drops are never silent - they increment a metric.

Metrics​

The crawler exposes Prometheus text at GET /metrics (plus /health, /ready). Metric families: discovery coverage/freshness, Helix + reputation alarms (429 / auth-fail / LOGIN_REQUIRED / ban), IRC connection/JOIN health, panels, bus/writer throughput + write-amplification (rows_upserted vs rows_skipped_noop), drop-policy activations, and opt-out enforcement (crawler_optout_*: registry size, roster-drop skips, PARTs, rows erased, refresh/erasure errors, honoured !optout commands). Arming alarms on these is the entry criterion for live ingestion.

Configuration​

Layered TOML (config/default.toml → {run_mode}.toml → local.toml → ENV CRAWLER_TWITCH__*). Secrets come from ENV / gitignored local.toml:

CRAWLER_TWITCH__TWITCH__CLIENT_ID=...
CRAWLER_TWITCH__TWITCH__CLIENT_SECRET=...
CRAWLER_TWITCH__TIMESCALE__URL=postgres://...
CRAWLER_TWITCH__DATABASE__URL=postgres://... # main DB, opt-out registry (read + !optout write)

Two DB pools: [timescale] (the write target for channel_history* and the crawled public_channel_{panels,metrics,profile} data tables - relocated to TSDB per ZAF-791/ZAF-843) and [database] (the main Postgres - the public_channel_settings opt-out registry; read by the enforcer, written only by the broadcaster !optout upsert).

Key [crawler] knobs: num_shards, owned_shards, discovery_interval_secs, stats_poll_interval_secs, page_size, max_channels_per_shard, join_bucket_capacity, join_refill_per_sec, writer_batch_size, optout_refresh_interval_secs, and scope (the crawl-scope fallback default - the live mode is the DB-backed crawler_settings row).

The scope selector adds two main-DB tables (owned by lumio-api migrations, not the crawler): crawler_settings (the singleton, admin-switchable scope mode) and crawl_watchlist (the curated list component, keyed by (platform, platform_channel_id)).