TypeScript 61.9%
HTML 37.2%
SQL 0.7%
1import { and, eq, isNull, lte, or, sql } from 'drizzle-orm';2import { connectors as connectorsTable } from '@rareindex/database';3import { listConnectorMeta, missingRequirements } from '@rareindex/connectors';4import { logger } from '@rareindex/shared';5import { db } from '../lib/db.ts';6import { JOBS, type Queue } from '../lib/queue.ts';78const PRIORITY: Record<string, number> = { high: 10, medium: 5, low: 1 };910/**11 * Dynamic crawl scheduling (§141, SPEC §17): every tick, enqueue a crawl for active connectors whose12 * next_run_at is due (or never ran). Gated connectors (missing API keys) are skipped without a run.13 * Active backfill campaigns (SPEC §9) are continued with `mode: 'backfill'` as long as the connector14 * has no crawl in flight, so a campaign survives worker restarts and time-boxed runs.15 */16export async function scheduleDueCrawls(queue: Queue): Promise<number> {17 const now = new Date();18 const gated = new Set(listConnectorMeta().filter((m) => missingRequirements(m).length > 0).map((m) => m.id));19 const due = await db()20 .select({ id: connectorsTable.id, priority: connectorsTable.priority })21 .from(connectorsTable)22 .where(and(eq(connectorsTable.status, 'active'), or(isNull(connectorsTable.nextRunAt), lte(connectorsTable.nextRunAt, now))));23 let n = 0;24 // pg-boss singleton keys do not dedupe across states reliably; skip connectors that already have a queued/active crawl.25 const queued = (await db().execute(sql`select distinct data->>'connectorId' as id from pgboss.job where name = ${JOBS.crawlRun} and state in ('created','retry','active')`)) as unknown as Array<{ id: string }>;26 const busy = new Set(queued.map((q) => q.id));27 const running = (await db().execute(sql`select distinct connector_id as id from connector_runs where status = 'running' and started_at > now() - interval '6 hours'`)) as unknown as Array<{ id: string }>;28 for (const r of running) busy.add(r.id);2930 // 1. backfill continuations31 const backfills = (await db().execute(sql`select b.connector_id as id, c.priority from connector_backfills b join connectors c on c.id = b.connector_id where b.status = 'running' and c.status = 'active' and b.updated_at < now() - interval '2 minutes'`)) as unknown as Array<{ id: string; priority: string }>;32 for (const b of backfills) {33 if (busy.has(b.id) || gated.has(b.id)) continue;34 const id = await queue.send(JOBS.crawlRun, { connectorId: b.id, mode: 'backfill', trigger: 'backfill' }, { singletonKey: `crawl:${b.id}`, priority: Math.max(1, (PRIORITY[b.priority] ?? 5) - 2), expireInSeconds: 6 * 3600, retryLimit: 1 });35 if (id) {36 n++;37 busy.add(b.id);38 }39 }4041 // 2. due incremental crawls42 for (const c of due) {43 if (busy.has(c.id) || gated.has(c.id)) continue;44 const id = await queue.send(JOBS.crawlRun, { connectorId: c.id, mode: 'incremental', trigger: 'schedule' }, { singletonKey: `crawl:${c.id}`, priority: PRIORITY[c.priority] ?? 5, expireInSeconds: 6 * 3600, retryLimit: 1 });45 if (id) n++;46 }47 // push next_run_at forward slightly so the same tick does not re-enqueue before the run starts48 if (due.length) {49 await db()50 .update(connectorsTable)51 .set({ nextRunAt: sql`now() + interval '10 minutes'` })52 .where(and(eq(connectorsTable.status, 'active'), or(isNull(connectorsTable.nextRunAt), lte(connectorsTable.nextRunAt, now))));53 }54 if (n) logger.info({ enqueued: n, backfills: backfills.length, gated: gated.size }, 'scheduled due crawls');55 return n;56}57