fix(cmpp): honor UTC inbox retry leases

This commit is contained in:
hectorzhao
2026-08-20 17:50:00 +08:00
parent 0b63bcd74e
commit 53073461e9
5 changed files with 44 additions and 6 deletions
@@ -530,6 +530,7 @@ async tryReserveDailySendQuota(applicationId: string, requestedCount: number, re
throw new BadRequestException('发送号码数量必须为正整数'); throw new BadRequestException('发送号码数量必须为正整数');
} }
const usageDate = shanghaiDateKey(); const usageDate = shanghaiDateKey();
const usageDateValue = new Date(`${usageDate}T00:00:00.000Z`);
const reserve = (client: Pick<Prisma.TransactionClient, '$queryRaw'>) => { const reserve = (client: Pick<Prisma.TransactionClient, '$queryRaw'>) => {
const reservationId = randomUUID(); const reservationId = randomUUID();
return client.$queryRaw<Array<{ tenantId: string; dailyLimit: number; usedCount: number | null }>>(Prisma.sql` return client.$queryRaw<Array<{ tenantId: string; dailyLimit: number; usedCount: number | null }>>(Prisma.sql`
@@ -583,7 +584,7 @@ async tryReserveDailySendQuota(applicationId: string, requestedCount: number, re
reservationKey: normalizedReservationKey, reservationKey: normalizedReservationKey,
tenantId: row.tenantId, tenantId: row.tenantId,
applicationId, applicationId,
usageDate, usageDate: usageDateValue,
requestedCount, requestedCount,
dailyLimit: Number(row.dailyLimit), dailyLimit: Number(row.dailyLimit),
usedCount: row.usedCount == null ? null : Number(row.usedCount), usedCount: row.usedCount == null ? null : Number(row.usedCount),
@@ -336,6 +336,14 @@ function createPrismaMock() {
count: jest.fn().mockResolvedValue(0), count: jest.fn().mockResolvedValue(0),
findFirst: jest.fn().mockResolvedValue(null), findFirst: jest.fn().mockResolvedValue(null),
}, },
smsApplicationDailyReservation: {
findUnique: jest.fn().mockResolvedValue(null),
create: jest.fn().mockResolvedValue({ id: 'daily-reservation-1' }),
},
phoneFrequencyReservation: {
findUnique: jest.fn().mockResolvedValue(null),
create: jest.fn().mockResolvedValue({ id: 'frequency-reservation-1' }),
},
smsBillingRecord: { smsBillingRecord: {
findFirst: jest.fn().mockResolvedValue(null), findFirst: jest.fn().mockResolvedValue(null),
create: jest.fn().mockResolvedValue({ id: 'bill-1' }), create: jest.fn().mockResolvedValue({ id: 'bill-1' }),
@@ -2194,6 +2202,34 @@ describe('SendChainService', () => {
} }
}); });
it('persists an idempotent daily quota reservation with a Prisma Date value', async () => {
const { service, prisma } = createService();
prisma.$queryRaw.mockResolvedValueOnce([{ tenantId: 'tenant-1', dailyLimit: 100000, usedCount: 1 }]);
await expect((service as any).tryReserveDailySendQuota('app-1', 1, 'workflow-1:daily-quota'))
.resolves.toEqual({ dailyLimit: 100000, usedCount: 1, reserved: true });
expect(prisma.smsApplicationDailyReservation.create).toHaveBeenCalledWith({
data: expect.objectContaining({
reservationKey: 'workflow-1:daily-quota',
usageDate: expect.any(Date),
}),
});
expect(prisma.smsApplicationDailyReservation.create.mock.calls[0][0].data.usageDate.toISOString())
.toMatch(/^\d{4}-\d{2}-\d{2}T00:00:00\.000Z$/);
});
it('claims Inbox leases against UTC for timestamp-without-time-zone columns', async () => {
const { service, prisma } = createService();
prisma.$queryRaw.mockResolvedValueOnce([]);
await (service as any).submission.inboundEntry.claimInboundWorkflows(5);
const sql = prisma.$queryRaw.mock.calls[0][0].strings.join(' ');
expect(sql).toContain("NOW() AT TIME ZONE 'UTC'");
expect(sql).toContain('FOR UPDATE SKIP LOCKED');
});
it('enqueues a freshly persisted inbound message without querying the task and message again', async () => { it('enqueues a freshly persisted inbound message without querying the task and message again', async () => {
const { service, prisma } = createService(); const { service, prisma } = createService();
const add = jest.fn().mockResolvedValue(undefined); const add = jest.fn().mockResolvedValue(undefined);
@@ -1086,10 +1086,10 @@ startInboundWorkflowWorker() {
FROM "CmppInboundSubmissionInbox" FROM "CmppInboundSubmissionInbox"
WHERE ( WHERE (
status = 'pending' status = 'pending'
AND "nextAttemptAt" <= NOW() AND "nextAttemptAt" <= (NOW() AT TIME ZONE 'UTC')
) OR ( ) OR (
status = 'processing' status = 'processing'
AND "lockedAt" <= NOW() - make_interval(secs => ${staleSeconds}) AND "lockedAt" <= (NOW() AT TIME ZONE 'UTC') - make_interval(secs => ${staleSeconds})
) )
-- Priority applications enter the same durable Inbox, but are claimed first while -- Priority applications enter the same durable Inbox, but are claimed first while
-- preserving FIFO within each class. This keeps the V5 priority contract effective -- preserving FIFO within each class. This keeps the V5 priority contract effective
@@ -1101,9 +1101,9 @@ startInboundWorkflowWorker() {
UPDATE "CmppInboundSubmissionInbox" AS inbox UPDATE "CmppInboundSubmissionInbox" AS inbox
SET status = 'processing', SET status = 'processing',
attempts = inbox.attempts + 1, attempts = inbox.attempts + 1,
"lockedAt" = NOW(), "lockedAt" = (NOW() AT TIME ZONE 'UTC'),
"lockedBy" = ${this.inboundWorkflowWorkerId}, "lockedBy" = ${this.inboundWorkflowWorkerId},
"updatedAt" = NOW() "updatedAt" = (NOW() AT TIME ZONE 'UTC')
FROM candidates FROM candidates
WHERE inbox.id = candidates.id WHERE inbox.id = candidates.id
RETURNING inbox.id, inbox."requestKey", inbox."applicationId", inbox.attempts, inbox.payload RETURNING inbox.id, inbox."requestKey", inbox."applicationId", inbox.attempts, inbox.payload
@@ -2105,6 +2105,7 @@
- V5 API入站数据库往返治理不得改变SubmitResp、余额冻结、号码频控、模板/签名、路由、队列或失败回执语义。同一Submit入口已取得的应用、企业和IP白名单快照必须复用于其全部目标号码,不得逐号码重复查询;任务风控与号码频控共用的默认规则完整性检查必须使用短TTL并发单飞,正常完整状态最多执行一次聚合数据库检查,实际生效规则仍逐次读取;刚持久化的单条CMPP内部任务和消息可凭已知ID、优先级直接入队,不得为入队重新查询同一任务和消息。缓存只覆盖默认规则“是否齐全”,检查失败必须立即失效,规则缺失须在短TTL到期后自动恢复;不得缓存余额、频控计数、应用启停或实际生效规则。 - V5 API入站数据库往返治理不得改变SubmitResp、余额冻结、号码频控、模板/签名、路由、队列或失败回执语义。同一Submit入口已取得的应用、企业和IP白名单快照必须复用于其全部目标号码,不得逐号码重复查询;任务风控与号码频控共用的默认规则完整性检查必须使用短TTL并发单飞,正常完整状态最多执行一次聚合数据库检查,实际生效规则仍逐次读取;刚持久化的单条CMPP内部任务和消息可凭已知ID、优先级直接入队,不得为入队重新查询同一任务和消息。缓存只覆盖默认规则“是否齐全”,检查失败必须立即失效,规则缺失须在短TTL到期后自动恢复;不得缓存余额、频控计数、应用启停或实际生效规则。
- 面向完整处理500条/秒目标的第一阶段,CMPP SubmitResp语义调整为“已完成最小协议/账号/IP/Src_Id校验且已持久化接收事实”,不再同步等待模板、风控、频控、日限、计费、路由和入队。Gateway必须为同一连接内同一Submit重试生成稳定请求键;PostgreSQL Inbox以唯一请求键和载荷摘要幂等保存原请求、稳定内部MessageId及响应。相同键相同载荷返回原响应,相同键不同载荷必须拒绝。 - 面向完整处理500条/秒目标的第一阶段,CMPP SubmitResp语义调整为“已完成最小协议/账号/IP/Src_Id校验且已持久化接收事实”,不再同步等待模板、风控、频控、日限、计费、路由和入队。Gateway必须为同一连接内同一Submit重试生成稳定请求键;PostgreSQL Inbox以唯一请求键和载荷摘要幂等保存原请求、稳定内部MessageId及响应。相同键相同载荷返回原响应,相同键不同载荷必须拒绝。
- Inbox Worker必须作为独立非root进程运行,使用独立数据库连接池和持续有界工作池;领取使用短事务、`FOR UPDATE SKIP LOCKED`和过期租约回收,实际业务处理在领取事务外执行。成功逐条完成,失败以有上限退避回到pending且不得静默丢弃。日发送配额、号码频控及余额冻结必须使用由Inbox请求派生的持久幂等键,确保进程崩溃或租约回收不会重复扣量、重复计频或重复冻结。 - Inbox Worker必须作为独立非root进程运行,使用独立数据库连接池和持续有界工作池;领取使用短事务、`FOR UPDATE SKIP LOCKED`和过期租约回收,实际业务处理在领取事务外执行。成功逐条完成,失败以有上限退避回到pending且不得静默丢弃。日发送配额、号码频控及余额冻结必须使用由Inbox请求派生的持久幂等键,确保进程崩溃或租约回收不会重复扣量、重复计频或重复冻结。
- Inbox的`DateTime`列沿用平台UTC无时区存储口径,领取、租约和退避SQL必须显式使用UTC时钟比较,不能受数据库会话或宿主机Asia/Shanghai时区影响;日期型日配额预留写Prisma时必须传合法Date对象,不能把`YYYY-MM-DD`字符串当作DateTime。
- 优先应用和普通应用共享同一耐久Inbox,但领取顺序必须保留优先级并在同一优先级内FIFO;进入BullMQ后继续沿用priority=1、normal=100。快路径成功只代表平台已可靠接收,异步业务拒绝仍必须落真实消息/任务状态并按既有CMPP失败回执链路通知客户,不得伪装为供应商最终送达。 - 优先应用和普通应用共享同一耐久Inbox,但领取顺序必须保留优先级并在同一优先级内FIFO;进入BullMQ后继续沿用priority=1、normal=100。快路径成功只代表平台已可靠接收,异步业务拒绝仍必须落真实消息/任务状态并按既有CMPP失败回执链路通知客户,不得伪装为供应商最终送达。
- 第一阶段的验收是入口可持续接收、Inbox不丢不重且最终可排空,并为后续500条/秒全链路扩容建立解耦边界;不能仅凭SubmitResp吞吐宣称完整500条/秒。压测必须同时报告SubmitResp成功率/延迟、Inbox pending/processing/最老等待、异步完成速率和排空时间,以及命令Stream、结果Outbox和数据库最终对账。 - 第一阶段的验收是入口可持续接收、Inbox不丢不重且最终可排空,并为后续500条/秒全链路扩容建立解耦边界;不能仅凭SubmitResp吞吐宣称完整500条/秒。压测必须同时报告SubmitResp成功率/延迟、Inbox pending/processing/最老等待、异步完成速率和排空时间,以及命令Stream、结果Outbox和数据库最终对账。
- 系统监控标题说明必须明确标注数据来自 Prometheus;“服务关键指标”位于趋势/核心服务区域之后、活动告警之前,并提供统一的“告警阈值设置”入口。安全检测页不重复渲染大号标题,说明文字必须明确标注使用 Fail2ban。 - 系统监控标题说明必须明确标注数据来自 Prometheus;“服务关键指标”位于趋势/核心服务区域之后、活动告警之前,并提供统一的“告警阈值设置”入口。安全检测页不重复渲染大号标题,说明文字必须明确标注使用 Fail2ban。
+1 -1
View File
@@ -4736,7 +4736,7 @@ npm run verify:phase8
| TC-CMPP-500-P1-003 | 幂等冲突 | 使用同一请求键提交不同载荷 | API拒绝冲突且不覆盖原Inbox/响应 | | TC-CMPP-500-P1-003 | 幂等冲突 | 使用同一请求键提交不同载荷 | API拒绝冲突且不覆盖原Inbox/响应 |
| TC-CMPP-500-P1-004 | 多号码与长短信 | 分别提交多号码Submit和完整长短信分片 | Inbox保留全部号码及稳定子MessageId;长短信Worker处理的是完整重组正文,不是最后一片正文 | | TC-CMPP-500-P1-004 | 多号码与长短信 | 分别提交多号码Submit和完整长短信分片 | Inbox保留全部号码及稳定子MessageId;长短信Worker处理的是完整重组正文,不是最后一片正文 |
| TC-CMPP-500-P1-005 | Worker领取与逐条完成 | 启动两个Worker并制造至少一个批次积压 | `SKIP LOCKED`领取不重复;每条独立完成,处理逻辑不占用领取事务 | | TC-CMPP-500-P1-005 | Worker领取与逐条完成 | 启动两个Worker并制造至少一个批次积压 | `SKIP LOCKED`领取不重复;每条独立完成,处理逻辑不占用领取事务 |
| TC-CMPP-500-P1-006 | 崩溃恢复 | 在领取后终止Worker,超过租约后重启 | processing记录被回收并完成日限频控、冻结、任务和消息均不重复 | | TC-CMPP-500-P1-006 | 崩溃恢复与时区 | 数据库会话使用Asia/Shanghai,在领取后终止Worker,超过租约后重启 | UTC无时区列按UTC时钟比较,退避不会立即重领;processing记录被回收并完成日限Date值合法,频控、冻结、任务和消息均不重复 |
| TC-CMPP-500-P1-007 | 异步业务拒绝 | 让已耐久受理短信命中真实模板/风控/余额拒绝 | SubmitResp仍表示已接收;后台生成真实失败状态和失败回执,不进入供应商发送 | | TC-CMPP-500-P1-007 | 异步业务拒绝 | 让已耐久受理短信命中真实模板/风控/余额拒绝 | SubmitResp仍表示已接收;后台生成真实失败状态和失败回执,不进入供应商发送 |
| TC-CMPP-500-P1-008 | 优先级 | 在普通Inbox积压期间持续混入priority应用 | priority先领取且类内FIFO;普通队列最终可排空;同时记录两类等待分位数 | | TC-CMPP-500-P1-008 | 优先级 | 在普通Inbox积压期间持续混入priority应用 | priority先领取且类内FIFO;普通队列最终可排空;同时记录两类等待分位数 |
| TC-CMPP-500-P1-009 | 进程隔离 | 检查systemd、进程、连接和指标端口 | API角色不运行发送后台任务;`cmpp-send-worker`独立非root运行,指标仅监听127.0.0.1:9465 | | TC-CMPP-500-P1-009 | 进程隔离 | 检查systemd、进程、连接和指标端口 | API角色不运行发送后台任务;`cmpp-send-worker`独立非root运行,指标仅监听127.0.0.1:9465 |