diff --git a/api/src/send-chain/send-chain.service.spec.ts b/api/src/send-chain/send-chain.service.spec.ts index b37de80..4f2aa74 100644 --- a/api/src/send-chain/send-chain.service.spec.ts +++ b/api/src/send-chain/send-chain.service.spec.ts @@ -924,6 +924,28 @@ describe('SendChainService', () => { })); }); + it('coalesces concurrent batch progress refreshes and keeps a trailing refresh', async () => { + const { service, prisma } = createService(); + let resolveFirst: ((value: Array<{ status: string; _count: { _all: number } }>) => void) | undefined; + prisma.smsMessageRecord.groupBy + .mockImplementationOnce(() => new Promise((resolve) => { resolveFirst = resolve; })) + .mockResolvedValue([{ status: 'delivered', _count: { _all: 2 } }]); + + const first = service['submission'].refreshTaskProgress('task-1'); + await Promise.resolve(); + const second = service['submission'].refreshTaskProgress('task-1'); + + expect(prisma.smsMessageRecord.groupBy).toHaveBeenCalledTimes(1); + resolveFirst?.([{ status: 'delivered', _count: { _all: 1 } }]); + await Promise.all([first, second]); + + expect(prisma.smsMessageRecord.groupBy).toHaveBeenCalledTimes(2); + expect(prisma.smsBatchTask.update).toHaveBeenLastCalledWith({ + where: { id: 'task-1' }, + data: expect.objectContaining({ progressTotal: 2, successTotal: 2, status: 'finished' }), + }); + }); + it('does not expose CMPP internal tasks through client task detail or messages', async () => { const { service, prisma } = createService(); prisma.smsBatchTask.findFirst.mockResolvedValue(null); diff --git a/api/src/send-chain/send-gateway-submit.service.ts b/api/src/send-chain/send-gateway-submit.service.ts index 8f61918..cb07e53 100644 --- a/api/src/send-chain/send-gateway-submit.service.ts +++ b/api/src/send-chain/send-gateway-submit.service.ts @@ -24,6 +24,8 @@ export class SendGatewaySubmitService { private sendQueue?: Queue; private gatewayQueue?: Queue; private worker?: Worker; + private readonly taskProgressRefreshes = new Map>(); + private readonly dirtyTaskProgressRefreshes = new Set(); constructor( private readonly prisma: PrismaService, @@ -449,6 +451,29 @@ async waitForChannelRateLimit(channelId: string, tps: number) { } async refreshTaskProgress(batchTaskId: string) { + const running = this.taskProgressRefreshes.get(batchTaskId); + if (running) { + // A state transition committed after the running aggregate may not be visible + // to its snapshot. Mark one trailing pass instead of issuing another identical + // GROUP BY concurrently for every message/result/receipt callback. + this.dirtyTaskProgressRefreshes.add(batchTaskId); + return running; + } + const refresh = this.refreshTaskProgressUntilClean(batchTaskId); + this.taskProgressRefreshes.set(batchTaskId, refresh); + try { + await refresh; + } finally { + if (this.taskProgressRefreshes.get(batchTaskId) === refresh) { + this.taskProgressRefreshes.delete(batchTaskId); + } + this.dirtyTaskProgressRefreshes.delete(batchTaskId); + } + } + +private async refreshTaskProgressUntilClean(batchTaskId: string) { + do { + this.dirtyTaskProgressRefreshes.delete(batchTaskId); const groups = await this.prisma.smsMessageRecord.groupBy({ by: ['status'], where: { batchTaskId }, @@ -468,6 +493,7 @@ async refreshTaskProgress(batchTaskId: string) { where: { id: batchTaskId }, data: { progressTotal, submittedTotal, successTotal, failedTotal, unknownTotal, timeoutTotal, status }, }); + } while (this.dirtyTaskProgressRefreshes.has(batchTaskId)); } getSendQueue(): Queue { diff --git a/docs/codebase-modularization-roadmap.md b/docs/codebase-modularization-roadmap.md index 5af6d33..1525e88 100644 --- a/docs/codebase-modularization-roadmap.md +++ b/docs/codebase-modularization-roadmap.md @@ -1172,3 +1172,4 @@ global,不能错误归入client。`client-signature-*`、发送页、企业认 - `send-chain.helpers`提供无副作用的优先级、主备和稳定weight选择;提交服务只提供消息稳定键与真实在线/报备候选,不保存进程内轮询状态。 - Gateway既有Submit池与结果Outbox继续独立有界;先以正价六通道实测定位容量,只有证据证明单条回调仍触发停止线时才扩展批量回调协议。 - Gateway上游模块对单分片使用聚合结果作为唯一Outbox边界,对多分片保留逐片持久化;结果Outbox不承担业务去重猜测,事件数量由上游已知分片数确定。 +- 批次任务进度聚合仍归`send-gateway-submit`统一实现;同进程、同批次并发调用使用单飞与尾随刷新收敛重复`GROUP BY`,不改变消息状态写入、跨进程幂等或PostgreSQL最终事实。 diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index b97a627..7483dc3 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -2136,4 +2136,5 @@ - 供应商接受后,释放冻结与正式扣费以一个原子SQL写入两条幂等审计流水;二者净余额变化为零,不得争抢企业账户锁或重复更新账户。只有恢复历史“已释放但未扣费”半完成状态时才取得账户锁补扣。零金额不创建AccountTransaction;拒绝、超时、重投和最终失败继续沿用逐短信幂等释放、扣费和退款。 - 同通道组内在线、已报备、同地域、同最低优先级且非备用的通道按weight和稳定消息键分流;备用、离线及更低优先级不抢占。只有显式配置为同优先级非备用的通道才参与主动双活。 - 单分片短短信的聚合SubmitResult已经携带完整分片信息,只发布聚合Outbox事件;不得再发布内容相同的分片事件造成双倍API回调和数据库写入。多分片长短信仍须在发送下一片前持久化当前分片事件,保持崩溃恢复边界。 +- 同一批次并发提交结果、最终回执和失败处理触发任务进度刷新时,进程内只允许一个PostgreSQL聚合查询执行;并发触发合并为一次尾随刷新,既避免逐消息并发扫描整个批次,也必须覆盖运行中查询快照之后已经提交的状态变化。消息状态、计费和幂等事实仍以PostgreSQL为准,不得缓存聚合结果替代最终刷新。 - 压测使用隔离测试企业、真实PostgreSQL账务、单价`0.0325元/计费条`、三运营商混合号码和六个模拟供应商账号;逐档核对消息、冻结/释放/扣费/退款、SmsBillingRecord、通道分布、队列和数据库。临时单价与主动双活配置须先快照、测试后恢复。 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index 90019a0..877ed38 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -4786,3 +4786,4 @@ npm run verify:phase8 | TC-CMPP-500-P4-007 | priority与normal并发积压 | priority保持明确服务能力且normal最终不饿死 | | TC-CMPP-500-P4-008 | 测试配置治理 | 单价、账户与组项目先快照;临时主动双活和单价测试后恢复;预生产和凭据不变 | | TC-CMPP-500-P4-009 | 单分片与多分片供应商Submit | 单分片只产生一个含segments的聚合Outbox事件;多分片逐片持久化并另有聚合事件,任何分片不丢失 | +| TC-CMPP-500-P4-010 | 同一批次并发任务进度刷新 | 并发调用共享一个运行中聚合查询,并在其间有新状态提交时只补一次尾随聚合;最终任务计数与消息终态一致 | diff --git a/docs/testing-progress.md b/docs/testing-progress.md index 13e19c5..fb5a530 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -3844,3 +3844,5 @@ git diff --check - 该轮最终双Stream归零,冻结/释放各2999、正式扣费2987、退款30、12条供应商最终拒绝只释放不扣费;SmsBillingRecord为charged2957/refunded30,余额恒等式精确成立,六通道与数据库最终无等待锁/idle in transaction。证据确认瓶颈是每个正价Submit结果在企业账户锁上串行,而不是单价0时可见的通道容量。 - 后续最小修复将正常“释放冻结+正式扣费”改为单个幂等SQL:冻结已在入口扣减,流水对净余额为零,因此正常结算不再取得账户锁或更新账户;历史已释放未扣费状态仍用原带锁路径补扣,并保留并发`ON CONFLICT`后一次MVCC恢复读。专项134项和TypeScript已通过,完整门禁、修复提交、新恢复资产与复测待补记。 - 无账户锁修复后100条/秒复测仍在注入结束出现命令Stream`128/1687`、结果Outbox`32/452`;相同观测点完成扣费从388提高至637、等待锁降为0,但单分片短信同时发布分片与聚合两个结果事件,仍造成约双倍API回调和分片审计写入。新增最小Gateway修复:单分片只发布已携带完整segments的聚合事件,多分片仍逐片先持久化,待新恢复资产和同口径复测。 +- 单分片回调收敛后第三轮100条/秒生成2999条,2999/2999受理,P50/P95/P99=`31/71/191ms`,三运营商`998/996/1005`、六通道`480至518`。注入期命令/结果Stream已排空但一度有2227条`submit_queued`,最终BullMQ wait/active/failed/delayed均0,消息2910 delivered、34 failed、55 submitted;冻结/释放各2999、charged 2987、refunded 21,SmsBillingRecord charged 2966/refunded 21,余额恒等式准确,数据库0等待锁/0 idle in transaction。该证据把剩余瓶颈定位到BullMQ发送Worker的逐消息数据库热路径,仍未达到全链100条/秒,因此没有升200档。 +- 热路径确认每条提交结果和回执都会并发执行整批`SmsMessageRecord GROUP BY status`并更新任务。最小修复对同进程同批次刷新做单飞合并,并在运行中查询后出现新状态时保留一次尾随刷新,不缓存业务结果。新增`TC-CMPP-500-P4-010`;SendChain专项122项与API正式TypeScript编译通过,部署和同口径100档复测待补记。