This commit is contained in:
Stephanos A
2026-06-17 18:12:58 +03:00
13 changed files with 397 additions and 199 deletions

View File

@@ -1,8 +1,57 @@
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
import { IsEmail, IsNotEmpty, IsOptional, IsString } from 'class-validator';
export class SendEmail {
@ApiProperty()
@IsEmail()
@IsNotEmpty()
to: string;
@ApiPropertyOptional()
@IsOptional()
@IsString()
sourceId?: string;
@ApiPropertyOptional()
@IsOptional()
@IsString()
sourceName?: string;
@ApiProperty()
@IsNotEmpty()
@IsString()
subject: string;
body: string;
@ApiPropertyOptional()
@IsOptional()
@IsString()
html?: string;
templateKey?: string;
context?: Record<string, unknown>;
@ApiPropertyOptional()
@IsOptional()
@IsString()
text?: string;
@ApiPropertyOptional()
@IsOptional()
body?: string;
@ApiPropertyOptional()
@IsOptional()
context?: Record<string, any>;
@ApiPropertyOptional()
@IsOptional()
@IsString()
templateName?: string;
@ApiPropertyOptional()
@IsOptional()
@IsEmail()
from?: string;
@ApiPropertyOptional()
@IsOptional()
@IsEmail()
replyTo?: string;
}

View File

@@ -1,9 +1,28 @@
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
import { IsArray, IsNotEmpty, IsOptional, IsString, ValidateNested } from 'class-validator';
import { Type } from 'class-transformer';
export class SendMessage {
@ApiProperty()
@IsNotEmpty()
@IsString()
to: string;
@ApiProperty()
@IsNotEmpty()
@IsString()
message: string;
@ApiPropertyOptional()
@IsOptional()
@IsString()
from?: string;
}
export class BulkMessagesDto {
@ApiProperty({ type: [SendMessage] })
@IsArray()
@ValidateNested({ each: true })
@Type(() => SendMessage)
messages: SendMessage[];
}

View File

@@ -7,7 +7,7 @@ import { TestNotificationDto } from './notifications.dto';
import { EmailClientService } from './email-client.service';
import { SmsClientService } from './sms-client.service';
import { SendEmail } from './dtos/email.dto';
import { SendMessage } from './dtos/sms.dto';
import { BulkMessagesDto, SendMessage } from './dtos/sms.dto';
@ApiTags('Notifications')
@Controller('notifications')
@@ -56,6 +56,15 @@ export class NotificationsController {
return this.smsClient.sendSms(dto);
}
@Post('send/sms/bulk')
@UseGuards(IamGuard)
@IamRoles('ADMIN', 'STAFF')
@ApiOperation({ summary: 'Send bulk SMS messages via the SMS microservice' })
@ApiBody({ type: BulkMessagesDto })
sendBulkSms(@Body() dto: BulkMessagesDto) {
return this.smsClient.sendBulkMessages(dto);
}
@Post('test')
@UseGuards(IamGuard)
@IamRoles('ADMIN', 'STAFF')

View File

