Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions prisma/migrations/20260929_add_notification/migration.sql
Original file line number Diff line number Diff line change
@@ -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;
16 changes: 16 additions & 0 deletions prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions src/app.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -124,6 +125,7 @@ import { BlockchainModule } from './blockchain/blockchain.module';
ClipsModule,
VideosModule,
JobsModule,
NotificationsModule,
StellarModule,
CsrfModule,
EncryptionModule,
Expand Down
2 changes: 1 addition & 1 deletion src/auth/auth.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {}
4 changes: 3 additions & 1 deletion src/auth/email-delivery.queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,14 +7,16 @@ 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;
subject: string;
template: EmailTemplate;
context: {
token: string;
/** Optional deep link included by job-completed notifications (#917). */
link?: string;
};
}

Expand Down
20 changes: 18 additions & 2 deletions src/auth/mail.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ export class MailService {
}

async sendTemplatedEmail(job: EmailDeliveryJobData): Promise<void> {
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" <noreply@clips.app>',
to: job.to,
Expand Down Expand Up @@ -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: `
<p>Your clips are ready! Click below to view them.</p>
<a href="${preview}" style="display:inline-block;padding:12px 24px;background:#6366f1;color:#fff;border-radius:6px;text-decoration:none;">View clips</a>
<p>Or copy this URL: ${preview}</p>
`,
};
}

if (template === 'magic-link') {
const link = `${baseUrl}/auth/verify-magic?token=${token}`;
return {
Expand Down
22 changes: 22 additions & 0 deletions src/clips/clips.events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
10 changes: 9 additions & 1 deletion src/clips/clips.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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: [
Expand All @@ -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: [
Expand All @@ -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 {}
98 changes: 98 additions & 0 deletions src/clips/nft-mint-enqueue.service.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
import { ConflictException, NotFoundException } from '@nestjs/common';
import { NftMintEnqueueService } from './nft-mint-enqueue.service';

function makeQueue() {
const jobs = new Map<string, any>();
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,
);
});
});
Loading
Loading