feat: add application queue priority routing

This commit is contained in:
hectorzhao
2026-07-07 14:49:10 +08:00
parent 72f2c010ce
commit 4c8f7dd73f
17 changed files with 392 additions and 46 deletions
@@ -0,0 +1,2 @@
ALTER TABLE "SmsApplication" ADD COLUMN "queuePriority" TEXT NOT NULL DEFAULT 'normal';
ALTER TABLE "SmsMessageRecord" ADD COLUMN "queuePriority" TEXT NOT NULL DEFAULT 'normal';
+2
View File
@@ -339,6 +339,7 @@ model SmsApplication {
secretHash String
dailyLimit Int?
customerUnitPrice Int @default(0)
queuePriority String @default("normal")
maxPhonesPerTask Int @default(1000000)
templateMismatchMode String @default("reject")
status String @default("active")
@@ -852,6 +853,7 @@ model SmsMessageRecord {
billingUnits Int @default(1)
unitPrice Int @default(0)
amountCents Int @default(0)
queuePriority String @default("normal")
channelId String?
submitId String?
gatewayMessageId String?
+21 -4
View File
@@ -17,6 +17,7 @@ function createPrismaMock() {
unitPrice: 3,
amountCents: 3,
status: 'queued',
queuePriority: 'normal',
submitId: 'SUB-1',
gatewayMessageId: 'GW-1',
channelId: 'channel-1',
@@ -58,7 +59,7 @@ function createPrismaMock() {
findUnique: jest.fn().mockResolvedValue({ id: 'tenant-1', status: 'active', certificationStatus: 'approved' }),
},
smsApplication: {
findUnique: jest.fn().mockResolvedValue({ id: 'app-1', tenantId: 'tenant-1', status: 'active', customerUnitPrice: 3 }),
findUnique: jest.fn().mockResolvedValue({ id: 'app-1', tenantId: 'tenant-1', status: 'active', customerUnitPrice: 3, queuePriority: 'normal' }),
},
smsTemplate: {
findUnique: jest.fn().mockResolvedValue({
@@ -183,8 +184,8 @@ describe('SendChainService', () => {
});
expect(prisma.smsMessageRecord.createMany).toHaveBeenCalledWith({
data: expect.arrayContaining([
expect.objectContaining({ phoneNumber: '13800000001', status: 'queued', billingUnits: 1, amountCents: 3 }),
expect.objectContaining({ phoneNumber: '13800000002', status: 'queued', billingUnits: 1, amountCents: 3 }),
expect.objectContaining({ phoneNumber: '13800000001', status: 'queued', billingUnits: 1, amountCents: 3, queuePriority: 'normal' }),
expect.objectContaining({ phoneNumber: '13800000002', status: 'queued', billingUnits: 1, amountCents: 3, queuePriority: 'normal' }),
]),
});
expect(billing.freeze).toHaveBeenCalledWith(expect.objectContaining({ amountCents: 6, smsUnits: 2, relatedId: 'task-1' }));
@@ -309,10 +310,25 @@ describe('SendChainService', () => {
service['getSendQueue'] = jest.fn().mockReturnValue({ add });
await expect(service.enqueueBatchTask('task-1')).resolves.toEqual({ taskId: 'task-1', enqueued: 1 });
expect(add).toHaveBeenCalledWith('send-message', { messageRecordId: 'record-1' }, { jobId: 'record-1', attempts: 3 });
expect(add).toHaveBeenCalledWith('send-message', { messageRecordId: 'record-1' }, { jobId: 'record-1', attempts: 3, priority: 100 });
expect(prisma.smsBatchTask.update).toHaveBeenCalledWith({ where: { id: 'task-1' }, data: { status: 'queued' } });
});
it('adds priority message jobs ahead of normal message jobs', async () => {
const { service, prisma } = createService();
const add = jest.fn().mockResolvedValue(undefined);
prisma.smsMessageRecord.findMany.mockResolvedValue([
{ id: 'record-priority', batchTaskId: 'task-1', queuePriority: 'priority' },
{ id: 'record-normal', batchTaskId: 'task-1', queuePriority: 'normal' },
]);
service['getSendQueue'] = jest.fn().mockReturnValue({ add });
await expect(service.enqueueBatchTask('task-1')).resolves.toEqual({ taskId: 'task-1', enqueued: 2 });
expect(add).toHaveBeenCalledWith('send-message', { messageRecordId: 'record-priority' }, { jobId: 'record-priority', attempts: 3, priority: 1 });
expect(add).toHaveBeenCalledWith('send-message', { messageRecordId: 'record-normal' }, { jobId: 'record-normal', attempts: 3, priority: 100 });
});
it('routes queued messages to gateway submit commands', async () => {
const { service, prisma } = createService();
const gatewayAdd = jest.fn().mockResolvedValue(undefined);
@@ -333,6 +349,7 @@ describe('SendChainService', () => {
messageType: 'SubmitCommand',
messageId: 'MSG-1',
channelId: 'channel-1',
queuePriority: 'normal',
phoneNumber: '13800000001',
route: expect.objectContaining({ channelCode: 'CMPP-A', rateLimitPerSecond: 100 }),
cmpp: expect.objectContaining({ serviceId: 'SMS', srcId: '10690000' }),
+35 -2
View File
@@ -80,6 +80,8 @@ interface SendJob {
messageRecordId: string;
}
type QueuePriority = 'normal' | 'priority';
type RoutedChannel = {
channel: {
id: string;
@@ -101,6 +103,10 @@ type RoutedChannel = {
const SEND_QUEUE = 'sms.send.queue';
const GATEWAY_SUBMIT_QUEUE = 'gateway.submit.queue';
const BULLMQ_PRIORITY: Record<QueuePriority, number> = {
priority: 1,
normal: 100,
};
@Injectable()
export class SendChainService implements OnModuleInit, OnModuleDestroy {
@@ -133,6 +139,7 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
const schedule = parseSchedule(data);
await this.validateSendResources(data.tenantId, data.applicationId, data.templateId);
const unitPrice = await this.resolveUnitPrice(data.tenantId, data.applicationId);
const queuePriority = await this.resolveQueuePriority(data.tenantId, data.applicationId);
const risk = await this.riskReview.evaluateTask({
tenantId: data.tenantId,
applicationId: data.applicationId,
@@ -223,6 +230,7 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
billingUnits: billing.billingUnitsPerMessage,
unitPrice: billing.unitPrice,
amountCents: billing.billingUnitsPerMessage * billing.unitPrice,
queuePriority,
status: batchStatus === 'ready' ? 'queued' : batchStatus === 'scheduled' ? 'scheduled' : batchStatus,
errorMessage: risk.status === 'rejected' ? risk.reason ?? undefined : undefined,
})),
@@ -365,12 +373,17 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
}
const messages = await this.prisma.smsMessageRecord.findMany({
where: { batchTaskId: taskId, status: 'queued' },
select: { id: true },
select: { id: true, queuePriority: true },
take: 100000,
});
const queue = this.getSendQueue();
for (const message of messages) {
await queue.add('send-message', { messageRecordId: message.id }, { jobId: message.id, attempts: 3 });
const queuePriority = normalizeQueuePriority(message.queuePriority);
await queue.add('send-message', { messageRecordId: message.id }, {
jobId: message.id,
attempts: 3,
priority: BULLMQ_PRIORITY[queuePriority],
});
}
await this.prisma.smsBatchTask.update({ where: { id: taskId }, data: { status: 'queued' } });
return { taskId, enqueued: messages.length };
@@ -653,6 +666,7 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
phoneNumber: string;
content: string;
billingUnits: number;
queuePriority?: string | null;
template?: { signature?: { id?: string | null; name?: string | null } | null } | null;
},
routed: RoutedChannel,
@@ -701,6 +715,7 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
applicationId: message.applicationId ?? 'unknown',
taskId: message.batchTaskId,
submitId,
queuePriority: normalizeQueuePriority(message.queuePriority),
phoneNumber: message.phoneNumber,
content: message.content,
signature: message.template?.signature?.name ?? 'SMS',
@@ -882,6 +897,20 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
return application.customerUnitPrice ?? 0;
}
private async resolveQueuePriority(tenantId: string, applicationId?: string): Promise<QueuePriority> {
if (!applicationId) {
return 'normal';
}
const application = await this.prisma.smsApplication.findUnique({
where: { id: applicationId },
select: { tenantId: true, queuePriority: true },
});
if (!application || application.tenantId !== tenantId) {
return 'normal';
}
return normalizeQueuePriority(application.queuePriority);
}
private async validateSendResources(tenantId: string, applicationId?: string, templateId?: string) {
const tenant = await this.prisma.tenant.findUnique({ where: { id: tenantId } });
if (!tenant || tenant.status !== 'active') {
@@ -1202,6 +1231,10 @@ function normalizeCarrier(carrier?: string | null) {
return value || 'mobile';
}
function normalizeQueuePriority(queuePriority?: string | null): QueuePriority {
return queuePriority === 'priority' ? 'priority' : 'normal';
}
function isCarrierCompatible(channelCarrier: string | null | undefined, targetCarrier: string) {
const normalized = normalizeCarrier(channelCarrier);
return normalized === 'all' || normalized === targetCarrier;
+39 -1
View File
@@ -8,6 +8,7 @@ function createPrismaMock() {
tenantId: 'tenant-1',
name: '应用A',
status: 'active',
queuePriority: 'normal',
tenant: { id: 'tenant-1', name: '租户A', code: 'TENANT-A' },
messageRecords: [{ status: 'delivered' }, { status: 'undelivered' }],
}]),
@@ -16,6 +17,7 @@ function createPrismaMock() {
tenantId: 'tenant-1',
name: '应用A',
status: 'active',
queuePriority: 'normal',
secretHash: 'secret-hash',
tenant: { id: 'tenant-1', name: '租户A', code: 'TENANT-A' },
}),
@@ -147,6 +149,7 @@ describe('SmsConfigService', () => {
expect.objectContaining({
id: 'app-1',
cmppStatus: 'connected',
queuePriority: 'normal',
sentToday: 2,
deliveryRate: 50,
cmppConnections: [expect.objectContaining({ connectionId: 'conn-a' })],
@@ -167,6 +170,40 @@ describe('SmsConfigService', () => {
}));
});
it('creates enterprise applications with persisted queue priority', async () => {
const prisma = createPrismaMock();
const service = new SmsConfigService(prisma as never);
await expect(service.createApplication({
tenantId: 'tenant-1',
name: '优先应用',
queuePriority: 'priority',
ipAllowlist: ['10.0.0.1/32'],
})).resolves.toEqual(expect.objectContaining({ id: 'app-new' }));
expect(prisma.smsApplication.create).toHaveBeenCalledWith(expect.objectContaining({
data: expect.objectContaining({
tenantId: 'tenant-1',
name: '优先应用',
queuePriority: 'priority',
ipAllowlist: { create: [{ ipCidr: '10.0.0.1/32' }] },
}),
}));
});
it('rejects invalid enterprise application queue priority', async () => {
const prisma = createPrismaMock();
const service = new SmsConfigService(prisma as never);
expect(() => service.createApplication({
tenantId: 'tenant-1',
name: '异常应用',
queuePriority: 'urgent',
})).toThrow('queuePriority must be normal or priority');
expect(prisma.smsApplication.create).not.toHaveBeenCalled();
});
it('updates enterprise application profile and allowlist through a transaction', async () => {
const prisma = createPrismaMock();
const tx = {
@@ -181,7 +218,7 @@ describe('SmsConfigService', () => {
prisma.$transaction.mockImplementationOnce((callback: (client: typeof tx) => unknown) => callback(tx));
const service = new SmsConfigService(prisma as never);
await expect(service.updateApplication('app-1', { name: '新应用', customerUnitPrice: 300, ipAllowlist: ['10.0.0.1/32'] }))
await expect(service.updateApplication('app-1', { name: '新应用', customerUnitPrice: 300, queuePriority: 'priority', ipAllowlist: ['10.0.0.1/32'] }))
.resolves.toEqual(expect.objectContaining({ id: 'app-1', name: '新应用' }));
expect(tx.smsApplicationIpAllowlist.deleteMany).toHaveBeenCalledWith({ where: { applicationId: 'app-1' } });
@@ -190,6 +227,7 @@ describe('SmsConfigService', () => {
data: expect.objectContaining({
name: '新应用',
customerUnitPrice: 300,
queuePriority: 'priority',
ipAllowlist: { create: [{ ipCidr: '10.0.0.1/32' }] },
}),
}));
+18
View File
@@ -10,6 +10,7 @@ export interface CreateSmsApplicationDto {
callbackUrl?: string;
dailyLimit?: number;
customerUnitPrice?: number;
queuePriority?: string;
maxPhonesPerTask?: number;
templateMismatchMode?: string;
ipAllowlist?: string[];
@@ -85,6 +86,9 @@ export interface ApplicationListQuery {
includeConnections?: boolean;
}
const APPLICATION_QUEUE_PRIORITIES = ['normal', 'priority'] as const;
type ApplicationQueuePriority = typeof APPLICATION_QUEUE_PRIORITIES[number];
@Injectable()
export class SmsConfigService {
constructor(private readonly prisma: PrismaService) {}
@@ -148,6 +152,7 @@ export class SmsConfigService {
createApplication(data: CreateSmsApplicationDto) {
const secret = randomBytes(24).toString('hex');
const queuePriority = normalizeApplicationQueuePriority(data.queuePriority);
return this.prisma.smsApplication.create({
data: {
tenantId: data.tenantId,
@@ -157,6 +162,7 @@ export class SmsConfigService {
secretHash: hashSecret(secret),
dailyLimit: data.dailyLimit,
customerUnitPrice: data.customerUnitPrice ?? 0,
queuePriority,
maxPhonesPerTask: data.maxPhonesPerTask ?? 1000000,
templateMismatchMode: data.templateMismatchMode ?? 'reject',
ipAllowlist: {
@@ -172,6 +178,9 @@ export class SmsConfigService {
if (!application) {
throw new NotFoundException('Application not found');
}
const queuePriority = data.queuePriority === undefined
? undefined
: normalizeApplicationQueuePriority(data.queuePriority);
return this.prisma.$transaction(async (tx) => {
if (data.ipAllowlist) {
@@ -185,6 +194,7 @@ export class SmsConfigService {
callbackUrl: data.callbackUrl,
dailyLimit: data.dailyLimit,
customerUnitPrice: data.customerUnitPrice,
queuePriority,
maxPhonesPerTask: data.maxPhonesPerTask,
templateMismatchMode: data.templateMismatchMode,
status: data.status,
@@ -738,6 +748,14 @@ function startOfToday() {
return date;
}
function normalizeApplicationQueuePriority(value?: string): ApplicationQueuePriority {
const queuePriority = value ?? 'normal';
if (!APPLICATION_QUEUE_PRIORITIES.includes(queuePriority as ApplicationQueuePriority)) {
throw new BadRequestException('queuePriority must be normal or priority');
}
return queuePriority as ApplicationQueuePriority;
}
function normalizeApplicationCmppStatus(connections: Array<{ status: string; currentConnections: number }>, applicationStatus: string) {
if (applicationStatus !== 'active') {
return 'inactive';