import { and, eq, inArray, isNull, lte, sql } from "drizzle-orm"; import { db } from "@/lib/db"; import logger from "@/lib/logger"; import { reminders } from "./schema"; import { notify } from "./notify"; const REMINDER_LOCK_KEY = 7_777_777; export async function scheduleReminder(input: { householdId: string; entityType: string; entityId: string; fireAt: Date; createdBy: string; channel?: string; title?: string; body?: string; }) { await db .insert(reminders) .values({ householdId: input.householdId, entityType: input.entityType, entityId: input.entityId, fireAt: input.fireAt, channel: input.channel ?? "auto", title: input.title ?? null, body: input.body ?? null, createdBy: input.createdBy, firedAt: null, }) .onConflictDoUpdate({ target: [reminders.entityType, reminders.entityId], set: { fireAt: input.fireAt, title: input.title ?? null, body: input.body ?? null, firedAt: null, createdBy: input.createdBy, }, }); } export async function cancelReminder(entityType: string, entityId: string) { await db .delete(reminders) .where(and(eq(reminders.entityType, entityType), eq(reminders.entityId, entityId))); } export async function listReminders(entityType: string, entityId: string) { return db .select() .from(reminders) .where(and(eq(reminders.entityType, entityType), eq(reminders.entityId, entityId))); } export async function tickReminders() { let dueReminders: (typeof reminders.$inferSelect)[] = []; try { await db.transaction(async (tx) => { const lockRows = await tx.execute<{ acquired: boolean }>( sql`SELECT pg_try_advisory_xact_lock(${REMINDER_LOCK_KEY}) AS acquired`, ); if (!lockRows[0]?.acquired) return; const now = new Date(); dueReminders = await tx .select() .from(reminders) .where(and(lte(reminders.fireAt, now), isNull(reminders.firedAt))); if (dueReminders.length > 0) { await tx .update(reminders) .set({ firedAt: now }) .where( inArray( reminders.id, dueReminders.map((r) => r.id), ), ); } }); } catch (err) { logger.error({ err }, "reminder tick error"); return; } await Promise.allSettled( dueReminders.map(async (reminder) => { if (!reminder.createdBy) return; try { await notify(reminder.createdBy, { title: reminder.title ?? "Reminder", body: reminder.body ?? "You have a reminder", url: reminder.entityType === "notes.note" ? `/notes/${reminder.entityId}` : "/", channels: ["push", "inapp"], }); } catch (err) { logger.error({ reminderId: reminder.id, err }, "reminder delivery failed"); } }), ); } let workerTimer: ReturnType | null = null; export function startReminderWorker() { if (workerTimer) return; workerTimer = setInterval(() => { tickReminders().catch((err) => logger.error({ err }, "reminder worker uncaught error")); }, 30_000); logger.info("reminder worker started (30s tick)"); }