perf: remove paid result account lock contention

This commit is contained in:
hectorzhao
2026-08-21 10:40:28 +08:00
parent 487b5282a6
commit fcf6d3e6b4
6 changed files with 58 additions and 48 deletions
+10 -8
View File
@@ -11,7 +11,7 @@ function createPrismaMock() {
findMany: jest.fn().mockResolvedValue([]), findMany: jest.fn().mockResolvedValue([]),
}, },
tenantAccount: { tenantAccount: {
findMany: jest.fn(), findMany: jest.fn().mockResolvedValue([]),
findUnique: jest.fn().mockImplementation(() => Promise.resolve({ ...accountState })), findUnique: jest.fn().mockImplementation(() => Promise.resolve({ ...accountState })),
findUniqueOrThrow: jest.fn().mockImplementation(() => Promise.resolve({ ...accountState })), findUniqueOrThrow: jest.fn().mockImplementation(() => Promise.resolve({ ...accountState })),
create: jest.fn(), create: jest.fn(),
@@ -30,7 +30,7 @@ function createPrismaMock() {
}), }),
}, },
accountTransaction: { accountTransaction: {
findMany: jest.fn(), findMany: jest.fn().mockResolvedValue([]),
findFirst: jest.fn(), findFirst: jest.fn(),
findUnique: jest.fn().mockResolvedValue(null), findUnique: jest.fn().mockResolvedValue(null),
findUniqueOrThrow: jest.fn().mockImplementation(({ where }) => Promise.resolve({ id: 'tx-charged', idempotencyKey: where.idempotencyKey })), 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' }), create: jest.fn().mockResolvedValue({ id: 'operation-1' }),
}, },
$executeRaw: jest.fn(), $executeRaw: jest.fn(),
$queryRaw: jest.fn(),
}; };
return Object.assign(prisma, { return Object.assign(prisma, {
$transaction: jest.fn((callback: (client: typeof prisma) => unknown) => callback(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 () => { it('settles a frozen SMS charge with one account lock and an idempotent ledger pair', async () => {
const prisma = createPrismaMock(); const prisma = createPrismaMock();
const rows = new Map<string, Record<string, unknown>>(); const rows = new Map<string, Record<string, unknown>>();
prisma.accountTransaction.findUnique.mockImplementation(({ where }) => Promise.resolve(rows.get(where.idempotencyKey) ?? null)); prisma.accountTransaction.findMany.mockImplementation(() => Promise.resolve([...rows.values()]));
prisma.accountTransaction.findUniqueOrThrow.mockImplementation(({ where }) => Promise.resolve(rows.get(where.idempotencyKey))); prisma.$queryRaw.mockImplementation(() => {
prisma.accountTransaction.createMany.mockImplementation(({ data }) => { const charged = { id: 'tx-charged', idempotencyKey: 'sms-charge:msg-paid-1', transactionType: 'charged', amountCents: -325 };
data.forEach((row: Record<string, unknown>) => rows.set(String(row.idempotencyKey), { id: `tx-${row.transactionType}`, ...row })); rows.set(String(charged.idempotencyKey), charged);
return Promise.resolve({ count: data.length }); return Promise.resolve([charged]);
}); });
const service = new BillingService(prisma as never); const service = new BillingService(prisma as never);
const input = { tenantId: 'tenant-1', amountCents: 325, messageId: 'msg-paid-1', taskId: 'task-paid-1' }; 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(first).toEqual(expect.objectContaining({ transactionType: 'charged', amountCents: -325 }));
expect(replay).toEqual(first); 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(); expect(prisma.tenantAccount.update).not.toHaveBeenCalled();
}); });
+42 -37
View File
@@ -441,44 +441,49 @@ export class BillingService {
if (amountCents <= 0) return null; if (amountCents <= 0) return null;
const releaseKey = `sms-charge-release:${data.messageId}`; const releaseKey = `sms-charge-release:${data.messageId}`;
const chargeKey = `sms-charge:${data.messageId}`; const chargeKey = `sms-charge:${data.messageId}`;
return this.prisma.$transaction(async (tx) => { const existing = await this.prisma.accountTransaction.findMany({
await tx.$executeRaw`SELECT pg_advisory_xact_lock(hashtextextended(${'tenant-account:' + data.tenantId}, 0))`; where: { idempotencyKey: { in: [releaseKey, chargeKey] } },
const [release, charge] = await Promise.all([ });
tx.accountTransaction.findUnique({ where: { idempotencyKey: releaseKey } }), const charge = existing.find((row) => row.idempotencyKey === chargeKey);
tx.accountTransaction.findUnique({ where: { idempotencyKey: chargeKey } }), if (charge) return charge;
]); const release = existing.find((row) => row.idempotencyKey === releaseKey);
if (charge) return charge; if (release) {
const account = await tx.tenantAccount.findUniqueOrThrow({ where: { tenantId: data.tenantId } }); // Recover the legacy two-transaction boundary: a crash may have committed
const balance = moneyToNumber(account.balanceCents); // release before charge, so this path must perform the missing balance debit.
if (release) { return this.applyAccountDelta({
const updated = await tx.tenantAccount.update({ tenantId: data.tenantId, transactionType: 'charged', idempotencyKey: chargeKey,
where: { tenantId: data.tenantId }, amountCents: -amountCents, relatedType: 'sms_message_record', relatedId: data.messageId,
data: { balanceCents: { decrement: amountCents } }, remark: '提交成功扣费(恢复既有已释放冻结)',
});
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: '提交成功扣费',
},
],
}); });
return tx.accountTransaction.findUniqueOrThrow({ where: { idempotencyKey: chargeKey } }); }
}, { isolationLevel: Prisma.TransactionIsolationLevel.ReadCommitted }); const rows = await this.prisma.$queryRaw<Array<{ id: string; idempotencyKey: string }>>(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) { release(data: BillingActionDto) {
+1 -1
View File
@@ -1168,6 +1168,6 @@ global,不能错误归入client。`client-signature-*`、发送页、企业认
## 2026-08-21 第四阶段模块边界 ## 2026-08-21 第四阶段模块边界
- `send-inbound-entry`负责编排正价批次短事务;账户固定锁序、总额覆盖、逐短信冻结幂等键留在PostgreSQL一致性边界,不迁入Redis。 - `send-inbound-entry`负责编排正价批次短事务;账户固定锁序、总额覆盖、逐短信冻结幂等键留在PostgreSQL一致性边界,不迁入Redis。
- `billing.service`集中维护账户锁和流水`settleFrozenCharge`释放/扣费合并为一个事务`send-accounting`只编排消息计费状态。 - `billing.service`集中维护余额锁和流水`settleFrozenCharge`利用冻结已完成扣减、释放扣费净变化为零的事实,以单个幂等SQL写成流水对且不争抢账户锁,只有历史半完成恢复才进入带锁余额补扣`send-accounting`只编排消息计费状态。
- `send-chain.helpers`提供无副作用的优先级、主备和稳定weight选择;提交服务只提供消息稳定键与真实在线/报备候选,不保存进程内轮询状态。 - `send-chain.helpers`提供无副作用的优先级、主备和稳定weight选择;提交服务只提供消息稳定键与真实在线/报备候选,不保存进程内轮询状态。
- Gateway既有Submit池与结果Outbox继续独立有界;先以正价六通道实测定位容量,只有证据证明单条回调仍触发停止线时才扩展批量回调协议。 - Gateway既有Submit池与结果Outbox继续独立有界;先以正价六通道实测定位容量,只有证据证明单条回调仍触发停止线时才扩展批量回调协议。
@@ -2133,6 +2133,6 @@
## 完整处理500条/秒第四阶段:正价计费与供应商并行(2026-08-21) ## 完整处理500条/秒第四阶段:正价计费与供应商并行(2026-08-21)
- 普通CMPP短短信批量路径必须支持正单价:按企业固定锁序一次校验批次总金额、一次更新余额,并为每短信写独立幂等冻结流水;任务、请求、消息与冻结流水同属一个短事务。余额加授信必须覆盖所需金额,不能只判断大于零。 - 普通CMPP短短信批量路径必须支持正单价:按企业固定锁序一次校验批次总金额、一次更新余额,并为每短信写独立幂等冻结流水;任务、请求、消息与冻结流水同属一个短事务。余额加授信必须覆盖所需金额,不能只判断大于零。
- 供应商接受后,释放冻结与正式扣费一个账户锁事务内保留两条审计流水;零金额不创建AccountTransaction拒绝、超时、重投和最终失败继续沿用逐短信幂等释放、扣费和退款。 - 供应商接受后,释放冻结与正式扣费一个原子SQL写入两条幂等审计流水;二者净余额变化为零,不得争抢企业账户锁或重复更新账户。只有恢复历史“已释放但未扣费”半完成状态时才取得账户锁补扣。零金额不创建AccountTransaction拒绝、超时、重投和最终失败继续沿用逐短信幂等释放、扣费和退款。
- 同通道组内在线、已报备、同地域、同最低优先级且非备用的通道按weight和稳定消息键分流;备用、离线及更低优先级不抢占。只有显式配置为同优先级非备用的通道才参与主动双活。 - 同通道组内在线、已报备、同地域、同最低优先级且非备用的通道按weight和稳定消息键分流;备用、离线及更低优先级不抢占。只有显式配置为同优先级非备用的通道才参与主动双活。
- 压测使用隔离测试企业、真实PostgreSQL账务、单价`0.0325元/计费条`、三运营商混合号码和六个模拟供应商账号;逐档核对消息、冻结/释放/扣费/退款、SmsBillingRecord、通道分布、队列和数据库。临时单价与主动双活配置须先快照、测试后恢复。 - 压测使用隔离测试企业、真实PostgreSQL账务、单价`0.0325元/计费条`、三运营商混合号码和六个模拟供应商账号;逐档核对消息、冻结/释放/扣费/退款、SmsBillingRecord、通道分布、队列和数据库。临时单价与主动双活配置须先快照、测试后恢复。
+1 -1
View File
@@ -4779,7 +4779,7 @@ npm run verify:phase8
| --- | --- | --- | | --- | --- | --- |
| TC-CMPP-500-P4-001 | 同企业正价短短信批量入站并重放 | 一次账户余额更新;每短信唯一`:freeze`流水;任务/请求/消息金额正确;重放不重复冻结 | | TC-CMPP-500-P4-001 | 同企业正价短短信批量入站并重放 | 一次账户余额更新;每短信唯一`:freeze`流水;任务/请求/消息金额正确;重放不重复冻结 |
| TC-CMPP-500-P4-002 | 余额加授信小于批次或单条金额 | 不得透支;批次安全回退且只发送余额覆盖的短信 | | 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-004 | 零价短信Submit接受 | 保留业务结果但无0金额AccountTransaction |
| TC-CMPP-500-P4-005 | 两个同优先级非备用通道与一个备用通道 | 稳定weighted分流只覆盖两个主动通道;排除已尝试通道后可安全切换;备用不抢占 | | TC-CMPP-500-P4-005 | 两个同优先级非备用通道与一个备用通道 | 稳定weighted分流只覆盖两个主动通道;排除已尝试通道后可安全切换;备用不抢占 |
| TC-CMPP-500-P4-006 | 三运营商、六通道、325金额单位阶梯压测 | 每档入口/Inbox/供应商/回执/账务对账一致,双Stream和数据库最终稳定排空;失败即停止升档 | | TC-CMPP-500-P4-006 | 三运营商、六通道、325金额单位阶梯压测 | 每档入口/Inbox/供应商/回执/账务对账一致,双Stream和数据库最终稳定排空;失败即停止升档 |
+3
View File
@@ -3840,3 +3840,6 @@ git diff --check
- 10个隔离应用仍为单价0、同企业余额1/授信0;移动主备各200条/秒,联通/电信主备各150条/秒,窗口32。组项目为主10/备用20,默认只能三主主动承载,六连接在线不等于六通道并行。 - 10个隔离应用仍为单价0、同企业余额1/授信0;移动主备各200条/秒,联通/电信主备各150条/秒,窗口32。组项目为主10/备用20,默认只能三主主动承载,六连接在线不等于六通道并行。
- 已实现正价Worker批次、覆盖本次金额的余额判断、单事务释放/扣费、零价免账户流水,以及同最低优先级非备用通道的稳定weighted分流;默认主备语义不变。API正式编译和Billing/纯策略/SendChain三套143项通过。 - 已实现正价Worker批次、覆盖本次金额的余额判断、单事务释放/扣费、零价免账户流水,以及同最低优先级非备用通道的稳定weighted分流;默认主备语义不变。API正式编译和Billing/纯策略/SendChain三套143项通过。
- 部署、恢复资产、测试充值、临时六通道主动双活、smoke及阶梯结果待补记;所有临时配置只允许在测试环境先快照后变更并在测试后恢复。 - 部署、恢复资产、测试充值、临时六通道主动双活、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已通过,完整门禁、修复提交、新恢复资产与复测待补记。