import { and, eq, sql } from "drizzle-orm"; import { auditEvents, emailDeliveries, events, getDb, guests } from "@album/database"; import { sendEmail } from "./index"; export async function processEmailDelivery() { const db = getDb(); const provider = process.env.EMAIL_PROVIDER ?? "resend"; await db.execute(sql`UPDATE email_deliveries SET status = CASE WHEN provider = 'resend' AND first_attempt_at > now() - interval '23 hours' AND attempts < 5 THEN 'pending' ELSE 'review' END, last_error = 'Worker interrupted; delivery needs reconciliation', updated_at = now() WHERE provider = ${provider} AND status = 'sending' AND updated_at < now() - interval '5 minutes'`); const rows = await db.execute<{ id: string; event_id: string }>(sql`UPDATE email_deliveries SET status = 'sending', attempts = attempts + 1, first_attempt_at = coalesce(first_attempt_at, now()), updated_at = now() WHERE id = (SELECT id FROM email_deliveries WHERE provider = ${provider} AND status = 'pending' AND next_attempt_at <= now() ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 1) RETURNING id, event_id`); const job = rows[0]; if (!job) return false; const scope = and(eq(emailDeliveries.id, job.id), eq(emailDeliveries.eventId, job.event_id)); const [delivery] = await db.select().from(emailDeliveries).where(scope); if (!delivery) return true; try { if (delivery.attempts > 1 && (!delivery.firstAttemptAt || Date.now() - delivery.firstAttemptAt.getTime() >= 23 * 3600000)) { await db.update(emailDeliveries).set({ status: "review", lastError: "Idempotency window expired", updatedAt: new Date() }).where(scope); return true; } const [event] = await db.select().from(events).where(eq(events.id, job.event_id)); const now = new Date(); const visible = event && (event.status !== "draft" || (event.publishAt && event.publishAt <= now)) && (!event.publishAt || event.publishAt <= now) && event.galleryPolicy !== "never" && (!event.galleryVisibleAt || event.galleryVisibleAt <= now) && (event.galleryPolicy === "automatic" || event.galleryReleasedAt || event.galleryVisibleAt) && delivery.payload.text.endsWith(`/e/${event.slug}`); const optedIn = await db.select({ id: guests.id }).from(guests).where(and(eq(guests.eventId, job.event_id), eq(guests.notifyWhenReady, true), sql`lower(${guests.email}) = ${delivery.recipient}`)).limit(1); if (!visible || !optedIn.length) { await db.update(emailDeliveries).set({ status: "review", lastError: "Gallery visibility or recipient consent changed", updatedAt: now }).where(scope); return true; } const result = await sendEmail(delivery.payload, { provider: delivery.provider, idempotencyKey: `gallery-ready/${delivery.id}` }); await db.transaction(async (tx) => { await tx.update(emailDeliveries).set({ status: "sent", providerId: result.id, lastError: null, updatedAt: new Date() }).where(scope); await tx.update(guests).set({ notifiedAt: new Date(), updatedAt: new Date() }).where(and(eq(guests.eventId, job.event_id), eq(guests.notifyWhenReady, true), sql`lower(${guests.email}) = ${delivery.recipient}`)); await tx.insert(auditEvents).values({ eventId: job.event_id, action: "guest.email.sent", subjectType: "email", subjectId: delivery.id, metadata: { providerId: result.id } }); }); } catch { const retry = delivery.provider === "resend" && delivery.attempts < 5 && delivery.firstAttemptAt && Date.now() - delivery.firstAttemptAt.getTime() < 23 * 3600000; await db.update(emailDeliveries).set({ status: retry ? "pending" : "review", nextAttemptAt: new Date(Date.now() + Math.min(3600, 30 * 2 ** delivery.attempts) * 1000), lastError: "Delivery could not be confirmed; no recipient data logged", updatedAt: new Date() }).where(scope); } return true; }