diff --git a/prisma/migrations/20260929_add_notification/migration.sql b/prisma/migrations/20260929_add_notification/migration.sql new file mode 100644 index 0000000..109910f --- /dev/null +++ b/prisma/migrations/20260929_add_notification/migration.sql @@ -0,0 +1,21 @@ +-- Create Notification table for job-completion notifications (#917) + +CREATE TABLE IF NOT EXISTS "Notification" ( + "id" INTEGER NOT NULL GENERATED ALWAYS AS IDENTITY, + "userId" INTEGER NOT NULL, + "jobId" TEXT NOT NULL, + "type" TEXT NOT NULL, + "title" TEXT NOT NULL, + "body" TEXT, + "link" TEXT, + "readAt" TIMESTAMP(3), + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "Notification_pkey" PRIMARY KEY ("id") +); + +CREATE UNIQUE INDEX IF NOT EXISTS "Notification_jobId_type_key" ON "Notification"("jobId", "type"); +CREATE INDEX IF NOT EXISTS "Notification_userId_createdAt_idx" ON "Notification"("userId", "createdAt"); + +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_userId_fkey" + FOREIGN KEY ("userId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/prisma/schema.prisma b/prisma/schema.prisma index db0409d..b3ab54a 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -216,6 +216,22 @@ model UserPlatform { user User @relation(fields: [userId], references: [id], onDelete: Cascade) } +model Notification { + id Int @id @default(autoincrement()) + userId Int + jobId String + type String + title String + body String? + link String? + readAt DateTime? + createdAt DateTime @default(now()) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + + @@unique([jobId, type]) + @@index([userId, createdAt]) +} + model Wallet { id Int @id @default(autoincrement()) userId Int diff --git a/src/app.module.ts b/src/app.module.ts index f3d0a63..cbedb5d 100644 --- a/src/app.module.ts +++ b/src/app.module.ts @@ -16,6 +16,7 @@ import { NftModule } from './nft/nft.module'; import { EventEmitterModule } from '@nestjs/event-emitter'; import { VideosModule } from './videos/videos.module'; import { JobsModule } from './jobs/jobs.module'; +import { NotificationsModule } from './notifications/notifications.module'; import { PayoutsModule } from './payouts/payouts.module'; import { StellarModule } from './stellar/stellar.module'; import { CsrfModule } from './csrf/csrf.module'; @@ -124,6 +125,7 @@ import { BlockchainModule } from './blockchain/blockchain.module'; ClipsModule, VideosModule, JobsModule, + NotificationsModule, StellarModule, CsrfModule, EncryptionModule, diff --git a/src/auth/auth.module.ts b/src/auth/auth.module.ts index 17089e5..faf6cc6 100644 --- a/src/auth/auth.module.ts +++ b/src/auth/auth.module.ts @@ -65,6 +65,6 @@ import { QueueOverflowService } from '../common/queue/queue-overflow.service'; AdminGuard, QueueOverflowService, ], - exports: [MailService, JwtModule, PassportModule, AdminGuard], + exports: [MailService, JwtModule, PassportModule, AdminGuard, EmailDeliveryService], }) export class AuthModule {} diff --git a/src/auth/email-delivery.queue.ts b/src/auth/email-delivery.queue.ts index 5e8f097..7bfd333 100644 --- a/src/auth/email-delivery.queue.ts +++ b/src/auth/email-delivery.queue.ts @@ -7,7 +7,7 @@ export const EMAIL_DELIVERY_JOB = 'deliver-email'; */ export const EMAIL_DELIVERY_QUEUE_PRIORITY = 5; -export type EmailTemplate = 'verification' | 'password-reset' | 'magic-link' | 'queue-alert'; +export type EmailTemplate = 'verification' | 'password-reset' | 'magic-link' | 'queue-alert' | 'job-completed'; export interface EmailDeliveryJobData { to: string; @@ -15,6 +15,8 @@ export interface EmailDeliveryJobData { template: EmailTemplate; context: { token: string; + /** Optional deep link included by job-completed notifications (#917). */ + link?: string; }; } diff --git a/src/auth/mail.service.ts b/src/auth/mail.service.ts index ef7ed4d..c3873cd 100644 --- a/src/auth/mail.service.ts +++ b/src/auth/mail.service.ts @@ -20,7 +20,7 @@ export class MailService { } async sendTemplatedEmail(job: EmailDeliveryJobData): Promise { - const content = this.buildTemplate(job.template, job.context.token); + const content = this.buildTemplate(job.template, job.context.token, job.context.link); const info = await this.transporter.sendMail({ from: process.env.SMTP_FROM || '"Clips App" ', to: job.to, @@ -80,8 +80,24 @@ export class MailService { }); } - private buildTemplate(template: EmailDeliveryJobData['template'], token: string) { + private buildTemplate( + template: EmailDeliveryJobData['template'], + token: string, + link?: string, + ) { const baseUrl = process.env.APP_BASE_URL || 'http://localhost:3000'; + if (template === 'job-completed') { + const preview = link ?? token; + return { + text: `Your clips are ready! View them here:\n\n${preview}`, + html: ` +

Your clips are ready! Click below to view them.

+ View clips +

Or copy this URL: ${preview}

+ `, + }; + } + if (template === 'magic-link') { const link = `${baseUrl}/auth/verify-magic?token=${token}`; return { diff --git a/src/clips/clips.events.ts b/src/clips/clips.events.ts index 8d6b43d..1e2dfff 100644 --- a/src/clips/clips.events.ts +++ b/src/clips/clips.events.ts @@ -84,3 +84,25 @@ export interface ClipFailedPayload { export const WS_CLIP_PROGRESS = 'clip.progress'; export const WS_CLIP_COMPLETED = 'clip.completed'; export const WS_CLIP_FAILED = 'clip.failed'; +/** Emitted to `user:{id}` room when a job-completion notification is created (#917). */ +export const WS_NOTIFICATION_CREATED = 'notification.created'; +/** Internal event bus topic fanned out to NotificationsService (#917). */ +export const JOB_COMPLETED_EVENT = 'job.completed'; + +export interface JobCompletedEvent { + jobId: string; + type: string; + userId: number; + title: string; + body?: string; + link?: string; +} + +export interface NotificationCreatedPayload { + id: number; + jobId: string; + type: string; + title: string; + body?: string; + link?: string; +} diff --git a/src/clips/clips.module.ts b/src/clips/clips.module.ts index d995c1a..779be90 100644 --- a/src/clips/clips.module.ts +++ b/src/clips/clips.module.ts @@ -13,6 +13,10 @@ import { NftMetadataService } from '../nft/nft-metadata.service'; import { RoyaltyConfigurationService } from '../nft/royalty-configuration.service'; import { registerQueue } from '../common'; import { CLIP_GENERATION_QUEUE } from './clip-generation.queue'; +import { NFT_MINT_QUEUE } from './nft-mint.queue'; +import { NftMintProcessor } from './nft-mint.processor'; +import { NftMintEnqueueService } from './nft-mint-enqueue.service'; +import { QueueOverflowService } from '../common/queue/queue-overflow.service'; @Module({ imports: [ @@ -21,6 +25,7 @@ import { CLIP_GENERATION_QUEUE } from './clip-generation.queue'; CircuitBreakerModule, IpfsUploadModule, registerQueue(CLIP_GENERATION_QUEUE), + registerQueue(NFT_MINT_QUEUE), ], controllers: [ClipsController], providers: [ @@ -30,8 +35,11 @@ import { CLIP_GENERATION_QUEUE } from './clip-generation.queue'; NftConfig, NftMetadataService, NftMintService, + NftMintEnqueueService, + NftMintProcessor, RoyaltyConfigurationService, + QueueOverflowService, ], - exports: [ClipsService, CloudinaryService, NftMintService], + exports: [ClipsService, CloudinaryService, NftMintService, NftMintEnqueueService], }) export class ClipsModule {} diff --git a/src/clips/nft-mint-enqueue.service.spec.ts b/src/clips/nft-mint-enqueue.service.spec.ts new file mode 100644 index 0000000..def89da --- /dev/null +++ b/src/clips/nft-mint-enqueue.service.spec.ts @@ -0,0 +1,98 @@ +import { ConflictException, NotFoundException } from '@nestjs/common'; +import { NftMintEnqueueService } from './nft-mint-enqueue.service'; + +function makeQueue() { + const jobs = new Map(); + return { + jobs, + getJob: jest.fn(async (id: string) => jobs.get(id)), + add: jest.fn(async (_name: string, _data: any, opts: any) => { + const job = { + id: opts.jobId, + attemptsMade: 0, + getState: async () => 'waiting', + }; + jobs.set(opts.jobId, job); + return job; + }), + }; +} + +function makeService() { + const queue = makeQueue(); + const overflowService = { + enqueue: jest.fn(async (o: any) => + queue.add(o.jobName, o.data, o.baseOptions), + ), + }; + const configService = { get: (_k: string, fb: string) => fb }; + const service = new NftMintEnqueueService( + queue as any, + overflowService as any, + configService as any, + ); + return { service, queue, overflowService }; +} + +describe('NftMintEnqueueService (#974)', () => { + it('uses a deterministic jobId per clip', async () => { + const { service, queue } = makeService(); + const res = await service.enqueueMint({ + clipId: 42, + walletAddress: 'GABC', + userId: 7, + }); + expect(res.jobId).toBe('nft-mint-clip-42'); + expect(queue.add).toHaveBeenCalledWith( + 'mint', + expect.objectContaining({ clipId: 42 }), + expect.objectContaining({ jobId: 'nft-mint-clip-42' }), + ); + }); + + it('prevents duplicate active mint jobs with 409', async () => { + const { service } = makeService(); + await service.enqueueMint({ clipId: 42, walletAddress: 'GABC', userId: 7 }); + await expect( + service.enqueueMint({ clipId: 42, walletAddress: 'GABC', userId: 7 }), + ).rejects.toBeInstanceOf(ConflictException); + }); + + it('allows re-enqueue after the previous job completed', async () => { + const { service, queue } = makeService(); + await service.enqueueMint({ clipId: 42, walletAddress: 'GABC', userId: 7 }); + const job = queue.jobs.get('nft-mint-clip-42'); + job.getState = async () => 'completed'; + const res = await service.enqueueMint({ + clipId: 42, + walletAddress: 'GABC', + userId: 7, + }); + expect(res.jobId).toBe('nft-mint-clip-42'); + }); + + it('maps job states to transaction states', async () => { + const { service, queue } = makeService(); + await service.enqueueMint({ clipId: 1, walletAddress: 'GABC', userId: 7 }); + const waiting = await service.getMintStatus('nft-mint-clip-1'); + expect(waiting.txState).toBe('pending'); + + const job = queue.jobs.get('nft-mint-clip-1'); + job.getState = async () => 'completed'; + job.returnvalue = { xdr: 'AAAA' }; + expect((await service.getMintStatus('nft-mint-clip-1')).txState).toBe('confirmed'); + + job.getState = async () => 'failed'; + job.failedReason = 'boom'; + const failed = await service.getMintStatus('nft-mint-clip-1'); + expect(failed.txState).toBe('failed'); + expect(failed.failedReason).toBe('boom'); + }); + + it('throws 404 for unknown jobId', async () => { + const { service } = makeService(); + await expect(service.getMintStatus('nft-mint-clip-999')).rejects.toBeInstanceOf( + NotFoundException, + ); + }); +}); diff --git a/src/clips/nft-mint-enqueue.service.ts b/src/clips/nft-mint-enqueue.service.ts new file mode 100644 index 0000000..09ce491 --- /dev/null +++ b/src/clips/nft-mint-enqueue.service.ts @@ -0,0 +1,121 @@ +import { + Injectable, + Logger, + ConflictException, + NotFoundException, +} from '@nestjs/common'; +import { InjectQueue } from '@nestjs/bullmq'; +import { Queue } from 'bullmq'; +import { ConfigService } from '@nestjs/config'; +import { NFT_MINT_QUEUE, NFT_MINT_JOB_OPTIONS } from './nft-mint.queue'; +import type { NftMintJob } from './nft-mint.processor'; +import { QueueOverflowService } from '../common/queue/queue-overflow.service'; +import { getBullMQRateLimitConfig } from '../config/bullmq.config'; + +export type NftMintTxState = + | 'pending' + | 'confirmed' + | 'failed' + | 'unknown'; + +/** + * NftMintEnqueueService (#974) + * + * Moves Soroban mint processing off the request path and into the dedicated + * `nft-mint` BullMQ queue (separate worker, own concurrency + retry policy), + * so video processing never blocks NFT minting and vice versa. + * + * Dedupe: the BullMQ jobId is deterministic (`nft-mint-{clipId}`), so a + * second enqueue for a clip with an active job returns 409 instead of + * double-minting. + */ +@Injectable() +export class NftMintEnqueueService { + private readonly logger = new Logger(NftMintEnqueueService.name); + + constructor( + @InjectQueue(NFT_MINT_QUEUE) private readonly mintQueue: Queue, + private readonly queueOverflowService: QueueOverflowService, + private readonly configService: ConfigService, + ) {} + + static jobIdForClip(clipId: number): string { + return `nft-mint-clip-${clipId}`; + } + + async enqueueMint(input: { + clipId: number; + walletAddress: string; + userId: number; + }): Promise<{ jobId: string; delayed: boolean; delayMs: number }> { + const jobId = NftMintEnqueueService.jobIdForClip(input.clipId); + + // Dedupe: an already-queued/active job for this clip is a 409. + const existing = await this.mintQueue.getJob(jobId); + if (existing) { + const state = await existing.getState().catch(() => 'unknown'); + if (state === 'waiting' || state === 'active' || state === 'delayed' || state === 'prioritized') { + throw new ConflictException( + `Mint job for clip ${input.clipId} is already ${state} (job ${jobId})`, + ); + } + } + + const rateLimits = getBullMQRateLimitConfig(this.configService); + const result = await this.queueOverflowService.enqueue({ + queue: this.mintQueue as Queue, + jobName: 'mint', + data: { + clipId: input.clipId, + walletAddress: input.walletAddress, + userId: input.userId, + }, + baseOptions: { + ...NFT_MINT_JOB_OPTIONS, + jobId, + } as Record, + rateLimitConfig: rateLimits.nftMint, + }); + + this.logger.log( + `Enqueued NFT mint job ${result.jobId} for clip ${input.clipId}` + + (result.delayed ? ` (delayed ${result.delayMs}ms by overflow guard)` : ''), + ); + return { + jobId: result.jobId ?? jobId, + delayed: result.delayed, + delayMs: result.delayMs, + }; + } + + async getMintStatus(jobId: string): Promise<{ + jobId: string; + state: string; + txState: NftMintTxState; + attemptsMade: number; + failedReason?: string; + result?: unknown; + }> { + const job = await this.mintQueue.getJob(jobId); + if (!job) { + throw new NotFoundException(`Mint job ${jobId} not found`); + } + const state = await job.getState(); + const txState: NftMintTxState = + state === 'completed' + ? 'confirmed' + : state === 'failed' + ? 'failed' + : state === 'waiting' || state === 'active' || state === 'delayed' || state === 'prioritized' + ? 'pending' + : 'unknown'; + return { + jobId, + state, + txState, + attemptsMade: job.attemptsMade, + failedReason: job.failedReason, + result: job.returnvalue, + }; + } +} diff --git a/src/clips/nft-mint.processor.ts b/src/clips/nft-mint.processor.ts index 3a55360..a937039 100644 --- a/src/clips/nft-mint.processor.ts +++ b/src/clips/nft-mint.processor.ts @@ -1,8 +1,10 @@ -import { Logger, OnModuleInit } from '@nestjs/common'; +import { Logger, OnModuleInit, Optional } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq'; +import { EventEmitter2 } from '@nestjs/event-emitter'; import { Job } from 'bullmq'; import { NFT_MINT_QUEUE } from './nft-mint.queue'; +import { JOB_COMPLETED_EVENT } from './clips.events'; import { NftMintService } from './nft-mint.service'; import { MetricsService } from '../metrics/metrics.service'; import { GracefulShutdownService } from '../common/shutdown/graceful-shutdown.service'; @@ -24,6 +26,7 @@ export class NftMintProcessor extends WorkerHost implements OnModuleInit { private readonly nftMintService: NftMintService, private readonly metricsService: MetricsService, private readonly shutdownService: GracefulShutdownService, + @Optional() private readonly eventEmitter?: EventEmitter2, ) { super(); } @@ -74,5 +77,13 @@ export class NftMintProcessor extends WorkerHost implements OnModuleInit { @OnWorkerEvent('completed') onCompleted(job: Job): void { this.logger.log(`NFT mint job ${job.id} completed for clip ${job.data.clipId}`); + this.eventEmitter?.emit(JOB_COMPLETED_EVENT, { + jobId: job.id ?? `nft-mint-clip-${job.data.clipId}`, + type: 'nft-mint', + userId: job.data.userId, + title: 'Your NFT mint completed', + body: 'Mint transaction finished', + link: `/clips/${job.data.clipId}`, + }); } } diff --git a/src/nft/nft.controller.ts b/src/nft/nft.controller.ts index e58e2d6..e6041e3 100644 --- a/src/nft/nft.controller.ts +++ b/src/nft/nft.controller.ts @@ -139,6 +139,7 @@ import { CollectionInfoResponseDto } from './dto/collection-info.dto'; import { TotalSupplyResponseDto } from './dto/total-supply.dto'; import { GasStatsResponseDto } from './dto/gas-stats.dto'; import { GasMetricsService } from './gas-metrics.service'; +import { NftMintEnqueueService } from '../clips/nft-mint-enqueue.service'; import { UpdateMetadataDto, UpdateMetadataResponseDto, @@ -195,6 +196,7 @@ export class NftController { private readonly nftMetadataRefreshService: NftMetadataRefreshService, private readonly claimRoyaltyService: ClaimRoyaltyService, private readonly royaltyClaimHistoryService: RoyaltyClaimHistoryService, + private readonly nftMintEnqueueService: NftMintEnqueueService, ) {} @Auth() @@ -758,6 +760,89 @@ export class NftController { ); } + @Auth() + @UseGuards(NftMintGuard, QueueRateLimitGuard) + @Post('enqueue-mint') + @HttpCode(HttpStatus.ACCEPTED) + @Throttle({ nftMint: { limit: 5, ttl: 60000 } }) + @QueueRateLimit({ queue: 'nft-mint', maxJobs: 5 }) + @ApiBearerAuth('access-token') + @ApiOperation({ + summary: 'Enqueue an NFT mint job on the dedicated nft-mint queue', + description: + 'Moves Soroban mint processing off the request path into the dedicated nft-mint BullMQ queue ' + + '(own worker, concurrency and retry policy), so video processing never blocks minting. ' + + 'Returns the BullMQ jobId immediately; use GET /nfts/mint-status/:jobId to track transaction states. ' + + 'Duplicate active jobs for the same clip return 409.', + }) + @ApiBody({ type: MintNftDto }) + @ApiResponse({ + status: 202, + description: 'Mint job enqueued; returns jobId for status tracking', + schema: { + example: { jobId: 'nft-mint-clip-42', delayed: false, delayMs: 0 }, + }, + }) + @ApiBadRequestResponse({ description: 'Invalid mint payload' }) + @ApiUnauthorizedResponse({ + description: 'Unauthorized — Bearer JWT required', + }) + @ApiForbiddenResponse({ + description: 'Caller does not own the clip', + }) + @ApiNotFoundResponse({ description: 'Clip not found' }) + @ApiConflictResponse({ + description: 'A mint job for this clip is already queued or running (duplicate prevented)', + type: NftMintConflictDto, + }) + @ApiTooManyRequestsResponse({ + description: + 'Too many active NFT mint jobs for this user. Retry after the number of seconds in the Retry-After header.', + }) + async enqueueMint( + @Body() dto: MintNftDto, + @Req() req: Request, + ): Promise<{ jobId: string; delayed: boolean; delayMs: number }> { + const userId = Number((req as any).user?.id ?? 0); + await this.nftMintService.validateClipOwner(dto.clipId, userId); + return this.nftMintEnqueueService.enqueueMint({ + clipId: dto.clipId, + walletAddress: dto.creatorWallet, + userId, + }); + } + + @Auth() + @Get('mint-status/:jobId') + @ApiBearerAuth('access-token') + @ApiOperation({ + summary: 'Get NFT mint job transaction status', + description: + 'Tracks a BullMQ mint job (returned by enqueue-mint) through pending → confirmed/failed states. ' + + 'pending covers waiting/active/delayed; confirmed means the mint transaction completed; ' + + 'failed includes the failure reason.', + }) + @ApiParam({ name: 'jobId', description: 'BullMQ job ID returned by enqueue-mint', example: 'nft-mint-clip-42' }) + @ApiResponse({ + status: 200, + description: 'Mint job status with transaction state', + schema: { + example: { + jobId: 'nft-mint-clip-42', + state: 'completed', + txState: 'confirmed', + attemptsMade: 1, + }, + }, + }) + @ApiUnauthorizedResponse({ + description: 'Unauthorized — Bearer JWT required', + }) + @ApiNotFoundResponse({ description: 'Mint job not found' }) + async mintJobStatus(@Param('jobId') jobId: string) { + return this.nftMintEnqueueService.getMintStatus(jobId); + } + @Auth() @Post('confirm-mint') @HttpCode(HttpStatus.OK) diff --git a/src/notifications/notifications.controller.ts b/src/notifications/notifications.controller.ts new file mode 100644 index 0000000..00c68a5 --- /dev/null +++ b/src/notifications/notifications.controller.ts @@ -0,0 +1,82 @@ +import { + Controller, + Get, + Patch, + Param, + Query, + Req, + UseGuards, + ParseIntPipe, +} from '@nestjs/common'; +import { + ApiTags, + ApiOperation, + ApiResponse, + ApiBearerAuth, + ApiQuery, + ApiParam, + ApiUnauthorizedResponse, + ApiInternalServerErrorResponse, + ApiNotFoundResponse, +} from '@nestjs/swagger'; +import type { Request } from 'express'; +import { LoginGuard } from '../auth/guards/login.guard'; +import { NotificationsService } from './notifications.service'; + +@ApiTags('notifications') +@ApiBearerAuth('access-token') +@ApiUnauthorizedResponse({ description: 'Unauthorized' }) +@ApiInternalServerErrorResponse({ description: 'Internal server error' }) +@UseGuards(LoginGuard) +@Controller('notifications') +export class NotificationsController { + constructor(private readonly notificationsService: NotificationsService) {} + + @Get() + @ApiOperation({ + summary: 'List job-completion notifications', + description: + 'Returns the authenticated user\u2019s completion notifications, newest first. ' + + 'Subscribe to the `notification.created` WebSocket event on the `/video-progress` namespace ' + + 'for real-time delivery with payload {id, jobId, type, title, body, link}.', + }) + @ApiQuery({ name: 'limit', required: false, description: 'Max items (default 20, max 100)', example: 20 }) + @ApiQuery({ name: 'unreadOnly', required: false, description: 'Only unread notifications', example: false }) + @ApiResponse({ + status: 200, + description: 'Notifications returned with read/unread state', + schema: { + example: [ + { + id: 1, + jobId: 'abc123', + type: 'clip-generation', + title: 'Your clips are ready', + body: 'Video processing finished', + link: '/clips/9', + readAt: null, + createdAt: '2026-09-29T00:00:00.000Z', + }, + ], + }, + }) + @ApiResponse({ status: 401, description: 'Unauthorized' }) + async list(@Req() req: Request, @Query('limit') limit?: string, @Query('unreadOnly') unreadOnly?: string) { + const userId = Number((req as any).user?.userId ?? 0); + return this.notificationsService.listForUser(userId, { + limit: limit ? parseInt(limit, 10) : undefined, + unreadOnly: unreadOnly === 'true', + }); + } + + @Patch(':id/read') + @ApiOperation({ summary: 'Mark a notification as read' }) + @ApiParam({ name: 'id', description: 'Notification ID' }) + @ApiResponse({ status: 200, description: 'Marked as read (count of updated rows)' }) + @ApiResponse({ status: 401, description: 'Unauthorized' }) + @ApiNotFoundResponse({ description: 'Notification not found' }) + async markRead(@Req() req: Request, @Param('id', ParseIntPipe) id: number) { + const userId = Number((req as any).user?.userId ?? 0); + return this.notificationsService.markRead(userId, id); + } +} diff --git a/src/notifications/notifications.module.ts b/src/notifications/notifications.module.ts new file mode 100644 index 0000000..787eebf --- /dev/null +++ b/src/notifications/notifications.module.ts @@ -0,0 +1,14 @@ +import { Module } from '@nestjs/common'; +import { NotificationsController } from './notifications.controller'; +import { NotificationsService } from './notifications.service'; +import { PrismaModule } from '../prisma/prisma.module'; +import { AuthModule } from '../auth/auth.module'; +import { VideoProgressGatewayModule } from '../videos/video-progress-gateway.module'; + +@Module({ + imports: [PrismaModule, AuthModule, VideoProgressGatewayModule], + controllers: [NotificationsController], + providers: [NotificationsService], + exports: [NotificationsService], +}) +export class NotificationsModule {} diff --git a/src/notifications/notifications.service.spec.ts b/src/notifications/notifications.service.spec.ts new file mode 100644 index 0000000..9438e83 --- /dev/null +++ b/src/notifications/notifications.service.spec.ts @@ -0,0 +1,95 @@ +import { NotificationsService } from './notifications.service'; + +function makeService() { + const store = new Map(); + let seq = 1; + const prisma = { + notification: { + findUnique: jest.fn(async ({ where }: any) => + store.get(`${where.jobId_type.jobId}:${where.jobId_type.type}`) ?? null, + ), + create: jest.fn(async ({ data }: any) => { + const row = { id: seq++, readAt: null, createdAt: new Date(), ...data }; + store.set(`${data.jobId}:${data.type}`, row); + return row; + }), + findMany: jest.fn(async () => [...store.values()]), + updateMany: jest.fn(async () => ({ count: 1 })), + }, + user: { + findUnique: jest.fn(async () => ({ email: 'user@example.com' })), + }, + }; + const emitted: Array<{ userId: number; payload: any }> = []; + const gateway = { + emitNotification: jest.fn((userId: number, payload: any) => { + emitted.push({ userId, payload }); + }), + }; + const enqueued: any[] = []; + const email = { + enqueue: jest.fn(async (data: any) => { + enqueued.push(data); + }), + }; + const service = new NotificationsService( + prisma as any, + gateway as any, + email as any, + ); + return { service, store, emitted, enqueued, email }; +} + +const input = { + jobId: 'job-1', + type: 'clip-generation', + userId: 7, + title: 'Your clips are ready', + body: 'Video processing finished', + link: '/videos/9', +}; + +describe('NotificationsService (#917)', () => { + it('creates in-app notification, emits WS event and enqueues email with clip link', async () => { + const { service, emitted, enqueued } = makeService(); + const { notification, duplicate } = await service.notifyJobCompleted(input); + expect(duplicate).toBe(false); + expect(notification.jobId).toBe('job-1'); + + expect(emitted).toHaveLength(1); + expect(emitted[0].userId).toBe(7); + expect(emitted[0].payload.link).toBe('/videos/9'); + + expect(enqueued).toHaveLength(1); + expect(enqueued[0].to).toBe('user@example.com'); + expect(enqueued[0].template).toBe('job-completed'); + expect(enqueued[0].context.link).toBe('/videos/9'); + }); + + it('suppresses duplicate notifications for the same job', async () => { + const { service, emitted, enqueued } = makeService(); + await service.notifyJobCompleted(input); + const second = await service.notifyJobCompleted(input); + expect(second.duplicate).toBe(true); + expect(emitted).toHaveLength(1); + expect(enqueued).toHaveLength(1); + }); + + it('still lands in-app + WS when email fails', async () => { + const { service, emitted, enqueued, email } = makeService(); + email.enqueue.mockRejectedValueOnce(new Error('smtp down')); + const { notification, duplicate } = await service.notifyJobCompleted(input); + expect(duplicate).toBe(false); + expect(notification.jobId).toBe('job-1'); + expect(emitted).toHaveLength(1); + expect(enqueued).toHaveLength(0); + }); + + it('lists and marks read', async () => { + const { service } = makeService(); + await service.notifyJobCompleted(input); + const list = await service.listForUser(7, {}); + expect(list).toHaveLength(1); + await service.markRead(7, list[0].id); + }); +}); diff --git a/src/notifications/notifications.service.ts b/src/notifications/notifications.service.ts new file mode 100644 index 0000000..77f2eb0 --- /dev/null +++ b/src/notifications/notifications.service.ts @@ -0,0 +1,117 @@ +import { Injectable, Logger } from '@nestjs/common'; +import { OnEvent } from '@nestjs/event-emitter'; +import { PrismaService } from '../prisma/prisma.service'; +import { VideoProgressGateway } from '../videos/video-progress.gateway'; +import { EmailDeliveryService } from '../auth/email-delivery.service'; +import { JOB_COMPLETED_EVENT } from '../clips/clips.events'; +import type { JobCompletedEvent } from '../clips/clips.events'; + +export interface JobCompletionInput { + jobId: string; + /** e.g. 'clip-generation' | 'clip-posting' | 'nft-mint' */ + type: string; + userId: number; + title: string; + body?: string; + /** Deep link to the generated clips, included in-app, over WS and by email. */ + link?: string; +} + +/** + * NotificationsService (#917) + * + * Single entry point for job-completion fan-out: persists an in-app + * notification, emits a per-user WebSocket event, and enqueues a + * completion email. Dedupe is enforced by the `@@unique([jobId, type])` + * constraint — a second completion for the same job is a no-op so users + * never receive duplicate notifications. + */ +@Injectable() +export class NotificationsService { + private readonly logger = new Logger(NotificationsService.name); + + constructor( + private readonly prisma: PrismaService, + private readonly progressGateway: VideoProgressGateway, + private readonly emailDeliveryService: EmailDeliveryService, + ) {} + + @OnEvent(JOB_COMPLETED_EVENT) + async handleJobCompletedEvent(event: JobCompletedEvent) { + await this.notifyJobCompleted(event); + } + + async notifyJobCompleted(input: JobCompletionInput) { + const existing = await this.prisma.notification.findUnique({ + where: { jobId_type: { jobId: input.jobId, type: input.type } }, + }); + if (existing) { + this.logger.debug( + `Duplicate completion for job ${input.jobId} (${input.type}) suppressed`, + ); + return { notification: existing, duplicate: true as const }; + } + + const notification = await this.prisma.notification.create({ + data: { + userId: input.userId, + jobId: input.jobId, + type: input.type, + title: input.title, + body: input.body, + link: input.link, + }, + }); + + this.progressGateway.emitNotification(input.userId, { + id: notification.id, + jobId: input.jobId, + type: input.type, + title: input.title, + body: input.body, + link: input.link, + }); + + try { + const user = await this.prisma.user.findUnique({ + where: { id: input.userId }, + select: { email: true }, + }); + if (user?.email) { + await this.emailDeliveryService.enqueue({ + to: user.email, + subject: input.title, + template: 'job-completed', + context: { token: input.link ?? input.jobId, link: input.link }, + }); + } + } catch (err) { + // Email is best-effort: the in-app notification + WS event already + // landed, so a mail failure must not fail the completion path. + this.logger.warn( + `Completion email for job ${input.jobId} failed (non-fatal): ${err instanceof Error ? err.message : String(err)}`, + ); + } + + return { notification, duplicate: false as const }; + } + + async listForUser(userId: number, opts: { limit?: number; unreadOnly?: boolean } = {}) { + const limit = Math.min(Math.max(opts.limit ?? 20, 1), 100); + return this.prisma.notification.findMany({ + where: { + userId, + ...(opts.unreadOnly ? { readAt: null } : {}), + }, + orderBy: { createdAt: 'desc' }, + take: limit, + }); + } + + async markRead(userId: number, id: number) { + return this.prisma.notification.updateMany({ + where: { id, userId, readAt: null }, + data: { readAt: new Date() }, + }); + } +} diff --git a/src/videos/clip-generation.processor.ts b/src/videos/clip-generation.processor.ts index 40f7468..dfac7ab 100644 --- a/src/videos/clip-generation.processor.ts +++ b/src/videos/clip-generation.processor.ts @@ -1,6 +1,8 @@ import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq'; import { Logger, OnModuleInit, Optional } from '@nestjs/common'; +import { EventEmitter2 } from '@nestjs/event-emitter'; import { Job } from 'bullmq'; +import { JOB_COMPLETED_EVENT } from '../clips/clips.events'; import { CLIP_GENERATION_QUEUE } from '../clips/clip-generation.queue'; import { VideoProgressGateway } from './video-progress.gateway'; import { PrismaService } from '../prisma/prisma.service'; @@ -39,11 +41,12 @@ export interface ClipGenerationJobData { @Processor(CLIP_GENERATION_QUEUE) export class ClipGenerationProcessor extends WorkerHost implements OnModuleInit { private readonly logger = new Logger(ClipGenerationProcessor.name); - constructor( private readonly prisma: PrismaService, @Optional() private readonly progressGateway?: VideoProgressGateway, - @Optional() private readonly shutdownService?: GracefulShutdownService, + @Optional() private readonly shutdownService?: +GracefulShutdownService, + @Optional() private readonly eventEmitter?: EventEmitter2, ) { super(); } @@ -121,6 +124,15 @@ export class ClipGenerationProcessor extends WorkerHost implements OnModuleInit this.progressGateway.emitCompleted(userId, videoId, clipsGenerated); } + this.eventEmitter?.emit(JOB_COMPLETED_EVENT, { + jobId: job.id ?? `clip-generation-${videoId}`, + type: 'clip-generation', + userId, + title: 'Your clips are ready', + body: 'Video processing finished', + link: `/videos/${videoId}`, + }); + this.logger.log( `Job ${job.id}: video ${videoId} processed — ${clipsGenerated} clip(s)`, ); diff --git a/src/videos/video-progress.gateway.ts b/src/videos/video-progress.gateway.ts index 5265e6c..2118a51 100644 --- a/src/videos/video-progress.gateway.ts +++ b/src/videos/video-progress.gateway.ts @@ -241,6 +241,30 @@ export class VideoProgressGateway } } + /** + * Emit a job-completion notification event to all of the user's sockets (#917). + */ + emitNotification( + userId: number, + payload: { + id: number; + jobId: string; + type: string; + title: string; + body?: string; + link?: string; + }, + ): void { + const sockets = this.userSockets.get(userId) ?? new Set(); + for (const socketId of sockets) { + try { + this.server?.to(socketId).emit('notification.created', payload); + } catch (err) { + this.logger.warn(`Failed to emit notification.created to socket ${socketId}`); + } + } + } + /** * Emit a failure event when a job fails (including timeout). */