feat: add NotificationsService with SSE subject pool

This commit is contained in:
2026-07-05 23:01:26 +08:00
parent bb11ca1880
commit 5f14b9ff15
6 changed files with 167 additions and 2 deletions

View File

@@ -0,0 +1,33 @@
import { IsString, IsNotEmpty, IsOptional, IsArray, IsInt } from 'class-validator';
export class CreateNotificationDto {
@IsArray()
@IsInt({ each: true })
recipientIds: number[];
@IsString()
@IsNotEmpty()
type: string;
@IsString()
@IsNotEmpty()
title: string;
@IsOptional()
@IsString()
content?: string;
@IsOptional()
@IsString()
link?: string;
}
export class NotificationQueryDto {
@IsOptional()
@IsInt()
after?: number;
@IsOptional()
@IsInt()
limit?: number;
}

View File

@@ -0,0 +1,11 @@
import { Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { Notification } from '../entities/notification.entity';
import { NotificationsService } from './notifications.service';
@Module({
imports: [TypeOrmModule.forFeature([Notification])],
providers: [NotificationsService],
exports: [NotificationsService],
})
export class NotificationsModule {}

View File

@@ -0,0 +1,90 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { Subject, Observable } from 'rxjs';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { Notification } from '../entities/notification.entity';
import { CreateNotificationDto } from './dto/notification.dto';
@Injectable()
export class NotificationsService {
private subjects = new Map<number, Subject<Notification>>();
constructor(
@InjectRepository(Notification)
private repo: Repository<Notification>,
private eventEmitter: EventEmitter2,
) {}
async create(dto: CreateNotificationDto): Promise<Notification[]> {
const notifications = dto.recipientIds.map((recipientId) => ({
recipientId,
type: dto.type,
title: dto.title,
content: dto.content ?? '',
link: dto.link,
}));
const saved = await this.repo.save(notifications);
// Push SSE + emit event
for (const n of saved) {
this.subjects.get(n.recipientId)?.next(n);
this.eventEmitter.emit('notification.created', n);
}
return saved;
}
async findByUser(
userId: number,
after?: number,
limit: number = 20,
): Promise<Notification[]> {
const qb = this.repo
.createQueryBuilder('n')
.where('n.recipientId = :userId', { userId })
.orderBy('n.createdAt', 'DESC')
.take(limit);
if (after) {
qb.andWhere('n.id < :after', { after });
}
return qb.getMany();
}
async getUnreadCount(userId: number): Promise<number> {
return this.repo.count({
where: { recipientId: userId, isRead: false },
});
}
async markRead(id: number, userId: number): Promise<void> {
await this.repo.update(
{ id, recipientId: userId },
{ isRead: true, readAt: new Date() },
);
}
async markAllRead(userId: number): Promise<void> {
await this.repo.update(
{ recipientId: userId, isRead: false },
{ isRead: true, readAt: new Date() },
);
}
subscribe(userId: number): Observable<Notification> {
if (!this.subjects.has(userId)) {
this.subjects.set(userId, new Subject<Notification>());
}
return this.subjects.get(userId)!.asObservable();
}
unsubscribe(userId: number): void {
const subj = this.subjects.get(userId);
if (subj) {
subj.complete();
this.subjects.delete(userId);
}
}
}