import { and, eq, isNull, lte, or, sql } from 'drizzle-orm'; import { connectors as connectorsTable } from '@rareindex/database'; import { listConnectorMeta, missingRequirements } from '@rareindex/connectors'; import { logger } from '@rareindex/shared'; import { db } from '../lib/db.ts'; import { JOBS, type Queue } from '../lib/queue.ts'; const PRIORITY: Record = { high: 10, medium: 5, low: 1 }; /** * Dynamic crawl scheduling (§141, SPEC §17): every tick, enqueue a crawl for active connectors whose * next_run_at is due (or never ran). Gated connectors (missing API keys) are skipped without a run. * Active backfill campaigns (SPEC §9) are continued with `mode: 'backfill'` as long as the connector * has no crawl in flight, so a campaign survives worker restarts and time-boxed runs. */ export async function scheduleDueCrawls(queue: Queue): Promise { const now = new Date(); const gated = new Set(listConnectorMeta().filter((m) => missingRequirements(m).length > 0).map((m) => m.id)); const due = await db() .select({ id: connectorsTable.id, priority: connectorsTable.priority }) .from(connectorsTable) .where(and(eq(connectorsTable.status, 'active'), or(isNull(connectorsTable.nextRunAt), lte(connectorsTable.nextRunAt, now)))); let n = 0; // pg-boss singleton keys do not dedupe across states reliably; skip connectors that already have a queued/active crawl. 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 }>; const busy = new Set(queued.map((q) => q.id)); 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 }>; for (const r of running) busy.add(r.id); // 1. backfill continuations 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 }>; for (const b of backfills) { if (busy.has(b.id) || gated.has(b.id)) continue; 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 }); if (id) { n++; busy.add(b.id); } } // 2. due incremental crawls for (const c of due) { if (busy.has(c.id) || gated.has(c.id)) continue; 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 }); if (id) n++; } // push next_run_at forward slightly so the same tick does not re-enqueue before the run starts if (due.length) { await db() .update(connectorsTable) .set({ nextRunAt: sql`now() + interval '10 minutes'` }) .where(and(eq(connectorsTable.status, 'active'), or(isNull(connectorsTable.nextRunAt), lte(connectorsTable.nextRunAt, now)))); } if (n) logger.info({ enqueued: n, backfills: backfills.length, gated: gated.size }, 'scheduled due crawls'); return n; }