diff --git a/api/src/billing/billing.service.spec.ts b/api/src/billing/billing.service.spec.ts index e082d79..2778902 100644 --- a/api/src/billing/billing.service.spec.ts +++ b/api/src/billing/billing.service.spec.ts @@ -11,7 +11,7 @@ function createPrismaMock() { findMany: jest.fn().mockResolvedValue([]), }, tenantAccount: { - findMany: jest.fn(), + findMany: jest.fn().mockResolvedValue([]), findUnique: jest.fn().mockImplementation(() => Promise.resolve({ ...accountState })), findUniqueOrThrow: jest.fn().mockImplementation(() => Promise.resolve({ ...accountState })), create: jest.fn(), @@ -30,7 +30,7 @@ function createPrismaMock() { }), }, accountTransaction: { - findMany: jest.fn(), + findMany: jest.fn().mockResolvedValue([]), findFirst: jest.fn(), findUnique: jest.fn().mockResolvedValue(null), findUniqueOrThrow: jest.fn().mockImplementation(({ where }) => Promise.resolve({ id: 'tx-charged', idempotencyKey: where.idempotencyKey })), @@ -55,6 +55,7 @@ function createPrismaMock() { create: jest.fn().mockResolvedValue({ id: 'operation-1' }), }, $executeRaw: jest.fn(), + $queryRaw: jest.fn(), }; return Object.assign(prisma, { $transaction: jest.fn((callback: (client: typeof prisma) => unknown) => callback(prisma)), @@ -319,11 +320,11 @@ describe('BillingService', () => { it('settles a frozen SMS charge with one account lock and an idempotent ledger pair', async () => { const prisma = createPrismaMock(); const rows = new Map>(); - prisma.accountTransaction.findUnique.mockImplementation(({ where }) => Promise.resolve(rows.get(where.idempotencyKey) ?? null)); - prisma.accountTransaction.findUniqueOrThrow.mockImplementation(({ where }) => Promise.resolve(rows.get(where.idempotencyKey))); - prisma.accountTransaction.createMany.mockImplementation(({ data }) => { - data.forEach((row: Record) => rows.set(String(row.idempotencyKey), { id: `tx-${row.transactionType}`, ...row })); - return Promise.resolve({ count: data.length }); + prisma.accountTransaction.findMany.mockImplementation(() => Promise.resolve([...rows.values()])); + prisma.$queryRaw.mockImplementation(() => { + const charged = { id: 'tx-charged', idempotencyKey: 'sms-charge:msg-paid-1', transactionType: 'charged', amountCents: -325 }; + rows.set(String(charged.idempotencyKey), charged); + return Promise.resolve([charged]); }); const service = new BillingService(prisma as never); const input = { tenantId: 'tenant-1', amountCents: 325, messageId: 'msg-paid-1', taskId: 'task-paid-1' }; @@ -333,7 +334,8 @@ describe('BillingService', () => { expect(first).toEqual(expect.objectContaining({ transactionType: 'charged', amountCents: -325 })); expect(replay).toEqual(first); - expect(prisma.accountTransaction.createMany).toHaveBeenCalledTimes(1); + expect(prisma.$queryRaw).toHaveBeenCalledTimes(1); + expect(prisma.$executeRaw).not.toHaveBeenCalled(); expect(prisma.tenantAccount.update).not.toHaveBeenCalled(); }); diff --git a/api/src/billing/billing.service.ts b/api/src/billing/billing.service.ts index 732c928..5a7d738 100644 --- a/api/src/billing/billing.service.ts +++ b/api/src/billing/billing.service.ts @@ -441,44 +441,49 @@ export class BillingService { if (amountCents <= 0) return null; const releaseKey = `sms-charge-release:${data.messageId}`; const chargeKey = `sms-charge:${data.messageId}`; - return this.prisma.$transaction(async (tx) => { - await tx.$executeRaw`SELECT pg_advisory_xact_lock(hashtextextended(${'tenant-account:' + data.tenantId}, 0))`; - const [release, charge] = await Promise.all([ - tx.accountTransaction.findUnique({ where: { idempotencyKey: releaseKey } }), - tx.accountTransaction.findUnique({ where: { idempotencyKey: chargeKey } }), - ]); - if (charge) return charge; - const account = await tx.tenantAccount.findUniqueOrThrow({ where: { tenantId: data.tenantId } }); - const balance = moneyToNumber(account.balanceCents); - if (release) { - const updated = await tx.tenantAccount.update({ - where: { tenantId: data.tenantId }, - data: { balanceCents: { decrement: amountCents } }, - }); - return tx.accountTransaction.create({ - data: { - tenantId: data.tenantId, transactionType: 'charged', idempotencyKey: chargeKey, - amountCents: -amountCents, balanceAfter: updated.balanceCents, - relatedType: 'sms_message_record', relatedId: data.messageId, remark: '提交成功扣费', - }, - }); - } - await tx.accountTransaction.createMany({ - data: [ - { - tenantId: data.tenantId, transactionType: 'released', idempotencyKey: releaseKey, - amountCents, balanceAfter: balance + amountCents, - relatedType: 'sms_batch_task', relatedId: data.taskId, remark: data.remark, - }, - { - tenantId: data.tenantId, transactionType: 'charged', idempotencyKey: chargeKey, - amountCents: -amountCents, balanceAfter: balance, - relatedType: 'sms_message_record', relatedId: data.messageId, remark: '提交成功扣费', - }, - ], + const existing = await this.prisma.accountTransaction.findMany({ + where: { idempotencyKey: { in: [releaseKey, chargeKey] } }, + }); + const charge = existing.find((row) => row.idempotencyKey === chargeKey); + if (charge) return charge; + const release = existing.find((row) => row.idempotencyKey === releaseKey); + if (release) { + // Recover the legacy two-transaction boundary: a crash may have committed + // release before charge, so this path must perform the missing balance debit. + return this.applyAccountDelta({ + tenantId: data.tenantId, transactionType: 'charged', idempotencyKey: chargeKey, + amountCents: -amountCents, relatedType: 'sms_message_record', relatedId: data.messageId, + remark: '提交成功扣费(恢复既有已释放冻结)', }); - return tx.accountTransaction.findUniqueOrThrow({ where: { idempotencyKey: chargeKey } }); - }, { isolationLevel: Prisma.TransactionIsolationLevel.ReadCommitted }); + } + const rows = await this.prisma.$queryRaw>(Prisma.sql` + WITH account AS ( + SELECT "balanceCents" FROM "TenantAccount" WHERE "tenantId" = ${data.tenantId} + ), inserted AS ( + INSERT INTO "AccountTransaction" ( + id, "tenantId", "transactionType", "idempotencyKey", "amountCents", + "balanceAfter", "relatedType", "relatedId", remark, "createdAt" + ) + SELECT gen_random_uuid()::text, ${data.tenantId}, ledger."transactionType", ledger."idempotencyKey", + ledger."amountCents", account."balanceCents" + ledger."balanceDelta", + ledger."relatedType", ledger."relatedId", ledger.remark, (NOW() AT TIME ZONE 'UTC') + FROM account + CROSS JOIN (VALUES + ('released', ${releaseKey}, ${amountCents}::bigint, ${amountCents}::bigint, 'sms_batch_task', ${data.taskId}, ${data.remark ?? null}), + ('charged', ${chargeKey}, ${-amountCents}::bigint, 0::bigint, 'sms_message_record', ${data.messageId}, '提交成功扣费') + ) AS ledger("transactionType", "idempotencyKey", "amountCents", "balanceDelta", "relatedType", "relatedId", remark) + ON CONFLICT ("idempotencyKey") DO NOTHING + RETURNING id, "idempotencyKey" + ) + SELECT id, "idempotencyKey" FROM inserted WHERE "idempotencyKey" = ${chargeKey} + UNION ALL + SELECT id, "idempotencyKey" FROM "AccountTransaction" WHERE "idempotencyKey" = ${chargeKey} + LIMIT 1 + `); + if (rows[0]) return rows[0]; + // A concurrent identical callback can win ON CONFLICT while remaining + // invisible to this statement's snapshot; one read repairs that MVCC edge. + return this.prisma.accountTransaction.findUniqueOrThrow({ where: { idempotencyKey: chargeKey } }); } release(data: BillingActionDto) { diff --git a/docs/codebase-modularization-roadmap.md b/docs/codebase-modularization-roadmap.md index c634ff5..499f44a 100644 --- a/docs/codebase-modularization-roadmap.md +++ b/docs/codebase-modularization-roadmap.md @@ -1168,6 +1168,6 @@ global,不能错误归入client。`client-signature-*`、发送页、企业认 ## 2026-08-21 第四阶段模块边界 - `send-inbound-entry`负责编排正价批次短事务;账户固定锁序、总额覆盖、逐短信冻结幂等键留在PostgreSQL一致性边界,不迁入Redis。 -- `billing.service`集中维护账户锁和流水,`settleFrozenCharge`把释放/扣费合并为一个事务;`send-accounting`只编排消息计费状态。 +- `billing.service`集中维护余额锁和流水;`settleFrozenCharge`利用冻结已完成扣减、释放与扣费净变化为零的事实,以单个幂等SQL写成流水对且不争抢账户锁,只有历史半完成恢复才进入带锁余额补扣;`send-accounting`只编排消息计费状态。 - `send-chain.helpers`提供无副作用的优先级、主备和稳定weight选择;提交服务只提供消息稳定键与真实在线/报备候选,不保存进程内轮询状态。 - Gateway既有Submit池与结果Outbox继续独立有界;先以正价六通道实测定位容量,只有证据证明单条回调仍触发停止线时才扩展批量回调协议。 diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index a0b5b31..cc0224d 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -2133,6 +2133,6 @@ ## 完整处理500条/秒第四阶段:正价计费与供应商并行(2026-08-21) - 普通CMPP短短信批量路径必须支持正单价:按企业固定锁序一次校验批次总金额、一次更新余额,并为每短信写独立幂等冻结流水;任务、请求、消息与冻结流水同属一个短事务。余额加授信必须覆盖所需金额,不能只判断大于零。 -- 供应商接受后,释放冻结与正式扣费在一个账户锁事务内保留两条审计流水;零金额不创建AccountTransaction。拒绝、超时、重投和最终失败继续沿用逐短信幂等释放、扣费和退款。 +- 供应商接受后,释放冻结与正式扣费以一个原子SQL写入两条幂等审计流水;二者净余额变化为零,不得争抢企业账户锁或重复更新账户。只有恢复历史“已释放但未扣费”半完成状态时才取得账户锁补扣。零金额不创建AccountTransaction;拒绝、超时、重投和最终失败继续沿用逐短信幂等释放、扣费和退款。 - 同通道组内在线、已报备、同地域、同最低优先级且非备用的通道按weight和稳定消息键分流;备用、离线及更低优先级不抢占。只有显式配置为同优先级非备用的通道才参与主动双活。 - 压测使用隔离测试企业、真实PostgreSQL账务、单价`0.0325元/计费条`、三运营商混合号码和六个模拟供应商账号;逐档核对消息、冻结/释放/扣费/退款、SmsBillingRecord、通道分布、队列和数据库。临时单价与主动双活配置须先快照、测试后恢复。 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index 91bd650..6f25e24 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -4779,7 +4779,7 @@ npm run verify:phase8 | --- | --- | --- | | TC-CMPP-500-P4-001 | 同企业正价短短信批量入站并重放 | 一次账户余额更新;每短信唯一`:freeze`流水;任务/请求/消息金额正确;重放不重复冻结 | | TC-CMPP-500-P4-002 | 余额加授信小于批次或单条金额 | 不得透支;批次安全回退且只发送余额覆盖的短信 | -| TC-CMPP-500-P4-003 | 正价短信Submit接受及Outbox重放 | 一个账户锁事务生成released/charged;余额净值不重复变化;SmsBillingRecord唯一charged | +| TC-CMPP-500-P4-003 | 正价短信Submit接受及Outbox重放 | 一个原子SQL生成released/charged且不取得账户锁;余额净值不重复变化;SmsBillingRecord唯一charged;历史半完成状态仍可锁定恢复 | | TC-CMPP-500-P4-004 | 零价短信Submit接受 | 保留业务结果但无0金额AccountTransaction | | TC-CMPP-500-P4-005 | 两个同优先级非备用通道与一个备用通道 | 稳定weighted分流只覆盖两个主动通道;排除已尝试通道后可安全切换;备用不抢占 | | TC-CMPP-500-P4-006 | 三运营商、六通道、325金额单位阶梯压测 | 每档入口/Inbox/供应商/回执/账务对账一致,双Stream和数据库最终稳定排空;失败即停止升档 | diff --git a/docs/testing-progress.md b/docs/testing-progress.md index a020a3f..606cd88 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -3840,3 +3840,6 @@ git diff --check - 10个隔离应用仍为单价0、同企业余额1/授信0;移动主备各200条/秒,联通/电信主备各150条/秒,窗口32。组项目为主10/备用20,默认只能三主主动承载,六连接在线不等于六通道并行。 - 已实现正价Worker批次、覆盖本次金额的余额判断、单事务释放/扣费、零价免账户流水,以及同最低优先级非备用通道的稳定weighted分流;默认主备语义不变。API正式编译和Billing/纯策略/SendChain三套143项通过。 - 部署、恢复资产、测试充值、临时六通道主动双活、smoke及阶梯结果待补记;所有临时配置只允许在测试环境先快照后变更并在测试后恢复。 +- 首轮正价三运营商100条/秒生成2999条,2999/2999 SubmitResp成功、P50/P95/P99=`32/68/130ms`,Inbox与唯一MessageId均2999;移动/联通/电信分别998/996/1005,六通道最终各承载488至514条,证明主动双活与跨运营商分流有效。注入结束时仅388条完成扣费,命令Stream`128/2131`、结果Outbox`32/192`、等待锁/idle in transaction各5,按停止线不升200档。 +- 该轮最终双Stream归零,冻结/释放各2999、正式扣费2987、退款30、12条供应商最终拒绝只释放不扣费;SmsBillingRecord为charged2957/refunded30,余额恒等式精确成立,六通道与数据库最终无等待锁/idle in transaction。证据确认瓶颈是每个正价Submit结果在企业账户锁上串行,而不是单价0时可见的通道容量。 +- 后续最小修复将正常“释放冻结+正式扣费”改为单个幂等SQL:冻结已在入口扣减,流水对净余额为零,因此正常结算不再取得账户锁或更新账户;历史已释放未扣费状态仍用原带锁路径补扣,并保留并发`ON CONFLICT`后一次MVCC恢复读。专项134项和TypeScript已通过,完整门禁、修复提交、新恢复资产与复测待补记。