Something went wrong. Try again.
Automations and webhooks for the AT Protocol airglow.run
automation webhook atproto atprotocol
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596// Partition source for the Jetstream manager: which sockets sync needs.
import { config } from "../config.js";import type { PartitionContribution, SyncEventHandler } from "../jetstream/consumer.js";import { isNsidAllowed, nsidRequiresWantedDids } from "../lexicons/match.js";import { handleSyncEvent } from "./events.js";import { listEnabledSets } from "./store.js";import { PROFILE_WATCH_APPS, followCollection, profileCollection } from "./targets.js";
export const SYNC_KEY_PREFIX = "sync:";/** Jetstream v2 accepts at most this many DIDs per subscription. */const MAX_DIDS_PER_SUBSCRIPTION = 10_000;
const handler: SyncEventHandler = (event) => handleSyncEvent(event);
/** Stable shard for a DID: FNV-1a over the DID string, modulo the count. */export function shardFor(did: string, shardCount: number): number { let hash = 0x811c9dc5; for (let i = 0; i < did.length; i++) { hash ^= did.charCodeAt(i); hash = Math.imul(hash, 0x01000193); } return (hash >>> 0) % Math.max(1, shardCount);}
export function syncPartitions(): PartitionContribution[] { const sets = listEnabledSets(); if (sets.length === 0) return [];
const allow = config.nsidAllowlist; const block = config.nsidBlocklist; const shardCount = config.syncShardCount ?? 1; const shards = new Map<number, { dids: Set<string>; collections: Set<string> }>(); const watchedProfileApps = new Set<string>();
// Paused users stay subscribed: the handler only touches the database, so // their events keep queuing in order and drain after re-authorization. for (const set of sets) { const collections = set.apps .map(followCollection) .filter((nsid) => isNsidAllowed(nsid, allow, block)); if (collections.length === 0) continue; const n = shardFor(set.did, shardCount); let shard = shards.get(n); if (!shard) { shard = { dids: new Set(), collections: new Set() }; shards.set(n, shard); } shard.dids.add(set.did); for (const nsid of collections) shard.collections.add(nsid); for (const app of set.apps) { if (PROFILE_WATCH_APPS.includes(app)) watchedProfileApps.add(app); } }
const contributions: PartitionContribution[] = []; for (const [n, shard] of [...shards].sort((a, b) => a[0] - b[0])) { if (shard.dids.size === 0) continue; if (shard.dids.size > MAX_DIDS_PER_SUBSCRIPTION) { console.warn( `sync: shard ${n} has ${shard.dids.size} DIDs, above Jetstream's ${MAX_DIDS_PER_SUBSCRIPTION}; raise SYNC_SHARD_COUNT`, ); } const cells = new Map<string, SyncEventHandler>(); for (const nsid of shard.collections) { cells.set(`${nsid}\0create`, handler); cells.set(`${nsid}\0delete`, handler); } contributions.push({ key: `${SYNC_KEY_PREFIX}${n}`, wantedDids: [...shard.dids].sort(), cells, lookbackUs: () => config.syncJetstreamMaxLookbackUs, resumeFallbackPrefix: SYNC_KEY_PREFIX, }); }
const profileCells = new Map<string, SyncEventHandler>(); for (const app of PROFILE_WATCH_APPS) { if (!watchedProfileApps.has(app)) continue; const nsid = profileCollection(app); if (!isNsidAllowed(nsid, allow, block)) continue; if (nsidRequiresWantedDids(nsid, config.nsidRequireDids)) { console.warn( `sync: not watching ${nsid} profile creates because it requires wantedDids; late joiners on that app wait for the sweep`, ); continue; } profileCells.set(`${nsid}\0create`, handler); } if (profileCells.size > 0) { contributions.push({ key: "", wantedDids: [], cells: profileCells }); } return contributions;}