diff --git a/api/prisma/migrations/20260712113000_add_cmpp_review_aggregation/migration.sql b/api/prisma/migrations/20260712113000_add_cmpp_review_aggregation/migration.sql new file mode 100644 index 0000000..7c304d0 --- /dev/null +++ b/api/prisma/migrations/20260712113000_add_cmpp_review_aggregation/migration.sql @@ -0,0 +1,21 @@ +ALTER TABLE "SmsSendTask" +ADD COLUMN "sourceType" TEXT NOT NULL DEFAULT 'risk', +ADD COLUMN "aggregationKey" TEXT, +ADD COLUMN "contentHash" TEXT, +ADD COLUMN "windowStartedAt" TIMESTAMP(3), +ADD COLUMN "windowEndsAt" TIMESTAMP(3); + +ALTER TABLE "SmsMessageRecord" +ADD COLUMN "signatureId" TEXT, +ADD COLUMN "reviewTaskId" TEXT; + +CREATE UNIQUE INDEX "SmsSendTask_aggregationKey_key" ON "SmsSendTask"("aggregationKey"); +CREATE INDEX "SmsMessageRecord_reviewTaskId_status_idx" ON "SmsMessageRecord"("reviewTaskId", "status"); + +ALTER TABLE "SmsMessageRecord" +ADD CONSTRAINT "SmsMessageRecord_signatureId_fkey" +FOREIGN KEY ("signatureId") REFERENCES "SmsSignature"("id") ON DELETE SET NULL ON UPDATE CASCADE; + +ALTER TABLE "SmsMessageRecord" +ADD CONSTRAINT "SmsMessageRecord_reviewTaskId_fkey" +FOREIGN KEY ("reviewTaskId") REFERENCES "SmsSendTask"("id") ON DELETE SET NULL ON UPDATE CASCADE; diff --git a/api/prisma/schema.prisma b/api/prisma/schema.prisma index 596ea73..416bbb1 100644 --- a/api/prisma/schema.prisma +++ b/api/prisma/schema.prisma @@ -411,6 +411,7 @@ model SmsSignature { templates SmsTemplate[] reportMaterials SignatureReportMaterial[] reportTasks ChannelSignatureReportTask[] + messageRecords SmsMessageRecord[] @@index([tenantId, auditStatus]) @@index([tenantId, reportStatus]) @@ -782,6 +783,11 @@ model SmsSendTask { applicationId String? templateId String? taskNo String @unique + sourceType String @default("risk") + aggregationKey String? @unique + contentHash String? + windowStartedAt DateTime? + windowEndsAt DateTime? content String category String? phoneTotal Int @@ -806,6 +812,7 @@ model SmsSendTask { createdBy User? @relation("SmsSendTaskCreator", fields: [createdById], references: [id]) reviewedBy User? @relation("SmsSendTaskReviewer", fields: [reviewedById], references: [id]) riskHits RiskHitRecord[] + messageRecords SmsMessageRecord[] @@index([tenantId, status, createdAt]) @@index([applicationId, createdAt]) @@ -899,6 +906,8 @@ model SmsMessageRecord { batchTaskId String? applicationId String? templateId String? + signatureId String? + reviewTaskId String? messageId String @unique phoneNumber String carrier String? @@ -926,6 +935,8 @@ model SmsMessageRecord { batchTask SmsBatchTask? @relation(fields: [batchTaskId], references: [id], onDelete: Cascade) application SmsApplication? @relation(fields: [applicationId], references: [id]) template SmsTemplate? @relation(fields: [templateId], references: [id]) + signature SmsSignature? @relation(fields: [signatureId], references: [id]) + reviewTask SmsSendTask? @relation(fields: [reviewTaskId], references: [id]) channel SmsChannel? @relation(fields: [channelId], references: [id]) submitRecords SmsSubmitRecord[] receiptRecords SmsReceiptRecord[] @@ -936,6 +947,7 @@ model SmsMessageRecord { @@index([tenantId, status, queuedAt]) @@index([batchTaskId, status]) + @@index([reviewTaskId, status]) @@index([phoneNumber]) @@index([gatewayMessageId]) } diff --git a/api/src/operations/operations.service.spec.ts b/api/src/operations/operations.service.spec.ts index 8366048..00f43d4 100644 --- a/api/src/operations/operations.service.spec.ts +++ b/api/src/operations/operations.service.spec.ts @@ -166,6 +166,17 @@ function createPrismaMock() { } describe('OperationsService', () => { + it('keeps task progress limited to client-created batch tasks', async () => { + const prisma = createPrismaMock(); + const service = new OperationsService(prisma as never); + + await service.listBatchTasks({ tenantId: 'tenant-1', status: 'queued' }); + + expect(prisma.smsBatchTask.findMany).toHaveBeenCalledWith(expect.objectContaining({ + where: { tenantId: 'tenant-1', status: 'queued', sourceType: 'client' }, + })); + }); + it('filters send-chain messages by tenant, application, channel, content, date, task, phone, and status', async () => { const prisma = createPrismaMock(); const service = new OperationsService(prisma as never); diff --git a/api/src/operations/operations.service.ts b/api/src/operations/operations.service.ts index 7cc4d6c..1300de4 100644 --- a/api/src/operations/operations.service.ts +++ b/api/src/operations/operations.service.ts @@ -78,7 +78,7 @@ export class OperationsService { listBatchTasks(query: { tenantId?: string; status?: string }) { return this.prisma.smsBatchTask.findMany({ - where: { tenantId: query.tenantId, status: query.status }, + where: { tenantId: query.tenantId, status: query.status, sourceType: 'client' }, include: { apiRequests: true }, orderBy: { createdAt: 'desc' }, }); diff --git a/api/src/risk-review/admin-risk-review.controller.ts b/api/src/risk-review/admin-risk-review.controller.ts index 9aac1b3..301219f 100644 --- a/api/src/risk-review/admin-risk-review.controller.ts +++ b/api/src/risk-review/admin-risk-review.controller.ts @@ -1,5 +1,6 @@ import { Body, Controller, Get, Param, Post, Query } from '@nestjs/common'; import { ApiTags } from '@nestjs/swagger'; +import { SendChainService } from '../send-chain/send-chain.service'; import { CreateRiskRuleDto, ReviewSmsTaskDto, @@ -9,7 +10,7 @@ import { @ApiTags('risk-review') @Controller('admin/risk-review') export class AdminRiskReviewController { - constructor(private readonly riskReview: RiskReviewService) {} + constructor(private readonly riskReview: RiskReviewService, private readonly sendChain: SendChainService) {} @Get('rules') listRules(@Query('tenantId') tenantId?: string) { @@ -37,13 +38,16 @@ export class AdminRiskReviewController { } @Post('tasks/:id/approve') - approveTask(@Param('id') taskId: string, @Body() body: ReviewSmsTaskDto) { - return this.riskReview.approveTask(taskId, body); + async approveTask(@Param('id') taskId: string, @Body() body: ReviewSmsTaskDto) { + const task = await this.riskReview.approveTask(taskId, body); + await this.sendChain.handleReviewDecision(taskId, 'approved', body.reason ?? '运营审核通过'); + return task; } @Post('tasks/:id/reject') - rejectTask(@Param('id') taskId: string, @Body() body: ReviewSmsTaskDto) { - return this.riskReview.rejectTask(taskId, body); + async rejectTask(@Param('id') taskId: string, @Body() body: ReviewSmsTaskDto) { + const task = await this.riskReview.rejectTask(taskId, body); + await this.sendChain.handleReviewDecision(taskId, 'rejected', body.reason ?? task.rejectReason ?? '运营审核驳回'); + return task; } } - diff --git a/api/src/risk-review/risk-review.module.ts b/api/src/risk-review/risk-review.module.ts index 58ba9f8..462166d 100644 --- a/api/src/risk-review/risk-review.module.ts +++ b/api/src/risk-review/risk-review.module.ts @@ -1,14 +1,14 @@ -import { Module } from '@nestjs/common'; +import { forwardRef, Module } from '@nestjs/common'; import { PrismaModule } from '../prisma/prisma.module'; import { AdminRiskReviewController } from './admin-risk-review.controller'; import { ClientRiskReviewController } from './client-risk-review.controller'; import { RiskReviewService } from './risk-review.service'; +import { SendChainModule } from '../send-chain/send-chain.module'; @Module({ - imports: [PrismaModule], + imports: [PrismaModule, forwardRef(() => SendChainModule)], controllers: [AdminRiskReviewController, ClientRiskReviewController], providers: [RiskReviewService], exports: [RiskReviewService], }) export class RiskReviewModule {} - diff --git a/api/src/risk-review/risk-review.service.spec.ts b/api/src/risk-review/risk-review.service.spec.ts index 10c3e30..d33a7e6 100644 --- a/api/src/risk-review/risk-review.service.spec.ts +++ b/api/src/risk-review/risk-review.service.spec.ts @@ -33,16 +33,65 @@ function createPrismaMock(overrides: Record = {}) { findUnique: jest.fn().mockResolvedValue({ id: 'risk-task-1', riskHits: [] }), update: jest.fn(), findMany: jest.fn(), + upsert: jest.fn().mockImplementation(({ create }: { create: Record }) => Promise.resolve({ id: 'review-task-1', ...create })), + }, + smsMessageRecord: { + update: jest.fn().mockResolvedValue({ id: 'message-1' }), }, riskHitRecord: { createMany: jest.fn(), findMany: jest.fn(), }, + $transaction: jest.fn(async (callback) => callback({ + smsSendTask: { + upsert: jest.fn().mockImplementation(({ create }: { create: Record }) => Promise.resolve({ id: 'review-task-1', ...create })), + }, + smsMessageRecord: { + update: jest.fn().mockResolvedValue({ id: 'message-1' }), + }, + })), ...overrides, }; } describe('RiskReviewService', () => { + it('groups identical CMPP template mismatches into a deterministic short review window', async () => { + const prisma = createPrismaMock(); + prisma.smsSendTask.findUnique.mockResolvedValue({ + id: 'review-task-1', + sourceType: 'cmpp_template_mismatch', + phoneTotal: 2, + _count: { messageRecords: 2 }, + riskHits: [], + }); + const service = new RiskReviewService(prisma as never); + + const result = await service.aggregateTemplateMismatch({ + tenantId: 'tenant-1', + applicationId: 'app-1', + account: '100001', + messageRecordId: 'message-1', + signatureId: 'sig-1', + content: '【签名】同一审核内容', + }); + + expect(prisma.$transaction).toHaveBeenCalledTimes(1); + expect(result).toEqual(expect.objectContaining({ sourceType: 'cmpp_template_mismatch', phoneTotal: 2 })); + }); + + it('does not allow an aggregated review task to be decided before its window closes', async () => { + const prisma = createPrismaMock(); + prisma.smsSendTask.findUnique.mockResolvedValue({ + id: 'review-task-1', + sourceType: 'cmpp_template_mismatch', + windowEndsAt: new Date(Date.now() + 10_000), + }); + const service = new RiskReviewService(prisma as never); + + await expect(service.approveTask('review-task-1', { reason: '通过' })).rejects.toThrow('聚合窗口尚未关闭'); + expect(prisma.smsSendTask.update).not.toHaveBeenCalled(); + }); + it('rejects tasks over the application max phone threshold', async () => { const prisma = createPrismaMock(); prisma.smsApplication.findUnique.mockResolvedValue({ id: 'app-1', maxPhonesPerTask: 2 }); diff --git a/api/src/risk-review/risk-review.service.ts b/api/src/risk-review/risk-review.service.ts index 87dc5ab..675d29d 100644 --- a/api/src/risk-review/risk-review.service.ts +++ b/api/src/risk-review/risk-review.service.ts @@ -1,6 +1,6 @@ import { BadRequestException, Injectable, NotFoundException } from '@nestjs/common'; import { Prisma } from '@prisma/client'; -import { randomUUID } from 'node:crypto'; +import { createHash, randomUUID } from 'node:crypto'; import { PrismaService } from '../prisma/prisma.service'; export interface CreateRiskRuleDto { @@ -33,6 +33,15 @@ export interface ReviewSmsTaskDto { reason?: string; } +export interface AggregateTemplateMismatchDto { + tenantId: string; + applicationId: string; + account: string; + messageRecordId: string; + signatureId: string; + content: string; +} + interface RuleEvaluation { ruleId?: string; ruleCode: string; @@ -153,8 +162,14 @@ export class RiskReviewService { where: { tenantId, status, + ...(status === 'pending_review' ? { + OR: [ + { sourceType: { not: 'cmpp_template_mismatch' } }, + { windowEndsAt: { lte: new Date() } }, + ], + } : {}), }, - include: { riskHits: true }, + include: { riskHits: true, _count: { select: { messageRecords: true } } }, orderBy: { createdAt: 'desc' }, }); } @@ -163,6 +178,61 @@ export class RiskReviewService { return this.listTasks(undefined, 'pending_review'); } + async aggregateTemplateMismatch(data: AggregateTemplateMismatchDto) { + const normalizedContent = data.content.replace(/\r\n/g, '\n').trim(); + const contentHash = createHash('sha256').update(normalizedContent, 'utf8').digest('hex'); + const windowMs = Math.max(1_000, Number(process.env.CMPP_TEMPLATE_REVIEW_WINDOW_MS ?? 10_000)); + const now = new Date(); + const windowStartedAt = new Date(Math.floor(now.getTime() / windowMs) * windowMs); + const windowEndsAt = new Date(windowStartedAt.getTime() + windowMs); + const aggregationKey = createHash('sha256') + .update(`${data.applicationId}|${data.account}|${contentHash}|${windowStartedAt.toISOString()}`, 'utf8') + .digest('hex'); + const task = await this.prisma.$transaction(async (tx) => { + const aggregated = await tx.smsSendTask.upsert({ + where: { aggregationKey }, + update: { + phoneTotal: { increment: 1 }, + uniquePhoneTotal: { increment: 1 }, + }, + create: { + tenantId: data.tenantId, + applicationId: data.applicationId, + taskNo: `CMPP-REVIEW-${windowStartedAt.getTime()}-${aggregationKey.slice(0, 8)}`, + sourceType: 'cmpp_template_mismatch', + aggregationKey, + contentHash, + windowStartedAt, + windowEndsAt, + content: normalizedContent, + phoneTotal: 1, + uniquePhoneTotal: 1, + status: 'pending_review', + riskDecision: 'manual_review', + reviewReason: '企业应用已配置模板不匹配进入人工审核', + variableIssues: { + sourceType: 'cmpp', + account: data.account, + aggregationWindowMs: windowMs, + } as Prisma.InputJsonValue, + }, + }); + await tx.smsMessageRecord.update({ + where: { id: data.messageRecordId }, + data: { + reviewTaskId: aggregated.id, + signatureId: data.signatureId, + status: 'pending_review', + }, + }); + return aggregated; + }); + return this.prisma.smsSendTask.findUnique({ + where: { id: task.id }, + include: { riskHits: true, _count: { select: { messageRecords: true } } }, + }); + } + async evaluateTask(data: EvaluateSmsTaskDto) { await this.ensureDefaultRules(); if (data.createdById) { @@ -261,6 +331,7 @@ export class RiskReviewService { if (!task) { throw new NotFoundException('SMS send task not found'); } + this.assertAggregationWindowClosed(task); return this.prisma.smsSendTask.update({ where: { id: taskId }, data: { @@ -280,6 +351,7 @@ export class RiskReviewService { if (!task) { throw new NotFoundException('SMS send task not found'); } + this.assertAggregationWindowClosed(task); const reason = data.reason ?? task.reviewReason ?? '审核拒绝'; return this.prisma.smsSendTask.update({ where: { id: taskId }, @@ -306,6 +378,12 @@ export class RiskReviewService { } } + private assertAggregationWindowClosed(task: { sourceType?: string | null; windowEndsAt?: Date | null }) { + if (task.sourceType === 'cmpp_template_mismatch' && task.windowEndsAt && task.windowEndsAt.getTime() > Date.now()) { + throw new BadRequestException('聚合窗口尚未关闭,请在窗口结束后审核'); + } + } + private async effectiveRules(tenantId: string) { const rules = await this.prisma.riskRule.findMany({ where: { diff --git a/api/src/send-chain/client-send-chain.controller.ts b/api/src/send-chain/client-send-chain.controller.ts index 9207143..8d2f4b2 100644 --- a/api/src/send-chain/client-send-chain.controller.ts +++ b/api/src/send-chain/client-send-chain.controller.ts @@ -1,4 +1,4 @@ -import { Body, Controller, Get, Param, Post, Query } from '@nestjs/common'; +import { BadRequestException, Body, Controller, Get, Param, Post, Query } from '@nestjs/common'; import { ApiTags } from '@nestjs/swagger'; import { TenantId } from '../common/tenant-id.decorator'; import { ConfirmImportDto, CreateBatchTaskDto, ImportPreviewDto, SendChainService } from './send-chain.service'; @@ -25,21 +25,28 @@ export class ClientSendChainController { @Get('batch-tasks') listBatchTasks(@TenantId() tenantId?: string, @Query('status') status?: string) { - return this.sendChain.listBatchTasks(tenantId, status); + return this.sendChain.listBatchTasks(requireTenantId(tenantId), status, 'client'); } @Get('batch-tasks/:id') - getBatchTask(@Param('id') taskId: string) { - return this.sendChain.getBatchTask(taskId); + getBatchTask(@TenantId() tenantId: string | undefined, @Param('id') taskId: string) { + return this.sendChain.getBatchTask(taskId, requireTenantId(tenantId), 'client'); } @Get('batch-tasks/:id/messages') - listTaskMessages(@Param('id') taskId: string) { - return this.sendChain.listMessages({ taskId }); + listTaskMessages(@TenantId() tenantId: string | undefined, @Param('id') taskId: string) { + return this.sendChain.listClientTaskMessages(taskId, requireTenantId(tenantId)); } @Post('batch-tasks/:id/cancel') - cancelBatchTask(@Param('id') taskId: string) { - return this.sendChain.cancelBatchTask(taskId); + cancelBatchTask(@TenantId() tenantId: string | undefined, @Param('id') taskId: string) { + return this.sendChain.cancelBatchTask(taskId, requireTenantId(tenantId), 'client'); } } + +function requireTenantId(tenantId?: string) { + if (!tenantId) { + throw new BadRequestException('Tenant context is required'); + } + return tenantId; +} diff --git a/api/src/send-chain/send-chain.module.ts b/api/src/send-chain/send-chain.module.ts index 87f4ec1..f4e1ed8 100644 --- a/api/src/send-chain/send-chain.module.ts +++ b/api/src/send-chain/send-chain.module.ts @@ -1,4 +1,4 @@ -import { Module } from '@nestjs/common'; +import { forwardRef, Module } from '@nestjs/common'; import { BillingModule } from '../billing/billing.module'; import { PrismaModule } from '../prisma/prisma.module'; import { RiskReviewModule } from '../risk-review/risk-review.module'; @@ -9,7 +9,7 @@ import { GatewayEventsController } from './gateway-events.controller'; import { SendChainService } from './send-chain.service'; @Module({ - imports: [PrismaModule, BillingModule, RiskReviewModule, SmsConfigModule], + imports: [PrismaModule, BillingModule, forwardRef(() => RiskReviewModule), SmsConfigModule], controllers: [AdminSendChainController, ClientSendChainController, GatewayEventsController], providers: [SendChainService], exports: [SendChainService], diff --git a/api/src/send-chain/send-chain.service.spec.ts b/api/src/send-chain/send-chain.service.spec.ts index f6f687a..b743a4b 100644 --- a/api/src/send-chain/send-chain.service.spec.ts +++ b/api/src/send-chain/send-chain.service.spec.ts @@ -95,9 +95,16 @@ function createPrismaMock() { signature: { auditStatus: 'approved', reportStatus: 'approved' }, }), }, + smsSignature: { + findFirst: jest.fn().mockResolvedValue({ id: 'sig-1', name: '签名', auditStatus: 'approved', reportStatus: 'approved' }), + }, + smsSendTask: { + findUnique: jest.fn().mockResolvedValue(null), + }, smsBatchTask: { create: jest.fn().mockResolvedValue(task), findUnique: jest.fn().mockResolvedValue(task), + findFirst: jest.fn().mockResolvedValue(task), findMany: jest.fn(), update: jest.fn().mockResolvedValue(task), }, @@ -296,6 +303,10 @@ function createService(prisma = createPrismaMock()) { reason: null, task: { id: 'risk-task-1' }, }), + aggregateTemplateMismatch: jest.fn().mockResolvedValue({ + id: 'review-task-1', + reviewReason: '企业应用已配置模板不匹配进入人工审核', + }), } as unknown as RiskReviewService; const service = new SendChainService(prisma as never, billing, riskReview); service['postGatewayControl'] = jest.fn().mockResolvedValue({ delivered: true }); @@ -372,7 +383,7 @@ describe('SendChainService', () => { it('cancels scheduled tasks before dispatch', async () => { const { service, prisma } = createService(); - prisma.smsBatchTask.findUnique.mockResolvedValue({ id: 'task-1', status: 'scheduled' }); + prisma.smsBatchTask.findFirst.mockResolvedValue({ id: 'task-1', status: 'scheduled' }); await service.cancelBatchTask('task-1'); @@ -386,6 +397,30 @@ describe('SendChainService', () => { }); }); + it('lists only client-created batch tasks for task progress', async () => { + const { service, prisma } = createService(); + prisma.smsBatchTask.findMany.mockResolvedValue([{ id: 'task-client', tenantId: 'tenant-1', sourceType: 'client' }]); + prisma.smsMessageRecord.groupBy.mockResolvedValue([]); + + await service.listBatchTasks('tenant-1', 'queued'); + + expect(prisma.smsBatchTask.findMany).toHaveBeenCalledWith(expect.objectContaining({ + where: { tenantId: 'tenant-1', status: 'queued', sourceType: 'client' }, + })); + }); + + it('does not expose CMPP internal tasks through client task detail or messages', async () => { + const { service, prisma } = createService(); + prisma.smsBatchTask.findFirst.mockResolvedValue(null); + + await expect(service.getBatchTask('task-cmpp', 'tenant-1', 'client')).rejects.toThrow('SMS batch task not found'); + await expect(service.listClientTaskMessages('task-cmpp', 'tenant-1')).rejects.toThrow('SMS batch task not found'); + expect(prisma.smsBatchTask.findFirst).toHaveBeenCalledWith(expect.objectContaining({ + where: { id: 'task-cmpp', tenantId: 'tenant-1', sourceType: 'client' }, + })); + expect(prisma.smsMessageRecord.findMany).not.toHaveBeenCalled(); + }); + it('terminates non-final tasks by canceling unsubmitted messages', async () => { const { service, prisma } = createService(); prisma.smsBatchTask.findUnique.mockResolvedValue({ id: 'task-1', status: 'sending' }); @@ -507,7 +542,7 @@ describe('SendChainService', () => { }); it('records an unreported CMPP message and returns success before delivering the template failure receipt', async () => { - const { service, prisma } = createService(); + const { service, prisma, riskReview } = createService(); prisma.smsTemplate.findFirst.mockResolvedValue(null); await expect(service.submitInboundMessage({ @@ -524,6 +559,74 @@ describe('SendChainService', () => { '/downstream/receipt', expect.objectContaining({ receiptStatus: 'undelivered', rawStatus: 'REJECTD', errorCode: 'TEMPLATE' }), ); + expect(riskReview.aggregateTemplateMismatch).not.toHaveBeenCalled(); + }); + + it('aggregates template-mismatched CMPP messages only when the application uses manual review', async () => { + const { service, prisma, riskReview } = createService(); + prisma.smsApplication.findFirst.mockResolvedValue({ + id: 'app-1', + tenantId: 'tenant-1', + cmppAccount: '100001', + status: 'active', + interfaceEnabled: true, + templateMismatchMode: 'manual_review', + customerUnitPrice: 3, + queuePriority: 'normal', + ipAllowlist: [{ ipCidr: '127.0.0.1/32' }], + tenant: { id: 'tenant-1', status: 'active', certificationStatus: 'approved' }, + }); + prisma.smsTemplate.findFirst.mockResolvedValue(null); + + await expect(service.submitInboundMessage({ + account: '100001', + phoneNumber: '13800000001', + content: '【签名】未匹配模板的内容', + remoteIp: '127.0.0.1', + })).resolves.toEqual(expect.objectContaining({ accepted: true, messageRecordId: 'record-1' })); + + expect(riskReview.aggregateTemplateMismatch).toHaveBeenCalledWith(expect.objectContaining({ + applicationId: 'app-1', + account: '100001', + messageRecordId: 'record-1', + signatureId: 'sig-1', + })); + expect(prisma.smsBatchTask.update).toHaveBeenCalledWith({ + where: { id: 'task-1' }, + data: expect.objectContaining({ status: 'pending_review', riskTaskId: 'review-task-1', auditStatus: 'pending' }), + }); + expect(prisma.smsReceiptRecord.create).not.toHaveBeenCalled(); + }); + + it('fans an approved aggregated review task back into each internal CMPP batch', async () => { + const { service, prisma } = createService(); + service.enqueueBatchTask = jest.fn().mockResolvedValue({ taskId: 'task-1', enqueued: 1 }); + prisma.smsSendTask.findUnique.mockResolvedValue({ + id: 'review-task-1', + messageRecords: [{ + id: 'record-1', + tenantId: 'tenant-1', + applicationId: 'app-1', + batchTaskId: 'task-1', + messageId: 'MSG-1', + phoneNumber: '13800000001', + amountCents: 3, + billingUnits: 1, + batchTask: { id: 'task-1', sourceType: 'cmpp' }, + }], + }); + + await expect(service.handleReviewDecision('review-task-1', 'approved', '审核通过')).resolves.toEqual({ + reviewTaskId: 'review-task-1', + decision: 'approved', + affected: 1, + }); + + expect(prisma.smsMessageRecord.update).toHaveBeenCalledWith({ + where: { id: 'record-1' }, + data: { status: 'queued', errorCode: null, errorMessage: null }, + }); + expect(service.enqueueBatchTask).toHaveBeenCalledWith('task-1'); }); it('previews imported phone files with duplicate, invalid, blacklist, and variable errors', async () => { diff --git a/api/src/send-chain/send-chain.service.ts b/api/src/send-chain/send-chain.service.ts index 3fd7782..90984f4 100644 --- a/api/src/send-chain/send-chain.service.ts +++ b/api/src/send-chain/send-chain.service.ts @@ -331,9 +331,9 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { return this.getBatchTask(task.id); } - async listBatchTasks(tenantId?: string, status?: string) { + async listBatchTasks(tenantId?: string, status?: string, sourceType = 'client') { const tasks = await this.prisma.smsBatchTask.findMany({ - where: { tenantId, status }, + where: { tenantId, status, sourceType }, include: { tenant: true, application: true, @@ -355,11 +355,20 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { })); } - getBatchTask(taskId: string) { - return this.prisma.smsBatchTask.findUnique({ - where: { id: taskId }, + async getBatchTask(taskId: string, tenantId?: string, sourceType = 'client') { + const task = await this.prisma.smsBatchTask.findFirst({ + where: { id: taskId, tenantId, sourceType }, include: { apiRequests: true, messages: { take: 20, orderBy: { queuedAt: 'asc' } } }, }); + if (!task) { + throw new NotFoundException('SMS batch task not found'); + } + return task; + } + + async listClientTaskMessages(taskId: string, tenantId: string) { + await this.getBatchTask(taskId, tenantId, 'client'); + return this.listMessages({ tenantId, taskId }); } listMessages(query: { @@ -506,8 +515,8 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { return { taskId, enqueued: messages.length }; } - async cancelBatchTask(taskId: string) { - const task = await this.prisma.smsBatchTask.findUnique({ where: { id: taskId } }); + async cancelBatchTask(taskId: string, tenantId?: string, sourceType = 'client') { + const task = await this.prisma.smsBatchTask.findFirst({ where: { id: taskId, tenantId, sourceType } }); if (!task) { throw new NotFoundException('SMS batch task not found'); } @@ -524,6 +533,50 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { }); } + async handleReviewDecision(reviewTaskId: string, decision: 'approved' | 'rejected', reason: string) { + const reviewTask = await this.prisma.smsSendTask.findUnique({ + where: { id: reviewTaskId }, + include: { + messageRecords: { + where: { status: 'pending_review' }, + include: { batchTask: true }, + }, + }, + }); + if (!reviewTask || reviewTask.messageRecords.length === 0) { + return { reviewTaskId, decision, affected: 0 }; + } + const batchTaskIds = new Set(); + for (const message of reviewTask.messageRecords) { + if (!message.tenantId || !message.applicationId || !message.batchTaskId) continue; + if (decision === 'approved') { + await this.prisma.smsMessageRecord.update({ + where: { id: message.id }, + data: { status: 'queued', errorCode: null, errorMessage: null }, + }); + await this.prisma.smsBatchTask.update({ + where: { id: message.batchTaskId }, + data: { status: 'ready', auditStatus: 'approved', reviewReason: reason, rejectReason: null }, + }); + batchTaskIds.add(message.batchTaskId); + } else { + await this.releaseMessageReservation( + message as typeof message & { tenantId: string; batchTaskId: string }, + '模板不匹配人工审核驳回释放冻结', + ); + await this.prisma.smsBatchTask.update({ + where: { id: message.batchTaskId }, + data: { status: 'rejected', auditStatus: 'rejected', rejectReason: reason }, + }); + await this.recordCmppFailureReceipt(message, 'REVIEW_REJECTED', reason); + } + } + for (const batchTaskId of batchTaskIds) { + await this.enqueueBatchTask(batchTaskId); + } + return { reviewTaskId, decision, affected: reviewTask.messageRecords.length }; + } + async terminateBatchTask(taskId: string) { const task = await this.prisma.smsBatchTask.findUnique({ where: { id: taskId } }); if (!task) { @@ -611,7 +664,7 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { async processSendJob(job: SendJob) { const message = await this.prisma.smsMessageRecord.findUnique({ where: { id: job.messageRecordId }, - include: { batchTask: true, template: { include: { signature: true } } }, + include: { batchTask: true, template: { include: { signature: true } }, signature: true }, }); if (!message || message.status !== 'queued') { return { skipped: true }; @@ -1513,6 +1566,60 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { await reject('INTERFACE', '短信应用 CMPP 接口已停用'); } else if (application.tenant.certificationStatus !== 'approved') { await reject('CERT', '企业认证未通过'); + } else if (!template && application.templateMismatchMode === 'manual_review') { + const signature = await this.resolveInboundSignatureCandidate(application.id, data.content); + if (!signature) { + await reject('SIGNATURE', '短信内容未识别到已审核且已报备的签名'); + } else { + const risk = await this.riskReview.evaluateTask({ + tenantId: application.tenantId, + applicationId: application.id, + content: data.content, + phones: [data.phoneNumber], + }); + if (risk.status === 'rejected') { + await reject('RISK', risk.reason || '短信被风控拒绝'); + } else { + const accountCheck = await this.billing.checkAccount({ + tenantId: application.tenantId, + amountCents: billing.amountCents, + smsUnits: billing.totalBillingUnits, + }); + if (!accountCheck.canSend) { + await reject('BALANCE', '企业账户余额、套餐余量或授信额度不足'); + } else { + if (billing.amountCents + billing.totalBillingUnits > 0) { + await this.billing.freeze({ + tenantId: application.tenantId, + amountCents: billing.amountCents, + smsUnits: billing.totalBillingUnits, + relatedType: 'sms_batch_task', + relatedId: task.id, + remark: 'CMPP 模板不匹配待审核短信冻结', + }); + } + const reviewTask = risk.status === 'pending_review' && risk.task + ? await this.attachMessageToReviewTask(risk.task.id, message.id, signature.id) + : await this.riskReview.aggregateTemplateMismatch({ + tenantId: application.tenantId, + applicationId: application.id, + account: data.account, + messageRecordId: message.id, + signatureId: signature.id, + content: data.content, + }); + await this.prisma.smsBatchTask.update({ + where: { id: task.id }, + data: { + status: 'pending_review', + riskTaskId: reviewTask?.id, + auditStatus: 'pending', + reviewReason: reviewTask?.reviewReason ?? '模板不匹配,等待人工审核', + }, + }); + } + } + } } else if (!template) { await reject('TEMPLATE', '短信内容未匹配到已报备模板'); } else if (template.auditStatus !== 'approved') { @@ -1616,6 +1723,7 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { billingUnits: number; queuePriority?: string | null; template?: { signature?: { id?: string | null; name?: string | null } | null } | null; + signature?: { id?: string | null; name?: string | null } | null; }, routed: RoutedChannel, attempt: number, @@ -1668,7 +1776,7 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { queuePriority: normalizeQueuePriority(message.queuePriority), phoneNumber: message.phoneNumber, content: message.content, - signature: message.template?.signature?.name ?? 'SMS', + signature: message.template?.signature?.name ?? message.signature?.name ?? 'SMS', templateId: message.templateId ?? 'unknown', billingUnits: message.billingUnits, route: { @@ -1897,6 +2005,28 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { }); } + private resolveInboundSignatureCandidate(applicationId: string, content: string) { + const match = content.match(/^【([^】]+)】/); + if (!match?.[1]) return null; + return this.prisma.smsSignature.findFirst({ + where: { + applicationId, + name: match[1], + auditStatus: 'approved', + reportStatus: 'approved', + }, + orderBy: { updatedAt: 'desc' }, + }); + } + + private async attachMessageToReviewTask(reviewTaskId: string, messageRecordId: string, signatureId: string) { + await this.prisma.smsMessageRecord.update({ + where: { id: messageRecordId }, + data: { reviewTaskId, signatureId, status: 'pending_review' }, + }); + return this.prisma.smsSendTask.findUnique({ where: { id: reviewTaskId } }); + } + private async recordCmppFailureReceipt( message: { id: string; @@ -2102,10 +2232,11 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { id: string; templateId?: string | null; template?: { signature?: { id?: string | null; name?: string | null } | null } | null; + signature?: { id?: string | null; name?: string | null } | null; }, channelId: string, ) { - let signatureId = message.template?.signature?.id ?? null; + let signatureId = message.template?.signature?.id ?? message.signature?.id ?? null; if (!signatureId && message.templateId) { const template = await this.prisma.smsTemplate.findUnique({ where: { id: message.templateId }, diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index 6b58496..903b7c4 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -139,11 +139,18 @@ 7. 若命中审核策略,运营端短信审核通过后进入发送队列,审核页面必须展示进入审核的原因。 8. 发送服务按通道组路由、通道限速、企业限速执行提交。 9. 平台批量任务只记录客户端创建的发送任务;API 调用和 CMPP 对接发送不进入批量任务。 + - 运营端和客户端“短信任务进度”只查询 `SmsBatchTask.sourceType=client`。 + - CMPP/API/通道测试可使用内部批次承载风控、计费、队列、重试和回执关联,但不得出现在客户批量任务列表。 + - 客户端批量任务列表、详情、短信明细和取消操作必须同时校验当前企业和 `sourceType=client`。 10. 所有来源的短信,包括平台批量任务、API 调用、CMPP 对接发送,全部按手机号维度进入短信记录。 11. 任务进度、发送详情和短信记录实时或准实时更新。 -12. 发送入队必须按短信应用的队列等级分流到优先队列或普通队列;同等条件下优先队列消息必须先于普通队列消息被 Send Worker 消费并提交 Gateway。 -13. 优先队列只能改变待发送消息的调度顺序,不得绕过企业/应用状态、签名模板审核、通道报备、余额/授信、黑名单、风控、通道组路由、通道限速和 Gateway 连接可用性校验。 -14. 同一队列内部按创建时间、任务顺序和手机号拆分顺序保持 FIFO 或可解释的稳定排序;优先队列插队时必须可在 trace 或任务日志中追踪队列等级和入队时间。 +12. 企业应用“不符合模板的短信”配置为 `manual_review` 时,合法的 CMPP Submit 在模板不匹配后进入人工审核;配置为 `reject` 时仍直接拒绝并返回 `REJECTD` Deliver Receipt,其他模式不得被人工审核聚合逻辑误接管。 +13. CMPP 模板不匹配审核支持短窗口内容指纹聚合:只有同一企业应用、同一 CMPP 账号、规范化后内容 SHA-256 完全一致且位于同一时间窗口的短信才能合并为一个审核任务。默认窗口 10 秒,可通过 `CMPP_TEMPLATE_REVIEW_WINDOW_MS` 调整。 +14. 聚合审核不合并短信记录、计费或回执:每个手机号仍有独立 `SmsMessageRecord/messageId/sequenceId`。审核通过后逐条进入真实路由和上游提交;审核驳回后逐条释放冻结并产生客户侧 `REJECTD` 回执。 +15. 人工审核只覆盖模板不匹配;签名必须能从短信前缀识别且已审核/报备通过。签名不合法、风控直接拒绝或余额不足不得因内容聚合而绕过。 +16. 发送入队必须按短信应用的队列等级分流到优先队列或普通队列;同等条件下优先队列消息必须先于普通队列消息被 Send Worker 消费并提交 Gateway。 +17. 优先队列只能改变待发送消息的调度顺序,不得绕过企业/应用状态、签名模板审核、通道报备、余额/授信、黑名单、风控、通道组路由、通道限速和 Gateway 连接可用性校验。 +18. 同一队列内部按创建时间、任务顺序和手机号拆分顺序保持 FIFO 或可解释的稳定排序;优先队列插队时必须可在 trace 或任务日志中追踪队列等级和入队时间。 ### 4.6 通道配置与路由 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index 93c4ce8..d4f61a4 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -1008,6 +1008,7 @@ - CMPP2.0 和 CMPP3.0 连接分别返回对应版本 ConnectResp,后续 Submit/Deliver 按该 TCP 连接协商版本解包和组包,不发生字段错位。 - 密码错误、应用停用、企业停用、短信接口关闭、IP 不在白名单时 connect/login 被拒绝。 - submit 被接受后返回 CMPP SubmitResp 成功,并在真实数据库创建 `sourceType=cmpp` 的发送记录,进入真实发送链路。 + - `sourceType=cmpp` 的内部批次不出现在运营端或客户端“短信任务进度”;客户端不能通过任务 ID 读取该内部批次的详情、短信明细或执行取消。 - Submit 应用身份使用 bind 已鉴权账号;`MsgSrc` 使用应用级企业代码并独立校验。企业代码与登录账号不同时仍能正确定位应用,企业代码不匹配时返回失败。 - 鉴权失败、IP 白名单不符、手机号等协议参数不合法时返回非零 SubmitResp,且不创建短信记录。 - 已鉴权且参数合法的 Submit 必须先返回成功 SubmitResp 和平台 Msg_Id;内容不匹配审核模板、签名/报备未通过、余额不足、应用在 bind 后停用、无可用通道、上游 Submit 最终失败时,均须真实创建短信记录、`SmsReceiptRecord` 和 `CmppDownstreamDelivery`,并向客户下发 `undelivered/REJECTD` Deliver Receipt,不得仅以 SubmitResp 失败替代回执。 @@ -1015,6 +1016,36 @@ - CMPP 包在进入 handler 前因长度、命令字、读包或 Unpack 失败时,Gateway 记录 `read/unpack packet failed`、远端地址、协议模式、错误类型和原始错误,不得静默断开。 - 企业应用列表和连接详情展示真实下游 CMPP 会话:bind 后当前连接数加一,显示客户 IP、企业代码、CMPP 版本、连接建立时间与最后心跳;连接持续未响应 `ACTIVE_TEST` 超过阈值后转为心跳超时/断开,不能继续显示为正常连接。 +### TC-SEND-038 批量任务与 CMPP 内部批次隔离 + +- 优先级:P0 +- 前置条件:同一企业已存在一个客户端批量发送任务,并通过下游 CMPP 提交一条短信。 +- 步骤: + 1. 查询运营端短信任务进度。 + 2. 查询该企业的客户端批量任务列表。 + 3. 使用 CMPP 内部批次 ID 请求客户端任务详情、短信明细和取消接口。 + 4. 查询运营端短信记录。 +- 预期结果: + - 两个任务进度列表只返回 `sourceType=client` 的客户端批量任务。 + - CMPP 内部批次的详情、短信明细和取消请求均返回不可见/不存在,不泄露内部任务。 + - CMPP 短信仍完整出现在短信记录、提交、回执和账务链路。 + +### TC-SEND-039 CMPP 模板不匹配短窗口聚合人工审核 + +- 优先级:P0 +- 前置条件:应用 A 配置 `templateMismatchMode=manual_review`,应用 B 配置 `reject`;两个应用均已配置审核和报备通过的签名、余额和通道组。 +- 步骤: + 1. 在 10 秒内用应用 A 的同一 CMPP 账号向不同手机号提交多条规范化后内容完全一致、签名合法但不匹配模板的短信。 + 2. 用应用 A 提交内容不同、账号不同或跨越聚合窗口的短信。 + 3. 用应用 B 提交同样的模板不匹配短信。 + 4. 分别审核通过和驳回应用 A 的聚合任务。 +- 预期结果: + - 应用 A 同账号、同内容指纹、同窗口的短信只创建一个 `sourceType=cmpp_template_mismatch` 审核任务,审核页展示真实聚合号码数。 + - 不同应用、账号、内容指纹或窗口的短信不合并。 + - 应用 B 不进入人工审核,继续逐条产生 `REJECTD` Deliver Receipt。 + - 审核通过后每条成员独立进入路由、提交和计费;审核驳回后每条成员独立释放冻结并向客户下发 `REJECTD`。 + - 签名无法识别/未报备、风控直接拒绝或余额不足时不进入聚合审核。 + ### TC-GW-007 CMPP 客户到上游 SMSC 完整闭环 - 优先级:P0 diff --git a/docs/testing-progress.md b/docs/testing-progress.md index 3d571e9..5eb2465 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -1600,3 +1600,20 @@ git diff --check - 保留的固定条数均属于明确的近期日志/监控窗口、超时扫描单批、导出保护、单消息重试尝试或唯一候选判定,这些响应不被页面展示为业务总数。 - 已执行 API 全量测试(13 suites、122 项通过),及受影响服务定向测试(5 suites、91 项通过)、API build、前端 build 与 `git diff --check`;前端仅有既有 chunk size warning。 - 已将 `390b9700` 部署生产,部署前完成 PostgreSQL 和当前发布源码备份,Prisma 确认 27 条 migration 均已应用。`cmpp-api`、`cmpp-gateway`、Nginx、MinIO 均为 active,`12026/17890/8090/3000` 监听、API/Gateway health 和外部 `12026` HTTP 均通过。使用生产平台管理员会话调用真实受保护 API:企业管理返回 105 条未删除企业,企业应用返回 204 条,不再封顶 100。 + +## 2026-07-12 客户批量任务与 CMPP 内部批次隔离 + +- 短信任务进度的产品口径回归为“客户在客户端提交的批量发送任务”。运营端和客户端任务列表统一强制 `SmsBatchTask.sourceType=client`,不再展示每条 CMPP Submit 创建的 `sourceType=cmpp` 内部批次或通道测试批次。 +- CMPP 内部批次仍保留在 PostgreSQL,继续承载模板/签名/报备校验、风控、冻结/计费、队列、补发和 Deliver Receipt 关联;所有 CMPP 号码仍进入短信记录。 +- 客户端任务详情、任务短信明细和取消接口同时校验当前 `tenantId` 和 `sourceType=client`,不能通过内部批次 ID 读取或操作 CMPP 内部任务。 +- 生产现状只读核对:`sourceType=cmpp` 2 个、`sourceType=admin_channel_test` 1 个、`sourceType=client` 0 个。部署后任务进度应显示 0 个客户批量任务,但不删除现有内部批次数据。 +- 已执行 `send-chain.service.spec.ts + operations.service.spec.ts`(2 suites、50 项通过)、API 全量测试(13 suites、125 项通过)、API build 和前端 build;前端仅有既有 chunk size warning。按用户要求本轮暂不部署生产。 + +## 2026-07-12 CMPP 模板不匹配短窗口聚合审核 + +- 只有企业应用配置 `templateMismatchMode=manual_review` 时,CMPP 模板不匹配短信才进入人工审核。`reject` 继续逐条失败并下发 `REJECTD`;其他模式不被聚合逻辑接管。 +- 新增 `SmsSendTask.sourceType/aggregationKey/contentHash/windowStartedAt/windowEndsAt`,以应用、CMPP 账号、规范化内容 SHA-256 和默认 10 秒窗口生成唯一聚合键。`SmsMessageRecord.reviewTaskId/signatureId` 保留每条成员与审核任务、真实签名的关联。 +- 聚合前仍校验签名审核/报备、风控直接拒绝和账户余额,并按每条短信独立冻结。审核通过后逐条恢复到各自 `sourceType=cmpp` 内部批次并入队;驳回后逐条释放冻结、写入失败记录并生成客户侧 Deliver Receipt。 +- 运营端短信审核页新增“审核来源”和“聚合号码数”,区分 CMPP 模板不匹配聚合与普通风控审核。 +- 聚合窗口关闭前,任务不返回到待审列表且审核接口拒绝提前操作,避免窗口内后到短信加入已完成任务。 +- Prisma migration:`20260712113000_add_cmpp_review_aggregation`。已执行 API 全量测试、API build、前端 build、Prisma validate 和 `git diff --check`;本轮按用户要求暂不部署生产。 diff --git a/src/api/adminApi.ts b/src/api/adminApi.ts index 2de61de..4b84beb 100644 --- a/src/api/adminApi.ts +++ b/src/api/adminApi.ts @@ -305,6 +305,10 @@ export type SmsBatchTask = { id: string; tenantId: string; taskNo: string; + sourceType?: string; + contentHash?: string | null; + windowStartedAt?: string | null; + windowEndsAt?: string | null; applicationId?: string | null; templateId?: string | null; content: string; @@ -559,6 +563,10 @@ export type RiskReviewTask = { applicationId?: string | null; templateId?: string | null; taskNo: string; + sourceType?: string; + contentHash?: string | null; + windowStartedAt?: string | null; + windowEndsAt?: string | null; content: string; category?: string | null; phoneTotal: number; @@ -574,6 +582,7 @@ export type RiskReviewTask = { createdAt: string; reviewedAt?: string | null; riskHits?: Array<{ id: string; ruleName: string; reason: string }>; + _count?: { messageRecords: number }; }; export type TenantAccount = { diff --git a/src/apps/admin/AdminSmsAuditPage.tsx b/src/apps/admin/AdminSmsAuditPage.tsx index bfbe753..caf3029 100644 --- a/src/apps/admin/AdminSmsAuditPage.tsx +++ b/src/apps/admin/AdminSmsAuditPage.tsx @@ -15,6 +15,10 @@ const statusTone: Record = { rejected: 'danger', }; +function sourceLabel(sourceType?: string) { + return sourceType === 'cmpp_template_mismatch' ? 'CMPP模板不匹配聚合' : '风控审核'; +} + export function AdminSmsAuditPage() { const [records, setRecords] = useState([]); const [keyword, setKeyword] = useState(''); @@ -69,8 +73,9 @@ export function AdminSmsAuditPage() { const columns: Array> = [ { key: 'taskNo', title: '任务编号', width: '180px', render: (record) => {record.taskNo} }, + { key: 'sourceType', title: '审核来源', width: '180px', render: (record) => {sourceLabel(record.sourceType)} }, { key: 'content', title: '短信内容', render: (record) => {record.content} }, - { key: 'phoneTotal', title: '号码数', width: '130px', render: (record) => record.phoneTotal.toLocaleString('zh-CN') }, + { key: 'phoneTotal', title: '聚合号码数', width: '140px', render: (record) => (record._count?.messageRecords ?? record.phoneTotal).toLocaleString('zh-CN') }, { key: 'createdAt', title: '提交时间', width: '190px', render: (record) => record.createdAt }, { key: 'reason', title: '审核原因', render: (record) => record.reviewReason ?? record.rejectReason ?? record.riskHits?.map((item) => item.reason).join(';') ?? '-' }, {