From 53073461e9b991978a3063a691d454b68c20ced9 Mon Sep 17 00:00:00 2001 From: hectorzhao Date: Thu, 20 Aug 2026 17:50:00 +0800 Subject: [PATCH] fix(cmpp): honor UTC inbox retry leases --- .../send-chain/send-batch-entry.service.ts | 3 +- api/src/send-chain/send-chain.service.spec.ts | 36 +++++++++++++++++++ .../send-chain/send-inbound-entry.service.ts | 8 ++--- .../first-version-development-requirements.md | 1 + docs/system-functional-test-cases.md | 2 +- 5 files changed, 44 insertions(+), 6 deletions(-) diff --git a/api/src/send-chain/send-batch-entry.service.ts b/api/src/send-chain/send-batch-entry.service.ts index c984fc6..7d90452 100644 --- a/api/src/send-chain/send-batch-entry.service.ts +++ b/api/src/send-chain/send-batch-entry.service.ts @@ -530,6 +530,7 @@ async tryReserveDailySendQuota(applicationId: string, requestedCount: number, re throw new BadRequestException('发送号码数量必须为正整数'); } const usageDate = shanghaiDateKey(); + const usageDateValue = new Date(`${usageDate}T00:00:00.000Z`); const reserve = (client: Pick) => { const reservationId = randomUUID(); return client.$queryRaw>(Prisma.sql` @@ -583,7 +584,7 @@ async tryReserveDailySendQuota(applicationId: string, requestedCount: number, re reservationKey: normalizedReservationKey, tenantId: row.tenantId, applicationId, - usageDate, + usageDate: usageDateValue, requestedCount, dailyLimit: Number(row.dailyLimit), usedCount: row.usedCount == null ? null : Number(row.usedCount), diff --git a/api/src/send-chain/send-chain.service.spec.ts b/api/src/send-chain/send-chain.service.spec.ts index ec3342c..60000f1 100644 --- a/api/src/send-chain/send-chain.service.spec.ts +++ b/api/src/send-chain/send-chain.service.spec.ts @@ -336,6 +336,14 @@ function createPrismaMock() { count: jest.fn().mockResolvedValue(0), 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: { findFirst: jest.fn().mockResolvedValue(null), 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 () => { const { service, prisma } = createService(); const add = jest.fn().mockResolvedValue(undefined); diff --git a/api/src/send-chain/send-inbound-entry.service.ts b/api/src/send-chain/send-inbound-entry.service.ts index 1080410..e5d59d2 100644 --- a/api/src/send-chain/send-inbound-entry.service.ts +++ b/api/src/send-chain/send-inbound-entry.service.ts @@ -1086,10 +1086,10 @@ startInboundWorkflowWorker() { FROM "CmppInboundSubmissionInbox" WHERE ( status = 'pending' - AND "nextAttemptAt" <= NOW() + AND "nextAttemptAt" <= (NOW() AT TIME ZONE 'UTC') ) OR ( 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 -- preserving FIFO within each class. This keeps the V5 priority contract effective @@ -1101,9 +1101,9 @@ startInboundWorkflowWorker() { UPDATE "CmppInboundSubmissionInbox" AS inbox SET status = 'processing', attempts = inbox.attempts + 1, - "lockedAt" = NOW(), + "lockedAt" = (NOW() AT TIME ZONE 'UTC'), "lockedBy" = ${this.inboundWorkflowWorkerId}, - "updatedAt" = NOW() + "updatedAt" = (NOW() AT TIME ZONE 'UTC') FROM candidates WHERE inbox.id = candidates.id RETURNING inbox.id, inbox."requestKey", inbox."applicationId", inbox.attempts, inbox.payload diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index d000a9a..aba3593 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -2105,6 +2105,7 @@ - V5 API入站数据库往返治理不得改变SubmitResp、余额冻结、号码频控、模板/签名、路由、队列或失败回执语义。同一Submit入口已取得的应用、企业和IP白名单快照必须复用于其全部目标号码,不得逐号码重复查询;任务风控与号码频控共用的默认规则完整性检查必须使用短TTL并发单飞,正常完整状态最多执行一次聚合数据库检查,实际生效规则仍逐次读取;刚持久化的单条CMPP内部任务和消息可凭已知ID、优先级直接入队,不得为入队重新查询同一任务和消息。缓存只覆盖默认规则“是否齐全”,检查失败必须立即失效,规则缺失须在短TTL到期后自动恢复;不得缓存余额、频控计数、应用启停或实际生效规则。 - 面向完整处理500条/秒目标的第一阶段,CMPP SubmitResp语义调整为“已完成最小协议/账号/IP/Src_Id校验且已持久化接收事实”,不再同步等待模板、风控、频控、日限、计费、路由和入队。Gateway必须为同一连接内同一Submit重试生成稳定请求键;PostgreSQL Inbox以唯一请求键和载荷摘要幂等保存原请求、稳定内部MessageId及响应。相同键相同载荷返回原响应,相同键不同载荷必须拒绝。 - 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不丢不重且最终可排空,并为后续500条/秒全链路扩容建立解耦边界;不能仅凭SubmitResp吞吐宣称完整500条/秒。压测必须同时报告SubmitResp成功率/延迟、Inbox pending/processing/最老等待、异步完成速率和排空时间,以及命令Stream、结果Outbox和数据库最终对账。 - 系统监控标题说明必须明确标注数据来自 Prometheus;“服务关键指标”位于趋势/核心服务区域之后、活动告警之前,并提供统一的“告警阈值设置”入口。安全检测页不重复渲染大号标题,说明文字必须明确标注使用 Fail2ban。 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index e913b65..97a5f66 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -4736,7 +4736,7 @@ npm run verify:phase8 | TC-CMPP-500-P1-003 | 幂等冲突 | 使用同一请求键提交不同载荷 | API拒绝冲突且不覆盖原Inbox/响应 | | TC-CMPP-500-P1-004 | 多号码与长短信 | 分别提交多号码Submit和完整长短信分片 | Inbox保留全部号码及稳定子MessageId;长短信Worker处理的是完整重组正文,不是最后一片正文 | | 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-008 | 优先级 | 在普通Inbox积压期间持续混入priority应用 | priority先领取且类内FIFO;普通队列最终可排空;同时记录两类等待分位数 | | TC-CMPP-500-P1-009 | 进程隔离 | 检查systemd、进程、连接和指标端口 | API角色不运行发送后台任务;`cmpp-send-worker`独立非root运行,指标仅监听127.0.0.1:9465 |