import { Injectable, Logger } from '@nestjs/common'; import { OnEvent } from '@nestjs/event-emitter'; import { InjectDataSource } from '@nestjs/typeorm'; import { DataSource } from 'typeorm'; import { PrismaService } from '../../common/prisma.service'; import { PushAdapter, NotificationChannel } from './notification.adapters'; import { EmailClientService } from './email-client.service'; import { SmsClientService } from './sms-client.service'; export type NotificationChannelType = 'EMAIL' | 'SMS' | 'PUSH' | 'IN_APP'; const UUID_RE = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i; @Injectable() export class NotificationsService { private readonly logger = new Logger(NotificationsService.name); private readonly channels: Map; constructor( private prisma: PrismaService, @InjectDataSource() private readonly dataSource: DataSource, private emailClient: EmailClientService, private smsClient: SmsClientService, private pushAdapter: PushAdapter, ) { this.channels = new Map([ ['EMAIL', { send: (to, subject, body) => this.emailClient.sendEmail({ to, subject, text: body }).then((r) => r.queued) }], ['SMS', { send: (to, _subject, body) => this.smsClient.sendSms({ to, message: body }).then((r) => r.queued) }], ['PUSH', this.pushAdapter as NotificationChannel], ]); } /** * Send notification using template key and context * @param templateKey - Template code from NotificationTemplate table * @param recipient - User/Passenger ID or email/phone * @param context - Variables to interpolate in template * @param channels - Optional array of channels to use (defaults to user preferences) */ async send( templateKey: string, recipient: string, context: Record, channels?: NotificationChannelType[], ): Promise<{ queued: boolean; channels: string[] }> { const template = await this.prisma.notificationTemplate.findUnique({ where: { code: templateKey }, }); if (!template || !template.active) { this.logger.warn(`Template ${templateKey} not found or inactive`); return { queued: false, channels: [] }; } const { subject, body } = this.interpolate(template, context); // Channel resolution: explicit argument wins; otherwise honor the template's declared // channel(s); otherwise fall back to the recipient's preferences. let targetChannels: NotificationChannelType[]; if (channels) { targetChannels = channels; } else if (template.channel) { targetChannels = this.parseTemplateChannels(template.channel); } else { targetChannels = await this.getUserPreferredChannels(recipient); } // Channels successfully handed off (in-app persisted / email+SMS enqueued to RabbitMQ). // NOTE: enqueue is fire-and-forget — this is NOT a delivery confirmation. const queuedChannels: string[] = []; if (targetChannels.includes('IN_APP')) { await this.createInAppNotification(recipient, subject, body, context); queuedChannels.push('IN_APP'); } for (const channelType of targetChannels) { if (channelType === 'IN_APP') continue; const adapter = this.channels.get(channelType); if (!adapter) { this.logger.warn(`No adapter for channel: ${channelType}`); continue; } const recipientAddress = await this.getRecipientAddress(recipient, channelType); if (!recipientAddress) { this.logger.warn(`No ${channelType} address for recipient: ${recipient}`); continue; } const queued = await adapter.send(recipientAddress, subject, body, context); if (queued) { queuedChannels.push(channelType); } } return { queued: queuedChannels.length > 0, channels: queuedChannels }; } /** * Parses a template's `channel` column (e.g. "EMAIL" or "EMAIL,SMS") into valid channel * types, always including IN_APP so an in-app record is created. */ private parseTemplateChannels(channel: string): NotificationChannelType[] { const valid: NotificationChannelType[] = ['EMAIL', 'SMS', 'PUSH', 'IN_APP']; const parsed = channel .split(',') .map((c) => c.trim().toUpperCase()) .filter((c): c is NotificationChannelType => valid.includes(c as NotificationChannelType)); return Array.from(new Set(['IN_APP', ...parsed])); } private async createInAppNotification( recipient: string, title: string, body: string, context: Record, ): Promise { let passengerId = recipient; if (!UUID_RE.test(recipient)) { const iamUserId = await this.resolveIamUserId(recipient); if (!iamUserId) { this.logger.warn(`Could not find passenger for recipient: ${recipient}`); return; } const passenger = await this.prisma.passenger.findUnique({ where: { iamUserId } }); if (!passenger) { this.logger.warn(`Could not find passenger for recipient: ${recipient}`); return; } passengerId = passenger.id; } await this.prisma.notification.create({ data: { passengerId, title, body, category: (context.category as any) || 'SYSTEM', deepLink: context.deepLink as string, metadata: context as any, }, }); } private interpolate( template: { subject?: string | null; bodyTemplate: string }, context: Record, ): { subject: string; body: string } { return { subject: this.applyVars(template.subject || 'Notification', context), body: this.applyVars(template.bodyTemplate, context), }; } /** Replaces {{variable}} placeholders in a string with values from the context. */ private applyVars(text: string, context: Record): string { let out = text; for (const [key, value] of Object.entries(context)) { const regex = new RegExp(`{{\\s*${key}\\s*}}`, 'g'); out = out.replace(regex, String(value)); } return out; } private async getUserPreferredChannels(recipient: string): Promise { const iamUserId = await this.resolveIamUserId(recipient); const preferences = iamUserId ? await this.prisma.userPreferences.findUnique({ where: { iamUserId } }) : null; if (!preferences) { return ['IN_APP', 'EMAIL']; } const channels: NotificationChannelType[] = ['IN_APP']; if (preferences.emailEnabled) channels.push('EMAIL'); if (preferences.smsEnabled) channels.push('SMS'); if (preferences.pushEnabled) channels.push('PUSH'); return channels; } private async getRecipientAddress( recipient: string, channel: NotificationChannelType, ): Promise { const iamUserId = await this.resolveIamUserId(recipient); if (!iamUserId) return null; const contact = await this.resolveContactInfo(iamUserId); switch (channel) { case 'EMAIL': return contact.email; case 'SMS': return contact.phone; case 'PUSH': return iamUserId; default: return null; } } private async resolveIamUserId(recipient: string): Promise { if (UUID_RE.test(recipient)) { const passenger = await this.prisma.passenger.findUnique({ where: { id: recipient } }); return passenger?.iamUserId ?? recipient; } const rows = await this.dataSource.query<{ id: string }[]>( `SELECT id FROM iam.users WHERE email = $1 OR phone_number = $1 LIMIT 1`, [recipient], ); return rows[0]?.id ?? null; } private async resolveContactInfo(iamUserId: string): Promise<{ email: string | null; phone: string | null }> { const rows = await this.dataSource.query<{ email: string; phone_number: string | null }[]>( `SELECT email, phone_number FROM iam.users WHERE id = $1 LIMIT 1`, [iamUserId], ); return { email: rows[0]?.email ?? null, phone: rows[0]?.phone_number ?? null }; } private sanitize(value: string): string { return value .replace(/[\r\n]/g, ' ') .replace(/[<>&"']/g, (c) => ({ '<': '<', '>': '>', '&': '&', '"': '"', "'": ''' }[c] ?? c)); } getForPassenger(passengerId: string) { return this.prisma.notification.findMany({ where: { passengerId }, orderBy: { createdAt: 'desc' }, take: 50, }); } markRead(id: string) { return this.prisma.notification.update({ where: { id }, data: { read: true } }); } async markAllRead(passengerId: string) { await this.prisma.notification.updateMany({ where: { passengerId, read: false }, data: { read: true }, }); return { updated: true }; } @OnEvent('booking.created') async onBookingCreated(payload: any) { const booking = payload.booking; await this.send( 'booking.created', booking.passengerId, { bookingRef: booking.bookingRef, amount: this.formatAmount(booking), currency: booking.displayCurrency ?? 'ETB', category: 'BOOKING', deepLink: `edr://bookings/${booking.bookingRef}`, }, // For now, always notify the travelling passenger on every channel. ['IN_APP', 'EMAIL', 'SMS'], ); } /** * Payment succeeded → one combined "payment successful, here is your ticket" notification. * Email carries the full ticket (HTML + QR); SMS is a short pointer to view it. The shallow * event payload is re-fetched with the relations needed to render the ticket. */ @OnEvent('payment.succeeded') async onPaymentSucceeded(payload: any) { const passengerId = payload.booking.passengerId; const bookingId = payload.booking.id; const booking = await this.prisma.booking.findUnique({ where: { id: bookingId }, include: { schedule: { include: { originStation: true, destinationStation: true, train: true } }, seats: { include: { seat: { include: { coach: { include: { coachType: true } } } } } }, }, }); const ticket = await this.prisma.ticket.findUnique({ where: { bookingId } }); const ref = booking?.bookingRef ?? payload.booking.bookingRef; const amount = this.formatAmount(booking ?? payload.booking); const currency = (booking ?? payload.booking).displayCurrency ?? 'ETB'; const ticketUrl = `${process.env.PORTAL_URL ?? 'http://localhost:5174'}/booking/confirmation?ref=${ref}`; // IN_APP — always created. await this.createInAppNotification( passengerId, 'Payment successful', `Your payment of ${amount} ${currency} for booking ${ref} was successful. Your ticket is ready.`, { category: 'PAYMENT', deepLink: `edr://tickets/${ref}` }, ); // Ticket not ready (generation failed/raced) — fall back to a payment-only confirmation. if (!ticket || !booking) { this.logger.warn(`payment.succeeded: ticket not ready for booking ${ref}; sending payment-only confirmation`); const text = `EDR: Payment of ${amount} ${currency} received for booking ${ref}. Your ticket is being prepared.`; await this.deliverEmail(passengerId, `Payment received — ${ref}`, text); await this.deliverSms(passengerId, text); return; } // SMS — short pointer (no HTML/QR over SMS). await this.deliverSms( passengerId, `EDR: Booking ${ref} confirmed, ${amount} ${currency} paid. Show ref ${ref} at the gate or view your ticket: ${ticketUrl}`, ); // EMAIL — rich HTML ticket with plain-text fallback. await this.deliverEmail( passengerId, `Your EDR ticket — ${ref}`, this.buildTicketEmailText(booking, amount, currency, ticketUrl), this.buildTicketEmailHtml(booking, ticket, amount, currency, ticketUrl), ); } private async deliverEmail(recipient: string, subject: string, text: string, html?: string): Promise { const to = await this.getRecipientAddress(recipient, 'EMAIL'); if (!to) { this.logger.warn(`No EMAIL address for recipient: ${recipient}`); return; } await this.emailClient.sendEmail({ to, subject, text, html }); } private async deliverSms(recipient: string, message: string): Promise { const to = await this.getRecipientAddress(recipient, 'SMS'); if (!to) { this.logger.warn(`No SMS address for recipient: ${recipient}`); return; } await this.smsClient.sendSms({ to, message }); } private buildTicketEmailText(booking: any, amount: string, currency: string, url: string): string { const s = booking.schedule ?? {}; const dep = s.departureAt ? new Date(s.departureAt).toLocaleString('en-GB') : 'TBD'; const passengers = (booking.seats ?? []).map((bs: any) => bs.passengerName).filter(Boolean).join(', '); return [ `Booking ${booking.bookingRef} confirmed.`, `${s.originStation?.name ?? ''} -> ${s.destinationStation?.name ?? ''}`, `Train: ${s.train?.name ?? s.train?.number ?? ''}`, `Departs: ${dep}`, passengers ? `Passengers: ${passengers}` : '', `Total paid: ${amount} ${currency}`, `View your ticket: ${url}`, ].filter(Boolean).join('\n'); } private buildTicketEmailHtml(booking: any, ticket: any, amount: string, currency: string, url: string): string { const s = booking.schedule ?? {}; const fmt = (d: any) => d ? new Date(d).toLocaleString('en-GB', { dateStyle: 'medium', timeStyle: 'short' }) : 'TBD'; const seatRows = (booking.seats ?? []) .map((bs: any) => { const coach = bs.seat?.coach?.number ?? '-'; const seatNo = bs.seat?.seatNumber ?? '-'; const cls = bs.seat?.coach?.coachType?.name ?? '-'; return ` ${bs.passengerName ?? ''} ${coach} ${seatNo} ${cls} `; }) .join(''); return `

Ethio-Djibouti Railway

Payment successful — your ticket is ready

Booking reference: ${booking.bookingRef}

From ${s.originStation?.name ?? ''} (${s.originStation?.code ?? ''})
To ${s.destinationStation?.name ?? ''} (${s.destinationStation?.code ?? ''})
Train ${s.train?.name ?? s.train?.number ?? ''}
Departs ${fmt(s.departureAt)}
Arrives ${fmt(s.arrivalAt)}

Passengers

${seatRows}
Name Coach Seat Class

Show this QR code at the gate

Ticket QR code
Total paid ${amount} ${currency}

© Ethio-Djibouti Railway. All rights reserved.

`; } @OnEvent('payment.failed') async onPaymentFailed(payload: any) { const booking = payload.booking; await this.send( 'payment.failed', booking.passengerId, { bookingRef: booking.bookingRef, category: 'PAYMENT', deepLink: `edr://bookings/${booking.bookingRef}`, }, ['IN_APP', 'EMAIL', 'SMS'], ); } @OnEvent('booking.cancelled') async onBookingCancelled(payload: any) { const booking = payload.booking; await this.send( 'booking.cancelled', booking.passengerId, { bookingRef: booking.bookingRef, // refundAmount is computed in ETB minor units in BookingsService.cancel(). refundAmount: ((payload.refundAmount ?? 0) / 100).toFixed(2), currency: 'ETB', category: 'BOOKING', deepLink: `edr://bookings/${booking.bookingRef}`, }, ['IN_APP', 'EMAIL', 'SMS'], ); } /** * Formats a booking's payable amount from minor units into a major-unit string. * Money is stored as integer minor units (e.g. 59600 santim) to avoid floating-point * drift; we divide by 100 only here, at the display edge. e.g. 59600 -> "596.00". */ private formatAmount(booking: any): string { const minor = booking.displayTotalMinor ?? booking.totalMinor ?? 0; return (minor / 100).toFixed(2); } }