From ead1e096a5ea2a0975729e9e8edf5081f3a5938f Mon Sep 17 00:00:00 2001 From: morteza-mortezai Date: Tue, 9 Dec 2025 10:19:28 +0330 Subject: [PATCH] notifcation --- src/modules/notifications/constants/queue.ts | 4 + .../dto/create-preference.dto.ts | 22 +- .../dto/update-preference.dto.ts | 4 + .../notification-preference.entity.ts | 20 +- .../entities/notification.entity.ts | 33 +-- .../notification-queue.interface.ts | 10 +- .../interfaces/notification.interface.ts | 6 + .../listeners/notification.listeners.ts | 29 +- .../notifications/notifications.controller.ts | 204 +++++++------- .../notifications/notifications.module.ts | 12 +- .../processors/notification.processor.ts | 248 ------------------ .../processors/push.processor.ts | 99 +++++++ .../notifications/processors/sms.processor.ts | 99 +++++++ .../notification-preference.service.ts | 81 +++--- .../services/notification-queue.service.ts | 77 ++++-- .../services/notification.service.ts | 168 +++++------- 16 files changed, 510 insertions(+), 606 deletions(-) create mode 100644 src/modules/notifications/constants/queue.ts create mode 100644 src/modules/notifications/dto/update-preference.dto.ts create mode 100644 src/modules/notifications/interfaces/notification.interface.ts delete mode 100644 src/modules/notifications/processors/notification.processor.ts create mode 100644 src/modules/notifications/processors/push.processor.ts create mode 100644 src/modules/notifications/processors/sms.processor.ts diff --git a/src/modules/notifications/constants/queue.ts b/src/modules/notifications/constants/queue.ts new file mode 100644 index 0000000..777a59e --- /dev/null +++ b/src/modules/notifications/constants/queue.ts @@ -0,0 +1,4 @@ +export enum NotificationQueueNameEnum { + SMS = 'sms', + PUSH = 'push', +} diff --git a/src/modules/notifications/dto/create-preference.dto.ts b/src/modules/notifications/dto/create-preference.dto.ts index a31edb7..b7b253c 100644 --- a/src/modules/notifications/dto/create-preference.dto.ts +++ b/src/modules/notifications/dto/create-preference.dto.ts @@ -1,6 +1,8 @@ -import { IsString, IsNotEmpty, IsBoolean, IsObject, IsOptional, ValidateNested } from 'class-validator'; +import { IsString, IsNotEmpty, IsBoolean, IsOptional, ValidateNested } from 'class-validator'; import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger'; -import { Type } from 'class-transformer'; +import { NotificationTitle } from '../entities/notification.entity'; +import { Enum } from '@mikro-orm/core'; +import { NotificationType } from '../interfaces/notification.interface'; class ChannelsDto { @ApiPropertyOptional({ description: 'Enable SMS channel' }) @@ -14,22 +16,16 @@ class ChannelsDto { push?: boolean; } -export class CreateNotificationPreferenceDto { +export class CreatePreferenceDto { @ApiProperty({ description: 'Notification type (e.g., ORDER_CONFIRMED, PAYMENT_FAILED)' }) @IsString() @IsNotEmpty() - notificationType: string; + @Enum(() => NotificationType) + notificationType!: NotificationType; @ApiProperty({ description: 'Notification channels configuration', type: ChannelsDto }) @ValidateNested() - @Type(() => ChannelsDto) - @IsObject() @IsNotEmpty() - channels: ChannelsDto; - - @ApiPropertyOptional({ description: 'Whether the preference is enabled', default: true }) - @IsBoolean() - @IsOptional() - enabled?: boolean; + @Enum(() => NotificationTitle) + title!: NotificationTitle; } - diff --git a/src/modules/notifications/dto/update-preference.dto.ts b/src/modules/notifications/dto/update-preference.dto.ts new file mode 100644 index 0000000..dbf2672 --- /dev/null +++ b/src/modules/notifications/dto/update-preference.dto.ts @@ -0,0 +1,4 @@ +import { PartialType } from '@nestjs/swagger'; +import { CreatePreferenceDto } from './create-preference.dto'; + +export class UpdatePreferenceDto extends PartialType(CreatePreferenceDto) {} diff --git a/src/modules/notifications/entities/notification-preference.entity.ts b/src/modules/notifications/entities/notification-preference.entity.ts index 1c29373..61d3b48 100644 --- a/src/modules/notifications/entities/notification-preference.entity.ts +++ b/src/modules/notifications/entities/notification-preference.entity.ts @@ -1,24 +1,18 @@ -import { Entity, Property, ManyToOne, Unique } from '@mikro-orm/core'; +import { Entity, Property, ManyToOne, Unique, Enum } from '@mikro-orm/core'; import { BaseEntity } from '../../../common/entities/base.entity'; import { Restaurant } from '../../restaurants/entities/restaurant.entity'; - -export interface NotificationChannels { - sms?: boolean; - push?: boolean; -} +import { NotificationTitle } from './notification.entity'; +import { NotificationType } from '../interfaces/notification.interface'; @Entity({ tableName: 'notification_preferences' }) -@Unique({ properties: ['restaurant', 'notificationType'] }) +@Unique({ properties: ['restaurant', 'title'] }) export class NotificationPreference extends BaseEntity { @ManyToOne(() => Restaurant) restaurant!: Restaurant; @Property() - notificationType!: string; // e.g., ORDER_CONFIRMED, PAYMENT_FAILED + title!: NotificationTitle; - @Property({ type: 'json' }) - channels!: NotificationChannels; // { sms: true, push: true } - - @Property({ default: true }) - enabled: boolean = true; + @Enum(() => NotificationType) + notificationType!: NotificationType; //sms, push, email } diff --git a/src/modules/notifications/entities/notification.entity.ts b/src/modules/notifications/entities/notification.entity.ts index 2c54ec2..84e5c81 100644 --- a/src/modules/notifications/entities/notification.entity.ts +++ b/src/modules/notifications/entities/notification.entity.ts @@ -3,16 +3,10 @@ import { BaseEntity } from '../../../common/entities/base.entity'; import { Restaurant } from '../../restaurants/entities/restaurant.entity'; import { User } from '../../users/entities/user.entity'; -export enum NotificationStatus { - PENDING = 'pending', - SENT = 'sent', - FAILED = 'failed', - SKIPPED = 'skipped', -} - -export enum NotificationChannel { - SMS = 'sms', - PUSH = 'push', +export enum NotificationTitle { + ORDER_CREATED = 'order.created', + REVIEW_CREATED = 'review.created', + ORDER_STATUS_CHANGED = 'order.status.changed', } @Entity({ tableName: 'notifications' }) @@ -23,24 +17,11 @@ export class Notification extends BaseEntity { @ManyToOne(() => User, { nullable: true }) user?: User; - + @Enum(() => NotificationTitle) + title!: NotificationTitle; @Property() - notificationType!: string; // e.g., ORDER_CONFIRMED, PAYMENT_FAILED - - @Enum(() => NotificationChannel) - channel!: NotificationChannel; - - @Property({ type: 'json' }) - payload!: Record; // raw outbound payload - - @Enum(() => NotificationStatus) - status: NotificationStatus = NotificationStatus.PENDING; - - - - @Property({ type: 'int', default: 0 }) - attemptCount: number = 0; + content!: string; @Property({ nullable: true, unique: true }) idempotencyKey?: string; diff --git a/src/modules/notifications/interfaces/notification-queue.interface.ts b/src/modules/notifications/interfaces/notification-queue.interface.ts index 5418b42..9a1909d 100644 --- a/src/modules/notifications/interfaces/notification-queue.interface.ts +++ b/src/modules/notifications/interfaces/notification-queue.interface.ts @@ -1,13 +1,14 @@ -import { NotificationChannel } from '../entities/notification.entity'; +import type { NotificationType } from './notification.interface'; +import type { NotificationTitle } from '../entities/notification.entity'; export interface NotificationQueueJob { restaurantId: string; userId?: string; - notificationType: string; - channel: NotificationChannel; - payload: Record; + title: NotificationTitle; + content: string; idempotencyKey?: string; notificationId?: string; // For retries + notificationType: NotificationType; } export interface NotificationQueueJobResult { @@ -16,4 +17,3 @@ export interface NotificationQueueJobResult { providerResponse?: Record; error?: string; } - diff --git a/src/modules/notifications/interfaces/notification.interface.ts b/src/modules/notifications/interfaces/notification.interface.ts new file mode 100644 index 0000000..9af5342 --- /dev/null +++ b/src/modules/notifications/interfaces/notification.interface.ts @@ -0,0 +1,6 @@ +export enum NotificationType { + NONE = 'none', + SMS = 'sms', + PUSH = 'push', + Both = 'both', +} diff --git a/src/modules/notifications/listeners/notification.listeners.ts b/src/modules/notifications/listeners/notification.listeners.ts index 57440e9..d703170 100644 --- a/src/modules/notifications/listeners/notification.listeners.ts +++ b/src/modules/notifications/listeners/notification.listeners.ts @@ -1,6 +1,7 @@ import { Injectable, Logger } from '@nestjs/common'; import { OnEvent } from '@nestjs/event-emitter'; import { NotificationService } from '../services/notification.service'; +import { NotificationTitle } from '../entities/notification.entity'; import { OrderCreatedEvent, ReviewCreatedEvent, OrderStatusChangedEvent } from '../events/notification.events'; @Injectable() @@ -17,12 +18,8 @@ export class NotificationListeners { await this.notificationService.sendNotification({ restaurantId: event.restaurantId, userId: event.userId, - notificationType: 'order.created', - payload: { - title: 'Order Confirmed', - message: `Your order #${event.orderId} has been confirmed`, - data: { orderId: event.orderId }, - }, + title: NotificationTitle.ORDER_CREATED, + content: `Your order #${event.orderId} has been confirmed`, }); } catch (error) { this.logger.error( @@ -40,12 +37,8 @@ export class NotificationListeners { await this.notificationService.sendNotification({ restaurantId: event.restaurantId, userId: event.userId, - notificationType: 'review.created', - payload: { - title: 'Review Received', - message: 'Thank you for your review!', - data: { reviewId: event.reviewId }, - }, + title: NotificationTitle.REVIEW_CREATED, + content: 'Thank you for your review!', }); } catch (error) { this.logger.error( @@ -63,16 +56,8 @@ export class NotificationListeners { await this.notificationService.sendNotification({ restaurantId: event.restaurantId, userId: event.userId, - notificationType: 'order.status.changed', - payload: { - title: 'Order Status Updated', - message: `Your order #${event.orderId} is now ${event.newStatus}`, - data: { - orderId: event.orderId, - oldStatus: event.oldStatus, - newStatus: event.newStatus, - }, - }, + title: NotificationTitle.ORDER_STATUS_CHANGED, + content: `Your order #${event.orderId} is now ${event.newStatus}`, }); } catch (error) { this.logger.error( diff --git a/src/modules/notifications/notifications.controller.ts b/src/modules/notifications/notifications.controller.ts index 069703d..245d9e1 100644 --- a/src/modules/notifications/notifications.controller.ts +++ b/src/modules/notifications/notifications.controller.ts @@ -2,12 +2,12 @@ import { Controller, Post, Get, Body, Param, UseGuards, Query, Patch, Delete } f import { ApiTags, ApiOperation, ApiBearerAuth, ApiParam, ApiQuery } from '@nestjs/swagger'; import { NotificationService } from './services/notification.service'; import { NotificationPreferenceService } from './services/notification-preference.service'; -import { SendNotificationDto } from './dto/send-notification.dto'; -import { CreateNotificationPreferenceDto } from './dto/create-preference.dto'; +import { CreatePreferenceDto } from './dto/create-preference.dto'; import { AuthGuard } from '../auth/guards/auth.guard'; import { UserId } from '../../common/decorators/user-id.decorator'; import { AdminAuthGuard } from '../auth/guards/adminAuth.guard'; import { RestId } from '../../common/decorators/rest-id.decorator'; +import { UpdatePreferenceDto } from './dto/update-preference.dto'; @ApiTags('notifications') @Controller() @@ -17,84 +17,92 @@ export class NotificationsController { private readonly preferenceService: NotificationPreferenceService, ) {} - @UseGuards(AdminAuthGuard) - @ApiBearerAuth() - @Post('admin/notifications/send') - @ApiOperation({ summary: 'Send a notification (queued)' }) - async sendNotification(@Body() dto: SendNotificationDto) { - return this.notificationService.sendNotification(dto); - } - - @UseGuards(AdminAuthGuard) - @ApiBearerAuth() - @Get('admin/notifications') - @ApiOperation({ summary: 'Get notifications for a restaurant' }) - @ApiQuery({ name: 'limit', required: false, type: Number }) - @ApiQuery({ name: 'notificationType', required: false, type: String }) - async getRestaurantNotifications( - @RestId() restaurantId: string, - @Query('limit') limit?: number, - @Query('notificationType') notificationType?: string, - ) { - if (notificationType) { - return this.notificationService.findByRestaurantAndType( - restaurantId, - notificationType, - limit ? parseInt(limit.toString(), 10) : 50, - ); - } - return this.notificationService.findByRestaurant(restaurantId, limit ? parseInt(limit.toString(), 10) : 50); - } - - @UseGuards(AdminAuthGuard) - @ApiBearerAuth() - @Get('admin/notifications/:id') - @ApiOperation({ summary: 'Get a notification by ID' }) - @ApiParam({ name: 'id', description: 'Notification ID' }) - async getNotification(@Param('id') id: string) { - return this.notificationService.findOne(id); - } - - @UseGuards(AdminAuthGuard) - @ApiBearerAuth() - @Post('admin/notifications/:id/retry') - @ApiOperation({ summary: 'Retry a failed notification' }) - @ApiParam({ name: 'id', description: 'Notification ID' }) - async retryNotification(@Param('id') id: string) { - await this.notificationService.retryFailedNotification(id); - return { message: 'Notification queued for retry' }; - } - @UseGuards(AuthGuard) @ApiBearerAuth() @Get('public/notifications') - @ApiOperation({ summary: 'Get user notifications' }) + @ApiOperation({ summary: 'Get user restaurant notifications' }) @ApiQuery({ name: 'limit', required: false, type: Number }) - async getUserNotifications(@UserId() userId: string, @Query('limit') limit?: number) { - return this.notificationService.findByUserId(userId, limit ? parseInt(limit.toString(), 10) : 50); + async getUserNotifications(@UserId() userId: string, @RestId() restaurantId: string, @Query('limit') limit?: number) { + return await this.notificationService.findByUserAndRestaurant( + userId, + restaurantId, + limit ? parseInt(limit.toString(), 10) : 50, + ); } - @UseGuards(AuthGuard) + @UseGuards(AdminAuthGuard) @ApiBearerAuth() - @Get('public/notifications/:id') - @ApiOperation({ summary: 'Get a notification by ID' }) - @ApiParam({ name: 'id', description: 'Notification ID' }) - async getUserNotification(@Param('id') id: string) { - return this.notificationService.findOne(id); + @Delete('admin/notification/:id') + @ApiOperation({ summary: 'Delete a notification ' }) + @ApiParam({ name: 'id', description: 'Notification ID' }) + async deleteNotification(@RestId() restaurantId: string, @Param('id') id: string) { + await this.notificationService.removeNotification(id, restaurantId); + return { message: 'Notification deleted successfully' }; } + // @UseGuards(AuthGuard) + // @ApiBearerAuth() + // @Get('public/notifications/:id') + // @ApiOperation({ summary: 'Get a notification by ID' }) + // @ApiParam({ name: 'id', description: 'Notification ID' }) + // async getUserNotification(@Param('id') id: string) { + // return this.notificationService.findOne(id); + // } + + // @UseGuards(AdminAuthGuard) + // @ApiBearerAuth() + // @Post('admin/notifications/send') + // @ApiOperation({ summary: 'Send a notification (queued)' }) + // async sendNotification(@Body() dto: SendNotificationDto) { + // return this.notificationService.sendNotification(dto); + // } + + // @UseGuards(AdminAuthGuard) + // @ApiBearerAuth() + // @Get('admin/notifications') + // @ApiOperation({ summary: 'Get notifications for a restaurant' }) + // @ApiQuery({ name: 'limit', required: false, type: Number }) + // @ApiQuery({ name: 'notificationType', required: false, type: String }) + // async getRestaurantNotifications( + // @RestId() restaurantId: string, + // @Query('limit') limit?: number, + // @Query('notificationType') notificationType?: string, + // ) { + // if (notificationType) { + // return this.notificationService.findByRestaurantAndType( + // restaurantId, + // notificationType, + // limit ? parseInt(limit.toString(), 10) : 50, + // ); + // } + // return this.notificationService.findByRestaurant(restaurantId, limit ? parseInt(limit.toString(), 10) : 50); + // } + + // @UseGuards(AdminAuthGuard) + // @ApiBearerAuth() + // @Get('admin/notifications/:id') + // @ApiOperation({ summary: 'Get a notification by ID' }) + // @ApiParam({ name: 'id', description: 'Notification ID' }) + // async getNotification(@Param('id') id: string) { + // return this.notificationService.findOne(id); + // } + + // @UseGuards(AdminAuthGuard) + // @ApiBearerAuth() + // @Post('admin/notifications/:id/retry') + // @ApiOperation({ summary: 'Retry a failed notification' }) + // @ApiParam({ name: 'id', description: 'Notification ID' }) + // async retryNotification(@Param('id') id: string) { + // await this.notificationService.retryFailedNotification(id); + // return { message: 'Notification queued for retry' }; + // } // Preference endpoints @UseGuards(AdminAuthGuard) @ApiBearerAuth() @Post('admin/notification-preferences') - @ApiOperation({ summary: 'Create or update notification preference' }) - async createPreference(@RestId() restaurantId: string, @Body() dto: CreateNotificationPreferenceDto) { - return this.preferenceService.createOrUpdate( - restaurantId, - dto.notificationType, - dto.channels, - dto.enabled ?? true, - ); + @ApiOperation({ summary: 'Create notification preference' }) + createPreference(@RestId() restaurantId: string, @Body() dto: CreatePreferenceDto) { + return this.preferenceService.createOrUpdate(restaurantId, dto); } @UseGuards(AdminAuthGuard) @@ -105,48 +113,48 @@ export class NotificationsController { return this.preferenceService.findByRestaurant(restaurantId); } - @UseGuards(AdminAuthGuard) - @ApiBearerAuth() - @Get('admin/notification-preferences/:type') - @ApiOperation({ summary: 'Get a specific notification preference' }) - @ApiParam({ name: 'type', description: 'Notification type' }) - async getPreference(@RestId() restaurantId: string, @Param('type') notificationType: string) { - return this.preferenceService.findByRestaurantAndType(restaurantId, notificationType); - } + // @UseGuards(AdminAuthGuard) + // @ApiBearerAuth() + // @Get('admin/notification-preferences/:type') + // @ApiOperation({ summary: 'Get a specific notification preference' }) + // @ApiParam({ name: 'type', description: 'Notification type' }) + // async getPreference(@RestId() restaurantId: string, @Param('type') notificationType: string) { + // return this.preferenceService.findByRestaurantAndType(restaurantId, notificationType); + // } + + // @UseGuards(AdminAuthGuard) + // @ApiBearerAuth() + // @Patch('admin/notification-preferences/:type/enabled') + // @ApiOperation({ summary: 'Enable or disable a notification preference' }) + // @ApiParam({ name: 'type', description: 'Notification type' }) + // async updatePreferenceEnabled( + // @RestId() restaurantId: string, + // @Param('type') notificationType: string, + // @Body('enabled') enabled: boolean, + // ) { + // return this.preferenceService.updateEnabled(restaurantId, notificationType, enabled); + // } @UseGuards(AdminAuthGuard) @ApiBearerAuth() - @Patch('admin/notification-preferences/:type/enabled') - @ApiOperation({ summary: 'Enable or disable a notification preference' }) - @ApiParam({ name: 'type', description: 'Notification type' }) - async updatePreferenceEnabled( - @RestId() restaurantId: string, - @Param('type') notificationType: string, - @Body('enabled') enabled: boolean, - ) { - return this.preferenceService.updateEnabled(restaurantId, notificationType, enabled); - } - - @UseGuards(AdminAuthGuard) - @ApiBearerAuth() - @Patch('admin/notification-preferences/:type/channels') + @Patch('admin/notification-preferences/:id') @ApiOperation({ summary: 'Update notification channels' }) - @ApiParam({ name: 'type', description: 'Notification type' }) - async updatePreferenceChannels( + @ApiParam({ name: 'id', description: 'Notification preference ID' }) + async updatePreference( @RestId() restaurantId: string, - @Param('type') notificationType: string, - @Body('channels') channels: { sms?: boolean; push?: boolean }, + @Param('id') preferenceId: string, + @Body() dto: UpdatePreferenceDto, ) { - return this.preferenceService.updateChannels(restaurantId, notificationType, channels); + return this.preferenceService.updatePreference(preferenceId, restaurantId, dto); } @UseGuards(AdminAuthGuard) @ApiBearerAuth() - @Delete('admin/notification-preferences/:type') + @Delete('admin/notification-preferences/:id') @ApiOperation({ summary: 'Delete a notification preference' }) - @ApiParam({ name: 'type', description: 'Notification type' }) - async deletePreference(@RestId() restaurantId: string, @Param('type') notificationType: string) { - await this.preferenceService.delete(restaurantId, notificationType); + @ApiParam({ name: 'id', description: 'Notification preference ID' }) + async deletePreference(@RestId() restaurantId: string, @Param('id') preferenceId: string) { + await this.preferenceService.delete(restaurantId, preferenceId); return { message: 'Preference deleted successfully' }; } } diff --git a/src/modules/notifications/notifications.module.ts b/src/modules/notifications/notifications.module.ts index 34629a7..45d1ea7 100644 --- a/src/modules/notifications/notifications.module.ts +++ b/src/modules/notifications/notifications.module.ts @@ -9,23 +9,22 @@ import { NotificationPreferenceService } from './services/notification-preferenc import { NotificationQueueService } from './services/notification-queue.service'; import { PushNotificationService } from './services/push-notification.service'; import { SmsNotificationService } from './services/sms-notification.service'; -import { NotificationProcessor } from './processors/notification.processor'; +import { SmsProcessor } from './processors/sms.processor'; +import { PushProcessor } from './processors/push.processor'; import { NotificationsController } from './notifications.controller'; import { User } from '../users/entities/user.entity'; import { Restaurant } from '../restaurants/entities/restaurant.entity'; import { NotificationListeners } from './listeners/notification.listeners'; import { AuthModule } from '../auth/auth.module'; import { ConfigService } from '@nestjs/config'; +import { NotificationQueueNameEnum } from './constants/queue'; @Module({ imports: [ AuthModule, JwtModule, MikroOrmModule.forFeature([Notification, NotificationPreference, User, Restaurant]), - BullModule.registerQueue({ - name: 'notification-queue', - }), - + BullModule.registerQueue({ name: NotificationQueueNameEnum.SMS }, { name: NotificationQueueNameEnum.PUSH }), BullModule.forRootAsync({ inject: [ConfigService], useFactory: (config: ConfigService) => ({ @@ -44,7 +43,8 @@ import { ConfigService } from '@nestjs/config'; NotificationQueueService, PushNotificationService, SmsNotificationService, - NotificationProcessor, + PushProcessor, + SmsProcessor, NotificationListeners, ], exports: [ diff --git a/src/modules/notifications/processors/notification.processor.ts b/src/modules/notifications/processors/notification.processor.ts deleted file mode 100644 index fbc78bf..0000000 --- a/src/modules/notifications/processors/notification.processor.ts +++ /dev/null @@ -1,248 +0,0 @@ -import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq'; -import { Logger } from '@nestjs/common'; -import { Job } from 'bullmq'; -import { InjectRepository } from '@mikro-orm/nestjs'; -import { EntityRepository, EntityManager } from '@mikro-orm/postgresql'; -import { Notification, NotificationStatus, NotificationChannel } from '../entities/notification.entity'; -import { NotificationQueueJob, NotificationQueueJobResult } from '../interfaces/notification-queue.interface'; -import { PushNotificationService } from '../services/push-notification.service'; -import { SmsNotificationService } from '../services/sms-notification.service'; -import { Restaurant } from '../../restaurants/entities/restaurant.entity'; -import { User } from '../../users/entities/user.entity'; - -@Processor('notification-queue') -export class NotificationProcessor extends WorkerHost { - private readonly logger = new Logger(NotificationProcessor.name); - - constructor( - @InjectRepository(Notification) - private readonly notificationRepository: EntityRepository, - private readonly em: EntityManager, - private readonly pushNotificationService: PushNotificationService, - private readonly smsNotificationService: SmsNotificationService, - ) { - super(); - } - - async process(job: Job): Promise { - const { restaurantId, userId, notificationType, channel, payload, idempotencyKey, notificationId } = job.data; - - this.logger.log( - `Processing notification job: ${job.id} - Type: ${notificationType}, Channel: ${channel}, Restaurant: ${restaurantId}`, - ); - - try { - // Find or create notification record - let notification: Notification | null = null; - - if (notificationId) { - // Retry case - find existing notification - notification = await this.notificationRepository.findOne({ id: notificationId }); - if (!notification) { - throw new Error(`Notification ${notificationId} not found for retry`); - } - } else { - // Check for duplicate using idempotency key - if (idempotencyKey) { - const existing = await this.notificationRepository.findOne({ idempotencyKey }); - if (existing && existing.status === NotificationStatus.SENT) { - this.logger.warn(`Duplicate notification skipped: ${idempotencyKey}`); - return { - success: true, - notificationId: existing.id, - }; - } - if (existing) { - notification = existing; - } - } - - // Create new notification if not found - if (!notification) { - const restaurant = await this.em.findOne(Restaurant, { id: restaurantId }); - if (!restaurant) { - throw new Error(`Restaurant ${restaurantId} not found`); - } - const user = userId ? await this.em.findOne(User, { id: userId }) : undefined; - - notification = this.em.create(Notification, { - restaurant, - user, - notificationType, - channel, - payload, - idempotencyKey, - status: NotificationStatus.PENDING, - attemptCount: 0, - }); - await this.em.persistAndFlush(notification); - } - } - - // Increment attempt count - notification.attemptCount += 1; - await this.em.persistAndFlush(notification); - - // Send notification based on channel - let providerResponse: Record | undefined; - let success = false; - - switch (channel) { - case NotificationChannel.PUSH: { - const result = await this.sendPushNotification(notification); - success = result.success; - providerResponse = result.providerResponse; - break; - } - case NotificationChannel.SMS: { - const result = await this.sendSmsNotification(notification); - success = result.success; - providerResponse = result.providerResponse; - break; - } - default: - throw new Error(`Unsupported notification channel: ${channel}`); - } - - // Update notification status - if (success) { - notification.status = NotificationStatus.SENT; - notification.sentAt = new Date(); - this.logger.log(`Notification ${notification.id} sent successfully`); - } else { - notification.status = NotificationStatus.FAILED; - this.logger.error(`Notification ${notification.id} failed to send`); - } - - await this.em.persistAndFlush(notification); - - return { - success, - notificationId: notification.id, - providerResponse, - }; - } catch (error) { - this.logger.error(`Error processing notification job ${job.id}:`, error); - - // Update notification status if it exists - if (notificationId) { - const notification = await this.notificationRepository.findOne({ id: notificationId }); - if (notification) { - notification.status = NotificationStatus.FAILED; - - await this.em.persistAndFlush(notification); - } - } - - throw error; - } - } - - private async sendPushNotification(notification: Notification): Promise<{ - success: boolean; - providerResponse?: Record; - }> { - try { - const pushToken = notification.payload.pushToken; - if (!pushToken) { - this.logger.warn(`No push token in payload for notification ${notification.id}`); - return { - success: false, - providerResponse: { error: 'No push token provided' }, - }; - } - - const response = await this.pushNotificationService.sendToToken({ - title: notification.payload.title || 'Notification', - body: notification.payload.message || notification.payload.body || '', - token: pushToken, - data: notification.payload.data, - }); - - if (response.success) { - return { - success: true, - providerResponse: { - messageId: response.messageId, - success: true, - }, - }; - } else { - return { - success: false, - providerResponse: { - error: response.error, - success: false, - }, - }; - } - } catch (error) { - return { - success: false, - providerResponse: { - error: error instanceof Error ? error.message : 'Unknown error', - success: false, - }, - }; - } - } - - private async sendSmsNotification(notification: Notification): Promise<{ - success: boolean; - providerResponse?: Record; - }> { - try { - const phoneNumber = notification.payload.phoneNumber; - if (!phoneNumber) { - this.logger.warn(`No phone number in payload for notification ${notification.id}`); - return { - success: false, - providerResponse: { error: 'No phone number provided' }, - }; - } - - const response = await this.smsNotificationService.send({ - phoneNumber, - message: notification.payload.message || notification.payload.body || '', - templateId: notification.payload.templateId, - parameters: notification.payload.parameters, - }); - - if (response.success) { - return { - success: true, - providerResponse: { - messageId: response.messageId, - success: true, - }, - }; - } else { - return { - success: false, - providerResponse: { - error: response.error, - success: false, - }, - }; - } - } catch (error) { - return { - success: false, - providerResponse: { - error: error instanceof Error ? error.message : 'Unknown error', - success: false, - }, - }; - } - } - - @OnWorkerEvent('completed') - onCompleted(job: Job) { - this.logger.log(`Notification job ${job.id} completed`); - } - - @OnWorkerEvent('failed') - onFailed(job: Job, error: Error) { - this.logger.error(`Notification job ${job.id} failed:`, error); - } -} diff --git a/src/modules/notifications/processors/push.processor.ts b/src/modules/notifications/processors/push.processor.ts new file mode 100644 index 0000000..dc6461a --- /dev/null +++ b/src/modules/notifications/processors/push.processor.ts @@ -0,0 +1,99 @@ +import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq'; +import { Logger } from '@nestjs/common'; +import { Job } from 'bullmq'; +import { InjectRepository } from '@mikro-orm/nestjs'; +import { EntityRepository, EntityManager } from '@mikro-orm/postgresql'; +import { Notification } from '../entities/notification.entity'; +import { NotificationQueueJob, NotificationQueueJobResult } from '../interfaces/notification-queue.interface'; +import { Restaurant } from '../../restaurants/entities/restaurant.entity'; +import { User } from '../../users/entities/user.entity'; +import { NotificationQueueNameEnum } from '../constants/queue'; + +@Processor(NotificationQueueNameEnum.PUSH) +export class PushProcessor extends WorkerHost { + private readonly logger = new Logger(PushProcessor.name); + + constructor( + @InjectRepository(Notification) + private readonly notificationRepository: EntityRepository, + private readonly em: EntityManager, + ) { + super(); + } + + async process(job: Job): Promise { + const { restaurantId, userId, title, content, idempotencyKey, notificationId } = job.data; + + this.logger.log(`Processing notification job: ${job.id} - Title: ${title}, Restaurant: ${restaurantId}`); + + try { + // Find or create notification record + let notification: Notification | null = null; + + if (notificationId) { + // Retry case - find existing notification + notification = await this.notificationRepository.findOne({ id: notificationId }); + if (!notification) { + throw new Error(`Notification ${notificationId} not found for retry`); + } + } else { + // Check for duplicate using idempotency key + if (idempotencyKey) { + const existing = await this.notificationRepository.findOne({ idempotencyKey }); + if (existing && existing.sentAt) { + this.logger.warn(`Duplicate notification skipped: ${idempotencyKey}`); + return { + success: true, + notificationId: existing.id, + }; + } + if (existing) { + notification = existing; + } + } + + // Create new notification if not found + if (!notification) { + const restaurant = await this.em.findOne(Restaurant, { id: restaurantId }); + if (!restaurant) { + throw new Error(`Restaurant ${restaurantId} not found`); + } + const user = userId ? await this.em.findOne(User, { id: userId }) : undefined; + + notification = this.em.create(Notification, { + restaurant, + user, + title, + content, + idempotencyKey, + }); + await this.em.persistAndFlush(notification); + } + } + + // Mark notification as sent + notification.sentAt = new Date(); + await this.em.persistAndFlush(notification); + + this.logger.log(`Notification ${notification.id} processed successfully`); + + return { + success: true, + notificationId: notification.id, + }; + } catch (error) { + this.logger.error(`Error processing notification job ${job.id}:`, error); + throw error; + } + } + + @OnWorkerEvent('completed') + onCompleted(job: Job) { + this.logger.log(`Notification job ${job.id} completed`); + } + + @OnWorkerEvent('failed') + onFailed(job: Job, error: Error) { + this.logger.error(`Notification job ${job.id} failed:`, error); + } +} diff --git a/src/modules/notifications/processors/sms.processor.ts b/src/modules/notifications/processors/sms.processor.ts new file mode 100644 index 0000000..681a49c --- /dev/null +++ b/src/modules/notifications/processors/sms.processor.ts @@ -0,0 +1,99 @@ +import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq'; +import { Logger } from '@nestjs/common'; +import { Job } from 'bullmq'; +import { InjectRepository } from '@mikro-orm/nestjs'; +import { EntityRepository, EntityManager } from '@mikro-orm/postgresql'; +import { Notification } from '../entities/notification.entity'; +import { NotificationQueueJob, NotificationQueueJobResult } from '../interfaces/notification-queue.interface'; +import { Restaurant } from '../../restaurants/entities/restaurant.entity'; +import { User } from '../../users/entities/user.entity'; +import { NotificationQueueNameEnum } from '../constants/queue'; + +@Processor(NotificationQueueNameEnum.SMS) +export class SmsProcessor extends WorkerHost { + private readonly logger = new Logger(SmsProcessor.name); + + constructor( + @InjectRepository(Notification) + private readonly notificationRepository: EntityRepository, + private readonly em: EntityManager, + ) { + super(); + } + + async process(job: Job): Promise { + const { restaurantId, userId, title, content, idempotencyKey, notificationId } = job.data; + + this.logger.log(`Processing notification job: ${job.id} - Title: ${title}, Restaurant: ${restaurantId}`); + + try { + // Find or create notification record + let notification: Notification | null = null; + + if (notificationId) { + // Retry case - find existing notification + notification = await this.notificationRepository.findOne({ id: notificationId }); + if (!notification) { + throw new Error(`Notification ${notificationId} not found for retry`); + } + } else { + // Check for duplicate using idempotency key + if (idempotencyKey) { + const existing = await this.notificationRepository.findOne({ idempotencyKey }); + if (existing && existing.sentAt) { + this.logger.warn(`Duplicate notification skipped: ${idempotencyKey}`); + return { + success: true, + notificationId: existing.id, + }; + } + if (existing) { + notification = existing; + } + } + + // Create new notification if not found + if (!notification) { + const restaurant = await this.em.findOne(Restaurant, { id: restaurantId }); + if (!restaurant) { + throw new Error(`Restaurant ${restaurantId} not found`); + } + const user = userId ? await this.em.findOne(User, { id: userId }) : undefined; + + notification = this.em.create(Notification, { + restaurant, + user, + title, + content, + idempotencyKey, + }); + await this.em.persistAndFlush(notification); + } + } + + // Mark notification as sent + notification.sentAt = new Date(); + await this.em.persistAndFlush(notification); + + this.logger.log(`Notification ${notification.id} processed successfully`); + + return { + success: true, + notificationId: notification.id, + }; + } catch (error) { + this.logger.error(`Error processing notification job ${job.id}:`, error); + throw error; + } + } + + @OnWorkerEvent('completed') + onCompleted(job: Job) { + this.logger.log(`Notification job ${job.id} completed`); + } + + @OnWorkerEvent('failed') + onFailed(job: Job, error: Error) { + this.logger.error(`Notification job ${job.id} failed:`, error); + } +} diff --git a/src/modules/notifications/services/notification-preference.service.ts b/src/modules/notifications/services/notification-preference.service.ts index 017d938..b8fe72c 100644 --- a/src/modules/notifications/services/notification-preference.service.ts +++ b/src/modules/notifications/services/notification-preference.service.ts @@ -1,40 +1,27 @@ -import { Injectable, NotFoundException, BadRequestException } from '@nestjs/common'; +import { Injectable, NotFoundException } from '@nestjs/common'; import { EntityManager } from '@mikro-orm/postgresql'; -import { NotificationPreference, NotificationChannels } from '../entities/notification-preference.entity'; +import { NotificationPreference } from '../entities/notification-preference.entity'; import { Restaurant } from '../../restaurants/entities/restaurant.entity'; +import { CreatePreferenceDto } from '../dto/create-preference.dto'; +import { UpdatePreferenceDto } from '../dto/update-preference.dto'; +import { NotificationTitle } from '../entities/notification.entity'; @Injectable() export class NotificationPreferenceService { constructor(private readonly em: EntityManager) {} - async createOrUpdate( - restaurantId: string, - notificationType: string, - channels: NotificationChannels, - enabled: boolean = true, - ): Promise { + async createOrUpdate(restaurantId: string, dto: CreatePreferenceDto): Promise { const restaurant = await this.em.findOne(Restaurant, { id: restaurantId }); if (!restaurant) { throw new NotFoundException('Restaurant not found'); } - let preference = await this.em.findOne(NotificationPreference, { - restaurant: { id: restaurantId }, - notificationType, + const preference = this.em.create(NotificationPreference, { + restaurant, + notificationType: dto.notificationType, + title: dto.title, }); - if (preference) { - preference.channels = channels; - preference.enabled = enabled; - } else { - preference = this.em.create(NotificationPreference, { - restaurant, - notificationType, - channels, - enabled, - }); - } - await this.em.persistAndFlush(preference); return preference; } @@ -47,47 +34,55 @@ export class NotificationPreferenceService { async findByRestaurantAndType( restaurantId: string, - notificationType: string, + title: NotificationTitle, ): Promise { return this.em.findOne(NotificationPreference, { restaurant: { id: restaurantId }, - notificationType, + title, }); } - async updateEnabled(restaurantId: string, notificationType: string, enabled: boolean): Promise { - const preference = await this.findByRestaurantAndType(restaurantId, notificationType); - if (!preference) { - throw new NotFoundException('Notification preference not found'); - } + // async updateEnabled( + // restaurantId: string, + // notificationType: string, + // enabled: boolean, + // ): Promise { + // const preference = await this.findByRestaurantAndType(restaurantId, notificationType); + // if (!preference) { + // throw new NotFoundException('Notification preference not found'); + // } - preference.enabled = enabled; - await this.em.persistAndFlush(preference); - return preference; - } + // preference.enabled = enabled; + // await this.em.persistAndFlush(preference); + // return preference; + // } - async updateChannels( + async updatePreference( + preferenceId: string, restaurantId: string, - notificationType: string, - channels: NotificationChannels, + dto: UpdatePreferenceDto, ): Promise { - const preference = await this.findByRestaurantAndType(restaurantId, notificationType); + const preference = await this.em.findOne(NotificationPreference, { + id: preferenceId, + restaurant: { id: restaurantId }, + }); if (!preference) { throw new NotFoundException('Notification preference not found'); } - preference.channels = channels; + this.em.assign(preference, dto); await this.em.persistAndFlush(preference); return preference; } - async delete(restaurantId: string, notificationType: string): Promise { - const preference = await this.findByRestaurantAndType(restaurantId, notificationType); + async delete(restaurantId: string, preferenceId: string): Promise { + const preference = await this.em.findOne(NotificationPreference, { + restaurant: { id: restaurantId }, + id: preferenceId, + }); if (!preference) { throw new NotFoundException('Notification preference not found'); } - await this.em.removeAndFlush(preference); } } - diff --git a/src/modules/notifications/services/notification-queue.service.ts b/src/modules/notifications/services/notification-queue.service.ts index 80e0cfe..cea17af 100644 --- a/src/modules/notifications/services/notification-queue.service.ts +++ b/src/modules/notifications/services/notification-queue.service.ts @@ -1,20 +1,23 @@ import { Injectable, Logger } from '@nestjs/common'; import { InjectQueue } from '@nestjs/bullmq'; import { Queue } from 'bullmq'; +import { NotificationQueueNameEnum } from '../constants/queue'; import { NotificationQueueJob } from '../interfaces/notification-queue.interface'; -import { NotificationChannel } from '../entities/notification.entity'; @Injectable() export class NotificationQueueService { private readonly logger = new Logger(NotificationQueueService.name); - constructor(@InjectQueue('notification-queue') private readonly notificationQueue: Queue) {} + constructor( + @InjectQueue(NotificationQueueNameEnum.SMS) private readonly smsQueue: Queue, + @InjectQueue(NotificationQueueNameEnum.PUSH) private readonly pushQueue: Queue, + ) {} - async addNotification(job: NotificationQueueJob): Promise { + async addSmsNotification(job: NotificationQueueJob): Promise { try { - const jobId = job.idempotencyKey || `${job.restaurantId}-${job.notificationType}-${job.channel}-${Date.now()}`; + const jobId = job.idempotencyKey || `${job.restaurantId}-${job.title}-${Date.now()}`; - await this.notificationQueue.add('send-notification', job, { + await this.smsQueue.add('sms-notification', job, { jobId, attempts: 3, backoff: { @@ -30,34 +33,60 @@ export class NotificationQueueService { }, }); - this.logger.log(`Notification job added to queue: ${jobId}`); + this.logger.log(`SMS notification job added to queue: ${jobId}`); } catch (error) { - this.logger.error(`Failed to add notification to queue:`, error); + this.logger.error(`Failed to add SMS notification to queue:`, error); throw error; } } - async addBulkNotifications(jobs: NotificationQueueJob[]): Promise { + async addPushNotification(job: NotificationQueueJob): Promise { try { - const queueJobs = jobs.map(job => ({ - name: 'send-notification', - data: job, - opts: { - jobId: job.idempotencyKey || `${job.restaurantId}-${job.notificationType}-${job.channel}-${Date.now()}`, - attempts: 3, - backoff: { - type: 'exponential', - delay: 2000, - }, - }, - })); + const jobId = job.idempotencyKey || `${job.restaurantId}-${job.title}-${Date.now()}`; - await this.notificationQueue.addBulk(queueJobs); - this.logger.log(`Added ${jobs.length} notification jobs to queue`); + await this.pushQueue.add('push-notification', job, { + jobId, + attempts: 3, + backoff: { + type: 'exponential', + delay: 2000, + }, + removeOnComplete: { + age: 24 * 3600, // Keep completed jobs for 24 hours + count: 1000, + }, + removeOnFail: { + age: 7 * 24 * 3600, // Keep failed jobs for 7 days + }, + }); + + this.logger.log(`Push notification job added to queue: ${jobId}`); } catch (error) { - this.logger.error(`Failed to add bulk notifications to queue:`, error); + this.logger.error(`Failed to add push notification to queue:`, error); throw error; } } -} + // async addBulkNotifications(jobs: NotificationQueueJob[]): Promise { + // try { + // const queueJobs = jobs.map(job => ({ + // name: 'send-notification', + // data: job, + // opts: { + // jobId: job.idempotencyKey || `${job.restaurantId}-${job.title}-${Date.now()}`, + // attempts: 3, + // backoff: { + // type: 'exponential', + // delay: 2000, + // }, + // }, + // })); + + // await this.notificationQueue.addBulk(queueJobs); + // this.logger.log(`Added ${jobs.length} notification jobs to queue`); + // } catch (error) { + // this.logger.error(`Failed to add bulk notifications to queue:`, error); + // throw error; + // } + // } +} diff --git a/src/modules/notifications/services/notification.service.ts b/src/modules/notifications/services/notification.service.ts index e5b6ecd..049d949 100644 --- a/src/modules/notifications/services/notification.service.ts +++ b/src/modules/notifications/services/notification.service.ts @@ -1,26 +1,18 @@ -import { Injectable, Logger, BadRequestException, NotFoundException } from '@nestjs/common'; +import { Injectable, Logger, NotFoundException } from '@nestjs/common'; import { EntityManager } from '@mikro-orm/postgresql'; -import { Notification, NotificationStatus, NotificationChannel } from '../entities/notification.entity'; +import { Notification, NotificationTitle } from '../entities/notification.entity'; import { NotificationPreferenceService } from './notification-preference.service'; import { NotificationQueueService } from './notification-queue.service'; import { NotificationQueueJob } from '../interfaces/notification-queue.interface'; import { Restaurant } from '../../restaurants/entities/restaurant.entity'; import { User } from '../../users/entities/user.entity'; +import { NotificationType } from '../interfaces/notification.interface'; export interface SendNotificationParams { restaurantId: string; userId?: string; - notificationType: string; - payload: { - title?: string; - message: string; - body?: string; - phoneNumber?: string; - pushToken?: string; - data?: Record; - templateId?: string; - parameters?: Array<{ name: string; value: string }>; - }; + title: NotificationTitle; + content: string; idempotencyKey?: string; } @@ -34,8 +26,8 @@ export class NotificationService { private readonly queueService: NotificationQueueService, ) {} - async sendNotification(params: SendNotificationParams): Promise { - const { restaurantId, userId, notificationType, payload, idempotencyKey } = params; + async sendNotification(params: SendNotificationParams): Promise { + const { restaurantId, userId, title, content, idempotencyKey } = params; // Verify restaurant exists const restaurant = await this.em.findOne(Restaurant, { id: restaurantId }); @@ -44,78 +36,54 @@ export class NotificationService { } // Get notification preferences - const preference = await this.preferenceService.findByRestaurantAndType(restaurantId, notificationType); + const preference = await this.preferenceService.findByRestaurantAndType(restaurantId, title); - if (!preference || !preference.enabled) { - this.logger.warn( - `Notification preference not found or disabled for restaurant ${restaurantId}, type ${notificationType}`, - ); - return []; + if (!preference) { + this.logger.warn(`Notification preference not found or disabled for restaurant ${restaurantId}, type ${title}`); + throw new NotFoundException(`Notification preference not found or disabled for type ${title}`); } - // Determine which channels to use - const channels: NotificationChannel[] = []; - if (preference.channels.sms && payload.phoneNumber) { - channels.push(NotificationChannel.SMS); + // Generate idempotency key if not provided + const finalIdempotencyKey = idempotencyKey || `${restaurantId}-${title}-${Date.now()}`; + + // Create notification record + const user = userId ? await this.em.findOne(User, { id: userId }) : undefined; + const notification = this.em.create(Notification, { + restaurant, + user, + title, + content, + idempotencyKey: finalIdempotencyKey, + }); + + await this.em.persistAndFlush(notification); + + if (preference.notificationType === NotificationType.NONE) { + this.logger.warn(`Notification type is NONE for restaurant ${restaurantId}, title ${title}`); + return notification; } - if (preference.channels.push && payload.pushToken) { - channels.push(NotificationChannel.PUSH); + // Add to queue + const job: NotificationQueueJob = { + restaurantId, + userId, + title, + content, + idempotencyKey: finalIdempotencyKey, + notificationId: notification.id, + notificationType: preference.notificationType, + }; + if (preference.notificationType === NotificationType.SMS) { + await this.queueService.addSmsNotification(job); + } else if (preference.notificationType === NotificationType.PUSH) { + await this.queueService.addPushNotification(job); + } else if (preference.notificationType === NotificationType.Both) { + await this.queueService.addSmsNotification(job); + await this.queueService.addPushNotification(job); } - if (channels.length === 0) { - this.logger.warn(`No enabled channels for notification type ${notificationType}`); - return []; - } + this.logger.log(`Queued notification for restaurant ${restaurantId}, title ${title}`); - // Create notification records and queue jobs - const notifications: Notification[] = []; - const jobs: NotificationQueueJob[] = []; - - for (const channel of channels) { - const channelIdempotencyKey = idempotencyKey - ? `${idempotencyKey}-${channel}` - : `${restaurantId}-${notificationType}-${channel}-${Date.now()}`; - - // Create notification record - const user = userId ? await this.em.findOne(User, { id: userId }) : undefined; - const notification = this.em.create(Notification, { - restaurant, - user, - notificationType, - channel, - payload, - status: NotificationStatus.PENDING, - idempotencyKey: channelIdempotencyKey, - attemptCount: 0, - }); - - await this.em.persistAndFlush(notification); - notifications.push(notification); - - // Add to queue - jobs.push({ - restaurantId, - userId, - notificationType, - channel, - payload, - idempotencyKey: channelIdempotencyKey, - notificationId: notification.id, - }); - } - - // Add all jobs to queue - if (jobs.length === 1) { - await this.queueService.addNotification(jobs[0]); - } else { - await this.queueService.addBulkNotifications(jobs); - } - - this.logger.log( - `Queued ${notifications.length} notification(s) for restaurant ${restaurantId}, type ${notificationType}`, - ); - - return notifications; + return notification; } async findOne(id: string): Promise { @@ -138,24 +106,31 @@ export class NotificationService { ); } - async findByUserId(userId: string, limit = 50): Promise { + async findByUserAndRestaurant(userId: string, restaurantId: string, limit = 50): Promise { return this.em.find( Notification, - { user: { id: userId } }, + { user: { id: userId }, restaurant: { id: restaurantId } }, { orderBy: { createdAt: 'DESC' }, limit, - populate: ['restaurant'], }, ); } - async findByRestaurantAndType(restaurantId: string, notificationType: string, limit = 50): Promise { + async removeNotification(id: string, restaurantId: string): Promise { + const notification = await this.em.findOne(Notification, { id, restaurant: { id: restaurantId } }); + if (!notification) { + throw new NotFoundException('Notification not found'); + } + await this.em.removeAndFlush(notification as any); + } + + async findByRestaurantAndType(restaurantId: string, title: NotificationTitle, limit = 50): Promise { return this.em.find( Notification, { restaurant: { id: restaurantId }, - notificationType, + title, }, { orderBy: { createdAt: 'DESC' }, @@ -164,27 +139,4 @@ export class NotificationService { }, ); } - - async retryFailedNotification(notificationId: string): Promise { - const notification = await this.findOne(notificationId); - - if (notification.status !== NotificationStatus.FAILED) { - throw new BadRequestException('Only failed notifications can be retried'); - } - - // Reset status and add to queue - notification.status = NotificationStatus.PENDING; - await this.em.persistAndFlush(notification); - - await this.queueService.addNotification({ - restaurantId: notification.restaurant.id, - notificationType: notification.notificationType, - channel: notification.channel, - payload: notification.payload, - idempotencyKey: notification.idempotencyKey, - notificationId: notification.id, - }); - - this.logger.log(`Retry queued for notification ${notificationId}`); - } }