From 14f993c1f8e4ef5075a16c64ee8e46987054e21e Mon Sep 17 00:00:00 2001 From: hectorzhao Date: Thu, 20 Aug 2026 15:29:45 +0800 Subject: [PATCH] perf(cmpp): reduce inbound database round trips --- .../risk-review/risk-review.service.spec.ts | 46 +++++++++++++++++++ api/src/risk-review/risk-review.service.ts | 33 +++++++++++-- api/src/send-chain/send-chain.service.spec.ts | 32 ++++++++++++- api/src/send-chain/send-chain.service.ts | 7 +-- .../send-chain/send-gateway-submit.service.ts | 14 +++++- .../send-chain/send-inbound-entry.service.ts | 16 +++---- api/src/send-chain/send-submission.service.ts | 7 +-- docs/codebase-modularization-roadmap.md | 1 + docs/contracts/send-chain-r9-submission.json | 6 +-- .../first-version-development-requirements.md | 1 + docs/system-functional-test-cases.md | 5 ++ docs/testing-progress.md | 9 ++++ 12 files changed, 153 insertions(+), 24 deletions(-) diff --git a/api/src/risk-review/risk-review.service.spec.ts b/api/src/risk-review/risk-review.service.spec.ts index 45b6579..4efc346 100644 --- a/api/src/risk-review/risk-review.service.spec.ts +++ b/api/src/risk-review/risk-review.service.spec.ts @@ -3,6 +3,7 @@ import { RiskReviewService } from './risk-review.service'; function createPrismaMock(overrides: Record = {}) { return { riskRule: { + count: jest.fn().mockResolvedValue(5), findFirst: jest.fn().mockResolvedValue({ id: 'default-rule' }), findUnique: jest.fn(), create: jest.fn(), @@ -61,6 +62,51 @@ function createPrismaMock(overrides: Record = {}) { } describe('RiskReviewService', () => { + it('coalesces concurrent default-rule checks and reuses the short completeness cache', async () => { + const prisma = createPrismaMock(); + let releaseCount: ((count: number) => void) | undefined; + prisma.riskRule.count.mockReturnValue(new Promise((resolve) => { releaseCount = resolve; })); + const service = new RiskReviewService(prisma as never); + + const first = service.ensureDefaultRules(); + const second = service.ensureDefaultRules(); + releaseCount?.(5); + await Promise.all([first, second]); + await service.ensureDefaultRules(); + + expect(prisma.riskRule.count).toHaveBeenCalledTimes(1); + expect(prisma.riskRule.findFirst).not.toHaveBeenCalled(); + }); + + it('clears a failed default-rule check so the next request can retry', async () => { + const prisma = createPrismaMock(); + prisma.riskRule.count + .mockRejectedValueOnce(new Error('database unavailable')) + .mockResolvedValueOnce(5); + const service = new RiskReviewService(prisma as never); + + await expect(service.ensureDefaultRules()).rejects.toThrow('database unavailable'); + await expect(service.ensureDefaultRules()).resolves.toBeUndefined(); + + expect(prisma.riskRule.count).toHaveBeenCalledTimes(2); + }); + + it('falls back to per-rule recovery when the completeness count finds a missing default', async () => { + const prisma = createPrismaMock(); + prisma.riskRule.count.mockResolvedValue(4); + prisma.riskRule.findFirst.mockImplementation(({ where }: { where: { code: string } }) => ( + Promise.resolve(where.code === 'PHONE_FREQUENCY_5M' ? null : { id: `rule-${where.code}` }) + )); + const service = new RiskReviewService(prisma as never); + service.createRule = jest.fn().mockResolvedValue({ id: 'restored-rule' }) as never; + + await service.ensureDefaultRules(); + + expect(prisma.riskRule.findFirst).toHaveBeenCalledTimes(5); + expect(service.createRule).toHaveBeenCalledTimes(1); + expect(service.createRule).toHaveBeenCalledWith(expect.objectContaining({ code: 'PHONE_FREQUENCY_5M' })); + }); + it('keeps phone-frequency periods fixed and rejects manual-review actions', () => { const service = new RiskReviewService(createPrismaMock() as never); diff --git a/api/src/risk-review/risk-review.service.ts b/api/src/risk-review/risk-review.service.ts index abf28bf..b6b015b 100644 --- a/api/src/risk-review/risk-review.service.ts +++ b/api/src/risk-review/risk-review.service.ts @@ -114,6 +114,8 @@ const RULE_DEFINITIONS = new Map(DEFAULT_RULES.map((rule) => [rule.code, rule])) @Injectable() export class RiskReviewService { + private defaultRulesCheck?: { expiresAt: number; promise: Promise }; + constructor(private readonly prisma: PrismaService) {} async listRules(applicationId?: string) { @@ -491,14 +493,39 @@ export class RiskReviewService { } async ensureDefaultRules() { + const now = Date.now(); + if (this.defaultRulesCheck && this.defaultRulesCheck.expiresAt > now) { + return this.defaultRulesCheck.promise; + } + + // Submit高并发时,任务风控和号码频控都会确认默认规则。短TTL只缓存“规则是否齐全”, + // 实际生效规则仍逐次查询;并发单飞避免每条短信重复执行5次存在性SQL,同时允许删除后自动恢复。 + const promise = this.ensureDefaultRulesFromDatabase(); + this.defaultRulesCheck = { expiresAt: now + 30_000, promise }; + try { + await promise; + } catch (error) { + if (this.defaultRulesCheck?.promise === promise) this.defaultRulesCheck = undefined; + throw error; + } + } + + private async ensureDefaultRulesFromDatabase() { + const existingCount = await this.prisma.riskRule.count({ + where: { + applicationId: null, + code: { in: DEFAULT_RULES.map((rule) => rule.code) }, + status: { not: 'deleted' }, + }, + }); + if (existingCount === DEFAULT_RULES.length) return; + for (const rule of DEFAULT_RULES) { const exists = await this.prisma.riskRule.findFirst({ where: { applicationId: null, code: rule.code, status: { not: 'deleted' } }, select: { id: true }, }); - if (!exists) { - await this.createRule(rule); - } + if (!exists) await this.createRule(rule); } } diff --git a/api/src/send-chain/send-chain.service.spec.ts b/api/src/send-chain/send-chain.service.spec.ts index 7040514..ce7756a 100644 --- a/api/src/send-chain/send-chain.service.spec.ts +++ b/api/src/send-chain/send-chain.service.spec.ts @@ -397,6 +397,7 @@ function createService( ); service['postGatewayControl'] = jest.fn().mockResolvedValue({ delivered: true }); service['publishGatewaySubmitCommand'] = jest.fn().mockResolvedValue(undefined); + service['getSendQueue'] = jest.fn().mockReturnValue({ add: jest.fn().mockResolvedValue(undefined) }); return { service, prisma, billing, riskReview, phoneFrequency }; } @@ -1627,6 +1628,7 @@ describe('SendChainService', () => { })).resolves.toEqual(expect.objectContaining({ accepted: true, phoneCount: 2 })); expect(prisma.smsMessageRecord.create).toHaveBeenCalledTimes(2); + expect(prisma.smsApplication.findFirst).toHaveBeenCalledTimes(1); expect(prisma.smsMessageRecord.update).toHaveBeenCalledWith({ where: { id: 'record-2' }, data: expect.objectContaining({ @@ -1766,7 +1768,10 @@ describe('SendChainService', () => { templateId: 'tpl-code', variables: { code: '715021' }, })); - expect(service.enqueueBatchTask).toHaveBeenCalledWith('task-1'); + expect(service.enqueueBatchTask).toHaveBeenCalledWith('task-1', { + messageRecordId: 'record-1', + queuePriority: 'normal', + }); expect(prisma.smsReceiptRecord.create).not.toHaveBeenCalled(); }); @@ -1814,7 +1819,10 @@ describe('SendChainService', () => { where: { id: 'record-1' }, data: { status: 'queued', signatureId: 'sig-1' }, }); - expect(service.enqueueBatchTask).toHaveBeenCalledWith('task-1'); + expect(service.enqueueBatchTask).toHaveBeenCalledWith('task-1', { + messageRecordId: 'record-1', + queuePriority: 'normal', + }); expect(prisma.smsReceiptRecord.create).not.toHaveBeenCalled(); }); @@ -2090,6 +2098,26 @@ describe('SendChainService', () => { }); }); + 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); + service['getSendQueue'] = jest.fn().mockReturnValue({ add }); + + await expect(service.enqueueBatchTask('task-1', { + messageRecordId: 'record-1', + queuePriority: 'priority', + })).resolves.toEqual({ taskId: 'task-1', enqueued: 1 }); + + expect(prisma.smsBatchTask.findUnique).not.toHaveBeenCalled(); + expect(prisma.smsMessageRecord.findMany).not.toHaveBeenCalled(); + expect(add).toHaveBeenCalledWith('send-message', { messageRecordId: 'record-1' }, { + jobId: 'record-1', + attempts: 3, + priority: 1, + }); + expect(prisma.smsBatchTask.update).toHaveBeenCalledWith({ where: { id: 'task-1' }, data: { status: 'queued' } }); + }); + it('reuses persisted carrier and province without querying routing dictionaries again', async () => { const { service, prisma } = createService(); service['identifyCarrier'] = jest.fn(); diff --git a/api/src/send-chain/send-chain.service.ts b/api/src/send-chain/send-chain.service.ts index 05b8d52..ddf3000 100644 --- a/api/src/send-chain/send-chain.service.ts +++ b/api/src/send-chain/send-chain.service.ts @@ -352,8 +352,8 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { return this.submission.confirmImport(data); } - async enqueueBatchTask(taskId: string) { - return this.submission.enqueueBatchTask(taskId); + async enqueueBatchTask(taskId: string, preparedMessage?: { messageRecordId: string; queuePriority?: string | null }) { + return this.submission.enqueueBatchTask(taskId, preparedMessage); } async cancelBatchTask(taskId: string, tenantId?: string, sourceType = 'client') { @@ -599,10 +599,11 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { data: GatewayInboundSubmitDto & { phoneNumber: string }, messageId: string, submitGroupMessageId: string, + application: NonNullable>>, synchronousRejection?: { code: string; reason: string }, receiptRejection?: { code: string; reason: string }, ) { - return this.submission.submitInboundSingleMessage(data, messageId, submitGroupMessageId, synchronousRejection, receiptRejection); + return this.submission.submitInboundSingleMessage(data, messageId, submitGroupMessageId, application, synchronousRejection, receiptRejection); } /** diff --git a/api/src/send-chain/send-gateway-submit.service.ts b/api/src/send-chain/send-gateway-submit.service.ts index 2ff97c4..15acdf4 100644 --- a/api/src/send-chain/send-gateway-submit.service.ts +++ b/api/src/send-chain/send-gateway-submit.service.ts @@ -67,7 +67,19 @@ export class SendGatewaySubmitService { } -async enqueueBatchTask(taskId: string) { +async enqueueBatchTask(taskId: string, preparedMessage?: { messageRecordId: string; queuePriority?: string | null }) { + if (preparedMessage) { + // CMPP内部任务在当前请求内刚完成持久化且不暴露取消入口,可安全复用已知ID; + // 普通批量任务仍走下方查询路径,以保留取消检查和多消息枚举语义。 + const queuePriority = normalizeQueuePriority(preparedMessage.queuePriority); + await this.facade.getSendQueue().add('send-message', { messageRecordId: preparedMessage.messageRecordId }, { + jobId: preparedMessage.messageRecordId, + attempts: 3, + priority: BULLMQ_PRIORITY[queuePriority], + }); + await this.prisma.smsBatchTask.update({ where: { id: taskId }, data: { status: 'queued' } }); + return { taskId, enqueued: 1 }; + } const task = await this.prisma.smsBatchTask.findUnique({ where: { id: taskId } }); if (!task) { throw new NotFoundException('SMS batch task not found'); diff --git a/api/src/send-chain/send-inbound-entry.service.ts b/api/src/send-chain/send-inbound-entry.service.ts index a8c8a33..a587a81 100644 --- a/api/src/send-chain/send-inbound-entry.service.ts +++ b/api/src/send-chain/send-inbound-entry.service.ts @@ -356,7 +356,7 @@ async submitCompleteInboundMessage( ...data, phoneNumber: submission.phoneNumber, phoneNumbers: undefined, - }, submission.messageId, submitGroupMessageId, submission.receiptRejection ? undefined : dailyLimitRejection, submission.receiptRejection)))); + }, submission.messageId, submitGroupMessageId, application, submission.receiptRejection ? undefined : dailyLimitRejection, submission.receiptRejection)))); } const first = results[0]; return { @@ -550,16 +550,11 @@ async submitInboundSingleMessage( data: GatewayInboundSubmitDto & { phoneNumber: string }, messageId: string, submitGroupMessageId: string, + application: NonNullable>>, synchronousRejection?: { code: string; reason: string }, receiptRejection?: { code: string; reason: string }, ) { - const application = await this.measureInboundStage( - 'application_lookup', - () => this.facade.findInboundApplication(data.account), - ); - if (!application) { - throw new BadRequestException('CMPP account is invalid'); - } + // 入口已按账号取得并校验同一个应用快照;复用它可避免每个目标号码再次查询应用、企业和IP白名单。 if (data.remoteIp && !isIpAllowed(data.remoteIp, application.ipAllowlist.map((item) => item.ipCidr))) { throw new BadRequestException('CMPP source IP is not in application allowlist'); } @@ -715,7 +710,10 @@ async submitInboundSingleMessage( where: { id: task.id }, data: { status: 'ready', riskTaskId: risk.task?.id, auditStatus: 'approved' }, }); - await this.facade.enqueueBatchTask(task.id); + await this.facade.enqueueBatchTask(task.id, { + messageRecordId: message.id, + queuePriority, + }); }); }; if (receiptRejection) { diff --git a/api/src/send-chain/send-submission.service.ts b/api/src/send-chain/send-submission.service.ts index e5e614e..4e63c5e 100644 --- a/api/src/send-chain/send-submission.service.ts +++ b/api/src/send-chain/send-submission.service.ts @@ -164,10 +164,11 @@ async submitInboundSingleMessage( data: GatewayInboundSubmitDto & { phoneNumber: string }, messageId: string, submitGroupMessageId: string, + application: NonNullable>>, synchronousRejection?: { code: string; reason: string }, receiptRejection?: { code: string; reason: string }, ) { - return this.inboundEntry.submitInboundSingleMessage(data, messageId, submitGroupMessageId, synchronousRejection, receiptRejection); + return this.inboundEntry.submitInboundSingleMessage(data, messageId, submitGroupMessageId, application, synchronousRejection, receiptRejection); } async evaluateRiskWithPhoneFrequency(input: { @@ -214,8 +215,8 @@ async runScheduledDispatchScan() { return this.scheduledDispatch.runScheduledDispatchScan(); } -async enqueueBatchTask(taskId: string) { - return this.gatewaySubmit.enqueueBatchTask(taskId); +async enqueueBatchTask(taskId: string, preparedMessage?: { messageRecordId: string; queuePriority?: string | null }) { + return this.gatewaySubmit.enqueueBatchTask(taskId, preparedMessage); } startWorker() { diff --git a/docs/codebase-modularization-roadmap.md b/docs/codebase-modularization-roadmap.md index 61f5dac..ee4f7be 100644 --- a/docs/codebase-modularization-roadmap.md +++ b/docs/codebase-modularization-roadmap.md @@ -1154,5 +1154,6 @@ global,不能错误归入client。`client-signature-*`、发送页、企业认 - V2提交工作池只归`gateway/internal/submitworker/`治理:Redis领取、全局槽位、在途消息ID和逐条ACK不能渗入上游连接池;`gateway/internal/upstream/`继续只负责通道连接、窗口和供应商协议往返。这样Worker吞吐调优不会改写CMPP连接状态机,连接池也不能自行确认Redis消息。 - V3客户入站窗口只归`gateway/internal/inbound/`与项目内受控的`third_party/gocmpp`服务循环治理:API认证只返回应用窗口,inbound会话负责收紧窗口,协议服务循环负责受限派发和断线等待;不得把客户入站槽位与`submitworker`供应商槽位或`upstream`供应商窗口合并成同一并发计数。 - V4供应商结果异步边界只归`gateway/internal/resultoutbox/`治理:`upstream`在每个真实分片SubmitResp后只调用持久化接口,`submitworker`只负责聚合结果入Outbox与命令ACK的原子边界,Outbox回调Worker独立控制API并发和PEL恢复。API的`SmsSubmitRecord.resultEventId`是跨重启持久幂等事实;不得把HTTP回调重新放回供应商连接池或Submit工作槽,也不得让Outbox承担业务计费、补发或状态机判断。 +- V5 API数据库往返优化仍归`send-inbound-entry`、`send-gateway-submit`和`risk-review`现有边界:入口应用快照沿稳定Facade显式传递,单条快速入队只接受已持久化ID和优先级,通用批量入队保留原查询与取消校验;默认规则并发单飞属于`RiskReviewService`内部完整性保障,不得在`SendChainService`新增第二套规则缓存,也不得把余额、频控状态或实际规则决策缓存进进程内存。 - `api/src/infrastructure-monitoring/`是运营端监控聚合与固定阈值应用边界:只消费代码白名单 PromQL,并通过版本化 PostgreSQL 单例、promtool 校验和原子规则热加载管理数值阈值;Exporter安装、端口隔离、固定规则模板和权限仍归`tools/monitoring/`治理。 - 活动告警已读也归该边界:Prometheus保留告警事实,Prisma仅持久化逐管理员、逐触发周期的阅读状态;全局布局只消费轻量未读汇总,不复制指纹、activeAt或用户隔离逻辑。 diff --git a/docs/contracts/send-chain-r9-submission.json b/docs/contracts/send-chain-r9-submission.json index 66e7d82..7a18486 100644 --- a/docs/contracts/send-chain-r9-submission.json +++ b/docs/contracts/send-chain-r9-submission.json @@ -31,7 +31,7 @@ }, { "name": "enqueueBatchTask", - "bodySha256": "909972eb0accfb6e31b9fd6118750005c4a6671577ef5b5c9222a341dac59800", + "bodySha256": "081fab5913288d7325a8c4d15cea5e7d37af1c6614154f2963eb635bbcc0f40e", "file": "send-gateway-submit.service.ts" }, { @@ -76,7 +76,7 @@ }, { "name": "submitCompleteInboundMessage", - "bodySha256": "10e2140580e4bae2055853981de65f4adf82c228b647dbb56fce7c39469acaf3", + "bodySha256": "14bd1ffd41b3e38d9ceeabd7f0e19d6bb356b4ff2b281acc31a8477c8bc18c13", "file": "send-inbound-entry.service.ts" }, { @@ -91,7 +91,7 @@ }, { "name": "submitInboundSingleMessage", - "bodySha256": "6fa4c577848eb06f635a8f186d1e1510b1213f0ba6d439d032ea08c593725d14", + "bodySha256": "6246c081fa2433ef5b3ffc00a2da7cfe32da9d02c310ef757b904800ffbd6732", "file": "send-inbound-entry.service.ts" }, { diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index 6a1b68d..dfc9b02 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -2102,6 +2102,7 @@ - V2有界工作池必须暴露配置槽位数和当前在途槽位数,使用固定`state=configured|in_flight`标签;该指标用于区分Worker容量耗尽与供应商窗口/限速等待,不得增加通道或消息标签。 - V3入站并发必须只并发Submit业务处理,连接认证保持串行先完成,心跳和Deliver ACK不得被长耗时Submit阻塞;每个SubmitResp继续使用原请求Sequence_Id关联,允许按实际完成顺序返回。同一连接关闭时必须先等待已接受的在途处理收尾,再清理会话和回执映射,避免迟到处理重新注册已断开的连接。`cmpp_gateway_inbound_submit_slots{state=configured|in_flight}`只暴露全部在线连接的聚合窗口与在途数量,不得增加账号、应用、连接或消息标签。 - V4结果Outbox必须暴露独立回调Worker的configured/in_flight槽位和Outbox pending/lag,仍只使用固定状态标签。`api_callback`从V4起只在Outbox回调Worker计时,不再混入供应商Submit工作槽;压测结束必须同时核对命令Stream和结果Outbox均`pending=0/lag=0`,并证明API回调故障时供应商槽继续释放、结果不丢失且恢复后只处理一次。 +- V5 API入站数据库往返治理不得改变SubmitResp、余额冻结、号码频控、模板/签名、路由、队列或失败回执语义。同一Submit入口已取得的应用、企业和IP白名单快照必须复用于其全部目标号码,不得逐号码重复查询;任务风控与号码频控共用的默认规则完整性检查必须使用短TTL并发单飞,正常完整状态最多执行一次聚合数据库检查,实际生效规则仍逐次读取;刚持久化的单条CMPP内部任务和消息可凭已知ID、优先级直接入队,不得为入队重新查询同一任务和消息。缓存只覆盖默认规则“是否齐全”,检查失败必须立即失效,规则缺失须在短TTL到期后自动恢复;不得缓存余额、频控计数、应用启停或实际生效规则。 - 系统监控标题说明必须明确标注数据来自 Prometheus;“服务关键指标”位于趋势/核心服务区域之后、活动告警之前,并提供统一的“告警阈值设置”入口。安全检测页不重复渲染大号标题,说明文字必须明确标注使用 Fail2ban。 - 告警阈值仅开放固定指标的警告/严重数值,不允许前端提交 PromQL、标签、文件路径或持续时间;必须满足警告值小于严重值。配置以 PostgreSQL 保存版本、期望值、生效值和应用状态,经 `promtool check rules` 校验、同目录原子替换和 Prometheus 热加载成功后才标记生效,失败保留上一生效规则并展示原因。 - 右上角预警中心增加“系统监控告警”,通过独立轻量接口统计 Prometheus 当前 firing/pending 告警及严重数,跳转系统监控活动告警区;任一预警域失败不得清空其他域。 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index c0c5aba..4b2cc5f 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -4690,6 +4690,11 @@ npm run verify:phase8 | TC-CMPP-PERF-V4-006 | Outbox有界指标 | 制造回调在途和积压后抓取Gateway metrics | 回调Worker configured/in_flight与真实槽位一致,Outbox pending/lag与Redis consumer group一致,指标不含手机号、消息、submit、企业、应用或通道标签 | | TC-CMPP-PERF-V4-007 | 双Stream发布后排空 | 完成一档隔离压测并等待异步处理结束 | `gateway.submit.commands`和`gateway.submit.results`均为`pending=0/lag=0`,结果Outbox成功事件已删除,无Gateway Submit死信;数据库业务数与客户端完全一致 | | TC-CMPP-PERF-V4-008 | V4阶梯容量复测 | 在V3相同隔离供应商模拟器与真实API/PostgreSQL/Redis/Gateway中依次执行10、20、30、40、50条/秒各60秒 | 对比V3的SubmitResp分位、供应商RTT、命令Stream、结果Outbox、API阶段和资源;任一档出现拒绝、连接错误、P95超过5秒、双Stream持续增长或服务异常立即停止升档 | +| TC-CMPP-PERF-V5-001 | 入站应用快照复用 | 分别提交单号码、多号码和完整长短信,并统计`SmsApplication`查询 | 每次Submit入口只查询一次账号对应应用、企业和IP白名单;全部目标号码复用该已校验快照,应用/IP/模板/风控业务结果不变 | +| TC-CMPP-PERF-V5-002 | 默认规则完整性检查单飞 | 同一API实例并发触发任务风控和号码频控,再在30秒内重复提交 | 并发调用共用一个检查Promise;默认5条规则完整时只执行一次聚合count,不再逐消息产生两轮各5次存在性查询;实际生效规则和频控状态仍逐消息读取 | +| TC-CMPP-PERF-V5-003 | 默认规则缓存失效与恢复 | 让聚合检查失败后重试;另在缓存期后模拟缺失一条默认规则 | 检查失败立即清除缓存并允许下次重试;短TTL内只缓存完整性,TTL后发现缺失规则并按既有创建校验恢复,不缓存应用规则内容或业务判断 | +| TC-CMPP-PERF-V5-004 | 已持久化单消息快速入队 | 创建CMPP单号码内部任务和消息,记录已知taskId、messageRecordId、queuePriority后进入队列 | BullMQ jobId仍为messageRecordId、attempts和优先级不变;不再回查刚创建的任务和消息,任务最终更新为queued;普通批量任务原通用入队路径保持兼容 | +| TC-CMPP-PERF-V5-005 | V5隔离环境阶梯复测 | 发布到虚拟机测试环境后,以V4相同模拟器、连接数、8槽Outbox和10/20/30/40/50条每秒阶梯执行 | 对比V4的SubmitResp分位及`application_lookup/risk_frequency/queue_publish/complete_submit`阶段;数据库业务数、冻结/计费、双Stream排空和回执幂等保持一致,遇拒绝、连接错误、P95超过5秒或持续积压立即停止 | | TC-GLOBAL-ALERT-001 | 铃铛分域预警菜单 | 准备签名清退未读消息和安全待处置告警后点击右上角铃铛 | 弹层分开显示“签名清退预警”和“安全检测与封禁”,分别展示真实数量和摘要,角标等于两项之和 | | TC-GLOBAL-ALERT-002 | 预警菜单跳转 | 分别点击铃铛中的两个菜单项 | 签名项跳转`/admin/signature-retirement`,安全项跳转`/admin/security-detection`,弹层关闭且对应页面读取真实后端数据 | | TC-GLOBAL-ALERT-003 | 域间故障隔离与轻量轮询 | 分别让一个汇总接口失败并观察30秒轮询请求 | 失败域显示0且另一域数据保留;安全预警使用专用汇总接口,不调用完整overview、规则、代理状态或告警大列表 | diff --git a/docs/testing-progress.md b/docs/testing-progress.md index 4144389..655d47e 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -3761,3 +3761,12 @@ git diff --check - 50条/秒触发停止线:因客户端背压只生成2790条而非约2999条,2787条成功、3条等待API满10秒后返回result=9,P50/P95/P99=`6765/7236/7789ms`,throttled ticks=3184。结束后命令Stream一度`pending=64/lag=489`、结果Outbox`pending=8/lag=1690`;命令约1分钟内、结果随后约20秒内归零。日志显示API饱和时协议遥测和连接状态回调超时,并暴露既有`CmppDownstreamConnection`创建/更新的P2002/P2025竞态;未出现`resultEventId`唯一键冲突或Outbox事件失败。停止后未继续上探。 - 整个窗口客户侧共9984条业务提交,真实PostgreSQL按`queuedAt`精确新增9984条;窗口内7922条实际供应商提交记录全部具有非空且互不重复的`resultEventId`,重复事件组0、Gateway Submit死信0。提交结果为accepted 7392、rejected 147、timeout 383;客户有限回执收集窗口和重连积压回执不替代数据库/Stream对账。最终确认40条/秒为当前测试环境安全档,50条/秒瓶颈转为API/数据库同步入站及遥测争用,后续应继续V1的重复查询/零散写入合并和P1分阶段指标分析。 - 本轮只连接隔离供应商模拟器`100.91.249.119:17900`,没有发送、补发或重投真实短信,没有修改真实通道账号、密码、启停状态、企业余额或客户连接。压测原始目录沿用既有`lg-v2-*`名称,但被测运行标识与结论均为V4;完整V4报告保存在短信平台测试项目。受保护的`api/tsconfig.build.tsbuildinfo`、`tsconfig.tsbuildinfo`、`outputs/`、`pnpm-lock.yaml`和空文件`=`继续排除提交、不删除、不归因。 + +# 2026-08-20 CMPP压测优化V5:API入站数据库往返治理(本地完成、待授权发布复测) + +- 接管复核确认本地`HEAD=67b760a5992d7ae10d7889fe6765e895d238ab3d`、`origin/main=c4f36fc50d7906dfb2f97c881e9ea43c6a64c370`,本地领先5个提交;测试环境只读标识仍为`485af688d21ee0bdb98281f4c20d151d69f7889f+workspace.v4.a63885439ed5`,预生产只读标识仍为`433b2ee56f6016ad8afff1bac73f510b8fd53083+gateway-v2.cd7bb8d05e7b`。本轮没有发布、回退、提交或推送。 +- V4高负载证据显示两次`application_lookup`合计约191.44ms、`risk_frequency`约677.38ms、`queue_publish`约242.43ms。代码复核确认每个目标号码会重复查询入口已取得的应用/企业/IP白名单;任务风控与号码频控各自调用默认规则检查,旧实现每次按5条规则逐条查询,单条短信可产生两轮共约10次存在性SQL;通用入队又回查刚创建的任务和消息。 +- V5复用Submit入口已经校验的应用快照并显式传入每个目标号码,删除第二次`application_lookup`;默认规则完整性改为30秒短TTL和并发单飞,正常完整状态使用一次聚合count,失败立即清除缓存,实际生效规则、应用规则和号码频控状态仍逐消息读取;缺失时继续走既有逐规则确认和创建逻辑。刚持久化的单条CMPP内部消息使用已知taskId、messageRecordId和queuePriority直接添加BullMQ任务,不再回查任务和消息;普通批量任务继续使用原通用入队路径。 +- 只读测试环境PostgreSQL `EXPLAIN (ANALYZE, BUFFERS)`确认应用查询使用`SmsApplication_cmppAccount_key`,执行约0.046ms;默认规则聚合检查使用`RiskRule_applicationId_code_key`并命中5行,执行约0.047ms。现有索引已经匹配查询,因此本轮不增加猜测性索引或migration;测试环境未安装`pg_stat_statements`,没有为诊断修改数据库扩展或配置。 +- 本地API TypeScript正式编译通过;API全量42套487项全部通过,其中`RiskReviewService`与`SendChainService`专项覆盖默认规则并发单飞/短缓存/失败重试/缺失恢复、单Submit只查一次应用、已知消息直接入队且不回查任务/消息,以及既有多号码、模板、频控、计费、补发和回执回归。5份Gateway队列契约、R0、R6、R7、R10和`git diff --check`通过;R9门禁在本轮已更新的`enqueueBatchTask`哈希通过后,被HEAD既有且本轮未修改的`dispatchDueScheduledTasks`哈希漂移阻断(契约期望`8fd154...`、当前方法`37db96...`),未为通过本需求门禁而改写或归因该并行历史。Jest使用`--forceExit`收尾仓库既有异步句柄;首次从仓库根目录误启动时扫描受保护`outputs/`并使用错误转换配置,未修改或删除其中任何文件,随后在`api`目录按项目配置重跑通过。 +- 当前尚未部署虚拟机测试环境或执行V5阶梯压测;`TC-CMPP-PERF-V5-005`需在取得明确发布授权、建立PostgreSQL/运行源码/环境恢复资产后执行。没有发送、补发或重投真实短信,没有修改通道账号、密码、启停状态、企业余额或客户连接;`api/tsconfig.build.tsbuildinfo`、`tsconfig.tsbuildinfo`、`outputs/`、`pnpm-lock.yaml`和空文件`=`继续保护、不归因。