diff --git a/src/modules/notifications/controllers/notifications.controller.ts b/src/modules/notifications/controllers/notifications.controller.ts index 71e5335..c50a7e0 100644 --- a/src/modules/notifications/controllers/notifications.controller.ts +++ b/src/modules/notifications/controllers/notifications.controller.ts @@ -177,6 +177,19 @@ export class NotificationsController { return { smsCount }; } + @UseGuards(AdminAuthGuard) + @ApiBearerAuth() + @Post('admin/notifications/test-notis') + @ApiOperation({ summary: 'Send a test in-app notification (socket) to the current admin' }) + async createTestNotis(@RestId() restaurantId: string, @AdminId() adminId: string) { + const notification = await this.notificationService.sendTestInAppNotification(adminId, restaurantId); + return { + message: 'Test in-app notification queued', + notification, + }; + } + + // super admin endpoints @UseGuards(SuperAdminAuthGuard) diff --git a/src/modules/notifications/notifications.gateway.ts b/src/modules/notifications/notifications.gateway.ts index 2bbd80d..e87411b 100644 --- a/src/modules/notifications/notifications.gateway.ts +++ b/src/modules/notifications/notifications.gateway.ts @@ -37,10 +37,19 @@ export class NotificationsGateway implements OnGatewayConnection, OnGatewayDisco ) {} async handleConnection(client: Socket) { + this.logger.log( + `[WS] connection attempt socketId=${client.id} transport=${client.conn.transport.name} ` + + `hasAuthToken=${Boolean(client.handshake.auth?.token || client.handshake.auth?.authorization)} ` + + `hasQueryToken=${Boolean(client.handshake.query?.token || client.handshake.query?.authorization)} ` + + `hasHeaderAuth=${Boolean(client.handshake.headers.authorization)}`, + ); + try { await this.authenticateConnection(client); const authenticatedClient = client as AuthenticatedSocket; - this.logger.log(`Admin connected: ${authenticatedClient.adminId}`); + this.logger.log( + `[WS] connected socketId=${client.id} adminId=${authenticatedClient.adminId} restId=${authenticatedClient.restId}`, + ); this.handleJoinRoom(authenticatedClient); // Start ping interval for this client @@ -51,7 +60,7 @@ export class NotificationsGateway implements OnGatewayConnection, OnGatewayDisco // this.pingIntervals.set(client.id, intervalId); } catch (error) { this.logger.error( - `Connection authentication failed: ${error instanceof Error ? error.message : 'Unknown error'}`, + `[WS] connection auth failed socketId=${client.id}: ${error instanceof Error ? error.message : 'Unknown error'}`, ); void client.emit('error', { message: 'Authentication failed' }); client.disconnect(); @@ -59,8 +68,11 @@ export class NotificationsGateway implements OnGatewayConnection, OnGatewayDisco } handleDisconnect(client: Socket) { - this.logger.log(`Client disconnected: ${client.id}`); - this.handleLeaveRoom(client); + const authenticatedClient = client as AuthenticatedSocket; + this.logger.log( + `[WS] disconnect socketId=${client.id} adminId=${authenticatedClient.adminId ?? 'n/a'} restId=${authenticatedClient.restId ?? 'n/a'}`, + ); + this.handleLeaveRoom(authenticatedClient); // Clear ping interval for this client // const intervalId = this.pingIntervals.get(client.id); @@ -72,13 +84,41 @@ export class NotificationsGateway implements OnGatewayConnection, OnGatewayDisco sendInAppNotification(repipient: { adminId: string; restaurantId: string }, payload: IInAppNotificationPayload) { const room = this.getRoom(repipient.adminId, repipient.restaurantId); + const adapter = this.server.sockets?.adapter ?? this.server.adapter; + const socketsMap = this.server.sockets?.sockets ?? this.server.sockets; + const socketsInRoom = adapter?.rooms?.get(room); + const roomSize = socketsInRoom?.size ?? 0; + const connectedSocketIds = + socketsMap instanceof Map ? [...socketsMap.keys()] : [...(Object.keys(socketsMap ?? {}) as string[])]; + const allRooms = adapter?.rooms ? [...adapter.rooms.keys()] : []; + + this.logger.log( + `[WS] emit start room=${room} adminId=${repipient.adminId} restaurantId=${repipient.restaurantId} ` + + `notificationId=${payload.notificationId} subject=${payload.subject} roomSize=${roomSize} ` + + `namespaceConnected=${connectedSocketIds.length}`, + ); + + if (roomSize === 0) { + this.logger.warn( + `[WS] emit to EMPTY room=${room} — no client joined this room. ` + + `connectedSockets=[${connectedSocketIds.join(', ') || 'none'}] ` + + `allRooms=[${allRooms.join(', ') || 'none'}]`, + ); + } else { + this.logger.log(`[WS] room members socketIds=[${[...(socketsInRoom ?? [])].join(', ')}]`); + } - this.logger.log(`Sending in app notification to admin: ${repipient.adminId}`); this.server.to(room).emit('notifications', payload, async () => { + this.logger.log( + `[WS] ACK received room=${room} notificationId=${payload.notificationId} (client acknowledged)`, + ); // 3. ACK (delivered) // await this.notificationService.markAsDelivered(payload.notificationId); }); - this.logger.log(`In app notification sent to admin: ${repipient.adminId}`); + + this.logger.log( + `[WS] emit queued event=notifications room=${room} roomSize=${roomSize} notificationId=${payload.notificationId}`, + ); } private getRoom(adminId: string, restaurantId: string) { @@ -86,16 +126,30 @@ export class NotificationsGateway implements OnGatewayConnection, OnGatewayDisco } handleLeaveRoom(client: AuthenticatedSocket) { - const room = this.getRoom(client.adminId!, client.restId!); + if (!client.adminId || !client.restId) { + this.logger.warn(`[WS] leave skipped socketId=${client.id} — missing adminId/restId`); + return; + } + + const room = this.getRoom(client.adminId, client.restId); void client.leave(room); - this.logger.log(`Admin ${client.adminId} (Client ${client.id}) left room: ${room}`); + this.logger.log(`[WS] left room=${room} adminId=${client.adminId} socketId=${client.id}`); void client.emit('left', { room, message: 'Successfully Admin left room' }); } handleJoinRoom(client: AuthenticatedSocket) { - const room = this.getRoom(client.adminId!, client.restId!); + if (!client.adminId || !client.restId) { + this.logger.warn(`[WS] join skipped socketId=${client.id} — missing adminId/restId`); + return; + } + + const room = this.getRoom(client.adminId, client.restId); void client.join(room); - this.logger.log(`Admin ${client.adminId} (Client ${client.id}) joined room: ${room}`); + const adapter = this.server.sockets?.adapter ?? this.server.adapter; + const roomSize = adapter?.rooms?.get(room)?.size ?? 0; + this.logger.log( + `[WS] joined room=${room} adminId=${client.adminId} restId=${client.restId} socketId=${client.id} roomSize=${roomSize}`, + ); void client.emit('joined', { room, message: 'Successfully Admin joined room' }); } diff --git a/src/modules/notifications/processors/in-app.processor.ts b/src/modules/notifications/processors/in-app.processor.ts index 1280fc9..fdc69b2 100644 --- a/src/modules/notifications/processors/in-app.processor.ts +++ b/src/modules/notifications/processors/in-app.processor.ts @@ -2,9 +2,7 @@ import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq'; import { Logger } from '@nestjs/common'; import { Job } from 'bullmq'; import { CreateRequestContext } from '@mikro-orm/core'; -import { InjectRepository } from '@mikro-orm/nestjs'; import { EntityManager } from '@mikro-orm/postgresql'; -import { Notification } from '../entities/notification.entity'; import { InAppNotificationQueueJob } from '../interfaces/jobs-queue.interface'; import { NotificationQueueNameEnum } from '../constants/queue'; import { NotificationsGateway } from '../notifications.gateway'; @@ -14,7 +12,6 @@ export class InAppProcessor extends WorkerHost { private readonly logger = new Logger(InAppProcessor.name); constructor( - @InjectRepository(Notification) private readonly em: EntityManager, private readonly notificationsGateway: NotificationsGateway, ) { @@ -25,10 +22,18 @@ export class InAppProcessor extends WorkerHost { async process(job: Job): Promise { const { recipient, subject, body, notificationId } = job.data; - this.logger.log(`Processing InApp notification - Recipient: ${JSON.stringify(recipient)}, subject: ${subject}`); + this.logger.log( + `[InAppProcessor] processing jobId=${job.id} attempt=${job.attemptsMade + 1} ` + + `recipient=${JSON.stringify(recipient)} subject=${subject} notificationId=${notificationId}`, + ); - if (!recipient.adminId) { - this.logger.warn(`Skipping in-app notification job ${job.id}: missing adminId`); + if (!recipient?.adminId) { + this.logger.warn(`[InAppProcessor] skipping job ${job.id}: missing adminId`); + return; + } + + if (!recipient.restaurantId) { + this.logger.warn(`[InAppProcessor] skipping job ${job.id}: missing restaurantId`); return; } @@ -38,20 +43,30 @@ export class InAppProcessor extends WorkerHost { { subject, body, notificationId }, ); - this.logger.log(`InApp notification sent successfully to ${recipient.adminId}`); + this.logger.log( + `[InAppProcessor] gateway emit done for admin=${recipient.adminId} restaurant=${recipient.restaurantId}`, + ); } catch (error) { - this.logger.error(`Error processing InApp notification job ${job.id}:`, error); + this.logger.error(`[InAppProcessor] error processing job ${job.id}:`, error); throw error; } } + @OnWorkerEvent('active') + onActive(job: Job) { + this.logger.log(`[InAppProcessor] job ${job.id} became active`); + } + @OnWorkerEvent('completed') onCompleted(job: Job) { - this.logger.log(`Notification job ${job.id} completed`); + this.logger.log(`[InAppProcessor] job ${job.id} completed`); } @OnWorkerEvent('failed') onFailed(job: Job, error: Error) { - this.logger.error(`Notification job ${job.id} failed:`, error); + this.logger.error( + `[InAppProcessor] job ${job.id} failed after ${job.attemptsMade} attempt(s): ${error.message}`, + error.stack, + ); } } diff --git a/src/modules/notifications/services/notification-queue.service.ts b/src/modules/notifications/services/notification-queue.service.ts index d556c51..6ef1732 100644 --- a/src/modules/notifications/services/notification-queue.service.ts +++ b/src/modules/notifications/services/notification-queue.service.ts @@ -119,6 +119,12 @@ export class NotificationQueueService { } async addBulkInAppNotifications(jobs: InAppNotificationQueueJob[]): Promise { try { + this.logger.log( + `[addBulkInAppNotifications] queueing ${jobs.length} job(s): ${jobs + .map(j => `admin=${j.recipient.adminId} rest=${j.recipient.restaurantId} notif=${j.notificationId}`) + .join(' | ')}`, + ); + const queueJobs = jobs.map(job => ({ name: 'in-app-notification', data: job, @@ -132,10 +138,12 @@ export class NotificationQueueService { }, })); - await this.inAppQueue.addBulk(queueJobs); - this.logger.log(`Added ${jobs.length} InApp notification jobs to queue`); + const added = await this.inAppQueue.addBulk(queueJobs); + this.logger.log( + `[addBulkInAppNotifications] added ${added.length} job(s) to queue "${NotificationQueueNameEnum.IN_APP}": ids=${added.map(j => j.id).join(', ')}`, + ); } catch (error) { - this.logger.error(`Failed to add bulk notifications to queue:`, error); + this.logger.error(`[addBulkInAppNotifications] failed to add jobs to queue:`, error); throw error; } } diff --git a/src/modules/notifications/services/notification.service.ts b/src/modules/notifications/services/notification.service.ts index 81c2e7f..29c682d 100644 --- a/src/modules/notifications/services/notification.service.ts +++ b/src/modules/notifications/services/notification.service.ts @@ -24,6 +24,11 @@ export class NotificationService { async sendNotification(params: NotifRequest): Promise { const { recipients, message, metadata, restaurantId } = params; + this.logger.log( + `[sendNotification] start restaurant=${restaurantId} title=${message.title} recipients=${recipients.length} ` + + `admins=${recipients.filter(r => 'adminId' in r).length} users=${recipients.filter(r => 'userId' in r).length}`, + ); + // create Database notifications const notifications = await this.createAdminBulkNotifications( recipients.map(recipient => ({ @@ -35,18 +40,35 @@ export class NotificationService { })), ); + this.logger.log(`[sendNotification] created ${notifications.length} DB notification(s)`); + // get admin prefrences const preference = await this.preferenceService.findByRestaurantAndType(restaurantId, message.title); - if (preference?.channels?.length === 0) { + this.logger.log( + `[sendNotification] preference=${preference ? `id=${preference.id} channels=${JSON.stringify(preference.channels)}` : 'NOT FOUND'}`, + ); + + if (!preference) { + this.logger.warn( + `[sendNotification] skipping channels — no preference for restaurant ${restaurantId}, title ${message.title}`, + ); + return notifications; + } + + if (preference.channels?.length === 0) { this.logger.warn(`Notification type is NONE for restaurant ${restaurantId}, title ${message.title}`); return notifications; } // send in app notification (admins only — user records must not emit to admin socket rooms) - if (preference?.channels?.includes(NotifChannelEnum.IN_APP)) { + if (preference.channels?.includes(NotifChannelEnum.IN_APP)) { + this.logger.log(`[sendNotification] IN_APP channel enabled, building jobs`); const inAppJobs = recipients.flatMap((recipient, index) => { if (!('adminId' in recipient) || !recipient.adminId) { + this.logger.debug( + `[sendNotification] skipping recipient index=${index} (not admin): ${JSON.stringify(recipient)}`, + ); return []; } @@ -61,16 +83,23 @@ export class NotificationService { }); if (inAppJobs.length > 0) { + this.logger.log( + `[sendNotification] queueing ${inAppJobs.length} in-app job(s): ${inAppJobs.map(j => j.recipient.adminId).join(', ')}`, + ); await this.queueService.addBulkInAppNotifications(inAppJobs); } else { this.logger.warn( `No admin recipients for in-app notification: restaurant ${restaurantId}, title ${message.title}`, ); } + } else { + this.logger.warn( + `[sendNotification] IN_APP channel NOT enabled for restaurant ${restaurantId}, title ${message.title}, channels=${JSON.stringify(preference.channels)}`, + ); } // add sms notifications to queue - if (preference?.channels?.includes(NotifChannelEnum.SMS)) { + if (preference.channels?.includes(NotifChannelEnum.SMS)) { await this.queueService.addBulkSmsNotifications( recipients.map(recipient => ({ recipient, @@ -293,4 +322,40 @@ export class NotificationService { async getSmsCountByRestaurantId(restaurantId: string): Promise { return this.smsLogRepository.getSmsCountByRestaurantId(restaurantId); } + + async sendTestInAppNotification(adminId: string, restaurantId: string): Promise { + const title = NotifTitleEnum.ORDER_CREATED; + const content = 'Test in-app notification'; + + this.logger.log(`[sendTestInAppNotification] start adminId=${adminId} restaurantId=${restaurantId}`); + + const [notification] = await this.createAdminBulkNotifications([ + { + restaurantId, + title, + content, + adminId, + userId: null, + }, + ]); + + this.logger.log( + `[sendTestInAppNotification] created notification id=${notification.id}, queueing in-app job`, + ); + + await this.queueService.addBulkInAppNotifications([ + { + recipient: { adminId, restaurantId }, + subject: title, + body: content, + notificationId: notification.id, + }, + ]); + + this.logger.log( + `[sendTestInAppNotification] queued successfully admin=${adminId} restaurant=${restaurantId} notificationId=${notification.id}`, + ); + + return notification; + } }