import type { Prisma } from "@b2bcall/database"; export interface ReservedLead { id: string; phoneNormalized: string; name: string | null; attemptCount: number; } /** * Reserva atômica de leads (agente.md secao 78): `FOR UPDATE SKIP LOCKED` * garante que dois workers concorrentes nunca pegam o mesmo lead — quem * chegar primeiro tranca a linha, o outro pula pra próxima em vez de * esperar (nada de fila de espera aqui, é melhor originar outro lead do * que travar o tick inteiro). `count` normalmente é pequeno (poucas * unidades por tick), então `SELECT ... LIMIT` é barato mesmo sem índice * dedicado além do já existente em `(campaign_id, status, next_attempt_at)`. * * Precisa rodar dentro do MESMO `withTenantContext` transaction que fez a * checagem de RLS — o lock de linha só vale até o commit/rollback da * transação atual. */ export async function reserveLeads( tx: Prisma.TransactionClient, tenantId: string, campaignId: string, count: number, ): Promise { if (count <= 0) return []; const rows = await tx.$queryRaw` SELECT id, phone_normalized AS "phoneNormalized", name, attempt_count AS "attemptCount" FROM leads WHERE tenant_id = ${tenantId}::uuid AND campaign_id = ${campaignId}::uuid AND status IN ('NEW', 'READY') AND (next_attempt_at IS NULL OR next_attempt_at <= now()) ORDER BY next_attempt_at ASC NULLS FIRST, created_at ASC LIMIT ${count} FOR UPDATE SKIP LOCKED `; if (rows.length === 0) return []; await tx.lead.updateMany({ where: { id: { in: rows.map((r) => r.id) } }, data: { status: "RESERVED" }, }); return rows; }