fix: reassemble inbound CMPP long messages

This commit is contained in:
hectorzhao
2026-07-23 21:04:35 +08:00
parent f186aee00b
commit b29576fcd1
9 changed files with 1102 additions and 22 deletions
+360 -9
View File
@@ -48,6 +48,12 @@ export interface GatewayInboundSubmitDto {
destId?: string;
sequenceId?: number;
remoteIp?: string;
longMessage?: {
reference: number;
total: number;
index: number;
format: number;
};
}
interface GatewayInboundSingleSubmitResult {
@@ -260,6 +266,9 @@ const DEFAULT_SCHEDULED_DISPATCH_STALE_MS = 2 * 60_000;
const SCHEDULED_DISPATCH_INITIAL_DELAY_MS = 1_000;
const DEFAULT_GATEWAY_SUBMIT_REQUEUE_STALE_MS = 2 * 60_000;
const DEFAULT_DOWNSTREAM_MANUAL_REQUEUE_STALE_MS = 2 * 60_000;
const DEFAULT_INBOUND_LONG_MESSAGE_SCAN_INTERVAL_MS = 60_000;
const INBOUND_LONG_MESSAGE_SCAN_INITIAL_DELAY_MS = 10_000;
const DEFAULT_INBOUND_LONG_MESSAGE_PROCESSING_STALE_SECONDS = 30;
const GATEWAY_SUBMIT_REQUEUE_IDEMPOTENCY_TTL_SECONDS = 30 * 24 * 60 * 60;
const BULLMQ_PRIORITY: Record<QueuePriority, number> = {
priority: 1,
@@ -279,6 +288,8 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
private scheduledDispatchInitialTimer?: ReturnType<typeof setTimeout>;
private scheduledDispatchIntervalTimer?: ReturnType<typeof setInterval>;
private scheduledDispatchScanRunning = false;
private inboundLongMessageInitialTimer?: ReturnType<typeof setTimeout>;
private inboundLongMessageIntervalTimer?: ReturnType<typeof setInterval>;
constructor(
private readonly prisma: PrismaService,
@@ -312,6 +323,25 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
);
this.scheduledDispatchIntervalTimer.unref?.();
}
if (process.env.CMPP_INBOUND_LONG_MESSAGE_SCAN_ENABLED !== 'false') {
this.inboundLongMessageInitialTimer = setTimeout(
() => void this.expireInboundLongMessages().catch((error) => {
this.logger.error(`Failed to expire inbound CMPP long messages: ${String(error)}`);
}),
INBOUND_LONG_MESSAGE_SCAN_INITIAL_DELAY_MS,
);
this.inboundLongMessageInitialTimer.unref?.();
this.inboundLongMessageIntervalTimer = setInterval(
() => void this.expireInboundLongMessages().catch((error) => {
this.logger.error(`Failed to expire inbound CMPP long messages: ${String(error)}`);
}),
positiveInteger(
process.env.CMPP_INBOUND_LONG_MESSAGE_SCAN_INTERVAL_MS,
DEFAULT_INBOUND_LONG_MESSAGE_SCAN_INTERVAL_MS,
),
);
this.inboundLongMessageIntervalTimer.unref?.();
}
}
async onModuleDestroy() {
@@ -319,6 +349,8 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
if (this.receiptTimeoutIntervalTimer) clearInterval(this.receiptTimeoutIntervalTimer);
if (this.scheduledDispatchInitialTimer) clearTimeout(this.scheduledDispatchInitialTimer);
if (this.scheduledDispatchIntervalTimer) clearInterval(this.scheduledDispatchIntervalTimer);
if (this.inboundLongMessageInitialTimer) clearTimeout(this.inboundLongMessageInitialTimer);
if (this.inboundLongMessageIntervalTimer) clearInterval(this.inboundLongMessageIntervalTimer);
await this.worker?.close();
await this.sendQueue?.close();
await this.gatewayQueue?.close();
@@ -2027,28 +2059,174 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
if (!application) {
throw new BadRequestException('CMPP account is invalid');
}
const dailyQuota = await this.tryReserveDailySendQuota(application.id, phoneNumbers.length);
if (data.longMessage) {
if (data.remoteIp && !isIpAllowed(data.remoteIp, application.ipAllowlist.map((item) => item.ipCidr))) {
throw new BadRequestException('CMPP source IP is not in application allowlist');
}
validateInboundApplicationSrcId(data.srcId, application);
const collection = await this.collectInboundLongMessageFragment(data, application, phoneNumbers);
if (collection.response) {
return collection.response;
}
if (!collection.complete) {
return {
accepted: true,
messageId: collection.messageId,
status: 'fragment_pending',
fragmentPending: true,
receivedSegments: collection.receivedSegments,
segmentTotal: data.longMessage.total,
phoneCount: phoneNumbers.length,
messages: phoneNumbers.map((phoneNumber) => ({
phoneNumber,
messageId: collection.messageId,
status: 'fragment_pending',
})),
};
}
try {
const response = await this.recoverCompletedInboundLongMessageResponse(
collection.messageId,
phoneNumbers,
) ?? await this.submitCompleteInboundMessage({
...data,
content: collection.content,
sequenceId: collection.sequenceId,
longMessage: undefined,
}, phoneNumbers, application, collection.messageId);
await this.prisma.cmppInboundLongMessage.update({
where: { id: collection.groupId },
data: {
status: 'completed',
response: JSON.parse(JSON.stringify(response)) as Prisma.InputJsonValue,
completedAt: new Date(),
},
});
return response;
} catch (error) {
await this.prisma.cmppInboundLongMessage.update({
where: { id: collection.groupId },
data: {
status: 'rejected',
completedAt: new Date(),
},
}).catch(() => undefined);
throw error;
}
}
return this.submitCompleteInboundMessage(data, phoneNumbers, application);
}
private async recoverCompletedInboundLongMessageResponse(messageId: string, phoneNumbers: string[]) {
const existing = await this.prisma.smsMessageRecord.findMany({
where: {
cmppSubmitGroupMessageId: messageId,
phoneNumber: { in: phoneNumbers },
},
select: {
id: true,
tenantId: true,
applicationId: true,
batchTaskId: true,
messageId: true,
phoneNumber: true,
status: true,
errorCode: true,
},
});
const byPhone = new Map(existing.map((item) => [item.phoneNumber, item]));
const ordered = phoneNumbers.map((phoneNumber) => byPhone.get(phoneNumber));
if (ordered.some((item) => !item)) {
return null;
}
const messages = ordered.map((item, index) => ({
phoneNumber: phoneNumbers[index],
messageId: item!.messageId,
messageRecordId: item!.id,
taskId: item!.batchTaskId ?? '',
status: item!.status,
}));
const first = ordered[0]!;
const dailyLimitRejected = ordered.every((item) => item!.errorCode === 'DAILY_LIMIT');
return {
accepted: !dailyLimitRejected,
tenantId: first.tenantId ?? '',
applicationId: first.applicationId ?? '',
taskId: first.batchTaskId ?? '',
messageId: first.messageId,
messageRecordId: first.id,
status: dailyLimitRejected ? 'rejected' : 'accepted',
result: dailyLimitRejected ? 8 : undefined,
phoneCount: messages.length,
messages,
};
}
private async submitCompleteInboundMessage(
data: GatewayInboundSubmitDto,
phoneNumbers: string[],
application: Awaited<ReturnType<SendChainService['findInboundApplication']>>,
requestedGroupMessageId?: string,
) {
if (!application) {
throw new BadRequestException('CMPP account is invalid');
}
const persisted = requestedGroupMessageId
? await this.prisma.smsMessageRecord.findMany({
where: {
cmppSubmitGroupMessageId: requestedGroupMessageId,
phoneNumber: { in: phoneNumbers },
},
select: {
id: true,
tenantId: true,
applicationId: true,
batchTaskId: true,
messageId: true,
phoneNumber: true,
status: true,
errorCode: true,
},
})
: [];
const persistedByPhone = new Map(persisted.map((item) => [item.phoneNumber, item]));
const missingPhoneCount = phoneNumbers.filter((phoneNumber) => !persistedByPhone.has(phoneNumber)).length;
const dailyQuota = missingPhoneCount > 0
? await this.tryReserveDailySendQuota(application.id, missingPhoneCount)
: { reserved: true, dailyLimit: application.dailyLimit ?? 100000 };
const dailyLimitRejection = dailyQuota.reserved
? undefined
: {
code: 'DAILY_LIMIT',
reason: `应用当日发送上限${dailyQuota.dailyLimit}条,本次${phoneNumbers.length}条超出剩余配额`,
reason: `应用当日发送上限${dailyQuota.dailyLimit}条,本次${missingPhoneCount}条超出剩余配额`,
};
const submitGroupMessageId = `MSG-${randomUUID()}`;
const submitGroupMessageId = requestedGroupMessageId ?? `MSG-${randomUUID()}`;
const submissions = phoneNumbers.map((phoneNumber, index) => ({
phoneNumber,
messageId: index === 0 ? submitGroupMessageId : `MSG-${randomUUID()}`,
persisted: persistedByPhone.get(phoneNumber),
messageId: persistedByPhone.get(phoneNumber)?.messageId
?? (index === 0 ? submitGroupMessageId : `MSG-${randomUUID()}`),
}));
const results: GatewayInboundSingleSubmitResult[] = [];
const concurrency = 10;
for (let offset = 0; offset < submissions.length; offset += concurrency) {
const batch = submissions.slice(offset, offset + concurrency);
results.push(...await Promise.all(batch.map((submission) => this.submitInboundSingleMessage({
...data,
phoneNumber: submission.phoneNumber,
phoneNumbers: undefined,
}, submission.messageId, submitGroupMessageId, dailyLimitRejection))));
results.push(...await Promise.all(batch.map((submission) => submission.persisted
? Promise.resolve({
accepted: submission.persisted.errorCode !== 'DAILY_LIMIT',
tenantId: submission.persisted.tenantId ?? application.tenantId,
applicationId: submission.persisted.applicationId ?? application.id,
taskId: submission.persisted.batchTaskId ?? '',
messageId: submission.persisted.messageId,
messageRecordId: submission.persisted.id,
status: submission.persisted.status,
})
: this.submitInboundSingleMessage({
...data,
phoneNumber: submission.phoneNumber,
phoneNumbers: undefined,
}, submission.messageId, submitGroupMessageId, dailyLimitRejection))));
}
const first = results[0];
return {
@@ -2065,6 +2243,173 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy {
};
}
private async collectInboundLongMessageFragment(
data: GatewayInboundSubmitDto,
application: NonNullable<Awaited<ReturnType<SendChainService['findInboundApplication']>>>,
phoneNumbers: string[],
) {
const fragment = data.longMessage;
if (!fragment || !Number.isInteger(fragment.reference) || fragment.reference < 0 || fragment.reference > 65535
|| !Number.isInteger(fragment.total) || fragment.total < 2 || fragment.total > 255
|| !Number.isInteger(fragment.index) || fragment.index < 1 || fragment.index > fragment.total
|| !Number.isInteger(fragment.format) || fragment.format < 0 || fragment.format > 255) {
throw new BadRequestException('CMPP long message fragment metadata is invalid');
}
const groupKey = createHash('sha256').update(JSON.stringify({
applicationId: application.id,
account: data.account,
srcId: data.srcId?.trim() ?? '',
phoneNumbers,
reference: fragment.reference,
total: fragment.total,
format: fragment.format,
})).digest('hex');
const contentHash = createHash('sha256').update(data.content).digest('hex');
const now = new Date();
const expiresAt = new Date(now.getTime() + positiveInteger(
process.env.CMPP_INBOUND_LONG_MESSAGE_TTL_SECONDS,
300,
) * 1000);
return this.prisma.$transaction(async (tx) => {
await tx.$executeRaw`SELECT pg_advisory_xact_lock(hashtextextended(${groupKey}, 0))`;
await tx.cmppInboundLongMessage.updateMany({
where: {
groupKey,
status: { in: ['collecting', 'processing'] },
expiresAt: { lte: now },
},
data: { status: 'expired', completedAt: now },
});
const recent = await tx.cmppInboundLongMessage.findFirst({
where: {
groupKey,
expiresAt: { gt: now },
},
include: { segments: { orderBy: { segmentIndex: 'asc' } } },
orderBy: { createdAt: 'desc' },
});
const matchingRecentSegment = recent?.segments.find((item) => item.segmentIndex === fragment.index);
if (recent && ['completed', 'rejected'].includes(recent.status)
&& matchingRecentSegment?.contentHash === contentHash
&& matchingRecentSegment.sequenceId === (data.sequenceId == null ? null : String(data.sequenceId))) {
return {
complete: recent.status === 'completed',
groupId: recent.id,
messageId: recent.messageId,
receivedSegments: recent.segments.length,
response: recent.response as any,
content: recent.segments.map((item) => item.content).join(''),
sequenceId: parseOptionalSequenceId(recent.segments[0]?.sequenceId),
};
}
let group = recent && ['collecting', 'processing'].includes(recent.status) ? recent : null;
if (!group) {
group = await tx.cmppInboundLongMessage.create({
data: {
tenantId: application.tenantId,
applicationId: application.id,
groupKey,
account: data.account,
srcId: data.srcId?.trim() || null,
phoneNumbers,
concatReference: fragment.reference,
segmentTotal: fragment.total,
msgFmt: fragment.format,
messageId: `MSG-${randomUUID()}`,
expiresAt,
},
include: { segments: { orderBy: { segmentIndex: 'asc' } } },
});
}
if (group.status === 'processing') {
const processingStaleMs = positiveInteger(
process.env.CMPP_INBOUND_LONG_MESSAGE_PROCESSING_STALE_SECONDS,
DEFAULT_INBOUND_LONG_MESSAGE_PROCESSING_STALE_SECONDS,
) * 1000;
const complete = group.segments.length === fragment.total
&& group.segments.every((item, index) => item.segmentIndex === index + 1);
if (complete && now.getTime() - group.updatedAt.getTime() >= processingStaleMs) {
await tx.cmppInboundLongMessage.update({
where: { id: group.id },
data: { status: 'processing', expiresAt },
});
return {
complete: true,
groupId: group.id,
messageId: group.messageId,
receivedSegments: group.segments.length,
response: null,
content: group.segments.map((item) => item.content).join(''),
sequenceId: parseOptionalSequenceId(group.segments[0]?.sequenceId),
};
}
return {
complete: false,
groupId: group.id,
messageId: group.messageId,
receivedSegments: group.segments.length,
response: group.response as any,
content: '',
sequenceId: undefined,
};
}
const existing = group.segments.find((item) => item.segmentIndex === fragment.index);
if (existing && (existing.contentHash !== contentHash
|| existing.sequenceId !== (data.sequenceId == null ? null : String(data.sequenceId)))) {
throw new BadRequestException(`CMPP long message fragment ${fragment.index} conflicts with the stored fragment`);
}
if (!existing) {
await tx.cmppInboundLongMessageSegment.create({
data: {
groupId: group.id,
segmentIndex: fragment.index,
sequenceId: data.sequenceId == null ? null : String(data.sequenceId),
content: data.content,
contentHash,
},
});
}
const segments = await tx.cmppInboundLongMessageSegment.findMany({
where: { groupId: group.id },
orderBy: { segmentIndex: 'asc' },
});
const complete = segments.length === fragment.total
&& segments.every((item, index) => item.segmentIndex === index + 1);
if (complete) {
await tx.cmppInboundLongMessage.update({
where: { id: group.id },
data: { status: 'processing', expiresAt },
});
}
return {
complete,
groupId: group.id,
messageId: group.messageId,
receivedSegments: segments.length,
response: null,
content: complete ? segments.map((item) => item.content).join('') : '',
sequenceId: parseOptionalSequenceId(segments[0]?.sequenceId),
};
});
}
async expireInboundLongMessages(now = new Date()) {
return this.prisma.cmppInboundLongMessage.updateMany({
where: {
status: { in: ['collecting', 'processing'] },
expiresAt: { lte: now },
},
data: {
status: 'expired',
completedAt: now,
},
});
}
private async submitInboundSingleMessage(
data: GatewayInboundSubmitDto & { phoneNumber: string },
messageId: string,
@@ -3782,6 +4127,12 @@ function positiveInteger(value: string | undefined, fallback: number) {
return Number.isInteger(parsed) && parsed > 0 ? parsed : fallback;
}
function parseOptionalSequenceId(value: string | null | undefined) {
if (!value) return undefined;
const parsed = Number(value);
return Number.isInteger(parsed) && parsed >= 0 && parsed <= 0xffffffff ? parsed : undefined;
}
function shanghaiDateKey(now = new Date()) {
const parts = new Intl.DateTimeFormat('en-CA', {
timeZone: 'Asia/Shanghai',