SPB Git forge

spb/rareindex

Public
54commits 1branches 0releases
7.1 MBsize
maindefault branch
10 days agolast push
TypeScript 61.9% HTML 37.2% SQL 0.7%
3.4 KB · 57 lines typescript
Raw Blame History
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