From c1a17699db29b7fcbde63c0b716cbf1abceb0cdf Mon Sep 17 00:00:00 2001 From: hectorzhao Date: Tue, 14 Jul 2026 18:45:51 +0800 Subject: [PATCH] feat: restore credit limits and harden SMS sending --- .env.example | 5 + .../migration.sql | 2 + api/prisma/schema.prisma | 1 + api/src/billing/billing.controller.ts | 9 +- api/src/billing/billing.service.spec.ts | 24 +- api/src/billing/billing.service.ts | 49 +- api/src/send-chain/send-chain.service.spec.ts | 41 +- api/src/send-chain/send-chain.service.ts | 70 ++- ...ccess-number-upstream-downstream-design.md | 492 ++++++++++++++++++ .../first-version-development-requirements.md | 11 +- docs/production-deployment.md | 8 + docs/system-functional-test-cases.md | 20 +- docs/testing-progress.md | 16 + gateway/internal/upstream/deliver_test.go | 11 + src/api/adminApi.ts | 5 +- src/apps/admin/AdminCustomerDetailPage.tsx | 1 + src/apps/admin/AdminCustomerFormPage.tsx | 50 +- src/apps/admin/AdminCustomersPage.tsx | 1 + src/apps/admin/AdminHome.tsx | 2 +- src/apps/client/ClientBillingPage.tsx | 8 +- src/apps/client/ClientHome.tsx | 6 +- tools/deploy/production-bootstrap.sh | 10 + tools/deploy/production-deploy.sh | 10 + 23 files changed, 796 insertions(+), 56 deletions(-) create mode 100644 api/prisma/migrations/20260714190000_restore_negative_credit_limit/migration.sql create mode 100644 docs/access-number-upstream-downstream-design.md diff --git a/.env.example b/.env.example index fe4691f..6e23f7a 100644 --- a/.env.example +++ b/.env.example @@ -7,6 +7,8 @@ DATABASE_URL=postgresql://cmpp:cmpp_password@localhost:5432/cmpp_platform?schema REDIS_HOST=127.0.0.1 REDIS_PORT=6379 REDIS_URL=redis://127.0.0.1:6379 +API_ENABLE_SEND_WORKER=true +API_SEND_WORKER_CONCURRENCY=50 ADMIN_SESSION_IDLE_TIMEOUT_MS=3600000 CLIENT_SESSION_IDLE_TIMEOUT_MS=7200000 SESSION_LOCK_RECOVERY_MS=14400000 @@ -17,6 +19,9 @@ OPERATION_LOG_RETENTION_DAYS=180 OPERATION_LOG_ARCHIVE_BATCH_SIZE=1000 OPERATION_LOG_ARCHIVE_MAX_BATCHES=20 OPERATION_LOG_ARCHIVE_INTERVAL_MS=86400000 +SMS_RECEIPT_TIMEOUT_SCAN_ENABLED=true +SMS_RECEIPT_TIMEOUT_HOURS=72 +SMS_RECEIPT_TIMEOUT_SCAN_INTERVAL_MS=300000 # Local HTTP development only. Production must use HTTPS and true. SESSION_COOKIE_SECURE=false MINIO_ENDPOINT=localhost:9000 diff --git a/api/prisma/migrations/20260714190000_restore_negative_credit_limit/migration.sql b/api/prisma/migrations/20260714190000_restore_negative_credit_limit/migration.sql new file mode 100644 index 0000000..6d58b3d --- /dev/null +++ b/api/prisma/migrations/20260714190000_restore_negative_credit_limit/migration.sql @@ -0,0 +1,2 @@ +ALTER TABLE "TenantAccount" + ADD COLUMN "creditCents" INTEGER NOT NULL DEFAULT 0; diff --git a/api/prisma/schema.prisma b/api/prisma/schema.prisma index c0fb8a8..d38f61b 100644 --- a/api/prisma/schema.prisma +++ b/api/prisma/schema.prisma @@ -267,6 +267,7 @@ model TenantAccount { id String @id @default(cuid()) tenantId String balanceCents Int @default(0) + creditCents Int @default(0) status String @default("active") createdAt DateTime @default(now()) updatedAt DateTime @updatedAt diff --git a/api/src/billing/billing.controller.ts b/api/src/billing/billing.controller.ts index 6228a09..a892836 100644 --- a/api/src/billing/billing.controller.ts +++ b/api/src/billing/billing.controller.ts @@ -1,4 +1,4 @@ -import { Body, Controller, Get, Post } from '@nestjs/common'; +import { Body, Controller, Get, Param, Post } from '@nestjs/common'; import { ApiTags } from '@nestjs/swagger'; import { TenantId } from '../common/tenant-id.decorator'; import { RequireRecentAuthentication } from '../auth/require-recent-authentication.decorator'; @@ -10,6 +10,7 @@ import { CreateSmsBillingRecordDto, CreateTenantAccountDto, EstimateSmsCostDto, + UpdateCreditLimitDto, } from './billing.service'; @ApiTags('billing') @@ -27,6 +28,12 @@ export class BillingController { return this.billing.createAccount(body); } + @Post('accounts/:tenantId/credit-limit') + @RequireRecentAuthentication() + updateCreditLimit(@Param('tenantId') tenantId: string, @Body() body: UpdateCreditLimitDto) { + return this.billing.updateCreditLimit(tenantId, body); + } + @Get('recharges') listRechargeOrders(@TenantId() tenantId?: string) { return this.billing.listRechargeOrders(tenantId); diff --git a/api/src/billing/billing.service.spec.ts b/api/src/billing/billing.service.spec.ts index e74eddc..d313cb4 100644 --- a/api/src/billing/billing.service.spec.ts +++ b/api/src/billing/billing.service.spec.ts @@ -1,7 +1,7 @@ import { BillingService } from './billing.service'; function createPrismaMock() { - const accountState = { tenantId: 'tenant-1', balanceCents: 1000 }; + const accountState = { id: 'account-1', tenantId: 'tenant-1', balanceCents: 1000, creditCents: 0 }; return { accountState, tenantAccount: { @@ -9,7 +9,8 @@ function createPrismaMock() { create: jest.fn(), upsert: jest.fn().mockImplementation(() => Promise.resolve({ ...accountState })), update: jest.fn().mockImplementation(({ data }) => { - accountState.balanceCents = data.balanceCents; + if (data.balanceCents !== undefined) accountState.balanceCents = data.balanceCents; + if (data.creditCents !== undefined) accountState.creditCents = data.creditCents; return Promise.resolve({ ...accountState }); }), }, @@ -62,7 +63,7 @@ describe('BillingService', () => { ); }); - it('checks only the cash balance before sending', async () => { + it('allows sending only when cash balance plus credit is greater than zero', async () => { const prisma = createPrismaMock(); const service = new BillingService(prisma as never); @@ -70,8 +71,23 @@ describe('BillingService', () => { expect.objectContaining({ availableAmount: 1000, canSend: true }), ); await expect(service.checkAccount({ tenantId: 'tenant-1', amountCents: 1001 })).resolves.toEqual( - expect.objectContaining({ canSend: false }), + expect.objectContaining({ availableAmount: 1000, canSend: true }), ); + await service.updateCreditLimit('tenant-1', { creditCents: -1000, operatorId: 'admin-1' }); + await expect(service.checkAccount({ tenantId: 'tenant-1', amountCents: 1 })).resolves.toEqual( + expect.objectContaining({ availableAmount: 0, balanceCents: 1000, creditCents: -1000, canSend: false }), + ); + await service.updateCreditLimit('tenant-1', { creditCents: 500, operatorId: 'admin-1' }); + await expect(service.checkAccount({ tenantId: 'tenant-1', amountCents: 999999 })).resolves.toEqual( + expect.objectContaining({ availableAmount: 1500, creditCents: 500, canSend: true }), + ); + await expect(service.updateCreditLimit('tenant-1', { creditCents: 1.5 })).rejects.toThrow('授信额度必须为整数金额'); + expect(prisma.operationLog.create).toHaveBeenCalledWith({ + data: expect.objectContaining({ + action: 'billing.credit_limit_updated', + detail: expect.objectContaining({ previousCreditCents: 0, creditCents: -1000 }), + }), + }); }); it('creates cash recharge orders and account transactions without plans', async () => { diff --git a/api/src/billing/billing.service.ts b/api/src/billing/billing.service.ts index a898ef7..be4c9c5 100644 --- a/api/src/billing/billing.service.ts +++ b/api/src/billing/billing.service.ts @@ -1,13 +1,20 @@ -import { Injectable } from '@nestjs/common'; +import { BadRequestException, Injectable } from '@nestjs/common'; import { Prisma } from '@prisma/client'; import { PrismaService } from '../prisma/prisma.service'; export interface CreateTenantAccountDto { tenantId: string; balanceCents?: number; + creditCents?: number; status?: string; } +export interface UpdateCreditLimitDto { + creditCents: number; + operatorId?: string; + remark?: string; +} + export interface CreateAccountTransactionDto { tenantId: string; transactionType: string; @@ -80,14 +87,40 @@ export class BillingService { } createAccount(data: CreateTenantAccountDto) { + assertCreditAmount(data.creditCents ?? 0); const createData: Prisma.TenantAccountUncheckedCreateInput = { tenantId: data.tenantId, balanceCents: data.balanceCents ?? 0, + creditCents: data.creditCents ?? 0, status: data.status ?? 'active', }; return this.prisma.tenantAccount.create({ data: createData }); } + async updateCreditLimit(tenantId: string, data: UpdateCreditLimitDto) { + assertCreditAmount(data.creditCents); + const account = await this.getAccountOrCreate(tenantId); + const updated = await this.prisma.tenantAccount.update({ + where: { tenantId }, + data: { creditCents: data.creditCents }, + }); + await this.prisma.operationLog.create({ + data: { + tenantId, + userId: data.operatorId, + action: 'billing.credit_limit_updated', + resource: 'tenant_account', + resourceId: account.id, + detail: { + previousCreditCents: account.creditCents, + creditCents: data.creditCents, + remark: data.remark, + } as Prisma.InputJsonValue, + }, + }); + return updated; + } + listRechargeOrders(tenantId?: string) { return this.prisma.rechargeOrder.findMany({ where: tenantId ? { tenantId } : undefined, @@ -195,12 +228,14 @@ export class BillingService { async checkAccount(data: BillingActionDto) { const account = await this.getAccountOrCreate(data.tenantId); const requiredAmount = data.amountCents ?? 0; - const availableAmount = account.balanceCents; + const availableAmount = account.balanceCents + account.creditCents; return { tenantId: data.tenantId, requiredAmount, availableAmount, - canSend: availableAmount >= requiredAmount, + balanceCents: account.balanceCents, + creditCents: account.creditCents, + canSend: availableAmount > 0, }; } @@ -296,7 +331,7 @@ export class BillingService { return this.prisma.tenantAccount.upsert({ where: { tenantId }, update: {}, - create: { tenantId, balanceCents: 0, status: 'active' }, + create: { tenantId, balanceCents: 0, creditCents: 0, status: 'active' }, }); } @@ -324,6 +359,12 @@ export class BillingService { } } +function assertCreditAmount(creditCents: number) { + if (!Number.isInteger(creditCents)) { + throw new BadRequestException('授信额度必须为整数金额(分)'); + } +} + function estimateBillingUnits(content: string) { const length = [...content].length; if (length <= 70) { diff --git a/api/src/send-chain/send-chain.service.spec.ts b/api/src/send-chain/send-chain.service.spec.ts index 2e7f3ca..6b80299 100644 --- a/api/src/send-chain/send-chain.service.spec.ts +++ b/api/src/send-chain/send-chain.service.spec.ts @@ -1654,18 +1654,49 @@ describe('SendChainService', () => { await expect(service.batchRequeueDownstreamDeliveries([])).rejects.toThrow('请选择至少一条下游投递记录'); }); - it('marks 72 hour unknown receipts as timeout', async () => { - const { service, prisma } = createService(); + it('marks submitted or unknown messages without a final receipt for 72 hours as timeout and refunds them', async () => { + const { service, prisma, billing } = createService(); prisma.smsMessageRecord.findMany.mockResolvedValue([ - { id: 'record-1', batchTaskId: 'task-1' }, - { id: 'record-2', batchTaskId: 'task-1' }, + { id: 'record-1', tenantId: 'tenant-1', batchTaskId: 'task-1', messageId: 'MSG-1', amountCents: 3, billingUnits: 1 }, + { id: 'record-2', tenantId: 'tenant-1', batchTaskId: 'task-1', messageId: 'MSG-2', amountCents: 3, billingUnits: 1 }, ]); + prisma.smsBillingRecord.findFirst + .mockResolvedValueOnce(null).mockResolvedValueOnce({ id: 'bill-1', billingStatus: 'charged' }) + .mockResolvedValueOnce(null).mockResolvedValueOnce({ id: 'bill-2', billingStatus: 'charged' }); await expect(service.markUnknownTimeout({ olderThanHours: 72 })).resolves.toEqual({ timeout: 2 }); + expect(prisma.smsMessageRecord.findMany).toHaveBeenCalledWith({ + where: { + tenantId: { not: null }, + status: { in: ['submitted', 'unknown'] }, + submittedAt: { lte: expect.any(Date) }, + }, + select: { id: true, tenantId: true, batchTaskId: true, messageId: true, amountCents: true, billingUnits: true }, + take: 10000, + }); expect(prisma.smsMessageRecord.updateMany).toHaveBeenCalledWith({ - where: { id: { in: ['record-1', 'record-2'] } }, + where: { id: 'record-1', status: { in: ['submitted', 'unknown'] } }, data: expect.objectContaining({ status: 'timeout', errorMessage: '72小时未收到明确回执,自动转超时' }), }); + expect(billing.refund).toHaveBeenCalledTimes(2); expect(prisma.smsBatchTask.update).toHaveBeenCalled(); }); + + it('starts the automatic receipt-timeout scan after application startup', async () => { + jest.useFakeTimers(); + const previousEnabled = process.env.SMS_RECEIPT_TIMEOUT_SCAN_ENABLED; + const { service } = createService(); + const scan = jest.spyOn(service, 'markUnknownTimeout').mockResolvedValue({ timeout: 0 }); + try { + process.env.SMS_RECEIPT_TIMEOUT_SCAN_ENABLED = 'true'; + service.onModuleInit(); + await jest.advanceTimersByTimeAsync(60_000); + expect(scan).toHaveBeenCalledWith({}); + await service.onModuleDestroy(); + } finally { + if (previousEnabled === undefined) delete process.env.SMS_RECEIPT_TIMEOUT_SCAN_ENABLED; + else process.env.SMS_RECEIPT_TIMEOUT_SCAN_ENABLED = previousEnabled; + jest.useRealTimers(); + } + }); }); diff --git a/api/src/send-chain/send-chain.service.ts b/api/src/send-chain/send-chain.service.ts index 8220fe9..0253073 100644 --- a/api/src/send-chain/send-chain.service.ts +++ b/api/src/send-chain/send-chain.service.ts @@ -1,4 +1,4 @@ -import { BadRequestException, Injectable, NotFoundException, OnModuleDestroy, OnModuleInit } from '@nestjs/common'; +import { BadRequestException, Injectable, Logger, NotFoundException, OnModuleDestroy, OnModuleInit } from '@nestjs/common'; import { Prisma } from '@prisma/client'; import { Queue, Worker } from 'bullmq'; import IORedis from 'ioredis'; @@ -218,6 +218,9 @@ const GATEWAY_SUBMIT_STREAM = 'gateway.submit.commands'; const DEFAULT_DOWNSTREAM_RETRY_DELAY_MS = 60_000; const DEFAULT_DOWNSTREAM_RETRY_MAX_DELAY_MS = 30 * 60_000; const DEFAULT_DOWNSTREAM_MAX_RETRIES = 10; +const DEFAULT_RECEIPT_TIMEOUT_HOURS = 72; +const DEFAULT_RECEIPT_TIMEOUT_SCAN_INTERVAL_MS = 5 * 60_000; +const RECEIPT_TIMEOUT_INITIAL_DELAY_MS = 60_000; const BULLMQ_PRIORITY: Record = { priority: 1, normal: 100, @@ -225,10 +228,14 @@ const BULLMQ_PRIORITY: Record = { @Injectable() export class SendChainService implements OnModuleInit, OnModuleDestroy { + private readonly logger = new Logger(SendChainService.name); private redis?: IORedis; private sendQueue?: Queue; private gatewayQueue?: Queue; private worker?: Worker; + private receiptTimeoutInitialTimer?: ReturnType; + private receiptTimeoutIntervalTimer?: ReturnType; + private receiptTimeoutScanRunning = false; constructor( private readonly prisma: PrismaService, @@ -240,9 +247,20 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { if (process.env.API_ENABLE_SEND_WORKER === 'true') { this.startWorker(); } + if (process.env.SMS_RECEIPT_TIMEOUT_SCAN_ENABLED !== 'false') { + this.receiptTimeoutInitialTimer = setTimeout(() => void this.runReceiptTimeoutScan(), RECEIPT_TIMEOUT_INITIAL_DELAY_MS); + this.receiptTimeoutInitialTimer.unref?.(); + this.receiptTimeoutIntervalTimer = setInterval( + () => void this.runReceiptTimeoutScan(), + positiveInteger(process.env.SMS_RECEIPT_TIMEOUT_SCAN_INTERVAL_MS, DEFAULT_RECEIPT_TIMEOUT_SCAN_INTERVAL_MS), + ); + this.receiptTimeoutIntervalTimer.unref?.(); + } } async onModuleDestroy() { + if (this.receiptTimeoutInitialTimer) clearTimeout(this.receiptTimeoutInitialTimer); + if (this.receiptTimeoutIntervalTimer) clearInterval(this.receiptTimeoutIntervalTimer); await this.worker?.close(); await this.sendQueue?.close(); await this.gatewayQueue?.close(); @@ -1807,30 +1825,47 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { } async markUnknownTimeout(data: TimeoutUnknownDto) { - const olderThanHours = data.olderThanHours ?? 72; + const olderThanHours = data.olderThanHours ?? positiveInteger(process.env.SMS_RECEIPT_TIMEOUT_HOURS, DEFAULT_RECEIPT_TIMEOUT_HOURS); const cutoff = new Date(Date.now() - olderThanHours * 60 * 60 * 1000); const candidates = await this.prisma.smsMessageRecord.findMany({ where: { - status: 'unknown', - deliveredAt: { lte: cutoff }, + tenantId: { not: null }, + status: { in: ['submitted', 'unknown'] }, + submittedAt: { lte: cutoff }, }, - select: { id: true, batchTaskId: true }, + select: { id: true, tenantId: true, batchTaskId: true, messageId: true, amountCents: true, billingUnits: true }, take: 10000, }); - await this.prisma.smsMessageRecord.updateMany({ - where: { id: { in: candidates.map((candidate) => candidate.id) } }, - data: { status: 'timeout', timeoutAt: new Date(), errorMessage: '72小时未收到明确回执,自动转超时' }, - }); + const timedOutTaskIds = new Set(); + let timeout = 0; for (const candidate of candidates) { - const message = await this.prisma.smsMessageRecord.findUnique({ where: { id: candidate.id } }); - if (message?.tenantId) { - await this.refundMessage(message as typeof message & { tenantId: string }, '72小时未收到明确回执,自动超时退款'); - } + if (!candidate.tenantId) continue; + const transitioned = await this.prisma.smsMessageRecord.updateMany({ + where: { id: candidate.id, status: { in: ['submitted', 'unknown'] } }, + data: { status: 'timeout', timeoutAt: new Date(), errorMessage: `${olderThanHours}小时未收到明确回执,自动转超时` }, + }); + if (transitioned.count !== 1) continue; + timeout += 1; + await this.refundMessage(candidate as typeof candidate & { tenantId: string }, `${olderThanHours}小时未收到明确回执,自动超时退款`); + if (candidate.batchTaskId) timedOutTaskIds.add(candidate.batchTaskId); } - for (const batchTaskId of new Set(candidates.map((candidate) => candidate.batchTaskId).filter((value): value is string => Boolean(value)))) { + for (const batchTaskId of timedOutTaskIds) { await this.refreshTaskProgress(batchTaskId); } - return { timeout: candidates.length }; + return { timeout }; + } + + private async runReceiptTimeoutScan() { + if (this.receiptTimeoutScanRunning) return; + this.receiptTimeoutScanRunning = true; + try { + const result = await this.markUnknownTimeout({}); + if (result.timeout > 0) this.logger.log(`Marked ${result.timeout} messages as receipt timeout and refunded charged messages`); + } catch (error) { + this.logger.error('Receipt timeout scan failed', error instanceof Error ? error.stack : String(error)); + } finally { + this.receiptTimeoutScanRunning = false; + } } private async submitMessageToGateway( @@ -2883,6 +2918,11 @@ function isProvinceChannel(item: { province?: string | null; channel: { sendRegi return itemProvince === target || sendRegion === target; } +function positiveInteger(value: string | undefined, fallback: number) { + const parsed = Number(value); + return Number.isInteger(parsed) && parsed > 0 ? parsed : fallback; +} + function bullmqConnection() { const redisUrl = new URL(process.env.REDIS_URL ?? 'redis://127.0.0.1:6379'); return { diff --git a/docs/access-number-upstream-downstream-design.md b/docs/access-number-upstream-downstream-design.md new file mode 100644 index 0000000..eaee2a8 --- /dev/null +++ b/docs/access-number-upstream-downstream-design.md @@ -0,0 +1,492 @@ +# CMPP 接入号上下游配置与匹配流程设计 + +## 1. 文档目的 + +本文梳理当前 CMPP 短信平台接入号的真实代码、生产数据和协议字段使用情况,并给出客户侧下游 CMPP 接入号、运营商侧上游 CMPP 接入号、号码映射、普通上行匹配和历史追溯的目标设计。 + +本文仅描述设计结论,不代表相关改造已经完成。后续实现仍必须基于真实 NestJS API、Prisma/PostgreSQL、Go Gateway、Redis 和生产 CMPP 链路验证,不允许使用前端静态数据、localStorage 或 mock 代替业务完成。 + +## 2. 结论摘要 + +当前系统已经具备“通道接入号下发、上行接收、模糊匹配、人工认领”的第一版链路,但尚未形成真正的应用级接入号分配和上下游映射模型。 + +现状核心问题如下: + +1. 上游通道接入号和客户应用接入号混用。 +2. 客户侧展示的接入地址、端口和接入号来自任意一条上游通道,而不是平台客户接入端点和本应用号码。 +3. 客户 CMPP SUBMIT 中的 `Src_Id` 虽已传给 NestJS,但未校验是否属于当前应用,也未参与后续通道号码映射。 +4. 多个应用绑定同一通道组时,上行 `Dest_Id` 会匹配到该组全部应用,无法唯一定位。 +5. `extensionDigits` 目前只是配置和队列元数据,Go Gateway 没有根据它生成或匹配应用扩展号。 +6. 上游 CMPP SUBMIT 的 `MsgSrc` 当前使用通道登录账号,而不是通道企业代码。 +7. 普通上行没有完整保留 `Service_Id`、`LinkID`、原始接入号和规范化接入号。 +8. 应用级 `cmppEnterpriseCode` 当前允许最长 32 字符,与 CMPP SUBMIT 中 `MsgSrc` 固定 6 字节的字段边界不一致。 + +生产只读检查显示: + +- 当前 3 条通道的 `srcId` 均为 `1069999999`。 +- 其中 2 条通道为 active,1 条为 disabled。 +- 生产通道未配置 `serviceId` 和有效扩展位数。 +- 4 个真实应用路由到同一个通道组和同一基础接入号。 +- 当前 2 条普通上行的 `destId` 都是 `1069999999`,匹配状态均为 `ambiguous`。 +- 生产存在大量批量验证应用,不应自动分配正式生产接入号。 + +因此,当前共享基础接入号只能支撑模糊候选和人工认领,不能作为多应用上行自动路由的生产方案。 + +## 3. CMPP 字段边界 + +接入号设计必须区分以下概念: + +| 概念 | CMPP 字段 | 含义 | 当前状态 | +| --- | --- | --- | --- | +| 登录账号 | CONNECT `Source_Addr` | 建立 CMPP TCP 会话的账号,协议字段为 6 字节 | 上下游已有账号字段 | +| 企业代码 | SUBMIT `MsgSrc` | SP 身份、地址翻译、计费和结算标识,协议字段为 6 字节 | 下游应用允许过长;上游错误使用登录账号 | +| 业务代码 | `Service_Id` | 业务类型,最长 10 字节 | 藏在通道 JSON 中,生产未配置 | +| 下发接入号 | SUBMIT `Src_Id` | 服务代码或以服务代码为前缀的长号码,最长 21 字节 | 当前只有通道级基础号码 | +| 上行接入号 | DELIVER `Dest_Id` | 用户实际回复到的服务代码或长号码 | 当前仅按通道基础号反查路由组 | +| 上行手机号 | DELIVER `Src_Terminal_Id` | 发送普通上行的用户手机号 | 当前已处理 | +| 点播关联标识 | `LinkID` | 点播类业务的关联标识 | 当前未传入 NestJS | + +CMPP 2.0/3.0 协议要求: + +- `MsgSrc` 和 CONNECT 登录账号是两个不同概念,即使运营商配置值相同,也必须在数据模型中分开保存。 +- `Src_Id` 可以是基础服务代码,也可以是以服务代码为前缀的长号码。 +- 普通上行 DELIVER 中,`Dest_Id` 是用户回复的接入号,`Src_Terminal_Id` 是用户手机号。 +- 状态报告 DELIVER 和普通上行 DELIVER 必须通过 `Registered_Delivery` 严格分流。 +- 普通上行自身的 `Msg_Id` 不能默认当作原下发消息 ID。 + +协议参考: + +- [CMPP 2.0 协议](https://www.kannel.org/~tolj/specs/CMPP2/CMPP-2.0.pdf) +- [CMPP 3.0 协议](https://www.kannel.org/~tolj/specs/CMPP2/CMPP-v30.pdf) + +行业监管要求端口类短信按照批准的码号结构、位长、用途、使用范围和期限使用,并保存发送时间、接收时间、发送端、接收端和内容等记录。因此系统不能因为协议字段最大长度为 21 字节,就允许任意拼接扩展码;扩展长度必须依据具体运营商合同和报备结果配置。 + +行业规则参考: + +- [通信短信息服务管理规定](https://www.miit.gov.cn/zwgk/zcwj/flfg/art/2020/art_77bc7219833c4a08b4563ba42ad23e1f.html) +- [电信网编号计划(2017 年版)](https://www.miit.gov.cn/jgsj/xgj/wjfb/art/2020/art_eb0adf5b6e7148cbb70802b264878b1e.html) +- [电信网码号资源审批服务指南](https://qhca.miit.gov.cn/zwgk/txfz/hmzy/art/2021/art_754e60bc98664ad4a91748a37d29c6bc.html) + +## 4. 当前真实实现 + +### 4.1 当前下发流程 + +1. 应用通过页面、HTTP API 或客户侧 CMPP 提交短信。 +2. NestJS 根据应用、运营商、省份和通道组选择上游通道。 +3. `SubmitCommand.cmpp.srcId` 直接取 `SmsChannel.srcId`。 +4. Go Gateway 构造运营商侧 CMPP SUBMIT。 +5. 上游 `MsgSrc` 当前取通道登录账号。 +6. 上游 `Src_Id` 取通道基础接入号。 +7. 所有走同一通道的应用对外使用同一基础号码。 + +相关代码: + +- `api/src/send-chain/send-chain.service.ts`:生成 `SubmitCommand`。 +- `gateway/internal/upstream/manager.go`:构造 CMPP 2.0/3.0 SUBMIT 报文。 + +### 4.2 当前客户 CMPP 入站流程 + +1. 客户通过 6 位账号 CONNECT。 +2. Gateway 根据会话账号定位应用,并校验应用企业代码。 +3. 客户 SUBMIT 的 `Src_Id` 被传给 NestJS。 +4. NestJS 未校验该号码是否分配给当前应用。 +5. NestJS 未把客户 `Src_Id` 保存到短信记录或提交记录。 +6. 后续上游发送重新取通道 `srcId`,客户提交号码被丢弃。 + +相关代码: + +- `gateway/internal/inbound/server.go`:客户 CMPP CONNECT/SUBMIT 解析。 +- `api/src/send-chain/send-chain.service.ts`:`submitInboundMessage`。 + +### 4.3 当前普通上行流程 + +1. 运营商通过 CMPP DELIVER 返回用户手机号和 `Dest_Id`。 +2. Gateway 将手机号、接入号和内容回调 NestJS。 +3. NestJS 尝试按 messageId 匹配。 +4. 未匹配时,根据“上行通道所属通道组”查询所有绑定该组的应用。 +5. 多应用共用通道组时生成 `ambiguous` 候选。 +6. 无接入号候选时,再按手机号和最近 72 小时下发记录匹配。 +7. 最后进入人工认领。 + +当前流程能避免多候选时误推客户,但不能提供生产级接入号唯一路由。 + +相关代码: + +- `gateway/internal/upstream/manager.go`:解析普通上行 DELIVER。 +- `api/src/send-chain/send-chain.service.ts`:`resolveUplinkMatch`。 + +### 4.4 当前客户 CMPP 参数接口问题 + +`SmsApplication` 没有应用接入号字段。`getApplicationCmppParams` 会查询最新一条未删除的 `SmsChannel`,将该上游通道的地址、端口和 `srcId` 返回给客户。 + +这会产生以下错误: + +- 客户看到运营商侧上游网关地址,而不是平台客户接入地址。 +- 不同应用得到同一个任意通道接入号。 +- 上游通道新增、删除或排序变化可能改变客户展示参数。 +- 应用接入参数与实际客户连接到 `8.160.169.106:17890` 的生产事实不一致。 + +相关代码:`api/src/sms-config/sms-config.service.ts` 的 `getApplicationCmppParams`。 + +## 5. 目标数据模型 + +目标模型分为平台接入端点、上游号码能力、应用虚拟接入号和上下游绑定四层。 + +### 5.1 `CmppDownstreamEndpoint` + +保存客户连接本平台的公共端点: + +- `id` +- `name` +- `gatewayHost` +- `gatewayPort` +- `supportedVersions` +- `heartbeatSeconds` +- `status` +- `effectiveFrom` +- `effectiveTo` + +生产客户参数应返回平台端点 `8.160.169.106:17890` 或后续正式域名/TLS 地址,不再读取任意上游通道。 + +### 5.2 `UpstreamAccessNumber` + +表示运营商或上游供应商在某条通道上实际开通的号码能力: + +- `id` +- `channelId` +- `enterpriseCode`:上游 `MsgSrc`,严格 6 字节。 +- `serviceId`:最长 10 字节。 +- `baseNumber` +- `extensionMode`:`none/fixed/allocated/shared`。 +- `allowedExtensionLengths`:由合同或报备决定。 +- `maxTotalLength`:不得超过 21。 +- `uplinkReturnMode`:`full/base_only/truncated/custom`。 +- `replySupported` +- `reportStatus` +- `effectiveFrom` +- `effectiveTo` +- `status` + +不能把扩展位数设成全平台统一常量。不同运营商、供应商和通道允许的扩展位数可能不同,应由通道号码能力明确配置。 + +### 5.3 `ApplicationAccessNumber` + +表示客户应用看到、提交和收到普通上行时使用的号码: + +- `id` +- `tenantId` +- `applicationId` +- `endpointId` +- `accessNumber` +- `serviceId` +- `isDefault` +- `direction`:`mt/mo/both`。 +- `replyEnabled` +- `shareMode`:`exclusive/shared`。 +- `effectiveFrom` +- `effectiveTo` +- `status` + +约束: + +1. 同一活动号码原则上只能属于一个应用。 +2. 允许共享时必须显式标记 `shared`。 +3. 共享号码不能仅按接入号直接自动投递。 +4. 应用开放 CMPP 提交前必须至少有一个默认活动号码。 +5. 客户 SUBMIT 中的 `Src_Id` 必须属于当前已鉴权应用。 + +### 5.4 `AccessNumberBinding` + +表示客户应用号码与各上游通道号码之间的真实映射: + +- `id` +- `applicationAccessNumberId` +- `channelId` +- `upstreamAccessNumberId` +- `upstreamFullNumber` +- `upstreamExtension` +- `serviceIdOverride` +- `carrier` +- `province` +- `priority` +- `matchMode` +- `effectiveFrom` +- `effectiveTo` +- `status` + +关键约束: + +- 可自动回复的号码,活动状态下 `(channelId, upstreamFullNumber)` 必须唯一。 +- 多应用共享同一上游号码时,必须标记共享并禁止仅凭接入号自动投递。 +- 激活绑定前校验号码前缀、长度、运营商、报备状态和通道状态。 +- 每个需要发送的运营商/省份路由必须有可用号码绑定,否则该通道不是有效候选。 + +### 5.5 短信和提交快照 + +在 `SmsMessageRecord` 或 `SmsSubmitRecord` 中增加不可变快照: + +- `applicationAccessNumberId` +- `accessNumberBindingId` +- `clientSrcId` +- `upstreamSrcId` +- `upstreamMsgSrc` +- `upstreamServiceId` +- `upstreamChannelId` + +这样即使以后修改应用号码、通道或绑定,历史回执、上行和审计仍可按提交当时的号码关系追溯。 + +## 6. 目标下发流程 + +```mermaid +flowchart LR + A["客户应用提交短信"] --> B["CONNECT账号定位应用"] + B --> C["校验MsgSrc等于应用企业代码"] + C --> D["校验或补全客户Src_Id"] + D --> E["运营商、省份和通道组路由"] + E --> F["查询应用号码到候选通道的有效绑定"] + F --> G["生成上游MsgSrc、Service_Id、Src_Id"] + G --> H["保存客户号和上游号快照"] + H --> I["Gateway发送上游CMPP SUBMIT"] +``` + +详细规则: + +1. CONNECT 账号只用于定位应用。 +2. 客户 `MsgSrc` 必须与应用企业代码严格匹配;建议统一为 6 位 ASCII,现有超长数据需迁移或重新分配。 +3. 客户 `Src_Id`: + - 未填且应用只有一个默认号码时,自动补全默认号码。 + - 已填写时,必须精确属于当前应用。 + - 号码停用、过期、未分配或不允许下发时,返回明确的 `Src_Id` 非法错误。 + - 校验失败不得进入计费、任务和上游发送队列。 +4. 选出上游通道后,再查找应用号码在该通道上的有效绑定。 +5. 主通道没有有效号码绑定时,可按通道组规则选择下一候选,但不能退化使用通道基础号码。 +6. 上游 CMPP 报文必须使用: + - `MsgSrc = UpstreamAccessNumber.enterpriseCode` + - `Service_Id = binding.serviceIdOverride` 或上游号码默认值 + - `Src_Id = AccessNumberBinding.upstreamFullNumber` + - `Dest_Terminal_Id = 用户手机号` +7. 最终客户号、上游号、企业代码、业务代码和绑定 ID 必须持久化为提交快照。 + +## 7. 目标普通上行流程 + +```mermaid +flowchart TD + A["收到上游普通DELIVER"] --> B["保存Channel、原始Dest_Id、Service_Id、LinkID"] + B --> C["按通道规则规范化Dest_Id"] + C --> D{"channelId和完整接入号唯一绑定?"} + D -- "是" --> E["定位应用和客户接入号"] + D -- "否" --> F{"通道明确配置截断或基础号回传?"} + F -- "是" --> G["手机号和近期下发号码快照辅助匹配"] + F -- "否" --> H["生成候选"] + G --> I{"得到唯一高置信结果?"} + I -- "是" --> E + I -- "否" --> H + H --> J["ambiguous或unmatched,进入人工认领"] + E --> K["转换为客户侧CMPP DELIVER"] +``` + +匹配优先级: + +1. `(channelId, rawDestId)` 精确匹配活动绑定。 +2. 按该通道明确配置的规范化、截断或别名规则匹配。 +3. 手机号、通道和近期下发记录中的 `upstreamSrcId` 快照匹配。 +4. 使用 `Service_Id`、`LinkID`、运营商、时间窗口作为辅助条件。 +5. 多候选进入人工认领。 +6. 完全无候选标记 `unmatched`。 + +禁止: + +- 仅因为多个应用绑定同一通道组,就把全部应用视为接入号候选。 +- 选择任意 active 应用兜底。 +- 对所有通道统一执行前缀匹配或截断。 +- 把普通上行自身的 `Msg_Id` 当成原下发消息 ID。 +- 在没有其他证据时将共享接入号自动投递给某个客户。 + +匹配成功后,客户侧普通上行 DELIVER 应使用: + +- `Dest_Id = ApplicationAccessNumber.accessNumber` +- `Src_Terminal_Id = 用户手机号` +- `Service_Id = 客户侧业务代码` +- `LinkID = 可用时保留关联值` + +原始上游号码仍需保存在数据库和运营审计中,但不能直接作为转换后的客户号码。 + +## 8. 配置界面设计 + +### 8.1 平台客户接入端点 + +独立配置: + +- 公网域名/IP +- 端口 +- CMPP 版本 +- TLS 状态 +- 心跳间隔 +- 启停状态 + +### 8.2 上游通道页面 + +通道连接配置和号码能力分区展示: + +- 登录账号 `Source_Addr` +- 企业代码 `MsgSrc` +- 默认 `Service_Id` +- 基础接入号 +- 扩展模式 +- 允许扩展长度 +- 最大总长度 +- 是否支持上行 +- 上行回传模式 +- 报备状态和生效时间 + +### 8.3 应用页面 + +应用 CMPP 参数区展示: + +- 平台接入地址和端口 +- 账号 +- 密码 +- 企业代码 +- 协议版本 +- 连接数和窗口 +- 已分配客户接入号列表 +- 默认号码 +- 上行/下发能力 +- 生效状态 + +### 8.4 上下游号码绑定矩阵 + +按应用号码、运营商、通道组和通道展示: + +- 客户号码 +- 上游基础号 +- 上游扩展码 +- 上游完整号码 +- `Service_Id` +- 报备状态 +- 是否支持上行 +- 主备优先级 +- 生效时间 + +绑定不完整时,不允许把对应通道标记为该应用的可用发送候选。 + +## 9. API 建议 + +建议新增或调整: + +- `GET /api/admin/cmpp-downstream-endpoints` +- `POST /api/admin/cmpp-downstream-endpoints` +- `GET /api/admin/channels/:id/access-numbers` +- `POST /api/admin/channels/:id/access-numbers` +- `GET /api/admin/applications/:id/access-numbers` +- `POST /api/admin/applications/:id/access-numbers` +- `GET /api/admin/applications/:id/access-number-bindings` +- `POST /api/admin/applications/:id/access-number-bindings` +- `POST /api/admin/access-number-bindings/:id/activate` +- `POST /api/admin/access-number-bindings/:id/disable` +- `GET /api/client/applications/:id/cmpp-params` + +客户参数接口返回结构建议: + +```json +{ + "gatewayHost": "8.160.169.106", + "gatewayPort": 17890, + "account": "123456", + "enterpriseCode": "900001", + "protocolVersion": "CMPP2.0", + "maxConnections": 1, + "windowSize": 16, + "accessNumbers": [ + { + "accessNumber": "10699999990001", + "serviceId": "NOTICE", + "isDefault": true, + "replyEnabled": true, + "status": "active" + } + ] +} +``` + +示例号码只用于说明结构,不能在未获得上游合同和报备确认时直接用于生产配置。 + +## 10. 生产数据迁移原则 + +1. 将现有 `SmsChannel.srcId=1069999999` 迁移为各通道的上游基础号码记录。 +2. 从上游供应商合同或后台确认: + - 上游企业代码。 + - `Service_Id`。 + - 允许扩展位数。 + - 是否完整回传扩展号。 + - 是否支持普通上行。 + - 多连接或多通道是否共享同一号段。 +3. 在未确认扩展规则前,不为多个应用自动生成扩展号码。 +4. 现有 2 条 `ambiguous` 上行继续保留人工认领,不自动修改历史归属。 +5. 大量批量验证应用不得自动分配正式号码。 +6. 修复或迁移超过 6 字节的 `cmppEnterpriseCode`。 +7. 修复客户参数接口,不再返回任意上游通道地址和号码。 +8. 对现有两条 active、配置相同的复制通道确认是否属于真实独立上游连接,避免重复通道与号码能力重复生效。 + +## 11. 实施优先级 + +### P0:数据正确性 + +- 新增上游号码、应用号码和绑定表。 +- 修复客户 CMPP 参数接口。 +- 严格校验客户 `MsgSrc` 和 `Src_Id`。 +- 上游 `MsgSrc` 改用通道企业代码。 +- 持久化客户号码和上游号码快照。 +- 上行事件补充 `Service_Id`、`LinkID` 和原始号码。 +- 删除“通道组内所有应用都是接入号候选”的匹配路径。 + +### P1:配置与运营 + +- 增加通道号码能力配置。 +- 增加应用接入号分配。 +- 增加应用号码到三网通道的绑定矩阵。 +- 激活前检查报备状态和三网映射完整性。 +- 人工认领可生成规则建议,但不得自动修改正式绑定。 + +### P2:治理与指标 + +- 接入号精确匹配率、歧义率、未匹配率。 +- 每个号码的下发量、上行量和最后活跃时间。 +- 共享号码积压告警。 +- 非法 `Src_Id`、未报备号码和异常截断告警。 +- 接入号到期、报备失效和停用后的发送阻断。 + +当前数据规模不需要为接入号表提前做分区;优先保证唯一约束、有效期索引、通道和号码组合索引,以及历史快照完整性。 + +## 12. 验收清单 + +后续实现至少覆盖以下真实链路用例: + +1. 客户参数接口只返回平台 CMPP 端点和当前应用号码。 +2. 客户提交未分配 `Src_Id` 时被拒绝,且不计费、不入队。 +3. 客户不填 `Src_Id` 且应用只有一个默认号时自动补全。 +4. 同一客户号码在移动、联通、电信映射为不同上游号码。 +5. 主备通道切换后使用各自绑定的上游 `Src_Id`。 +6. 上游 `MsgSrc` 与通道登录账号不同时仍能正常提交。 +7. 完整扩展号上行能唯一匹配应用。 +8. 运营商只返回基础号或截断号时,严格按该通道规则辅助匹配。 +9. 两个应用共享号码时不得自动误投。 +10. 普通上行和状态报告严格分流。 +11. 长上行完成重组后仍按完整接入号匹配。 +12. 修改号码配置后,历史消息仍按提交快照关联。 +13. 接入号停用、过期或报备失效后阻断新发送,但历史记录仍可查询。 +14. 上行匹配失败时真实写入 PostgreSQL 候选或未匹配记录,并支持人工认领。 +15. Gateway、NestJS、PostgreSQL 和客户 CMPP 连接上的号码值可以相互核对。 + +## 13. 后续落地建议 + +建议后续另开实现批次,按以下顺序完成: + +1. 先取得上游供应商的真实企业代码、业务代码、扩展位数和上行回传规则。 +2. 建立 Prisma 模型、迁移和唯一约束。 +3. 修复客户 CMPP 参数接口和应用企业代码校验。 +4. 改造下发路由,在选中通道后解析号码绑定并保存快照。 +5. 扩展 Gateway 上行事件字段。 +6. 重写普通上行匹配优先级。 +7. 补齐运营端配置页面和人工认领页面。 +8. 使用真实 PostgreSQL、Redis、Gateway 和生产同类 CMPP 测试对端执行完整验收。 diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index 0e30ae9..7f35942 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -4,7 +4,7 @@ 本文基于当前前端设计原型整理,用于交给 Codex 或开发团队执行第一版落地开发。 -当前确认:第一版保留短信业务,排除彩信功能;账户按现金余额计费,人工充值和充值记录进入第一版开发范围,套餐、短信余量、授信额度、账单流水页面和公开交易查询 API 不进入第一版。彩信服务、彩信应用/签名/模板 Tab,以及运营端彩信相关菜单标记为“待开发”;业务性能指标为“平台可稳定入队并调度 500 条短信/秒,实际向通道 submit 受通道限速配置控制”。 +当前确认:第一版保留短信业务,排除彩信功能;账户按现金余额和授信额度计费,人工充值和充值记录进入第一版开发范围,套餐、短信余量、账单流水页面和公开交易查询 API 不进入第一版。彩信服务、彩信应用/签名/模板 Tab,以及运营端彩信相关菜单标记为“待开发”;业务性能指标为“平台可稳定入队并调度 500 条短信/秒,实际向通道 submit 受通道限速配置控制”。 ## 1. 项目目标 @@ -17,7 +17,7 @@ - 支持通道级限速、失败重试、回执同步和发送记录追踪。 - 支持企业、应用、签名、模板、通道、通道组、报备任务等核心配置数据的后台维护。 - 彩信功能仅保留菜单占位或隐藏,不进入第一版开发范围。 -- 账户现金余额计费、人工充值、充值记录进入第一版范围,并与发送记录形成可对账闭环;不提供套餐、短信余量和授信额度。 +- 账户现金余额、授信额度、人工充值和充值记录进入第一版范围,并与发送记录形成可对账闭环;不提供套餐或短信余量。 ## 2. 角色与权限 @@ -329,13 +329,14 @@ 1. 客户端可查看现金余额和充值记录;充值由运营端人工入账。 2. 发送创建时按短信内容计费条数和企业应用客户单价生成预估费用,计费条数只按 70/67 字规则拆分;不按移动、联通、电信配置不同客户价。 -3. 平台发送前只检查企业现金余额,余额大于等于预估费用即可发送。 +3. 平台发送前只判断 `现金余额 + 授信额度 > 0`;授信额度可为正数、负数或 0,和小于等于 0 时禁止发送。 4. 发送链路需记录计费条数、计费单价、计费金额、账务状态。 5. 账单流水与短信记录可追溯关联,支持按企业、应用、任务、手机号、时间对账。 6. 最终失败、超时失败需要退费。 7. 三网通道成本只用于平台内部成本核算,不影响客户扣费金额。 8. 当前版本计费口径固定为提交 accepted 扣费、最终 failed receipt/timeout 退款。 9. 所有面向用户展示的金额、余额、充值金额和单价统一以人民币元展示并固定保留三位小数;内部仍使用分或最小计费单位持久化,不以展示精度改变账务计算。 +10. API 必须定时扫描提交成功但超过 72 小时仍未收到明确最终回执的短信,转为 timeout 并退还已扣金额;扫描需覆盖 `submitted` 和 `unknown`,且用条件更新避免多实例重复退款。 ## 5. 功能需求 @@ -503,7 +504,7 @@ - 充值记录:支持运营人员人工充值和负数冲正。 - 账单流水:支持冻结、扣费、退费、解冻、人工调整、失败返还。 -- 计费规则:支持按短信计费条数和企业应用单价计算费用,发送额度仅取现金余额。 +- 计费规则:支持按短信计费条数和企业应用单价计算费用;发送准入额度为现金余额加授信额度。 - 账务流水必须与短信记录形成可追溯关系,支持对账导出。 - 最终失败、超时失败需要退费。 - 计费口径可配置为按提交成功计费或按回执成功计费。 @@ -1225,7 +1226,7 @@ 验收标准: - 客户端账户余额和充值记录进入第一版。 -- 发送前仅校验企业现金余额。 +- 发送前校验企业现金余额与授信额度之和必须大于 0。 - 每条短信记录可追溯到账务流水。 ### 阶段 6:风控与审核 diff --git a/docs/production-deployment.md b/docs/production-deployment.md index 472d786..f56ec9a 100644 --- a/docs/production-deployment.md +++ b/docs/production-deployment.md @@ -29,6 +29,8 @@ REPO_URL=http://175.27.255.91:3000/hectorzhao/lislgosms.git BRANCH=main PUBLIC_HTTP_PORT=12026 API_PORT=3000 +API_ENABLE_SEND_WORKER=true +API_SEND_WORKER_CONCURRENCY=50 ADMIN_SESSION_IDLE_TIMEOUT_MS=3600000 CLIENT_SESSION_IDLE_TIMEOUT_MS=7200000 SESSION_LOCK_RECOVERY_MS=14400000 @@ -40,6 +42,9 @@ OPERATION_LOG_RETENTION_DAYS=180 OPERATION_LOG_ARCHIVE_BATCH_SIZE=1000 OPERATION_LOG_ARCHIVE_MAX_BATCHES=20 OPERATION_LOG_ARCHIVE_INTERVAL_MS=86400000 +SMS_RECEIPT_TIMEOUT_SCAN_ENABLED=true +SMS_RECEIPT_TIMEOUT_HOURS=72 +SMS_RECEIPT_TIMEOUT_SCAN_INTERVAL_MS=300000 CMPP_DOWNSTREAM_ACK_TIMEOUT_SECONDS=30 GATEWAY_CMPP_ADDR=0.0.0.0:17890 OBJECT_STORAGE_DRIVER=minio @@ -55,6 +60,8 @@ PROD_ADMIN_PASSWORD='change-me' 脚本会安装 Node.js、Go、PostgreSQL、Redis、MinIO、Nginx,创建 systemd 服务,执行 Prisma migrate,构建前端/API/Gateway,并创建平台管理员。Node.js、Go 和 MinIO 下载会按服务器架构自动选择 x64/amd64 或 arm64。 +`API_ENABLE_SEND_WORKER=true` 是生产发送链路必填项。后续发布脚本会在构建和迁移前校验该开关以及正整数 `API_SEND_WORKER_CONCURRENCY`;缺失时直接终止发布,防止 API/Gateway 健康但 BullMQ 短信队列无人消费。 + 如生产验证服务器临时无法稳定下载 MinIO,可显式传入 `OBJECT_STORAGE_DRIVER=local`,文件会通过真实 API 保存到服务器本地目录 `OBJECT_STORAGE_LOCAL_ROOT`,`cmpp-minio` 服务会跳过安装和启动。该模式只建议用于验证环境;正式生产建议恢复 `OBJECT_STORAGE_DRIVER=minio`。 ## 后续发布 @@ -91,6 +98,7 @@ curl http://127.0.0.1:8090/health curl http://127.0.0.1:12026/ redis-cli -h 127.0.0.1 -p 6379 ping pg_isready -d "$(grep '^DATABASE_URL=' /etc/cmpp-platform/cmpp-platform.env | cut -d= -f2-)" +grep -E '^(API_ENABLE_SEND_WORKER|API_SEND_WORKER_CONCURRENCY)=' /etc/cmpp-platform/cmpp-platform.env ``` ## 回滚 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index 08c5312..a3bd279 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -21,7 +21,7 @@ | 模板 | 验证码模板 `验证码为 ${code}`,营销模板 `尊敬的${name},优惠活动开始`,分别准备草稿、待审核、通过、驳回。 | | 通道 | active CMPP 通道、disabled 通道、备用通道;通道组包含主备优先级。 | | 号码 | 合法号码、重复号码、非法号码、企业黑名单号码、全局黑名单号码。 | -| 账户 | 现金余额充足、现金余额不足;不配置套餐余量或授信额度。 | +| 账户 | 现金余额与授信额度组合后的和为正数、0、负数;授信额度覆盖正数、负数和 0;不配置套餐余量。 | | 企业认证 | 未认证、待审核、已通过、已驳回四类企业认证资料。 | | 客户 | 正常客户、停用客户、欠费客户、未认证客户、跨租户客户、客户联系人和开票资料。 | | 导入文件 | UTF-8 CSV、GBK CSV、TXT、超 20 MB 文件、含空行/重复/非法号码/非法字符文件。 | @@ -805,6 +805,19 @@ - jobId 使用 messageRecordId,避免重复入队。 - 批量任务状态变为 queued。 +### TC-SEND-002A 生产发送 Worker 配置门禁 + +- 优先级:P0 +- 前置条件:使用生产环境文件执行发布脚本。 +- 步骤: + 1. 删除或关闭 `API_ENABLE_SEND_WORKER`,执行发布脚本。 + 2. 配置 `API_ENABLE_SEND_WORKER=true` 和正整数 `API_SEND_WORKER_CONCURRENCY`,重新发布。 + 3. 创建一条 queued 短信并观察 Redis BullMQ 与数据库提交记录。 +- 预期结果: + - 步骤 1 在构建、迁移和服务重启前终止,明确提示发送 Worker 未启用。 + - 步骤 2 发布成功,API 启动发送 Worker。 + - 步骤 3 的 job 不长时停留在 `prioritized/wait`,且真实生成 `SmsSubmitRecord` 并进入 Gateway。 + ### TC-SEND-003 通道路由和 Gateway SubmitCommand - 优先级:P0 @@ -3205,12 +3218,13 @@ npm run verify:phase8 | 用例 | 细化执行点 | 必查断言 | | --- | --- | --- | | TC-BILLING-006 | 运营端人工充值现金金额,客户端查看 Dashboard 和充值记录。 | RechargeOrder 状态为 paid/manual_topup;TenantAccount 现金余额同步增加;AccountTransaction 类型 recharge;充值记录“充值后余额”必须等于该订单关联 AccountTransaction.balanceAfter,不能用当前账户余额替代;运营日志可追溯。 | -| TC-BILLING-007 | 分别填写正数金额和负数金额执行充值、冲正;提交 0 或非法金额。 | 正负金额方向正确;0 和非法金额被拒绝;数据模型和接口不存在短信条数、套餐或授信字段。 | +| TC-BILLING-007 | 分别填写正数金额和负数金额执行充值、冲正;提交 0 或非法金额;将授信额度配置为正数、负数和 0。 | 充值正负金额方向正确;0 和非法充值金额被拒绝;三种授信值都能保存并写操作日志;数据模型和接口不存在短信套餐条数字段。 | | TC-BILLING-008 | 对已充值记录执行撤销/冲正,分别覆盖未消费和已部分消费。 | 未消费可全额回退;已消费按规则拒绝或生成人工调整;原订单状态和反向流水清晰;日志记录原因。 | | TC-BILLING-009 | 无权限用户、审核员、管理员分别执行充值;大额人工充值不走审批。 | 权限不足被拒绝并写失败日志;有权限用户确认后立即入账;不产生 pending 审批态;充值订单、账户余额、流水和日志同步完成。 | | TC-BILLING-010 | 余额不足发送失败,人工充值后重试发送并模拟 delivered。 | 充值前不扣费;充值后发送成功;冻结、扣费、短信计费记录完整;reconciliation diff 为 0。 | -| TC-BILLING-011 | 账户现金余额 100 分,预估费用 5 分,且数据库不存在套餐和授信数据;再将余额改为 4 分重试。 | 100 分时可发送,4 分时提示“企业账户余额不足”;判断只依赖 TenantAccount.balanceCents。 | +| TC-BILLING-011 | 分别准备 `余额+授信` 为正数、0 和负数的账户,使用相同短信费用发起发送。 | 和为正数时允许发送;和为 0 或负数时提示余额不足。判断公式为 `balanceCents + creditCents > 0`,与本次费用和套餐无关。 | | TC-BILLING-012 | 已扣费短信收到最终失败回执;另一个未提交成功任务只释放冻结。 | 前者只生成一条 refunded 流水并计入“今日返还”;重复回执不重复退款;后者的 released 流水不计入“今日返还”。 | +| TC-BILLING-013 | 准备已提交扣费但 72 小时完全无回执的 `submitted` 短信,以及有 `UNKNOWN` 回执且超过 72 小时的短信;启动 API 定时扫描并模拟重复扫描。 | 两类短信都转为 timeout 并退款;任务进度刷新;同一短信只退款一次;定时扫描默认启用且每 5 分钟执行。 | ### 17.6 系统日志细化 diff --git a/docs/testing-progress.md b/docs/testing-progress.md index f668915..f06139e 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -1,5 +1,21 @@ # 第一版系统化测试进度 +## 2026-07-14 生产发送 Worker 配置缺失修复(待部署验收) + +- 生产号码 `18821203795` 的最新短信于 18:24:59 审核通过后恢复为 queued,BullMQ 已生成 job,但一直停留在 `bull:sms.send.queue:prioritized`,无通道、`submitId` 和 `SmsSubmitRecord`。 +- 生产 API 进程环境缺少 `API_ENABLE_SEND_WORKER=true`,而 `SendChainService.onModuleInit()` 只在该值严格为 `true` 时启动 BullMQ Worker;Gateway 健康,但任务尚未进入 Gateway。 +- 修复为生产初始化默认写入 `API_ENABLE_SEND_WORKER=true` 和并发数 50,发布脚本在构建、迁移和重启前强制校验开关与正整数并发数,避免再次带病发布。 +- 与本批工作区改动合并回归:Prisma validate/generate 通过,API 完整 17 suites、163 项通过,API build、前端 build 和 Gateway 全量 Go 测试通过;前端仅有既有 Vite chunk size warning。 + +## 2026-07-14 恢复授信额度并补齐 72 小时无回执退款(待部署验收) + +- 产品口径修正:套餐和短信余量继续移除,但企业账户保留授信额度;授信额度可为正数、负数或 0。 +- 发送校验公式调整为 `balanceCents + creditCents > 0`,和大于 0 才允许发送,和小于等于 0 时禁止发送;不扣除本次预估费用后再判断。 +- 新增授信调整真实 API、操作日志、企业新建/编辑页输入和运营端企业列表展示,并新增后续迁移恢复 `TenantAccount.creditCents`;客户端账户页不展示授信额度字段,只展示现金余额和合并后的可用发送额度。 +- 退款时点审查:明确失败回执在补发不可用或耗尽后,于同一次回执处理内同步退款;Submit 拒绝/超时发生在扣费前,只释放冻结而非退款。已修复原 72 小时接口无自动调度且漏掉 `submitted` 的缺口:API 启动 60 秒后首次扫描,之后默认每 5 分钟扫描 `submitted/unknown`,超过 72 小时即转 timeout 并退款,单条条件更新避免并发扫描重复处理。 +- CMPP 2.0/3.0 的 `Stat` 定义包含标准值 `UNKNOWN`,不是平台自造状态;当前 Gateway 对 `DELIVRD` 以外的非空最终状态(包括 `UNKNOWN`)统一映射为 `undelivered`,因此会进入失败补发,无法补发时直接失败退款。只有空 `Stat` 才映射为平台 `unknown`,并由 72 小时扫描兜底。 +- 本地验证通过:新增迁移已应用于本地 PostgreSQL,Prisma validate/generate、完整 API 17 suites/163 项、API build、前端 build 和 Go Gateway 全量测试通过;前端仅保留既有 Vite chunk size warning。 + ## 2026-07-14 计费收敛为现金余额、失败退款展示 - 生产只读核查确认 17:01 的人工充值已将目标企业现金余额从 0 增加到 100 分;17:04 发送预估费用仅 5 分却被拒绝,根因是旧逻辑同时要求 `TenantAccount.smsUnits >= billingUnits`,充值只增加现金时短信余量仍为 0。 diff --git a/gateway/internal/upstream/deliver_test.go b/gateway/internal/upstream/deliver_test.go index b13d44f..ce0a4ba 100644 --- a/gateway/internal/upstream/deliver_test.go +++ b/gateway/internal/upstream/deliver_test.go @@ -86,3 +86,14 @@ func TestHandleCMPP2DeliverReceiptPostsReceiptEvent(t *testing.T) { t.Fatal("timed out waiting for receipt event") } } + +func TestReceiptStatusTreatsNonDeliveredFinalStatesAsUndelivered(t *testing.T) { + for _, stat := range []string{"UNKNOWN", "UNDELIV", "EXPIRED", "DELETED", "REJECTD"} { + if got := receiptStatus(stat); got != "undelivered" { + t.Fatalf("receiptStatus(%q) = %q, want undelivered", stat, got) + } + } + if got := receiptStatus(""); got != "unknown" { + t.Fatalf("receiptStatus(empty) = %q, want unknown", got) + } +} diff --git a/src/api/adminApi.ts b/src/api/adminApi.ts index 7c419b8..eb96ee9 100644 --- a/src/api/adminApi.ts +++ b/src/api/adminApi.ts @@ -255,7 +255,7 @@ export type DashboardResponse = { recentFailed: number; alertCount: number; }; - accounts: Array<{ id: string; tenantId: string; balanceCents: number; status: string; tenant?: TenantOption }>; + accounts: Array<{ id: string; tenantId: string; balanceCents: number; creditCents: number; status: string; tenant?: TenantOption }>; recentTasks: Array>; recentRecharges: Array; }; @@ -661,6 +661,7 @@ export type TenantAccount = { id: string; tenantId: string; balanceCents: number; + creditCents: number; status: string; tenant?: TenantOption; }; @@ -953,6 +954,8 @@ export const adminApi = { listSystemLogs: (query: { tenantId?: string; keyword?: string; level?: string; module?: string; range?: string; page?: number; pageSize?: number }) => request(withQuery('/admin/system-logs', query)), listAccounts: () => request('/admin/billing/accounts'), + updateCreditLimit: (tenantId: string, body: { creditCents: number; operatorId?: string; remark?: string }) => + request(`/admin/billing/accounts/${tenantId}/credit-limit`, { method: 'POST', body: JSON.stringify(body) }), listManualRecharges: (tenantId?: string) => request(withQuery('/admin/billing/manual-recharges', { tenantId })), createManualRecharge: (body: { tenantId: string; amountCents: number; operatorId?: string; remark?: string }) => request('/admin/billing/manual-recharges', { method: 'POST', body: JSON.stringify(body) }), diff --git a/src/apps/admin/AdminCustomerDetailPage.tsx b/src/apps/admin/AdminCustomerDetailPage.tsx index 37f689c..1b942cb 100644 --- a/src/apps/admin/AdminCustomerDetailPage.tsx +++ b/src/apps/admin/AdminCustomerDetailPage.tsx @@ -69,6 +69,7 @@ export function AdminCustomerDetailPage() {
企业编码{tenant?.code ?? '-'}{tenant?.status ?? '-'}
计费方式按量计费仅从现金余额扣费
现金余额¥{formatCents(account?.balanceCents)}真实账户余额
+
授信额度¥{formatCents(account?.creditCents)}可配置为正数、负数或 0
diff --git a/src/apps/admin/AdminCustomerFormPage.tsx b/src/apps/admin/AdminCustomerFormPage.tsx index 4b7c1c8..f41995b 100644 --- a/src/apps/admin/AdminCustomerFormPage.tsx +++ b/src/apps/admin/AdminCustomerFormPage.tsx @@ -7,6 +7,7 @@ import { Breadcrumb, Button, FileActions, Input, Select, Textarea } from '@/comp type EnterpriseForm = { name: string; creditCode: string; + creditLimit: string; province: string; city: string; address: string; @@ -44,6 +45,7 @@ const cityOptionsByProvince: Record { - setForm(formFromTenant(tenant)); + Promise.all([adminApi.getTenant(enterpriseId), adminApi.listAccounts()]) + .then(([tenant, accounts]) => { + setForm(formFromTenant(tenant, accounts.find((account) => account.tenantId === enterpriseId)?.creditCents ?? 0)); setError(''); }) .catch((failure: Error) => setError(failure.message || '企业信息加载失败')); @@ -117,21 +120,30 @@ export function AdminCustomerFormPage() { if (!form.creditCode.trim()) nextErrors.creditCode = '请填写统一社会信用代码'; if (!form.contactName.trim()) nextErrors.contactName = '请填写联系人姓名'; if (!form.contactPhone.trim()) nextErrors.contactPhone = '请填写手机号'; + const creditLimit = Number(form.creditLimit); + if (!Number.isFinite(creditLimit)) nextErrors.creditLimit = '请填写有效的授信额度'; setErrors(nextErrors); return Object.keys(nextErrors).length === 0; } - function submitForm() { + async function submitForm() { if (!validateForm()) return; setSaving(true); - const { photoContentType, photoFileName, ...payload } = form; - const request = isEdit && enterpriseId - ? adminApi.updateTenant(enterpriseId, payload) - : adminApi.createTenant(payload); - request - .then(() => navigate('/admin/customers')) - .catch((failure: Error) => setError(failure.message || '企业保存失败')) - .finally(() => setSaving(false)); + const { creditLimit, photoContentType, photoFileName, ...payload } = form; + try { + const tenant = isEdit && enterpriseId + ? await adminApi.updateTenant(enterpriseId, payload) + : await adminApi.createTenant(payload); + await adminApi.updateCreditLimit(tenant.id, { + creditCents: Math.round(Number(creditLimit) * 100), + remark: isEdit ? '企业编辑页调整授信额度' : '创建企业初始化授信额度', + }); + navigate('/admin/customers'); + } catch (failure) { + setError(failure instanceof Error ? failure.message : '企业保存失败'); + } finally { + setSaving(false); + } } function uploadEnterprisePhoto(file: File | undefined) { @@ -204,6 +216,18 @@ export function AdminCustomerFormPage() { />
+
+ updateForm('creditLimit', event.target.value)} + required + type="number" + value={form.creditLimit} + /> +
+
updateForm('city', event.target.value)} options={cityOptions} value={form.city} /> diff --git a/src/apps/admin/AdminCustomersPage.tsx b/src/apps/admin/AdminCustomersPage.tsx index 314b610..af27a1b 100644 --- a/src/apps/admin/AdminCustomersPage.tsx +++ b/src/apps/admin/AdminCustomersPage.tsx @@ -86,6 +86,7 @@ export function AdminCustomersPage({ basePath = '/admin/customers' }: AdminCusto ); }, }, + { key: 'creditLimit', title: '授信额度', width: '150px', align: 'right', render: (record) => `¥${formatCents(record.account?.creditCents ?? 0)}` }, { key: 'todaySpend', title: '今日消费', width: '150px', align: 'right', render: (record) => `¥${formatCents(record.todaySpendCents)}` }, { key: 'todayRefund', title: '今日返还', width: '150px', align: 'right', render: (record) => `¥${formatCents(record.todayRefundCents)}` }, { key: 'status', title: '企业状态', width: '130px', render: (record) => {record.status === 'active' ? '正常' : '已禁用'} }, diff --git a/src/apps/admin/AdminHome.tsx b/src/apps/admin/AdminHome.tsx index eab9564..de5d5d3 100644 --- a/src/apps/admin/AdminHome.tsx +++ b/src/apps/admin/AdminHome.tsx @@ -69,7 +69,7 @@ export function AdminHome() { const todaySpend = Math.abs(dashboard?.recentRecharges .filter((item) => item.tenantId === account.tenantId) .reduce((sum, item) => sum + item.amountCents, 0) ?? 0) / 100; - const availableBalance = account.balanceCents / 100; + const availableBalance = (account.balanceCents + account.creditCents) / 100; return { id: account.tenantId, enterprise: account.tenant?.name ?? account.tenantId, diff --git a/src/apps/client/ClientBillingPage.tsx b/src/apps/client/ClientBillingPage.tsx index fdf6f92..4cb75fc 100644 --- a/src/apps/client/ClientBillingPage.tsx +++ b/src/apps/client/ClientBillingPage.tsx @@ -7,6 +7,7 @@ import { formatDateTime } from '@/utils/dateTime'; export function ClientBillingPage() { const [balanceCents, setBalanceCents] = useState(0); + const [creditCents, setCreditCents] = useState(0); const [orders, setOrders] = useState([]); const [page, setPage] = useState(1); const [loading, setLoading] = useState(true); @@ -21,6 +22,7 @@ export function ClientBillingPage() { Promise.all([clientApi.getDashboard(), clientApi.listOrders()]) .then(([dashboard, nextOrders]) => { setBalanceCents(dashboard.accounts[0]?.balanceCents ?? 0); + setCreditCents(dashboard.accounts[0]?.creditCents ?? 0); setOrders(nextOrders); setError(''); }) @@ -36,7 +38,11 @@ export function ClientBillingPage() {
-
当前可用余额¥{formatCents(balanceCents)}短信发送仅按现金余额校验。
+
现金余额¥{formatCents(balanceCents)}充值和扣费后的账面余额。
+
+
+ +
可用发送额度¥{formatCents(balanceCents + creditCents)}大于 0 时可以发送短信。
diff --git a/src/apps/client/ClientHome.tsx b/src/apps/client/ClientHome.tsx index 6b28b48..d231460 100644 --- a/src/apps/client/ClientHome.tsx +++ b/src/apps/client/ClientHome.tsx @@ -50,7 +50,7 @@ export function ClientHome() { }, []); const account = dashboard?.accounts[0]; - const availableBalance = (account?.balanceCents ?? 0) / 100; + const availableBalance = ((account?.balanceCents ?? 0) + (account?.creditCents ?? 0)) / 100; const todaySpend = (dashboard?.today.spendCents ?? 0) / 100; const todayRefund = Math.abs((dashboard?.transactions._sum.amountCents ?? 0) < 0 ? 0 : dashboard?.transactions._sum.amountCents ?? 0) / 100; const balanceBaseline = Math.max(availableBalance + todaySpend - todayRefund, availableBalance, 1); @@ -102,9 +102,9 @@ export function ClientHome() {
- 账户剩余余额 + 可用发送额度 ¥{formatAmount(availableBalance)} - 今日消费 ¥{formatAmount(todaySpend)} + 现金余额 ¥{formatCents(account?.balanceCents)}
今日发送 diff --git a/tools/deploy/production-bootstrap.sh b/tools/deploy/production-bootstrap.sh index 6f30a1c..f816dca 100644 --- a/tools/deploy/production-bootstrap.sh +++ b/tools/deploy/production-bootstrap.sh @@ -6,6 +6,8 @@ REPO_URL="${REPO_URL:-http://175.27.255.91:3000/hectorzhao/lislgosms.git}" BRANCH="${BRANCH:-main}" PUBLIC_HTTP_PORT="${PUBLIC_HTTP_PORT:-12026}" API_PORT="${API_PORT:-3000}" +API_ENABLE_SEND_WORKER="${API_ENABLE_SEND_WORKER:-true}" +API_SEND_WORKER_CONCURRENCY="${API_SEND_WORKER_CONCURRENCY:-50}" GATEWAY_CONTROL_ADDR="${GATEWAY_CONTROL_ADDR:-127.0.0.1:8090}" GATEWAY_CMPP_ADDR="${GATEWAY_CMPP_ADDR:-0.0.0.0:17890}" DB_NAME="${DB_NAME:-cmpp_platform}" @@ -24,6 +26,9 @@ OPERATION_LOG_RETENTION_DAYS="${OPERATION_LOG_RETENTION_DAYS:-180}" OPERATION_LOG_ARCHIVE_BATCH_SIZE="${OPERATION_LOG_ARCHIVE_BATCH_SIZE:-1000}" OPERATION_LOG_ARCHIVE_MAX_BATCHES="${OPERATION_LOG_ARCHIVE_MAX_BATCHES:-20}" OPERATION_LOG_ARCHIVE_INTERVAL_MS="${OPERATION_LOG_ARCHIVE_INTERVAL_MS:-86400000}" +SMS_RECEIPT_TIMEOUT_SCAN_ENABLED="${SMS_RECEIPT_TIMEOUT_SCAN_ENABLED:-true}" +SMS_RECEIPT_TIMEOUT_HOURS="${SMS_RECEIPT_TIMEOUT_HOURS:-72}" +SMS_RECEIPT_TIMEOUT_SCAN_INTERVAL_MS="${SMS_RECEIPT_TIMEOUT_SCAN_INTERVAL_MS:-300000}" if [[ "$(id -u)" -ne 0 ]]; then echo "Run as root." >&2 @@ -172,6 +177,8 @@ write_env() { cat >/etc/cmpp-platform/cmpp-platform.env <&2 + exit 1 +fi + +if [[ ! "${API_SEND_WORKER_CONCURRENCY:-}" =~ ^[1-9][0-9]*$ ]]; then + echo "API_SEND_WORKER_CONCURRENCY must be a positive integer in $ENV_FILE." >&2 + exit 1 +fi + cd "$APP_DIR" echo "[deploy] Installing dependencies"