chore: refactor the announcement module and service and add queue

This commit is contained in:
mahyargdz
2025-03-04 10:51:16 +03:30
parent e234f98998
commit bf991914e4
9 changed files with 265 additions and 105 deletions
@@ -38,17 +38,23 @@ export class CreateAnnouncementDto {
@ApiProperty({ description: "Is this announcement important?", example: false })
isImportant: boolean;
// @IsOptional()
// @IsNotEmpty({ message: AnnouncementMessage.SERVICE_MUST_BE_ID })
// @ArrayMinSize(1)
// @IsUUID("4", { message: AnnouncementMessage.SERVICE_MUST_BE_UUID, each: true })
// @ApiProperty({ type: [String], description: "Service IDs", example: ["d290f1ee-6c54-4b01-90e6-d701748f0851"] })
// serviceIds: string[];
@IsOptional()
@IsNotEmpty({ message: AnnouncementMessage.SERVICE_MUST_BE_ID })
@ArrayMinSize(1)
@IsUUID("4", { message: AnnouncementMessage.SERVICE_MUST_BE_UUID, each: true })
@ApiProperty({ type: [String], description: "Service IDs", example: ["d290f1ee-6c54-4b01-90e6-d701748f0851"] })
serviceIds: string[];
@IsUUID("4", { message: AnnouncementMessage.SERVICE_MUST_BE_UUID })
@ApiProperty({ description: "Service ID", example: "d290f1ee-6c54-4b01-90e6-d701748f0851" })
serviceId?: string;
@IsOptional()
@ArrayMinSize(1)
@IsArray({ message: AnnouncementMessage.USER_IDS_MUST_BE_ARRAY })
@IsUUID("4", { each: true, message: AnnouncementMessage.USER_MUST_BE_UUID })
@ApiProperty({ type: [String], description: "Array of users ids", example: ["f7b1b3b1-1b1b-4b1b-8b1b-1b1b1b1b1b1b"] })
userIds: string[];
userIds?: string[];
}
@@ -1,13 +1,13 @@
import { Controller, Get, Param, Query } from "@nestjs/common";
import { Body, Controller, Get, Param, Post, Query } from "@nestjs/common";
import { ApiProperty, ApiTags } from "@nestjs/swagger";
import { CreateAnnouncementDto } from "./DTO/create-announcement.dto";
import { SearchAnnouncementQueryDto } from "./DTO/search-announcement-query.dto";
import { AnnouncementService } from "./providers/announcement.service";
import { AuthGuards } from "../../common/decorators/auth-guard.decorator";
import { PermissionsDec } from "../../common/decorators/permission.decorator";
import { UserDec } from "../../common/decorators/user.decorator";
import { ParamDto } from "../../common/DTO/param.dto";
import { SearchCriticismQueryDto } from "../criticisms/DTO/search-criticism-query.dto";
import { User } from "../users/entities/user.entity";
import { PermissionEnum } from "../users/enums/permission.enum";
@ApiTags("Announcements")
@@ -15,47 +15,47 @@ import { PermissionEnum } from "../users/enums/permission.enum";
export class AnnouncementController {
constructor(private readonly announcementService: AnnouncementService) {}
// @AuthGuards()
// @PermissionsDec(PermissionEnum.ANNOUNCEMENTS)
// @ApiProperty({ description: "Create a new announcement ===> login as admin" })
// @Post()
// async create(@Body() createAnnouncementDto: CreateAnnouncementDto) {
// return await this.announcementService.createAnnouncement(createAnnouncementDto);
// }
@AuthGuards()
@PermissionsDec(PermissionEnum.ANNOUNCEMENTS)
@ApiProperty({ description: "Create a new announcement ===> login as admin" })
@Post()
async create(@Body() createAnnouncementDto: CreateAnnouncementDto) {
return await this.announcementService.createAnnouncement(createAnnouncementDto);
}
@AuthGuards()
@ApiProperty({ description: "Get all announcements ===> login as admin" })
@PermissionsDec(PermissionEnum.ANNOUNCEMENTS)
@Get()
async getAllAnnouncements(@Query() queryDto: SearchCriticismQueryDto, @UserDec() user: User) {
return await this.announcementService.getAllAnnouncements(queryDto, user.id);
async getAnnouncements(@Query() queryDto: SearchAnnouncementQueryDto) {
return await this.announcementService.getAnnouncements(queryDto);
}
@AuthGuards()
@ApiProperty({ description: "Get all announcements ===> login as user" })
@ApiProperty({ description: "Get all announcements for user ===> login as user" })
@Get("user")
async getAllAnnouncementsByUser(@Query() queryDto: SearchCriticismQueryDto, @UserDec() user: User) {
return await this.announcementService.getAllAnnouncementsByUser(queryDto, user.id);
async getAnnouncementsByUser(@Query() queryDto: SearchAnnouncementQueryDto, @UserDec("id") userId: string) {
return await this.announcementService.getAnnouncementsByUser(queryDto, userId);
}
@AuthGuards()
@ApiProperty({ description: "Get one announcements with id" })
@Get(":id")
getAnnouncement(@Param() paramDto: ParamDto, @UserDec() user: User) {
return this.announcementService.getOneAnnouncement(paramDto.id, user.id);
getAnnouncementById(@Param() paramDto: ParamDto, @UserDec("id") userId: string) {
return this.announcementService.getAnnouncementById(paramDto.id, userId);
}
@AuthGuards()
@ApiProperty({ description: "Get all public announcements" })
@Get("public")
async getPublicAnnouncements() {
return await this.announcementService.getPublicAnnouncements();
}
// @AuthGuards()
// @ApiProperty({ description: "Get all public announcements" })
// @Get("public")
// async getPublicAnnouncements() {
// return await this.announcementService.getPublicAnnouncements();
// }
@AuthGuards()
@ApiProperty({ description: "Get all public announcements" })
@Get("/service/:serviceId")
async getAnnouncementsByService(@Query() queryDto: SearchCriticismQueryDto, @Param("serviceId") serviceId: string) {
return await this.announcementService.getAnnouncementByService(queryDto, serviceId);
}
// @AuthGuards()
// @ApiProperty({ description: "Get all public announcements" })
// @Get("/service/:serviceId")
// async getAnnouncementsByService(@Query() queryDto: SearchAnnouncementQueryDto, @Param("serviceId") serviceId: string) {
// return await this.announcementService.getAnnouncementByService(queryDto, serviceId);
// }
}
@@ -1,9 +1,12 @@
import { BullModule } from "@nestjs/bullmq";
import { Module } from "@nestjs/common";
import { TypeOrmModule } from "@nestjs/typeorm";
import { AnnouncementController } from "./announcement.controller";
import { ANNOUNCEMENT } from "./constants";
import { Announcement } from "./entities/announcement.entity";
import { AnnouncementService } from "./providers/announcement.service";
import { AnnouncementProcessor } from "./queue/announcment.processor";
import { AnnouncementRepository } from "./repositories/announcement.repository";
import { DanakServicesModule } from "../danak-services/danak-services.module";
import { UsersModule } from "../users/users.module";
@@ -11,8 +14,13 @@ import { UserAnnouncement } from "./entities/user-announcement.entity";
import { UserAnnouncementRepository } from "./repositories/user-announcement.repository";
@Module({
imports: [TypeOrmModule.forFeature([Announcement, UserAnnouncement]), DanakServicesModule, UsersModule],
providers: [AnnouncementService, AnnouncementRepository, UserAnnouncementRepository],
imports: [
TypeOrmModule.forFeature([Announcement, UserAnnouncement]),
BullModule.registerQueue({ name: ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME }),
DanakServicesModule,
UsersModule,
],
providers: [AnnouncementService, AnnouncementRepository, UserAnnouncementRepository, AnnouncementProcessor],
controllers: [AnnouncementController],
exports: [AnnouncementService],
})
+12
View File
@@ -0,0 +1,12 @@
export const ANNOUNCEMENT = Object.freeze({
ANNOUNCEMENT_QUEUE_NAME: "announcement",
ANNOUNCEMENT_SEND_JOB_NAME: "sendAnnouncement",
ANNOUNCEMENT_PUBLISH_JOB_NAME: "publishAnnouncement",
//
ANNOUNCEMENT_SEND_JOB_PRIORITY: 1,
ANNOUNCEMENT_SEND_JOB_ATTEMPTS: 3,
ANNOUNCEMENT_SEND_JOB_BACKOFF: 1000,
ANNOUNCEMENT_SEND_JOB_TIMEOUT: 10000,
//
ANNOUNCEMENT_PUBLISH_JOB_PRIORITY: 1,
});
@@ -16,7 +16,7 @@ export class Announcement extends BaseEntity {
@Column({ type: "timestamp", nullable: true, default: null })
publishAt: Date | null;
@ManyToOne(() => DanakService, (service) => service.announcements, { eager: true, onDelete: "CASCADE" })
@ManyToOne(() => DanakService, (service) => service.announcements, { onDelete: "CASCADE" })
targetService: DanakService;
@Column({ type: "boolean", default: false })
@@ -27,8 +27,4 @@ export class Announcement extends BaseEntity {
@Column({ type: "boolean", default: false })
isPublic: boolean;
isPublicAnnouncement(): boolean {
return this.isPublic;
}
}
@@ -0,0 +1,6 @@
export interface ISendAnnouncement {
announcementId: string;
isPublic: boolean;
serviceId: string;
userIds?: string[];
}
@@ -1,7 +1,13 @@
import { InjectQueue } from "@nestjs/bullmq";
import { BadRequestException, Injectable } from "@nestjs/common";
import { Queue } from "bullmq";
import dayjs from "dayjs";
import { AnnouncementMessage } from "../../../common/enums/message.enum";
import { DanakServicesService } from "../../danak-services/providers/danak-services.service";
import { PaginationUtils } from "../../utils/providers/pagination.utils";
import { ANNOUNCEMENT } from "../constants";
import { CreateAnnouncementDto } from "../DTO/create-announcement.dto";
import { SearchAnnouncementQueryDto } from "../DTO/search-announcement-query.dto";
import { AnnouncementRepository } from "../repositories/announcement.repository";
import { UserAnnouncementRepository } from "../repositories/user-announcement.repository";
@@ -9,33 +15,78 @@ import { UserAnnouncementRepository } from "../repositories/user-announcement.re
@Injectable()
export class AnnouncementService {
constructor(
@InjectQueue(ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME) private readonly announcementQueue: Queue,
private readonly announcementRepository: AnnouncementRepository,
private readonly userAnnouncementRepository: UserAnnouncementRepository,
private readonly danakServicesService: DanakServicesService,
) {}
/***************************************** */
async getAllAnnouncements(queryDto: SearchAnnouncementQueryDto, userId: string) {
async createAnnouncement(createAnnouncementDto: CreateAnnouncementDto) {
const { serviceId, userIds } = createAnnouncementDto;
let isPublic = true;
const announcement = this.announcementRepository.create({ ...createAnnouncementDto });
if (!announcement.publishAt) {
announcement.publishAt = new Date();
}
if (dayjs(announcement.publishAt).isBefore(dayjs().subtract(1, "day")))
throw new BadRequestException(AnnouncementMessage.PUBLISH_AT_MUST_BE_FUTURE_DATE);
if (serviceId) {
announcement.targetService = await this.danakServicesService.findServiceById(serviceId);
isPublic = false;
}
if (userIds && userIds.length > 0) isPublic = false;
await this.announcementRepository.save(announcement);
await this.announcementQueue.add(
ANNOUNCEMENT.ANNOUNCEMENT_SEND_JOB_NAME,
{
announcementId: announcement.id,
isPublic,
serviceId,
userIds: userIds,
},
{
delay: dayjs(announcement.publishAt).diff(dayjs(), "millisecond"),
attempts: ANNOUNCEMENT.ANNOUNCEMENT_SEND_JOB_ATTEMPTS,
backoff: ANNOUNCEMENT.ANNOUNCEMENT_SEND_JOB_BACKOFF,
},
);
return {
message: AnnouncementMessage.CREATED_AND_SENT_TO_USERS_AFTER_PUBLISH,
announcement,
};
}
/***************************************** */
async getAnnouncements(queryDto: SearchAnnouncementQueryDto) {
const { limit, skip } = PaginationUtils(queryDto);
const queryBuilder = this.announcementRepository
.createQueryBuilder("announcement")
.leftJoin("announcement.service", "danak_service")
.addSelect(["danak_service.id", "danak_service.name"])
.leftJoin("announcement.targetGroups", "group")
.addSelect(["group.id", "group.name"])
.leftJoinAndSelect("announcement.userAnnouncements", "userAnnouncements", "userAnnouncements.userId = :userId", { userId });
.leftJoin("announcement.targetService", "danak_service")
.addSelect(["danak_service.id", "danak_service.name"]);
// .leftJoinAndSelect("announcement.userAnnouncements", "userAnnouncements", "userAnnouncements.userId = :userId", { userId });
if (queryDto.since) {
queryBuilder.andWhere("announcement.createdAt >= :since", { since: queryDto.since });
queryBuilder.andWhere("announcement.createdAt >= :since", { since: dayjs(queryDto.since).toDate() });
}
if (queryDto.publishAt) {
queryBuilder.andWhere("announcement.publishAt >= :publishAt", { publishAt: queryDto.publishAt });
queryBuilder.andWhere("announcement.publishAt >= :publishAt", { publishAt: dayjs(queryDto.publishAt).toDate() });
}
if (queryDto.q) {
queryBuilder.andWhere("announcement.title LIKE :q", { q: `%${queryDto.q}%` });
queryBuilder.andWhere("announcement.title ILIKE :q", { q: `%${queryDto.q}%` });
}
queryBuilder.orderBy("announcement.createdAt", "DESC");
@@ -47,89 +98,60 @@ export class AnnouncementService {
/***************************************** */
async getAllAnnouncementsByUser(queryDto: SearchAnnouncementQueryDto, userId: string) {
async getAnnouncementsByUser(queryDto: SearchAnnouncementQueryDto, userId: string) {
const { limit, skip } = PaginationUtils(queryDto);
const queryBuilder = this.announcementRepository
.createQueryBuilder("announcement")
.leftJoin("announcement.service", "service")
const queryBuilder = this.userAnnouncementRepository
.createQueryBuilder("userAnnouncement")
.leftJoin("userAnnouncement.announcement", "announcement")
.addSelect([
"announcement.id",
"announcement.title",
"announcement.isImportant",
"announcement.createdAt",
"announcement.publishAt",
"announcement.content",
])
.leftJoin("userAnnouncement.user", "user")
.leftJoin("announcement.targetService", "service")
.addSelect(["service.id", "service.name"])
.leftJoin("announcement.users", "announcementUser")
.leftJoin("announcement.targetGroups", "group")
.addSelect(["group.id", "group.name"])
.leftJoin("group.users", "groupUser")
.leftJoin("service.subscriptionPlans", "subscriptionPlan")
.leftJoin(
"UserSubscription",
"userSubscription",
"userSubscription.planId = subscriptionPlan.id AND userSubscription.userId = :userId",
{ userId },
)
.leftJoinAndSelect("announcement.userAnnouncements", "userAnnouncements", "userAnnouncements.userId = :userId", { userId })
.where("announcement.isPublic = :isPublic", { isPublic: true })
.orWhere("groupUser.id = :userId", { userId })
.orWhere("userSubscription.userId = :userId", { userId })
.orWhere("announcementUser.id = :userId", { userId });
.where("user.id = :userId", { userId });
if (queryDto.since) {
queryBuilder.andWhere("announcement.createdAt >= :since", { since: queryDto.since });
queryBuilder.andWhere("announcement.createdAt >= :since", { since: dayjs(queryDto.since).toDate() });
}
if (queryDto.publishAt) {
queryBuilder.andWhere("announcement.publishAt >= :publishAt", { publishAt: queryDto.publishAt });
queryBuilder.andWhere("announcement.publishAt >= :publishAt", { publishAt: dayjs(queryDto.publishAt).toDate() });
}
if (queryDto.q) {
queryBuilder.andWhere("announcement.title LIKE :q", { q: `%${queryDto.q}%` });
queryBuilder.andWhere("announcement.title ILIKE :q", { q: `%${queryDto.q}%` });
}
queryBuilder.orderBy("announcement.createdAt", "DESC");
const rawResults = await queryBuilder.getRawMany();
const [, count] = await queryBuilder.skip(skip).take(limit).getManyAndCount();
//TODO: fix later
const announcements = rawResults.map((raw) => ({
id: raw.announcement_id,
createdAt: raw.announcement_createdAt,
updatedAt: raw.announcement_updatedAt,
title: raw.announcement_title,
content: raw.announcement_content,
publishAt: raw.announcement_publishAt,
service: {
id: raw.service_id,
name: raw.service_name,
},
isImportant: raw.announcement_isImportant,
targetGroups: raw.group_id ? [{ id: raw.group_id, name: raw.group_name }] : [],
isPublic: raw.announcement_isPublic,
isRead: raw.userAnnouncements_isRead ? true : false,
}));
const [announcements, count] = await queryBuilder.skip(skip).take(limit).getManyAndCount();
return { announcements, count, paginate: true };
}
/***************************************** */
async getOneAnnouncement(id: string, userId: string) {
async getAnnouncementById(id: string, userId: string) {
const announcement = await this.announcementRepository.findOneBy({ id });
if (!announcement) throw new BadRequestException(AnnouncementMessage.NOT_FOUND);
const userAnnouncement = this.userAnnouncementRepository.create({
announcement: { id: announcement.id },
user: { id: userId },
isRead: true,
});
const userAnnouncement = await this.userAnnouncementRepository.findOneBy({ announcement: { id }, user: { id: userId } });
if (!userAnnouncement) throw new BadRequestException(AnnouncementMessage.NOT_FOUND);
userAnnouncement.isRead = true;
await this.userAnnouncementRepository.save(userAnnouncement);
return { announcement };
}
/***************************************** */
async getPublicAnnouncements() {
const announcements = await this.announcementRepository.find({ where: { isPublic: true }, relations: ["service", "targetGroups"] });
return { announcements };
}
/***************************************** */
async getAnnouncementByService(queryDto: SearchAnnouncementQueryDto, serviceId: string) {
const { limit, skip } = PaginationUtils(queryDto);
+107
View File
@@ -0,0 +1,107 @@
import { Processor } from "@nestjs/bullmq";
import { Logger } from "@nestjs/common";
import { Job } from "bullmq";
import { DataSource } from "typeorm";
import { WorkerProcessor } from "../../../common/queues/worker.processor";
import { SubscriptionPlan } from "../../subscriptions/entities/subscription.entity";
import { UserSubscription } from "../../subscriptions/entities/user-subscription.entity";
import { User } from "../../users/entities/user.entity";
import { ANNOUNCEMENT } from "../constants";
import { UserAnnouncement } from "../entities/user-announcement.entity";
import { ISendAnnouncement } from "../interfaces/ISendAnnouncement";
@Processor(ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME)
export class AnnouncementProcessor extends WorkerProcessor {
protected readonly logger = new Logger(AnnouncementProcessor.name);
constructor(private readonly dataSource: DataSource) {
super();
}
async process(job: Job, token?: string) {
this.logger.log(job, token);
switch (job.name) {
case ANNOUNCEMENT.ANNOUNCEMENT_SEND_JOB_NAME:
return this.sendAnnouncement(job, token);
case ANNOUNCEMENT.ANNOUNCEMENT_PUBLISH_JOB_NAME:
return this.publishAnnouncement(job, token);
default:
this.logger.error(`Unknown job name: ${job.name}`);
return;
}
}
//********************************** */
private async sendAnnouncement(job: Job, token?: string) {
this.logger.log(`Sending announcement: ${job.data.announcementId} to users. ${token}`);
const queryRunner = this.dataSource.createQueryRunner();
await queryRunner.connect();
await queryRunner.startTransaction();
const data = job.data as ISendAnnouncement;
try {
// Create a query builder to fetch unique users
const userQueryBuilder = queryRunner.manager.createQueryBuilder(User, "user").select("user.id").distinct(true);
// Add public users condition
if (data.isPublic) {
// No additional condition needed for public - it will select all users
}
// Add service subscription condition
if (data.serviceId) {
userQueryBuilder
.leftJoin(UserSubscription, "us", "us.userId = user.id")
.leftJoin(SubscriptionPlan, "sp", "us.planId = sp.id AND sp.serviceId = :serviceId", { serviceId: data.serviceId })
.andWhere("sp.id IS NOT NULL");
}
// Add specific users condition
if (data.userIds && data.userIds.length > 0) {
if (!data.isPublic && !data.serviceId) {
// Only apply direct user filtering if it's the only condition
userQueryBuilder.andWhere("user.id IN (:...userIds)", { userIds: data.userIds });
} else {
// Use UNION to add specific users
userQueryBuilder.orWhere("user.id IN (:...userIds)", { userIds: data.userIds });
}
}
// Execute the query to get distinct user IDs
const userResult = await userQueryBuilder.getRawMany();
const users = userResult.map((row) => row.user_id);
if (users.length === 0) {
this.logger.error(`No users found`);
await queryRunner.rollbackTransaction();
return;
}
const announcementId = job.data.announcementId;
this.logger.debug(`Sending announcement to ${users.length} users`);
const userAnnouncements = users.map((userId) => ({
user: { id: userId },
announcement: { id: announcementId },
isRead: false,
}));
await queryRunner.manager.insert(UserAnnouncement, userAnnouncements);
this.logger.log(`Successfully sent announcement ${announcementId} to ${users.length} users`);
await queryRunner.commitTransaction();
} catch (error) {
this.logger.error(`Failed to send announcement:`, error);
await queryRunner.rollbackTransaction();
throw error;
} finally {
await queryRunner.release();
}
}
//********************************** */
private async publishAnnouncement(job: Job, token?: string) {
this.logger.log(`Publishing announcement: ${job.data.announcementId} ${token}`);
}
}