@@ -1,6 +1,5 @@
import { DynamicModule, Module } from '@nestjs/common';
import { Module } from '@nestjs/common';
import { HttpModule } from '@nestjs/axios';
import { ConfigModule, ConfigService } from '@nestjs/config';
import { ClientsModule, Transport } from '@nestjs/microservices';
import { NotificationsController } from './notifications.controller';
import { NotificationsService } from './notifications.service';
@@ -8,64 +7,39 @@ import { EmailAdapter, SmsAdapter, PushAdapter } from './notification.adapters';
import { EmailClientService } from './email-client.service';
import { SmsClientService } from './sms-client.service';
const rmqClientsModule = ClientsModule.registerAsync([
{
name: 'EMAIL_SERVICE',
imports: [ConfigModule],
inject: [ConfigService],
useFactory: (config: ConfigService) => ({
transport: Transport.RMQ,
options: {
urls: [config.get<string>('RABBITMQ_URL') ?? 'amqp://localhost:5672'],
queue: config.get<string>('EMAIL_QUEUE') ?? 'email_queue',
queueOptions: { durable: true },
noAck: true,
@Module({
imports: [
HttpModule.register({ timeout: 10_000 }),
ClientsModule.register([
{
name: 'EMAIL_SERVICE',
transport: Transport.RMQ,
options: {
urls: [process.env.RABBITMQ_URL as string],
queue: process.env.EMAIL_QUEUE ?? 'email_queue',
queueOptions: { durable: true },
},
},
}),
},
{
name: 'SMS_SERVICE',
imports: [ConfigModule],
inject: [ConfigService],
useFactory: (config: ConfigService) => ({
transport: Transport.RMQ,
options: {
urls: [config.get<string>('RABBITMQ_URL') ?? 'amqp://localhost:5672'],
queue: config.get<string>('SMS_QUEUE') ?? 'sms_queue',
queueOptions: { durable: true },
noAck: true,
{
name: 'SMS_SERVICE',
transport: Transport.RMQ,
options: {
urls: [process.env.RABBITMQ_URL as string],
queue: process.env.SMS_QUEUE ?? 'sms_queue',
queueOptions: { durable: true },
},
},
}),
},
]);
@Module({})
export class NotificationsModule {
static register(): DynamicModule {
const rmqEnabled = process.env.RABBITMQ_ENABLED !== 'false';
return {
module: NotificationsModule,
imports: [
HttpModule.register({ timeout: 10_000 }),
...(rmqEnabled ? [rmqClientsModule] : []),
],
controllers: [NotificationsController],
providers: [
NotificationsService,
EmailAdapter,
SmsAdapter,
PushAdapter,
...(rmqEnabled
? [EmailClientService, SmsClientService]
: [
{ provide: 'EMAIL_SERVICE', useValue: null },
{ provide: 'SMS_SERVICE', useValue: null },
EmailClientService,
SmsClientService,
]),
],
exports: [NotificationsService, EmailClientService, SmsClientService],
};
}
}
]),
],
controllers: [NotificationsController],
providers: [
NotificationsService,
EmailAdapter,
SmsAdapter,
PushAdapter,
EmailClientService,
SmsClientService,
],
exports: [NotificationsService, EmailClientService, SmsClientService],
})
export class NotificationsModule {}

View File

@@ -20,7 +20,7 @@ export class NotificationsService {
private pushAdapter: PushAdapter,
) {
this.channels = new Map<NotificationChannelType, NotificationChannel>([
['EMAIL', { send: (to, subject, body) => this.emailClient.sendEmail({ to, subject, body }).then(() => true) }],
['EMAIL', { send: (to, subject, body) => this.emailClient.sendEmail({ to, subject, text: body }).then(() => true) }],
['SMS', { send: (to, _subject, body) => this.smsClient.sendSms({ to, message: body }).then(() => true) }],
['PUSH', this.pushAdapter as NotificationChannel],
]);
@@ -107,7 +107,7 @@ export class NotificationsService {
await this.emailClient.sendEmail({
to: passenger.user.email,
subject: this.sanitize(dto.title),
body: this.sanitize(dto.body),
text: this.sanitize(dto.body),
});
}

View File

@@ -1,4 +1,9 @@
import { Inject, Injectable, Logger, OnApplicationBootstrap } from '@nestjs/common';
import {
Inject,
Injectable,
Logger,
OnApplicationBootstrap,
} from '@nestjs/common';
import { ClientProxy } from '@nestjs/microservices';
import { BulkMessagesDto, SendMessage } from './dtos/sms.dto';
@@ -8,7 +13,7 @@ export class SmsClientService implements OnApplicationBootstrap {
constructor(
@Inject('SMS_SERVICE')
private readonly smsClient: ClientProxy,
private smsClient: ClientProxy,
) {}
private readonly enabled = process.env.RABBITMQ_ENABLED !== 'false';
@@ -17,8 +22,12 @@ export class SmsClientService implements OnApplicationBootstrap {
if (!this.enabled) return;
this.smsClient
.connect()
.then(() => this.logger.log('Connected to SMS service'))
.catch((err) => this.logger.error('Error connecting to SMS service', err));
.then(() => {
this.logger.log('connected to SMS service');
})
.catch((err) => {
console.error('Error happened at SMS service', err);
});
}
async sendSms(dto: SendMessage) {

View File

@@ -79,6 +79,28 @@ export class PaymentsController {
return this.service.getIntentByBookingId(bookingId);
}
@Get("waafi/return")
@ApiOperation({
summary:
"DEMO ONLY — confirm a Waafi payment from the browser-return params and return JSON for the " +
"UI to display. The frontend success page forwards the Waafi query params here. Gated by " +
"WAAFI_DEMO_TRUST_RETURN (INSECURE; real confirmation is the webhook/HPP_GETTRANINFO).",
})
@ApiQuery({ name: "referenceId", required: true })
@ApiQuery({ name: "state", required: true })
@ApiQuery({ name: "transactionId", required: false })
waafiReturn(
@Query("referenceId") referenceId: string,
@Query("state") state: string,
@Query("transactionId") transactionId: string,
) {
return this.service.confirmWaafiReturnDemo({
referenceId,
state,
transactionId,
});
}
@Post("refund")
@UseGuards(JwtGuard, RolesGuard)
@Roles(UserRole.ADMIN, UserRole.STAFF, UserRole.AGENT)

View File

@@ -44,6 +44,8 @@ export class PaymentsService {
private readonly logger = new Logger(PaymentsService.name);
private readonly walletDemoAutoSucceed = true;
private readonly waafiDemoTrustReturn = true;
constructor(
private prisma: PrismaService,
private seatsService: SeatsService,
@@ -192,6 +194,41 @@ export class PaymentsService {
return { returnUrl, failureUrl };
}
async confirmWaafiReturnDemo(params: {
referenceId?: string;
state?: string;
transactionId?: string;
}): Promise<{ confirmed: boolean; bookingId?: string; reason?: string }> {
if (!this.waafiDemoTrustReturn) {
return { confirmed: false, reason: "demo-disabled" };
}
if ((params.state ?? "").toUpperCase() !== "APPROVED") {
return { confirmed: false, reason: `not-approved (${params.state})` };
}
if (!params.referenceId) {
return { confirmed: false, reason: "missing-referenceId" };
}
const intent = await this.prisma.paymentIntent.findFirst({
where: { merchantOrderId: params.referenceId },
});
if (!intent) {
this.logger.warn(
`waafi demo return: no local intent for referenceId ${params.referenceId}`,
);
return { confirmed: false, reason: "intent-not-found" };
}
this.logger.warn(
`WAAFI_DEMO_TRUST_RETURN enabled — confirming booking ${intent.bookingId} from browser return (INSECURE, demo only)`,
);
await this.finalizePaymentSuccess({
intentId: intent.id,
providerTxnId: params.transactionId,
});
return { confirmed: true, bookingId: intent.bookingId };
}
private async syncIntentProjection(
bookingId: string,
snapshot: PaymentIntentSnapshot,
@@ -453,6 +490,30 @@ export class PaymentsService {
});
}
/**
* Guard against an implausible paidAt from a provider event (e.g. a Telebirr epoch parsed as
* ms×1000 → year 58429), which Prisma/Postgres rejects and would otherwise dead-letter the
* whole confirmation. Falls back to "now" for missing/invalid/far-future/ancient values so the
* booking still confirms.
*/
private sanitizePaidAt(value?: Date): Date {
const now = new Date();
if (!value) return now;
const t = value.getTime();
const oneDayMs = 86_400_000;
if (
Number.isNaN(t) ||
t > now.getTime() + oneDayMs ||
t < Date.UTC(2000, 0, 1)
) {
this.logger.warn(
`finalizePaymentSuccess: implausible paidAt (epoch=${t}); using current time instead`,
);
return now;
}
return value;
}
async finalizePaymentSuccess(input: {
intentId: string;
providerTxnId?: string;
@@ -477,7 +538,7 @@ export class PaymentsService {
});
if (!booking) throw new NotFoundException("Booking not found");
const paidAt = input.paidAt ?? new Date();
const paidAt = this.sanitizePaidAt(input.paidAt);
await this.prisma.$transaction(async (tx) => {
await tx.paymentIntent.update({
where: { id: intent.id },