From 25f923f530aa12e8be76ab901d2392ad739a2131 Mon Sep 17 00:00:00 2001 From: weeihan Date: Thu, 23 Jul 2026 16:37:51 +0800 Subject: [PATCH] =?UTF-8?q?feat(db):=20phase=204=20group=201=20=E2=80=94?= =?UTF-8?q?=20lib/=20settings=20+=20notifications=20to=20Drizzle?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Convert lib/settings.ts, lib/notifications/{in-app,email,capa-escalation, effectiveness-recheck}.ts from Supabase PostgREST to Drizzle asAdmin queries. Drop supabase arg from all call sites in app/api/ and cron routes. Rewrite notification unit tests to mock @/lib/db/with-user instead of SupabaseClient. Co-Authored-By: Claude Sonnet 4.6 --- app/api/capa/[id]/verify/route.ts | 2 +- app/api/capa/route.ts | 2 +- app/api/cron/capa-escalation/route.ts | 4 +- app/api/cron/effectiveness-recheck/route.ts | 4 +- app/api/dashboard/ai/risk-flags/route.ts | 2 +- app/api/incidents/[id]/ai/rca-draft/route.ts | 2 +- .../incidents/[id]/ai/triage-suggest/route.ts | 2 +- app/api/incidents/[id]/close/route.ts | 2 +- app/api/incidents/[id]/similar/route.ts | 2 +- app/api/incidents/ai/quality-check/route.ts | 2 +- app/api/incidents/route.ts | 7 +- lib/notifications/capa-escalation.ts | 123 ++++++++------- lib/notifications/effectiveness-recheck.ts | 83 +++++----- lib/notifications/email.ts | 49 +++--- lib/notifications/in-app.ts | 34 +++-- lib/settings.ts | 20 ++- .../lib/notifications/capa-escalation.test.ts | 143 +++++++----------- tests/lib/notifications/email.test.ts | 59 ++++---- tests/lib/notifications/in-app.test.ts | 54 ++++--- 19 files changed, 296 insertions(+), 300 deletions(-) diff --git a/app/api/capa/[id]/verify/route.ts b/app/api/capa/[id]/verify/route.ts index d00d7d8..f29eb20 100644 --- a/app/api/capa/[id]/verify/route.ts +++ b/app/api/capa/[id]/verify/route.ts @@ -56,7 +56,7 @@ export async function POST( }) if (capa.owner_user_id) { - await createInAppNotifications(supabase, [{ + await createInAppNotifications([{ userId: capa.owner_user_id, title: body.verdict === 'verified' ? 'Your CAPA action was verified' diff --git a/app/api/capa/route.ts b/app/api/capa/route.ts index c7880e5..66e4aeb 100644 --- a/app/api/capa/route.ts +++ b/app/api/capa/route.ts @@ -67,7 +67,7 @@ export async function POST(request: NextRequest) { p_new_value: { incident_id, description, owner_user_id, department, due_date }, }) - await createInAppNotifications(supabase, [{ + await createInAppNotifications([{ userId: owner_user_id, title: `CAPA assigned to you, due ${due_date}`, link: `/hse/capa/${capa.id}`, diff --git a/app/api/cron/capa-escalation/route.ts b/app/api/cron/capa-escalation/route.ts index cbc4b41..5772ed3 100644 --- a/app/api/cron/capa-escalation/route.ts +++ b/app/api/cron/capa-escalation/route.ts @@ -2,7 +2,6 @@ export const dynamic = 'force-dynamic' import { timingSafeEqual } from 'crypto' import { NextRequest, NextResponse } from 'next/server' -import { createClient } from '@/lib/supabase/server' import { escalateOverdueCapa } from '@/lib/notifications/capa-escalation' export async function GET(request: NextRequest) { @@ -15,7 +14,6 @@ export async function GET(request: NextRequest) { return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) } - const supabase = await createClient() - const { notified } = await escalateOverdueCapa(supabase) + const { notified } = await escalateOverdueCapa() return NextResponse.json({ ok: true, notified }) } diff --git a/app/api/cron/effectiveness-recheck/route.ts b/app/api/cron/effectiveness-recheck/route.ts index fb49fd1..1861193 100644 --- a/app/api/cron/effectiveness-recheck/route.ts +++ b/app/api/cron/effectiveness-recheck/route.ts @@ -2,7 +2,6 @@ export const dynamic = 'force-dynamic' import { timingSafeEqual } from 'crypto' import { NextRequest, NextResponse } from 'next/server' -import { createClient } from '@/lib/supabase/server' import { sendEffectivenessRecheckNotifications } from '@/lib/notifications/effectiveness-recheck' export async function GET(request: NextRequest) { @@ -15,7 +14,6 @@ export async function GET(request: NextRequest) { return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) } - const supabase = await createClient() - const { notified } = await sendEffectivenessRecheckNotifications(supabase) + const { notified } = await sendEffectivenessRecheckNotifications() return NextResponse.json({ ok: true, notified }) } diff --git a/app/api/dashboard/ai/risk-flags/route.ts b/app/api/dashboard/ai/risk-flags/route.ts index a72f8c4..c027420 100644 --- a/app/api/dashboard/ai/risk-flags/route.ts +++ b/app/api/dashboard/ai/risk-flags/route.ts @@ -92,7 +92,7 @@ export async function POST() { return NextResponse.json({ error: 'Rate limit: wait 60 seconds between AI requests' }, { status: 429 }) } - const deepseekKey = await getApiKey(supabase, 'DEEPSEEK_API_KEY') + const deepseekKey = await getApiKey('DEEPSEEK_API_KEY') const client = createDeepSeekClient(deepseekKey) let res: Awaited> diff --git a/app/api/incidents/[id]/ai/rca-draft/route.ts b/app/api/incidents/[id]/ai/rca-draft/route.ts index 19da493..864d422 100644 --- a/app/api/incidents/[id]/ai/rca-draft/route.ts +++ b/app/api/incidents/[id]/ai/rca-draft/route.ts @@ -28,7 +28,7 @@ export async function POST( if ((recentCount ?? 0) > 0) return NextResponse.json({ error: 'Rate limited — please wait 60 seconds' }, { status: 429 }) - const deepseekKey = await getApiKey(supabase, 'DEEPSEEK_API_KEY') + const deepseekKey = await getApiKey('DEEPSEEK_API_KEY') const client = createDeepSeekClient(deepseekKey) const { data: incident } = await supabase diff --git a/app/api/incidents/[id]/ai/triage-suggest/route.ts b/app/api/incidents/[id]/ai/triage-suggest/route.ts index fe308b4..445686e 100644 --- a/app/api/incidents/[id]/ai/triage-suggest/route.ts +++ b/app/api/incidents/[id]/ai/triage-suggest/route.ts @@ -28,7 +28,7 @@ export async function POST( if ((recentCount ?? 0) > 0) return NextResponse.json({ error: 'Rate limited — please wait 60 seconds' }, { status: 429 }) - const deepseekKey = await getApiKey(supabase, 'DEEPSEEK_API_KEY') + const deepseekKey = await getApiKey('DEEPSEEK_API_KEY') const client = createDeepSeekClient(deepseekKey) const { data: incident } = await supabase diff --git a/app/api/incidents/[id]/close/route.ts b/app/api/incidents/[id]/close/route.ts index 3ce7595..64b975c 100644 --- a/app/api/incidents/[id]/close/route.ts +++ b/app/api/incidents/[id]/close/route.ts @@ -51,7 +51,7 @@ export async function POST( }) if (incident.reported_by) { - await createInAppNotifications(supabase, [{ + await createInAppNotifications([{ userId: incident.reported_by, title: `Your incident report ${incident.reference_no ?? ''} has been closed`, link: '/reporter', diff --git a/app/api/incidents/[id]/similar/route.ts b/app/api/incidents/[id]/similar/route.ts index da3ba94..6bac5af 100644 --- a/app/api/incidents/[id]/similar/route.ts +++ b/app/api/incidents/[id]/similar/route.ts @@ -18,7 +18,7 @@ export async function GET( const supabase = await createClient() - const googleAiKey = await getApiKey(supabase, 'GOOGLE_AI_API_KEY') + const googleAiKey = await getApiKey('GOOGLE_AI_API_KEY') const { data: incident } = await supabase .from('incidents') diff --git a/app/api/incidents/ai/quality-check/route.ts b/app/api/incidents/ai/quality-check/route.ts index 4373925..0493f83 100644 --- a/app/api/incidents/ai/quality-check/route.ts +++ b/app/api/incidents/ai/quality-check/route.ts @@ -22,7 +22,7 @@ export async function POST(request: NextRequest) { if ((recentCount ?? 0) > 0) return NextResponse.json({ error: 'Rate limited — please wait 60 seconds' }, { status: 429 }) - const deepseekKey = await getApiKey(supabase, 'DEEPSEEK_API_KEY') + const deepseekKey = await getApiKey('DEEPSEEK_API_KEY') const client = createDeepSeekClient(deepseekKey) let body: { description?: string; incident_type?: string } diff --git a/app/api/incidents/route.ts b/app/api/incidents/route.ts index d3bcd53..80f9f65 100644 --- a/app/api/incidents/route.ts +++ b/app/api/incidents/route.ts @@ -165,8 +165,8 @@ async function handlePost(request: Request) { // WhatsApp alert — fire-and-forget alongside email ;(async () => { try { - const phoneNumberId = await getApiKey(supabase, 'META_WHATSAPP_PHONE_NUMBER_ID') - const accessToken = await getApiKey(supabase, 'META_WHATSAPP_ACCESS_TOKEN') + const phoneNumberId = await getApiKey('META_WHATSAPP_PHONE_NUMBER_ID') + const accessToken = await getApiKey('META_WHATSAPP_ACCESS_TOKEN') const { data: siteData } = await supabase .from('sites').select('name').eq('id', zone.site_id).single() const siteName = (siteData as { name: string } | null)?.name ?? 'Unknown' @@ -177,7 +177,6 @@ async function handlePost(request: Request) { .eq('site_id', zone.site_id) await createInAppNotifications( - supabase, (recipients ?? []).map((r: { id: string }) => ({ userId: r.id, title: `New ${input.incident_type.replace(/_/g, ' ')} incident ${incident.reference_no ?? ''} at ${siteName}`, @@ -205,7 +204,7 @@ async function handlePost(request: Request) { // Embed description asynchronously for future similarity search const supabaseForEmbed = supabase import('@/lib/settings').then(({ getApiKey }) => - getApiKey(supabaseForEmbed, 'GOOGLE_AI_API_KEY').then(googleAiKey => + getApiKey('GOOGLE_AI_API_KEY').then(googleAiKey => import('@/lib/claude/embed').then(({ embedText }) => embedText(input.description.trim(), googleAiKey).then(embedding => supabase.from('incidents').update({ diff --git a/lib/notifications/capa-escalation.ts b/lib/notifications/capa-escalation.ts index bec2499..422d877 100644 --- a/lib/notifications/capa-escalation.ts +++ b/lib/notifications/capa-escalation.ts @@ -1,5 +1,8 @@ +import 'server-only' import { Resend } from 'resend' -import type { SupabaseClient } from '@supabase/supabase-js' +import { asAdmin } from '@/lib/db/with-user' +import { capaActions, incidents, users, notificationsLog } from '@/lib/db/schema' +import { and, eq, inArray, not } from 'drizzle-orm' import { sendWhatsAppMessage } from '@/lib/notifications/whatsapp' import { createInAppNotifications } from '@/lib/notifications/in-app' import { getApiKey } from '@/lib/settings' @@ -27,17 +30,25 @@ const THRESHOLD_SUBJECT: Record = { overdue_7d: '[IMS] URGENT: CAPA action 7 days OVERDUE', } -export async function escalateOverdueCapa( - supabase: SupabaseClient -): Promise<{ notified: number }> { - const { data: capas } = await supabase - .from('capa_actions') - .select(` - id, description, due_date, incident_id, owner_user_id, - incidents (reference_no, site_id), - owner:users!owner_user_id (email, name, phone) - `) - .not('status', 'in', '(verified,closed)') +export async function escalateOverdueCapa(): Promise<{ notified: number }> { + const capas = await asAdmin(db => + db.select({ + id: capaActions.id, + description: capaActions.description, + dueDate: capaActions.dueDate, + incidentId: capaActions.incidentId, + ownerUserId: capaActions.ownerUserId, + incidentRef: incidents.referenceNo, + siteId: incidents.siteId, + ownerEmail: users.email, + ownerName: users.name, + ownerPhone: users.phone, + }) + .from(capaActions) + .leftJoin(incidents, eq(capaActions.incidentId, incidents.id)) + .leftJoin(users, eq(capaActions.ownerUserId, users.id)) + .where(not(inArray(capaActions.status, ['verified', 'closed']))) + ) if (!capas || capas.length === 0) return { notified: 0 } @@ -46,49 +57,52 @@ export async function escalateOverdueCapa( const siteUrl = process.env.NEXT_PUBLIC_SITE_URL ?? 'http://localhost:3000' let notified = 0 - const whatsappPhoneNumberId = await getApiKey(supabase, 'META_WHATSAPP_PHONE_NUMBER_ID') - const whatsappAccessToken = await getApiKey(supabase, 'META_WHATSAPP_ACCESS_TOKEN') + const whatsappPhoneNumberId = await getApiKey('META_WHATSAPP_PHONE_NUMBER_ID') + const whatsappAccessToken = await getApiKey('META_WHATSAPP_ACCESS_TOKEN') for (const capa of capas) { - const threshold = getEscalationThreshold((capa as { due_date: string }).due_date) + const threshold = getEscalationThreshold(capa.dueDate) if (!threshold) continue - const { data: alreadySent } = await supabase - .from('notifications_log') - .select('id') - .eq('capa_id', capa.id) - .eq('channel', 'email') - .eq('status', threshold) - .limit(1) - .maybeSingle() + const alreadySent = await asAdmin(db => + db.select({ id: notificationsLog.id }) + .from(notificationsLog) + .where(and( + eq(notificationsLog.capaId, capa.id), + eq(notificationsLog.channel, 'email'), + eq(notificationsLog.status, threshold), + )) + .limit(1) + ) + if (alreadySent.length > 0) continue - if (alreadySent) continue - - const ownerEmail = (capa.owner as unknown as { email: string } | null)?.email - const ownerName = (capa.owner as unknown as { name: string } | null)?.name ?? 'Owner' + const ownerEmail = capa.ownerEmail + const ownerName = capa.ownerName ?? 'Owner' if (!ownerEmail) continue - const incidentRef = (capa.incidents as unknown as { reference_no: string | null } | null)?.reference_no ?? capa.incident_id + const incidentRef = capa.incidentRef ?? capa.incidentId const capaUrl = `${siteUrl}/ims/hse/capa/${capa.id}` const subject = THRESHOLD_SUBJECT[threshold] const html = `

Hi ${ownerName},

CAPA action for incident ${incidentRef} requires attention.

-

Action: ${(capa as { description: string }).description}

-

Due date: ${(capa as { due_date: string }).due_date}

+

Action: ${capa.description}

+

Due date: ${capa.dueDate}

View CAPA

` - const text = `CAPA ${incidentRef}: ${(capa as { description: string }).description}\nDue: ${(capa as { due_date: string }).due_date}\n${capaUrl}` + const text = `CAPA ${incidentRef}: ${capa.description}\nDue: ${capa.dueDate}\n${capaUrl}` const to = [ownerEmail] - const siteId = (capa.incidents as unknown as { site_id: string } | null)?.site_id - if (siteId && ['due_today', 'overdue_3d', 'overdue_7d'].includes(threshold)) { - const { data: hseUsers } = await supabase - .from('users') - .select('email') - .eq('site_id', siteId) - .in('role', ['hse', 'supervisor', 'management']) - if (hseUsers) to.push(...hseUsers.map((u: { email: string }) => u.email)) + if (capa.siteId && ['due_today', 'overdue_3d', 'overdue_7d'].includes(threshold)) { + const hseUsers = await asAdmin(db => + db.select({ email: users.email }) + .from(users) + .where(and( + eq(users.siteId, capa.siteId!), + inArray(users.role, ['hse', 'supervisor', 'management']), + )) + ) + to.push(...hseUsers.map(u => u.email).filter(Boolean) as string[]) } const { error } = await resend.emails.send({ from, to: [...new Set(to)], subject, html, text }) @@ -99,13 +113,13 @@ export async function escalateOverdueCapa( // WhatsApp for urgent thresholds only (owner must have a phone number) if (['overdue_3d', 'overdue_7d'].includes(threshold)) { - const ownerPhone = (capa.owner as unknown as { phone: string | null } | null)?.phone ?? '' + const ownerPhone = capa.ownerPhone ?? '' if (ownerPhone) { try { await sendWhatsAppMessage( ownerPhone, 'ims_capa_overdue', - [incidentRef, (capa as { description: string }).description, (capa as { due_date: string }).due_date], + [incidentRef ?? '', capa.description, capa.dueDate], whatsappPhoneNumberId, whatsappAccessToken, ) @@ -115,21 +129,22 @@ export async function escalateOverdueCapa( } } - await supabase.from('notifications_log').insert({ - capa_id: capa.id, - channel: 'email', - recipient: to.join(','), - status: threshold, - }) + await asAdmin(db => + db.insert(notificationsLog).values({ + capaId: capa.id, + channel: 'email', + recipient: [...new Set(to)].join(','), + status: threshold, + }) + ) - const ownerUserId = (capa as { owner_user_id: string | null }).owner_user_id - if (ownerUserId) { - await createInAppNotifications(supabase, [{ - userId: ownerUserId, - title: `${THRESHOLD_SUBJECT[threshold].replace('[IMS] ', '')} — ${incidentRef}`, + if (capa.ownerUserId) { + await createInAppNotifications([{ + userId: capa.ownerUserId, + title: `${THRESHOLD_SUBJECT[threshold].replace('[IMS] ', '')} — ${capa.incidentRef ?? capa.incidentId}`, link: `/hse/capa/${capa.id}`, - incidentId: capa.incident_id as string, - capaId: capa.id as string, + incidentId: capa.incidentId, + capaId: capa.id, }]) } diff --git a/lib/notifications/effectiveness-recheck.ts b/lib/notifications/effectiveness-recheck.ts index 20faaf0..f042e8a 100644 --- a/lib/notifications/effectiveness-recheck.ts +++ b/lib/notifications/effectiveness-recheck.ts @@ -1,5 +1,8 @@ +import 'server-only' import { Resend } from 'resend' -import type { SupabaseClient } from '@supabase/supabase-js' +import { asAdmin } from '@/lib/db/with-user' +import { capaActions, incidents, users } from '@/lib/db/schema' +import { and, eq, inArray, isNotNull, lt } from 'drizzle-orm' import { sendWhatsAppMessage } from '@/lib/notifications/whatsapp' import { getApiKey } from '@/lib/settings' @@ -27,22 +30,31 @@ export function shouldSendRecheck( return recheckDate <= today } -export async function sendEffectivenessRecheckNotifications( - supabase: SupabaseClient, -): Promise<{ notified: number }> { +export async function sendEffectivenessRecheckNotifications(): Promise<{ notified: number }> { const today = new Date().toISOString().split('T')[0] - const { data: capas } = await supabase - .from('capa_actions') - .select(` - id, description, effectiveness_recheck_date, effectiveness_recheck_round, - verified_at, incident_id, - incidents (reference_no), - verifier:users!verified_by (email, name, phone) - `) - .in('status', ['verified', 'closed']) - .not('effectiveness_recheck_date', 'is', null) - .lt('effectiveness_recheck_round', 3) + const capas = await asAdmin(db => + db.select({ + id: capaActions.id, + description: capaActions.description, + effectivenessRecheckDate: capaActions.effectivenessRecheckDate, + effectivenessRecheckRound: capaActions.effectivenessRecheckRound, + verifiedAt: capaActions.verifiedAt, + incidentId: capaActions.incidentId, + incidentRef: incidents.referenceNo, + verifierEmail: users.email, + verifierName: users.name, + verifierPhone: users.phone, + }) + .from(capaActions) + .leftJoin(incidents, eq(capaActions.incidentId, incidents.id)) + .leftJoin(users, eq(capaActions.verifiedBy, users.id)) + .where(and( + inArray(capaActions.status, ['verified', 'closed']), + isNotNull(capaActions.effectivenessRecheckDate), + lt(capaActions.effectivenessRecheckRound, 3), + )) + ) if (!capas || capas.length === 0) return { notified: 0 } @@ -53,8 +65,8 @@ export async function sendEffectivenessRecheckNotifications( let whatsappPhoneId: string | null = null let whatsappToken: string | null = null try { - whatsappPhoneId = await getApiKey(supabase, 'META_WHATSAPP_PHONE_NUMBER_ID') - whatsappToken = await getApiKey(supabase, 'META_WHATSAPP_ACCESS_TOKEN') + whatsappPhoneId = await getApiKey('META_WHATSAPP_PHONE_NUMBER_ID') + whatsappToken = await getApiKey('META_WHATSAPP_ACCESS_TOKEN') } catch { // WhatsApp not configured — email only } @@ -62,21 +74,20 @@ export async function sendEffectivenessRecheckNotifications( let notified = 0 for (const capa of capas) { - const recheckDate = (capa as { effectiveness_recheck_date: string | null }).effectiveness_recheck_date - const round = (capa as { effectiveness_recheck_round: number }).effectiveness_recheck_round + const recheckDate = capa.effectivenessRecheckDate + const round = capa.effectivenessRecheckRound if (!shouldSendRecheck(recheckDate, round, today)) continue - const verifierEmail = (capa.verifier as unknown as { email: string } | null)?.email - const verifierName = (capa.verifier as unknown as { name: string } | null)?.name ?? 'HSE Officer' - const verifierPhone = (capa.verifier as unknown as { phone: string | null } | null)?.phone ?? '' + const verifierEmail = capa.verifierEmail + const verifierName = capa.verifierName ?? 'HSE Officer' + const verifierPhone = capa.verifierPhone ?? '' if (!verifierEmail) continue - const incidentRef = (capa.incidents as unknown as { reference_no: string | null } | null)?.reference_no - ?? (capa as { incident_id: string }).incident_id + const incidentRef = capa.incidentRef ?? capa.incidentId const roundLabel = getRoundLabel(round) - const capaDesc = (capa as { description: string }).description - const capaUrl = `${siteUrl}/ims/hse/capa/${(capa as { id: string }).id}` + const capaDesc = capa.description + const capaUrl = `${siteUrl}/ims/hse/capa/${capa.id}` const subject = `[IMS] ${roundLabel} effectiveness check — ${incidentRef}` const html = ` @@ -99,22 +110,22 @@ export async function sendEffectivenessRecheckNotifications( sendWhatsAppMessage( verifierPhone, 'ims_effectiveness_recheck', - [roundLabel, incidentRef, capaDesc], + [roundLabel, incidentRef ?? '', capaDesc], whatsappPhoneId, whatsappToken, ).catch(err => console.error('WhatsApp recheck error:', err)) } // Advance round - const verifiedAt = (capa as { verified_at: string }).verified_at - const nextDate = getNextRecheckDate(verifiedAt, round) - await supabase - .from('capa_actions') - .update({ - effectiveness_recheck_round: round + 1, - effectiveness_recheck_date: nextDate ?? null, - }) - .eq('id', (capa as { id: string }).id) + const nextDate = getNextRecheckDate(capa.verifiedAt!.toISOString(), round) + await asAdmin(db => + db.update(capaActions) + .set({ + effectivenessRecheckRound: round + 1, + effectivenessRecheckDate: nextDate ?? null, + }) + .where(eq(capaActions.id, capa.id)) + ) notified++ } diff --git a/lib/notifications/email.ts b/lib/notifications/email.ts index 08a232b..3c10948 100644 --- a/lib/notifications/email.ts +++ b/lib/notifications/email.ts @@ -1,5 +1,8 @@ +import 'server-only' import { Resend } from 'resend' -import { createClient } from '@/lib/supabase/server' +import { asAdmin } from '@/lib/db/with-user' +import { users, incidents, sites } from '@/lib/db/schema' +import { and, eq, inArray } from 'drizzle-orm' import { newIncidentTemplate } from '@/lib/notifications/templates/new-incident' export async function sendNewIncidentEmail( @@ -8,29 +11,37 @@ export async function sendNewIncidentEmail( reference_no: string, incidentType: string, ): Promise { - const supabase = await createClient() - - const { data: recipients } = await supabase - .from('users') - .select('email, name, role') - .in('role', ['supervisor', 'hse']) - .eq('site_id', siteId) + const recipients = await asAdmin(db => + db.select({ email: users.email }) + .from(users) + .where(and( + inArray(users.role, ['supervisor', 'hse']), + eq(users.siteId, siteId), + )) + ) if (!recipients || recipients.length === 0) return - const to = recipients.map((r: { email: string }) => r.email).filter(Boolean) + const to = recipients.map(r => r.email).filter(Boolean) as string[] if (to.length === 0) return - const { data: incident } = await supabase - .from('incidents') - .select('reported_at, sites (name), reporter:users!reported_by (name)') - .eq('id', incidentId) - .single() - - const siteName = (incident?.sites as unknown as { name: string } | null)?.name ?? 'Unknown Site' - const reporterName = (incident?.reporter as unknown as { name: string } | null)?.name ?? 'Unknown' - const reportedAt = incident?.reported_at - ? new Date(incident.reported_at).toLocaleString('en-MY', { timeZone: 'Asia/Kuala_Lumpur' }) + const incidentRows = await asAdmin(db => + db.select({ + reportedAt: incidents.reportedAt, + siteName: sites.name, + reporterName: users.name, + }) + .from(incidents) + .leftJoin(sites, eq(incidents.siteId, sites.id)) + .leftJoin(users, eq(incidents.reportedBy, users.id)) + .where(eq(incidents.id, incidentId)) + .limit(1) + ) + const incident = incidentRows[0] + const siteName = incident?.siteName ?? 'Unknown Site' + const reporterName = incident?.reporterName ?? 'Unknown' + const reportedAt = incident?.reportedAt + ? new Date(incident.reportedAt).toLocaleString('en-MY', { timeZone: 'Asia/Kuala_Lumpur' }) : '-' const siteUrl = process.env.NEXT_PUBLIC_SITE_URL ?? 'http://localhost:3000' diff --git a/lib/notifications/in-app.ts b/lib/notifications/in-app.ts index 4ce6bb2..d002d92 100644 --- a/lib/notifications/in-app.ts +++ b/lib/notifications/in-app.ts @@ -1,4 +1,6 @@ -import type { SupabaseClient } from '@supabase/supabase-js' +import 'server-only' +import { asAdmin } from '@/lib/db/with-user' +import { notificationsLog } from '@/lib/db/schema' export interface InAppNotification { userId: string @@ -8,11 +10,7 @@ export interface InAppNotification { capaId?: string } -// Inserts go through the create_in_app_notification SECURITY DEFINER RPC: -// notifications_log INSERT is RLS-restricted to elevated roles, but reporters -// must still be able to trigger alerts to supervisors/HSE. export async function createInAppNotifications( - supabase: SupabaseClient, notifications: InAppNotification[], ): Promise<{ created: number }> { const seen = new Set() @@ -24,18 +22,22 @@ export async function createInAppNotifications( if (seen.has(key)) continue seen.add(key) - const { error } = await supabase.rpc('create_in_app_notification', { - p_recipient: n.userId, - p_title: n.title, - p_link: n.link ?? null, - p_incident_id: n.incidentId ?? null, - p_capa_id: n.capaId ?? null, - }) - if (error) { - console.error('in-app notification error:', error) - continue + try { + await asAdmin(db => + db.insert(notificationsLog).values({ + channel: 'in_app', + recipient: n.userId, + recipientUserId: n.userId, + title: n.title, + link: n.link ?? null, + incidentId: n.incidentId ?? null, + capaId: n.capaId ?? null, + }) + ) + created++ + } catch (e) { + console.error('in-app notification error:', e) } - created++ } return { created } diff --git a/lib/settings.ts b/lib/settings.ts index 8dbe0f1..8cf5652 100644 --- a/lib/settings.ts +++ b/lib/settings.ts @@ -1,13 +1,17 @@ -import type { SupabaseClient } from '@supabase/supabase-js' +import 'server-only' +import { asAdmin } from '@/lib/db/with-user' +import { appSettings } from '@/lib/db/schema' +import { eq } from 'drizzle-orm' -export async function getApiKey(supabase: SupabaseClient, key: string): Promise { +export async function getApiKey(key: string): Promise { try { - const { data } = await supabase - .from('app_settings') - .select('value') - .eq('key', key) - .single() - if (data?.value) return data.value + const rows = await asAdmin(db => + db.select({ value: appSettings.value }) + .from(appSettings) + .where(eq(appSettings.key, key)) + .limit(1) + ) + if (rows[0]?.value) return rows[0].value } catch { // fall through to env } diff --git a/tests/lib/notifications/capa-escalation.test.ts b/tests/lib/notifications/capa-escalation.test.ts index 96138f0..db55aaa 100644 --- a/tests/lib/notifications/capa-escalation.test.ts +++ b/tests/lib/notifications/capa-escalation.test.ts @@ -1,3 +1,4 @@ +// @vitest-environment node import { describe, it, expect, vi, beforeEach } from 'vitest' import { getEscalationThreshold, escalateOverdueCapa } from '@/lib/notifications/capa-escalation' @@ -17,6 +18,16 @@ vi.mock('@/lib/settings', () => ({ getApiKey: vi.fn().mockResolvedValue('test-cred'), })) +vi.mock('@/lib/notifications/in-app', () => ({ + createInAppNotifications: vi.fn().mockResolvedValue({ created: 1 }), +})) + +// At top of file, before other vi.mock calls +let mockAsAdminImpl: (fn: (db: unknown) => unknown) => unknown +vi.mock('@/lib/db/with-user', () => ({ + asAdmin: vi.fn().mockImplementation((fn: (db: unknown) => unknown) => mockAsAdminImpl(fn)), +})) + describe('getEscalationThreshold', () => { function daysFromNow(n: number): string { const d = new Date() @@ -49,106 +60,60 @@ describe('getEscalationThreshold', () => { }) }) -describe('escalateOverdueCapa — WhatsApp', () => { +describe('escalateOverdueCapa', () => { function daysFromNow(n: number): string { const d = new Date() d.setDate(d.getDate() + n) return d.toISOString().split('T')[0] } - function makeSupabaseMock(overrides: { - capas?: unknown[] - alreadySent?: unknown - hseUsers?: unknown[] - }) { - const { capas = [], alreadySent = null, hseUsers = [] } = overrides + beforeEach(() => { vi.clearAllMocks() }) - // Each .from() call returns a fresh chainable builder - // We track call order: first from('capa_actions'), then from('notifications_log'), then from('users'), then from('notifications_log').insert - let fromCallIndex = 0 + it('returns notified:0 when no active capas', async () => { + let call = 0 + mockAsAdminImpl = (fn) => { + call++ + if (call === 1) return fn({ select: () => ({ from: () => ({ leftJoin: () => ({ leftJoin: () => ({ where: () => [] }) }) }) }) }) + return fn({}) + } + const { notified } = await escalateOverdueCapa() + expect(notified).toBe(0) + }) - const makeChain = (resolvedValue: unknown) => { - const chain: Record = {} - const methods = ['select', 'not', 'eq', 'in', 'limit', 'maybeSingle', 'insert'] - for (const m of methods) { - chain[m] = vi.fn(() => chain) + it('skips capa with no threshold match', async () => { + let call = 0 + mockAsAdminImpl = (fn) => { + call++ + if (call === 1) { + // Return capa due in 2 days (no threshold) + return Promise.resolve([{ + id: 'c1', description: 'Fix it', dueDate: daysFromNow(2), + incidentId: 'i1', ownerUserId: 'u1', + incidentRef: 'SITE-202507-0001', siteId: 's1', + ownerEmail: 'owner@example.com', ownerName: 'Owner', ownerPhone: null, + }]) } - // Terminal: awaiting the chain resolves to resolvedValue - Object.defineProperty(chain, 'then', { - get() { - return (resolve: (v: unknown) => unknown) => Promise.resolve(resolvedValue).then(resolve) - }, - }) - return chain + return Promise.resolve([]) } - - const supabase = { - from: vi.fn(() => { - const index = fromCallIndex++ - if (index === 0) return makeChain({ data: capas, error: null }) // capa_actions - if (index === 1) return makeChain({ data: alreadySent, error: null }) // notifications_log check - if (index === 2) return makeChain({ data: hseUsers, error: null }) // hse users for CC - return makeChain({ data: null, error: null }) // notifications_log insert - }), - } - - return supabase as unknown as import('@supabase/supabase-js').SupabaseClient - } - - beforeEach(() => { - vi.clearAllMocks() + const { notified } = await escalateOverdueCapa() + expect(notified).toBe(0) }) - it('sends WhatsApp when threshold is overdue_7d and owner has phone', async () => { - const { sendWhatsAppMessage: mockWA } = await import('@/lib/notifications/whatsapp') - - const supabase = makeSupabaseMock({ - capas: [ - { - id: 'capa-1', - description: 'Fix safety barrier', - due_date: daysFromNow(-7), - incident_id: 'inc-1', - incidents: { reference_no: 'KL-202501-0001', site_id: 'site-1' }, - owner: { email: 'owner@test.com', name: 'Alice', phone: '60123456789' }, - }, - ], - alreadySent: null, - hseUsers: [], - }) - - await escalateOverdueCapa(supabase) - - expect(mockWA).toHaveBeenCalledWith( - expect.stringMatching(/^\d+$/), - 'ims_capa_overdue', - expect.arrayContaining([expect.any(String)]), - 'test-cred', - 'test-cred', - ) - }) - - it('does NOT send WhatsApp when threshold is warning_3d', async () => { - const { sendWhatsAppMessage: mockWA } = await import('@/lib/notifications/whatsapp') - vi.clearAllMocks() - - const supabase = makeSupabaseMock({ - capas: [ - { - id: 'capa-2', - description: 'Inspect equipment', - due_date: daysFromNow(3), - incident_id: 'inc-2', - incidents: { reference_no: 'KL-202501-0002', site_id: 'site-1' }, - owner: { email: 'owner@test.com', name: 'Bob', phone: '60129876543' }, - }, - ], - alreadySent: null, - hseUsers: [], - }) - - await escalateOverdueCapa(supabase) - - expect(mockWA).not.toHaveBeenCalled() + it('sends email and WhatsApp for overdue_3d capa with phone', async () => { + const { sendWhatsAppMessage } = await import('@/lib/notifications/whatsapp') + let call = 0 + mockAsAdminImpl = (fn) => { + call++ + if (call === 1) return Promise.resolve([{ + id: 'c1', description: 'Fix it', dueDate: daysFromNow(-3), + incidentId: 'i1', ownerUserId: 'u1', + incidentRef: 'SITE-202507-0001', siteId: 's1', + ownerEmail: 'owner@example.com', ownerName: 'Owner', ownerPhone: '+60123456789', + }]) + return Promise.resolve([]) // alreadySent = [], hseUsers = [], insert void + } + const { notified } = await escalateOverdueCapa() + expect(notified).toBe(1) + expect(sendWhatsAppMessage).toHaveBeenCalled() }) }) diff --git a/tests/lib/notifications/email.test.ts b/tests/lib/notifications/email.test.ts index 8033430..22c900c 100644 --- a/tests/lib/notifications/email.test.ts +++ b/tests/lib/notifications/email.test.ts @@ -1,8 +1,8 @@ +// @vitest-environment node import { describe, it, expect, vi, beforeEach } from 'vitest' -const { mockSend, mockFrom } = vi.hoisted(() => ({ +const { mockSend } = vi.hoisted(() => ({ mockSend: vi.fn(), - mockFrom: vi.fn(), })) vi.mock('resend', () => { @@ -11,42 +11,36 @@ vi.mock('resend', () => { return { Resend: MockResend } }) -vi.mock('@/lib/supabase/server', () => ({ - createClient: vi.fn().mockResolvedValue({ from: mockFrom }), +// Mock asAdmin: first call returns recipients, second returns incident+join rows +let asAdminCallCount = 0 +const mockRecipients = [ + { email: 'supervisor@setiacorp.com' }, + { email: 'hse@setiacorp.com' }, +] +const mockIncidentRows = [{ + reportedAt: new Date('2026-07-10T09:00:00Z'), + siteName: 'SCW1', + reporterName: 'John Doe', +}] + +vi.mock('@/lib/db/with-user', () => ({ + asAdmin: vi.fn().mockImplementation((fn: (db: unknown) => unknown) => { + asAdminCallCount++ + if (asAdminCallCount % 2 === 1) { + // first call: recipients query + return Promise.resolve(mockRecipients) + } + // second call: incident join query + return Promise.resolve(mockIncidentRows) + }), })) import { sendNewIncidentEmail } from '@/lib/notifications/email' -const mockIncidentRow = { - reported_at: '2026-07-10T09:00:00Z', - sites: { name: 'SCW1' }, - reporter: { name: 'John Doe' }, -} - -function setupMocks(recipientData: object[]) { - mockFrom.mockImplementation((table: string) => { - if (table === 'users') { - return { - select: vi.fn().mockReturnThis(), - in: vi.fn().mockReturnThis(), - eq: vi.fn().mockResolvedValue({ data: recipientData, error: null }), - } - } - return { - select: vi.fn().mockReturnThis(), - eq: vi.fn().mockReturnThis(), - single: vi.fn().mockResolvedValue({ data: mockIncidentRow, error: null }), - } - }) -} - beforeEach(() => { vi.clearAllMocks() + asAdminCallCount = 0 mockSend.mockResolvedValue({ data: { id: 'email-123' }, error: null }) - setupMocks([ - { email: 'supervisor@setiacorp.com', name: 'Ahmad', role: 'supervisor' }, - { email: 'hse@setiacorp.com', name: 'Priya', role: 'hse' }, - ]) }) describe('sendNewIncidentEmail', () => { @@ -59,7 +53,8 @@ describe('sendNewIncidentEmail', () => { }) it('does not send if no recipients', async () => { - setupMocks([]) + const { asAdmin } = await import('@/lib/db/with-user') + vi.mocked(asAdmin).mockResolvedValueOnce([]) await sendNewIncidentEmail('inc-001', 'site-001', 'SCW1-202607-0001', 'near_miss') expect(mockSend).not.toHaveBeenCalled() }) diff --git a/tests/lib/notifications/in-app.test.ts b/tests/lib/notifications/in-app.test.ts index 037b633..0be4e64 100644 --- a/tests/lib/notifications/in-app.test.ts +++ b/tests/lib/notifications/in-app.test.ts @@ -1,60 +1,58 @@ +// @vitest-environment node import { describe, it, expect, vi, beforeEach } from 'vitest' -import { createInAppNotifications } from '@/lib/notifications/in-app' -import type { SupabaseClient } from '@supabase/supabase-js' -function makeSupabaseMock(rpcResult: { error: unknown } = { error: null }) { - return { - rpc: vi.fn().mockResolvedValue(rpcResult), - } as unknown as SupabaseClient -} +// Mock asAdmin before importing the module under test +const mockInsert = vi.fn() +const mockValues = vi.fn().mockResolvedValue([]) +vi.mock('@/lib/db/with-user', () => ({ + asAdmin: vi.fn().mockImplementation(fn => + fn({ + insert: mockInsert.mockReturnValue({ values: mockValues }), + }) + ), +})) + +import { createInAppNotifications } from '@/lib/notifications/in-app' describe('createInAppNotifications', () => { beforeEach(() => { vi.clearAllMocks() + mockInsert.mockReturnValue({ values: mockValues }) + mockValues.mockResolvedValue([]) }) - it('calls create_in_app_notification RPC once per recipient', async () => { - const supabase = makeSupabaseMock() - const { created } = await createInAppNotifications(supabase, [ + it('inserts once per recipient', async () => { + const { created } = await createInAppNotifications([ { userId: 'u1', title: 'New incident', link: '/hse/incidents/i1', incidentId: 'i1' }, { userId: 'u2', title: 'New incident', link: '/hse/incidents/i1', incidentId: 'i1' }, ]) expect(created).toBe(2) - expect(supabase.rpc).toHaveBeenCalledTimes(2) - expect(supabase.rpc).toHaveBeenCalledWith('create_in_app_notification', { - p_recipient: 'u1', - p_title: 'New incident', - p_link: '/hse/incidents/i1', - p_incident_id: 'i1', - p_capa_id: null, - }) + expect(mockValues).toHaveBeenCalledTimes(2) }) it('skips entries missing userId or title', async () => { - const supabase = makeSupabaseMock() - const { created } = await createInAppNotifications(supabase, [ + const { created } = await createInAppNotifications([ { userId: '', title: 'x' }, { userId: 'u1', title: '' }, ]) expect(created).toBe(0) - expect(supabase.rpc).not.toHaveBeenCalled() + expect(mockValues).not.toHaveBeenCalled() }) - it('counts only successful inserts when RPC errors', async () => { - const supabase = makeSupabaseMock({ error: { message: 'boom' } }) - const { created } = await createInAppNotifications(supabase, [ + it('counts only successful inserts; failed insert is caught', async () => { + mockValues.mockRejectedValueOnce(new Error('DB error')) + const { created } = await createInAppNotifications([ { userId: 'u1', title: 'x' }, ]) expect(created).toBe(0) }) - it('deduplicates recipients for the same notification', async () => { - const supabase = makeSupabaseMock() - const { created } = await createInAppNotifications(supabase, [ + it('deduplicates same userId+title+incidentId', async () => { + const { created } = await createInAppNotifications([ { userId: 'u1', title: 'same', incidentId: 'i1' }, { userId: 'u1', title: 'same', incidentId: 'i1' }, ]) expect(created).toBe(1) - expect(supabase.rpc).toHaveBeenCalledTimes(1) + expect(mockValues).toHaveBeenCalledTimes(1) }) })