diff --git a/api/prisma/migrations/20260813150000_harden_downstream_requeue_tasks/migration.sql b/api/prisma/migrations/20260813150000_harden_downstream_requeue_tasks/migration.sql new file mode 100644 index 0000000..7f46654 --- /dev/null +++ b/api/prisma/migrations/20260813150000_harden_downstream_requeue_tasks/migration.sql @@ -0,0 +1,19 @@ +ALTER TABLE "DownstreamRequeueTask" + ADD COLUMN "applicationFailures" JSONB NOT NULL DEFAULT '{}', + ADD COLUMN "scanLeaseOwner" TEXT, + ADD COLUMN "scanLeaseUntil" TIMESTAMP(3); + +CREATE TABLE "DownstreamRequeueRateWindow" ( + "id" TEXT NOT NULL, + "applicationId" TEXT NOT NULL, + "windowStartedAt" TIMESTAMP(3) NOT NULL, + "consumed" INTEGER NOT NULL DEFAULT 0, + "updatedAt" TIMESTAMP(3) NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + CONSTRAINT "DownstreamRequeueRateWindow_pkey" PRIMARY KEY ("id") +); + +CREATE UNIQUE INDEX "DownstreamRequeueRateWindow_applicationId_windowStartedAt_key" + ON "DownstreamRequeueRateWindow"("applicationId", "windowStartedAt"); +CREATE INDEX "DownstreamRequeueRateWindow_windowStartedAt_idx" + ON "DownstreamRequeueRateWindow"("windowStartedAt"); diff --git a/api/prisma/schema.prisma b/api/prisma/schema.prisma index aeacd92..e9d8ee7 100644 --- a/api/prisma/schema.prisma +++ b/api/prisma/schema.prisma @@ -2137,7 +2137,10 @@ model DownstreamRequeueTask { skippedCount Int @default(0) waitingCount Int @default(0) consecutiveFailures Int @default(0) + applicationFailures Json @default("{}") lastError String? + scanLeaseOwner String? + scanLeaseUntil DateTime? createdById String? startedAt DateTime? pausedAt DateTime? @@ -2155,6 +2158,18 @@ model DownstreamRequeueTask { @@index([tenantId, createdAt]) } +model DownstreamRequeueRateWindow { + id String @id @default(cuid()) + applicationId String + windowStartedAt DateTime + consumed Int @default(0) + updatedAt DateTime @updatedAt + createdAt DateTime @default(now()) + + @@unique([applicationId, windowStartedAt]) + @@index([windowStartedAt]) +} + model DownstreamRequeueTaskItem { id String @id @default(cuid()) taskId String diff --git a/api/src/operations/admin-operations.controller.ts b/api/src/operations/admin-operations.controller.ts index 5bb505a..7267f54 100644 --- a/api/src/operations/admin-operations.controller.ts +++ b/api/src/operations/admin-operations.controller.ts @@ -364,18 +364,20 @@ export class AdminOperationsController { } @Post('downstream-requeue-tasks/preview') - previewDownstreamRequeueTask(@Body() body: { filter?: Record }) { - return this.sendChain.previewDownstreamRequeueTask(body.filter ?? {}); + previewDownstreamRequeueTask( + @Body() body: { filter?: Record }, + @CurrentSessionUserId() operatorId?: string, + ) { + return this.sendChain.previewDownstreamRequeueTask(body.filter ?? {}, operatorId); } @Post('downstream-requeue-tasks') createDownstreamRequeueTask( - @Body() body: { filter?: Record; snapshotAt?: string; reason?: string; ratePerSecond?: number; consecutiveFailureLimit?: number }, + @Body() body: { previewToken?: string; reason?: string; ratePerSecond?: number; consecutiveFailureLimit?: number }, @CurrentSessionUserId() operatorId?: string, ) { return this.sendChain.createDownstreamRequeueTask({ - filter: body.filter ?? {}, - snapshotAt: body.snapshotAt ?? '', + previewToken: body.previewToken ?? '', reason: body.reason ?? '', ratePerSecond: body.ratePerSecond, consecutiveFailureLimit: body.consecutiveFailureLimit, @@ -392,6 +394,17 @@ export class AdminOperationsController { return this.sendChain.getDownstreamRequeueTask(id); } + @Get('downstream-requeue-tasks/:id/items') + listDownstreamRequeueTaskItems( + @Param('id') id: string, + @Query('status') status?: string, + @Query('keyword') keyword?: string, + @Query('page') page?: string, + @Query('pageSize') pageSize?: string, + ) { + return this.sendChain.listDownstreamRequeueTaskItems(id, { status, keyword, page: Number(page), pageSize: Number(pageSize) }); + } + @Post('downstream-requeue-tasks/:id/:action') changeDownstreamRequeueTaskStatus( @Param('id') id: string, diff --git a/api/src/send-chain/send-chain.service.ts b/api/src/send-chain/send-chain.service.ts index 48a1073..0cda40c 100644 --- a/api/src/send-chain/send-chain.service.ts +++ b/api/src/send-chain/send-chain.service.ts @@ -513,11 +513,11 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { return this.completion.batchRequeueDownstreamDeliveries(ids); } - previewDownstreamRequeueTask(filter: DownstreamRequeueFilter) { - return this.downstreamRequeueTasks.preview(filter); + previewDownstreamRequeueTask(filter: DownstreamRequeueFilter, operatorId?: string) { + return this.downstreamRequeueTasks.preview(filter, operatorId); } - createDownstreamRequeueTask(data: { filter: DownstreamRequeueFilter; snapshotAt: string; reason: string; ratePerSecond?: number; consecutiveFailureLimit?: number }, operatorId?: string) { + createDownstreamRequeueTask(data: { previewToken: string; reason: string; ratePerSecond?: number; consecutiveFailureLimit?: number }, operatorId?: string) { return this.downstreamRequeueTasks.create(data, operatorId); } @@ -529,6 +529,10 @@ export class SendChainService implements OnModuleInit, OnModuleDestroy { return this.downstreamRequeueTasks.get(id); } + listDownstreamRequeueTaskItems(id: string, query: { status?: string; keyword?: string; page?: number; pageSize?: number }) { + return this.downstreamRequeueTasks.listItems(id, query); + } + changeDownstreamRequeueTaskStatus(id: string, action: 'pause' | 'resume' | 'terminate', operatorId?: string) { return this.downstreamRequeueTasks.changeStatus(id, action, operatorId); } diff --git a/api/src/send-chain/send-downstream-requeue-task.service.spec.ts b/api/src/send-chain/send-downstream-requeue-task.service.spec.ts index e5cfa07..996b0f3 100644 --- a/api/src/send-chain/send-downstream-requeue-task.service.spec.ts +++ b/api/src/send-chain/send-downstream-requeue-task.service.spec.ts @@ -3,66 +3,103 @@ import { SendDownstreamRequeueTaskService } from './send-downstream-requeue-task function prismaMock(): Record { const result: Record = { - cmppDownstreamDelivery: { - count: jest.fn(), groupBy: jest.fn(), findFirst: jest.fn(), findMany: jest.fn(), findUnique: jest.fn(), - }, - downstreamRequeueTask: { - findFirst: jest.fn(), findMany: jest.fn(), count: jest.fn(), findUnique: jest.fn(), create: jest.fn(), update: jest.fn(), - }, - downstreamRequeueTaskItem: { - createMany: jest.fn(), groupBy: jest.fn(), findMany: jest.fn(), findFirst: jest.fn(), updateMany: jest.fn(), update: jest.fn(), - }, + cmppDownstreamDelivery: { count: jest.fn(), groupBy: jest.fn(), findFirst: jest.fn(), findMany: jest.fn(), findUnique: jest.fn() }, + cmppDownstreamConnection: { count: jest.fn() }, + downstreamRequeueTask: { findFirst: jest.fn(), findMany: jest.fn(), count: jest.fn(), findUnique: jest.fn(), create: jest.fn(), update: jest.fn(), updateMany: jest.fn() }, + downstreamRequeueTaskItem: { createMany: jest.fn(), count: jest.fn(), groupBy: jest.fn(), findMany: jest.fn(), findFirst: jest.fn(), updateMany: jest.fn(), update: jest.fn() }, operationLog: { create: jest.fn() }, + downstreamRequeueRateWindow: { deleteMany: jest.fn() }, + $queryRaw: jest.fn(), }; result.$transaction = jest.fn(async (callback: (tx: unknown) => unknown) => callback(result)); return result; } let mock: Record; +let service: SendDownstreamRequeueTaskService; describe('SendDownstreamRequeueTaskService', () => { - beforeEach(() => { mock = prismaMock(); }); + beforeEach(() => { + process.env.DATABASE_URL = 'postgresql://test:test@127.0.0.1/test'; + mock = prismaMock(); + service = new SendDownstreamRequeueTaskService(mock as never, { requeueDownstreamDelivery: jest.fn() }); + }); - it('previews all matches separately from replayable records', async () => { - mock.cmppDownstreamDelivery.count.mockResolvedValueOnce(12).mockResolvedValueOnce(8); - mock.cmppDownstreamDelivery.groupBy - .mockResolvedValueOnce([{ status: 'pending', _count: { _all: 8 } }, { status: 'delivered', _count: { _all: 4 } }]) - .mockResolvedValueOnce([{ applicationId: 'app-1', _count: { _all: 12 } }]); + it('keeps the selected status in matched counts and returns a signed preview token', async () => { + mock.cmppDownstreamDelivery.count.mockResolvedValueOnce(8).mockResolvedValueOnce(8); + mock.cmppDownstreamDelivery.groupBy.mockResolvedValueOnce([{ status: 'pending', _count: { _all: 8 } }]).mockResolvedValueOnce([{ applicationId: 'app-1', _count: { _all: 8 } }]); mock.cmppDownstreamDelivery.findFirst.mockResolvedValue({ createdAt: new Date('2026-08-11T00:00:00Z') }); - const service = new SendDownstreamRequeueTaskService(mock as never, { requeueDownstreamDelivery: jest.fn() }); - const result = await service.preview({ status: 'all' }); - expect(result).toEqual(expect.objectContaining({ matchedCount: 12, replayableCount: 8, skippedCount: 4, applicationCount: 1 })); + const result = await service.preview({ status: 'pending' }, 'user-1'); + expect(result).toEqual(expect.objectContaining({ matchedCount: 8, replayableCount: 8, skippedCount: 0, previewToken: expect.stringContaining('.') })); + expect(mock.cmppDownstreamDelivery.count).toHaveBeenNthCalledWith(1, expect.objectContaining({ where: expect.objectContaining({ status: 'pending' }) })); }); - it('rejects delivered filters and short reasons', async () => { - const service = new SendDownstreamRequeueTaskService(mock as never, { requeueDownstreamDelivery: jest.fn() }); - await expect(service.create({ filter: { status: 'delivered' }, snapshotAt: new Date().toISOString(), reason: '事故恢复' })).rejects.toBeInstanceOf(BadRequestException); - await expect(service.create({ filter: { status: 'pending' }, snapshotAt: new Date().toISOString(), reason: '短' })).rejects.toBeInstanceOf(BadRequestException); + it('rejects a tampered preview token and short reasons', async () => { + await expect(service.create({ previewToken: 'invalid.token', reason: '处理事故积压' }, 'user-1')).rejects.toBeInstanceOf(BadRequestException); + await expect(service.create({ previewToken: 'invalid.token', reason: '短' }, 'user-1')).rejects.toBeInstanceOf(BadRequestException); }); - it('rejects a new task when the same application scope already has an unfinished task', async () => { - mock.downstreamRequeueTask.findFirst.mockResolvedValue({ taskNo: 'DRT-EXISTING' }); - const service = new SendDownstreamRequeueTaskService(mock as never, { requeueDownstreamDelivery: jest.fn() }); - await expect(service.create({ filter: { applicationId: 'app-1', status: 'pending' }, snapshotAt: new Date().toISOString(), reason: '处理历史回执积压' })).rejects.toThrow('DRT-EXISTING'); + it('materializes only the server-signed preview range', async () => { + mock.cmppDownstreamDelivery.count.mockResolvedValue(2); + mock.cmppDownstreamDelivery.groupBy.mockResolvedValueOnce([{ status: 'failed', _count: { _all: 2 } }]).mockResolvedValueOnce([{ applicationId: 'app-1', _count: { _all: 2 } }]); + mock.cmppDownstreamDelivery.findFirst.mockResolvedValue({ createdAt: new Date() }); + const preview = await service.preview({ applicationId: 'app-1', status: 'failed' }, 'user-1'); + mock.downstreamRequeueTask.findFirst.mockResolvedValue(null); + mock.cmppDownstreamDelivery.findMany.mockResolvedValue([{ id: 'd-1', applicationId: 'app-1', status: 'failed' }]); + mock.downstreamRequeueTask.create.mockResolvedValue({ id: 'task-1' }); + mock.downstreamRequeueTask.findUnique.mockResolvedValue({ id: 'task-1' }); + mock.downstreamRequeueTaskItem.groupBy.mockResolvedValue([]); + await service.create({ previewToken: preview.previewToken, reason: '处理历史回执积压' }, 'user-1'); + expect(mock.cmppDownstreamDelivery.findMany).toHaveBeenCalledWith(expect.objectContaining({ where: expect.objectContaining({ AND: expect.any(Array) }) })); + expect(mock.downstreamRequeueTaskItem.createMany).toHaveBeenCalledWith({ data: [expect.objectContaining({ deliveryId: 'd-1', applicationId: 'app-1' })] }); }); - it('skips a delivery that automatic recovery already confirmed before task execution', async () => { + it('paginates all task items with status and keyword filters', async () => { + mock.downstreamRequeueTask.findUnique.mockResolvedValue({ id: 'task-1' }); + mock.downstreamRequeueTaskItem.findMany.mockResolvedValue([{ id: 'item-1' }]); + mock.downstreamRequeueTaskItem.count.mockResolvedValue(21); + await expect(service.listItems('task-1', { status: 'failed', keyword: 'MSG-1', page: 2, pageSize: 20 })).resolves.toEqual({ items: [{ id: 'item-1' }], total: 21, page: 2, pageSize: 20 }); + expect(mock.downstreamRequeueTaskItem.findMany).toHaveBeenCalledWith(expect.objectContaining({ skip: 20, take: 20, where: expect.objectContaining({ status: 'failed', OR: expect.any(Array) }) })); + }); + + it('recovers an expired processing claim before scanning queued work', async () => { mock.downstreamRequeueTask.findMany.mockResolvedValue([{ id: 'task-1' }]); - mock.downstreamRequeueTask.findUnique - .mockResolvedValueOnce({ id: 'task-1', taskNo: 'DRT-1', status: 'queued', startedAt: null, ratePerSecond: 10, consecutiveFailures: 0, consecutiveFailureLimit: 10 }) - .mockResolvedValueOnce({ status: 'running' }) - .mockResolvedValue({ status: 'running' }); - mock.downstreamRequeueTaskItem.findMany - .mockResolvedValueOnce([]) - .mockResolvedValueOnce([{ id: 'item-1', deliveryId: 'delivery-1' }]) - .mockResolvedValueOnce([]); + mock.downstreamRequeueTask.updateMany.mockResolvedValue({ count: 1 }); + mock.downstreamRequeueTask.findUnique.mockResolvedValueOnce({ id: 'task-1', taskNo: 'DRT-1', status: 'running', startedAt: new Date(), ratePerSecond: 10, consecutiveFailureLimit: 10, applicationFailures: {} }).mockResolvedValue({ status: 'running' }); mock.downstreamRequeueTaskItem.updateMany.mockResolvedValue({ count: 1 }); - mock.cmppDownstreamDelivery.findUnique.mockResolvedValue({ status: 'delivered', payload: {}, deliveryType: 'receipt', application: { status: 'active', interfaceEnabled: true } }); - mock.downstreamRequeueTaskItem.groupBy.mockResolvedValue([{ status: 'skipped', _count: { _all: 1 } }]); + mock.downstreamRequeueTaskItem.findMany.mockResolvedValue([]); + mock.downstreamRequeueTaskItem.groupBy.mockResolvedValue([]); + await service.runScan(); + expect(mock.downstreamRequeueTaskItem.updateMany).toHaveBeenCalledWith(expect.objectContaining({ where: expect.objectContaining({ status: 'processing', claimedAt: expect.any(Object) }), data: expect.objectContaining({ status: 'queued', claimedAt: null }) })); + }); + + it('moves an offline application to waiting_connection without calling Gateway', async () => { const requeue = jest.fn(); - const service = new SendDownstreamRequeueTaskService(mock as never, { requeueDownstreamDelivery: requeue }); + service = new SendDownstreamRequeueTaskService(mock as never, { requeueDownstreamDelivery: requeue }); + mock.downstreamRequeueTask.findMany.mockResolvedValue([{ id: 'task-1' }]); + mock.downstreamRequeueTask.updateMany.mockResolvedValue({ count: 1 }); + mock.downstreamRequeueTask.findUnique.mockResolvedValueOnce({ id: 'task-1', taskNo: 'DRT-1', status: 'queued', startedAt: null, ratePerSecond: 10, consecutiveFailureLimit: 10, applicationFailures: {} }).mockResolvedValue({ status: 'running' }); + mock.downstreamRequeueTaskItem.findMany.mockResolvedValueOnce([]).mockResolvedValueOnce([]).mockResolvedValueOnce([{ id: 'item-1', deliveryId: 'd-1', applicationId: 'app-1' }]).mockResolvedValueOnce([]).mockResolvedValueOnce([]); + mock.downstreamRequeueTaskItem.updateMany.mockResolvedValue({ count: 1 }); + mock.downstreamRequeueTaskItem.groupBy.mockResolvedValue([{ status: 'waiting_connection', _count: { _all: 1 } }]); + mock.$queryRaw.mockResolvedValue([{ consumed: 1 }]); + mock.cmppDownstreamDelivery.findUnique.mockResolvedValue({ id: 'd-1', applicationId: 'app-1', status: 'failed', payload: {}, deliveryType: 'receipt', application: { status: 'active', interfaceEnabled: true } }); + mock.downstreamRequeueTaskItem.findFirst.mockResolvedValue(null); + mock.cmppDownstreamConnection.count.mockResolvedValue(0); await service.runScan(); expect(requeue).not.toHaveBeenCalled(); - expect(mock.downstreamRequeueTaskItem.update).toHaveBeenCalledWith(expect.objectContaining({ data: expect.objectContaining({ status: 'skipped', skipReason: '已被客户确认' }) })); + expect(mock.downstreamRequeueTaskItem.update).toHaveBeenCalledWith(expect.objectContaining({ data: expect.objectContaining({ status: 'waiting_connection' }) })); + }); + + it('counts ACK failures by application and auto-pauses at the threshold', async () => { + mock.downstreamRequeueTask.findMany.mockResolvedValue([{ id: 'task-1' }]); + mock.downstreamRequeueTask.updateMany.mockResolvedValue({ count: 1 }); + mock.downstreamRequeueTask.findUnique.mockResolvedValueOnce({ id: 'task-1', taskNo: 'DRT-1', status: 'running', startedAt: new Date(), ratePerSecond: 10, consecutiveFailureLimit: 1, applicationFailures: {} }).mockResolvedValue({ status: 'running' }); + mock.downstreamRequeueTaskItem.findMany.mockResolvedValueOnce([]).mockResolvedValueOnce([{ id: 'item-1', applicationId: 'app-1', status: 'waiting_ack', delivery: { status: 'rejected', ackResult: 1, ackDeadlineAt: new Date(), lastError: 'ACK Result=1' } }]).mockResolvedValueOnce([]).mockResolvedValueOnce([]).mockResolvedValueOnce([]); + mock.downstreamRequeueTaskItem.groupBy.mockResolvedValue([{ status: 'failed', _count: { _all: 1 } }]); + await service.runScan(); + expect(mock.downstreamRequeueTaskItem.update).toHaveBeenCalledWith(expect.objectContaining({ data: expect.objectContaining({ status: 'failed' }) })); + expect(mock.downstreamRequeueTask.updateMany).toHaveBeenCalledWith(expect.objectContaining({ data: expect.objectContaining({ status: 'paused', lastError: expect.stringContaining('app-1') }) })); + expect(mock.operationLog.create).toHaveBeenCalledWith(expect.objectContaining({ data: expect.objectContaining({ action: 'gateway.downstream_requeue_task_auto_paused' }) })); }); }); diff --git a/api/src/send-chain/send-downstream-requeue-task.service.ts b/api/src/send-chain/send-downstream-requeue-task.service.ts index 324e9cf..fd7f407 100644 --- a/api/src/send-chain/send-downstream-requeue-task.service.ts +++ b/api/src/send-chain/send-downstream-requeue-task.service.ts @@ -1,5 +1,6 @@ import { BadRequestException, NotFoundException } from '@nestjs/common'; import { Prisma } from '@prisma/client'; +import { createHash, createHmac, randomUUID, timingSafeEqual } from 'node:crypto'; import { PrismaService } from '../prisma/prisma.service'; import { parseDateBoundary } from '../operations/operations.helpers'; @@ -15,68 +16,117 @@ export type DownstreamRequeueFilter = { type RequeueFacade = { requeueDownstreamDelivery(id: string): Promise }; const REPLAYABLE_STATUSES = ['pending', 'failed', 'unconfirmed', 'rejected']; +const ACTIVE_TASK_STATUSES = ['queued', 'running', 'paused']; +const PROCESSING_LEASE_MS = 2 * 60_000; +const SCAN_LEASE_MS = 15_000; +const PREVIEW_TOKEN_TTL_MS = 15 * 60_000; + +function normalizedFilter(filter: DownstreamRequeueFilter): DownstreamRequeueFilter { + return { + tenantId: filter.tenantId || 'all', + applicationId: filter.applicationId || 'all', + deliveryType: filter.deliveryType || 'all', + status: filter.status || 'all', + keyword: filter.keyword?.trim() || undefined, + createdAtFrom: filter.createdAtFrom || undefined, + createdAtTo: filter.createdAtTo || undefined, + }; +} function taskWhere(filter: DownstreamRequeueFilter, snapshotAt: Date, replayableByDefault = true): Prisma.CmppDownstreamDeliveryWhereInput { - const from = parseDateBoundary(filter.createdAtFrom, false); - const to = parseDateBoundary(filter.createdAtTo, true); + const normalized = normalizedFilter(filter); + const from = parseDateBoundary(normalized.createdAtFrom, false); + const to = parseDateBoundary(normalized.createdAtTo, true); return { - tenantId: filter.tenantId && filter.tenantId !== 'all' ? filter.tenantId : undefined, - applicationId: filter.applicationId && filter.applicationId !== 'all' ? filter.applicationId : undefined, - deliveryType: filter.deliveryType && filter.deliveryType !== 'all' ? filter.deliveryType : undefined, - status: filter.status && filter.status !== 'all' ? filter.status : replayableByDefault ? { in: REPLAYABLE_STATUSES } : undefined, + tenantId: normalized.tenantId !== 'all' ? normalized.tenantId : undefined, + applicationId: normalized.applicationId !== 'all' ? normalized.applicationId : undefined, + deliveryType: normalized.deliveryType !== 'all' ? normalized.deliveryType : undefined, + status: normalized.status !== 'all' ? normalized.status : replayableByDefault ? { in: REPLAYABLE_STATUSES } : undefined, createdAt: { ...(from ? { gte: from } : {}), lte: to && to < snapshotAt ? to : snapshotAt }, - OR: filter.keyword ? [ - { messageId: { contains: filter.keyword } }, - { payload: { path: ['account'], string_contains: filter.keyword } }, - { payload: { path: ['phoneNumber'], string_contains: filter.keyword } }, - { lastError: { contains: filter.keyword } }, - { tenant: { name: { contains: filter.keyword } } }, - { application: { name: { contains: filter.keyword } } }, + OR: normalized.keyword ? [ + { messageId: { contains: normalized.keyword } }, + { payload: { path: ['account'], string_contains: normalized.keyword } }, + { payload: { path: ['phoneNumber'], string_contains: normalized.keyword } }, + { lastError: { contains: normalized.keyword } }, + { tenant: { name: { contains: normalized.keyword } } }, + { application: { name: { contains: normalized.keyword } } }, ] : undefined, }; } +function previewSecret() { + const source = process.env.DOWNSTREAM_REQUEUE_PREVIEW_SECRET || process.env.DATABASE_URL; + if (!source) throw new BadRequestException('后台重投预检签名密钥未配置'); + return createHash('sha256').update(`cmpp-downstream-requeue-preview\0${source}`).digest(); +} + +function signPreview(payload: Record) { + const encoded = Buffer.from(JSON.stringify(payload)).toString('base64url'); + const signature = createHmac('sha256', previewSecret()).update(encoded).digest('base64url'); + return `${encoded}.${signature}`; +} + +function verifyPreview(token: string, operatorId?: string) { + const [encoded, supplied] = String(token || '').split('.'); + if (!encoded || !supplied) throw new BadRequestException('预检凭证无效,请重新预检'); + const expected = createHmac('sha256', previewSecret()).update(encoded).digest(); + let actual: Buffer; + try { actual = Buffer.from(supplied, 'base64url'); } catch { throw new BadRequestException('预检凭证无效,请重新预检'); } + if (expected.length !== actual.length || !timingSafeEqual(expected, actual)) throw new BadRequestException('预检凭证无效,请重新预检'); + const payload = JSON.parse(Buffer.from(encoded, 'base64url').toString('utf8')) as { filter: DownstreamRequeueFilter; snapshotAt: string; operatorId?: string; expiresAt: number }; + if (payload.expiresAt < Date.now()) throw new BadRequestException('预检凭证已过期,请重新预检'); + if ((payload.operatorId || '') !== (operatorId || '')) throw new BadRequestException('预检凭证与当前操作人不一致'); + return payload; +} + +function jsonFailures(value: unknown): Record { + if (!value || typeof value !== 'object' || Array.isArray(value)) return {}; + return Object.fromEntries(Object.entries(value).map(([key, count]) => [key, Math.max(0, Number(count) || 0)])); +} + export class SendDownstreamRequeueTaskService { constructor(private readonly prisma: PrismaService, private readonly facade: RequeueFacade) {} - async preview(filter: DownstreamRequeueFilter) { + async preview(filter: DownstreamRequeueFilter, operatorId?: string) { const snapshotAt = new Date(); - const base = taskWhere({ ...filter, status: 'all' }, snapshotAt, false); - const where = taskWhere(filter, snapshotAt); + const normalized = normalizedFilter(filter); + const base = taskWhere(normalized, snapshotAt, false); + const replayableWhere = { AND: [base, { status: { in: REPLAYABLE_STATUSES } }] } as Prisma.CmppDownstreamDeliveryWhereInput; const [matchedCount, replayableCount, statusGroups, appGroups, oldest] = await Promise.all([ this.prisma.cmppDownstreamDelivery.count({ where: base }), - this.prisma.cmppDownstreamDelivery.count({ where: { AND: [where, { status: { in: REPLAYABLE_STATUSES } }] } }), + this.prisma.cmppDownstreamDelivery.count({ where: replayableWhere }), this.prisma.cmppDownstreamDelivery.groupBy({ by: ['status'], where: base, _count: { _all: true } }), this.prisma.cmppDownstreamDelivery.groupBy({ by: ['applicationId'], where: base, _count: { _all: true } }), this.prisma.cmppDownstreamDelivery.findFirst({ where: base, orderBy: { createdAt: 'asc' }, select: { createdAt: true } }), ]); + const tokenPayload = { filter: normalized, snapshotAt: snapshotAt.toISOString(), operatorId: operatorId || '', expiresAt: Date.now() + PREVIEW_TOKEN_TTL_MS }; return { snapshotAt, + previewToken: signPreview(tokenPayload), matchedCount, replayableCount, skippedCount: matchedCount - replayableCount, applicationCount: appGroups.length, oldestCreatedAt: oldest?.createdAt ?? null, statusCounts: Object.fromEntries(statusGroups.map((item) => [item.status, item._count._all])), - filter: { ...filter, status: filter.status ?? 'all' }, + filter: normalized, }; } - async create(data: { filter: DownstreamRequeueFilter; snapshotAt: string; reason: string; ratePerSecond?: number; consecutiveFailureLimit?: number }, createdById?: string) { + async create(data: { previewToken: string; reason: string; ratePerSecond?: number; consecutiveFailureLimit?: number }, createdById?: string) { const reason = data.reason?.trim(); if (!reason || reason.length < 5) throw new BadRequestException('任务原因至少填写5个字'); - const snapshotAt = new Date(data.snapshotAt); - if (Number.isNaN(snapshotAt.getTime()) || snapshotAt.getTime() > Date.now() + 10_000) throw new BadRequestException('预检快照时间无效'); - if (data.filter.status === 'delivered' || data.filter.status === 'awaiting_ack') throw new BadRequestException('第一版后台任务不支持已确认或等待ACK记录'); + const preview = verifyPreview(data.previewToken, createdById); + const filter = normalizedFilter(preview.filter); + const snapshotAt = new Date(preview.snapshotAt); + if (filter.status === 'delivered' || filter.status === 'awaiting_ack') throw new BadRequestException('第一版后台任务不支持已确认或等待ACK记录'); const activeTask = await this.prisma.downstreamRequeueTask.findFirst({ where: { - status: { in: ['queued', 'running', 'paused'] }, - ...(data.filter.applicationId && data.filter.applicationId !== 'all' - ? { OR: [{ applicationId: data.filter.applicationId }, { applicationId: null }] } - : {}), + status: { in: ACTIVE_TASK_STATUSES }, + ...(filter.applicationId !== 'all' ? { OR: [{ applicationId: filter.applicationId }, { applicationId: null }] } : {}), }, select: { taskNo: true } }); if (activeTask) throw new BadRequestException(`当前应用范围已有未结束任务 ${activeTask.taskNo}`); - const where = { AND: [taskWhere(data.filter, snapshotAt), { status: { in: REPLAYABLE_STATUSES } }] } as Prisma.CmppDownstreamDeliveryWhereInput; - const deliveries = await this.prisma.cmppDownstreamDelivery.findMany({ where, orderBy: [{ createdAt: 'asc' }, { id: 'asc' }], take: 100001, select: { id: true, tenantId: true, applicationId: true, status: true } }); + const where = { AND: [taskWhere(filter, snapshotAt), { status: { in: REPLAYABLE_STATUSES } }] } as Prisma.CmppDownstreamDeliveryWhereInput; + const deliveries = await this.prisma.cmppDownstreamDelivery.findMany({ where, orderBy: [{ createdAt: 'asc' }, { id: 'asc' }], take: 100001, select: { id: true, applicationId: true, status: true } }); if (!deliveries.length) throw new BadRequestException('当前筛选条件下没有可重投记录'); if (deliveries.length > 100000) throw new BadRequestException('单个任务最多处理100000条,请缩小日期范围'); const ratePerSecond = Math.min(50, Math.max(1, Number(data.ratePerSecond ?? 10))); @@ -85,14 +135,14 @@ export class SendDownstreamRequeueTaskService { const task = await this.prisma.$transaction(async (tx) => { const created = await tx.downstreamRequeueTask.create({ data: { taskNo, - tenantId: data.filter.tenantId && data.filter.tenantId !== 'all' ? data.filter.tenantId : null, - applicationId: data.filter.applicationId && data.filter.applicationId !== 'all' ? data.filter.applicationId : null, - filterSnapshot: data.filter as Prisma.InputJsonValue, + tenantId: filter.tenantId !== 'all' ? filter.tenantId : null, + applicationId: filter.applicationId !== 'all' ? filter.applicationId : null, + filterSnapshot: filter as Prisma.InputJsonValue, snapshotAt, reason, ratePerSecond, consecutiveFailureLimit: failureLimit, totalCount: deliveries.length, createdById, } }); await tx.downstreamRequeueTaskItem.createMany({ data: deliveries.map((item) => ({ taskId: created.id, deliveryId: item.id, applicationId: item.applicationId, previousStatus: item.status })) }); - await tx.operationLog.create({ data: { userId: createdById, action: 'gateway.downstream_requeue_task_created', resource: 'downstream_requeue_task', resourceId: created.id, detail: { taskNo, reason, totalCount: deliveries.length, snapshotAt, filter: data.filter, ratePerSecond } } }); + await tx.operationLog.create({ data: { userId: createdById, action: 'gateway.downstream_requeue_task_created', resource: 'downstream_requeue_task', resourceId: created.id, detail: { taskNo, reason, totalCount: deliveries.length, snapshotAt, filter, ratePerSecond, consecutiveFailureLimit: failureLimit } } }); return created; }); return this.get(task.id); @@ -112,11 +162,30 @@ export class SendDownstreamRequeueTaskService { async get(id: string) { const task = await this.prisma.downstreamRequeueTask.findUnique({ where: { id }, include: { tenant: true, application: true, createdBy: { select: { id: true, displayName: true, username: true } } } }); if (!task) throw new NotFoundException('后台重投任务不存在'); - const [itemGroups, recentItems] = await Promise.all([ - this.prisma.downstreamRequeueTaskItem.groupBy({ by: ['status'], where: { taskId: id }, _count: { _all: true } }), - this.prisma.downstreamRequeueTaskItem.findMany({ where: { taskId: id }, include: { delivery: { select: { messageId: true, deliveryType: true, status: true, lastError: true } } }, orderBy: { updatedAt: 'desc' }, take: 50 }), + const itemGroups = await this.prisma.downstreamRequeueTaskItem.groupBy({ by: ['status'], where: { taskId: id }, _count: { _all: true } }); + return { ...task, itemCounts: Object.fromEntries(itemGroups.map((item) => [item.status, item._count._all])) }; + } + + async listItems(id: string, query: { status?: string; keyword?: string; page?: number; pageSize?: number }) { + const task = await this.prisma.downstreamRequeueTask.findUnique({ where: { id }, select: { id: true } }); + if (!task) throw new NotFoundException('后台重投任务不存在'); + const page = Math.max(1, Number(query.page ?? 1)); + const pageSize = Math.min(100, Math.max(1, Number(query.pageSize ?? 20))); + const keyword = query.keyword?.trim(); + const where: Prisma.DownstreamRequeueTaskItemWhereInput = { + taskId: id, + status: query.status && query.status !== 'all' ? query.status : undefined, + OR: keyword ? [ + { delivery: { messageId: { contains: keyword } } }, + { skipReason: { contains: keyword } }, + { errorMessage: { contains: keyword } }, + ] : undefined, + }; + const [items, total] = await Promise.all([ + this.prisma.downstreamRequeueTaskItem.findMany({ where, include: { delivery: { select: { messageId: true, deliveryType: true, status: true, lastError: true } } }, orderBy: [{ updatedAt: 'desc' }, { id: 'desc' }], skip: (page - 1) * pageSize, take: pageSize }), + this.prisma.downstreamRequeueTaskItem.count({ where }), ]); - return { ...task, itemCounts: Object.fromEntries(itemGroups.map((item) => [item.status, item._count._all])), recentItems }; + return { items, total, page, pageSize }; } async changeStatus(id: string, action: 'pause' | 'resume' | 'terminate', operatorId?: string) { @@ -126,107 +195,156 @@ export class SendDownstreamRequeueTaskService { const allowed = action === 'pause' ? ['queued', 'running'] : action === 'resume' ? ['paused'] : ['queued', 'running', 'paused']; if (!allowed.includes(task.status)) throw new BadRequestException('当前任务状态不允许此操作'); const status = action === 'pause' ? 'paused' : action === 'resume' ? 'queued' : 'terminated'; - const updated = await this.prisma.downstreamRequeueTask.update({ where: { id }, data: { status, pausedAt: status === 'paused' ? new Date() : null, finishedAt: status === 'terminated' ? new Date() : undefined } }); - if (status === 'terminated') await this.prisma.downstreamRequeueTaskItem.updateMany({ where: { taskId: id, status: 'queued' }, data: { status: 'unprocessed', skipReason: '任务已终止', completedAt: new Date() } }); + const updated = await this.prisma.downstreamRequeueTask.update({ where: { id }, data: { status, pausedAt: status === 'paused' ? new Date() : null, finishedAt: status === 'terminated' ? new Date() : undefined, scanLeaseOwner: null, scanLeaseUntil: null } }); + if (status === 'terminated') await this.prisma.downstreamRequeueTaskItem.updateMany({ where: { taskId: id, status: { in: ['queued', 'waiting_connection'] } }, data: { status: 'unprocessed', skipReason: '任务已终止', completedAt: new Date() } }); await this.prisma.operationLog.create({ data: { userId: operatorId, action: `gateway.downstream_requeue_task_${action}`, resource: 'downstream_requeue_task', resourceId: id, detail: { taskNo: task.taskNo, previousStatus: task.status, status } } }); return updated; } async runScan() { - const tasks = await this.prisma.downstreamRequeueTask.findMany({ where: { status: { in: ['queued', 'running'] } }, orderBy: { createdAt: 'asc' }, take: 3 }); + await this.prisma.downstreamRequeueRateWindow.deleteMany({ where: { windowStartedAt: { lt: new Date(Date.now() - 5 * 60_000) } } }); + const tasks = await this.prisma.downstreamRequeueTask.findMany({ where: { status: { in: ['queued', 'running'] } }, orderBy: { createdAt: 'asc' }, take: 3, select: { id: true } }); for (const task of tasks) await this.processTask(task.id); } private async processTask(taskId: string) { - const task = await this.prisma.downstreamRequeueTask.findUnique({ where: { id: taskId } }); - if (!task || !['queued', 'running'].includes(task.status)) return; - await this.reconcileWaiting(taskId); - await this.prisma.downstreamRequeueTask.update({ where: { id: taskId }, data: { status: 'running', startedAt: task.startedAt ?? new Date() } }); - const items = await this.prisma.downstreamRequeueTaskItem.findMany({ where: { taskId, status: 'queued' }, orderBy: { createdAt: 'asc' }, take: Math.min(10, task.ratePerSecond), select: { id: true, deliveryId: true } }); - let consecutiveFailures = task.consecutiveFailures; - for (const item of items) { - const latestTask = await this.prisma.downstreamRequeueTask.findUnique({ where: { id: taskId }, select: { status: true } }); - if (latestTask?.status === 'paused' || latestTask?.status === 'terminated') break; - const claimed = await this.prisma.downstreamRequeueTaskItem.updateMany({ where: { id: item.id, status: 'queued' }, data: { status: 'processing', claimedAt: new Date() } }); - if (!claimed.count) continue; - try { - const delivery = await this.prisma.cmppDownstreamDelivery.findUnique({ - where: { id: item.deliveryId }, - include: { application: { select: { status: true, interfaceEnabled: true } } }, - }); - if (!delivery) { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'skipped', skipReason: '投递记录已不存在', completedAt: new Date() } }); - continue; - } - if (!REPLAYABLE_STATUSES.includes(delivery.status)) { - if (delivery.status === 'awaiting_ack') { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'waiting_external_ack', skipReason: null } }); - } else { - const skipReason = delivery.status === 'delivered' ? '已被客户确认' : '执行前状态已变化'; - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'skipped', skipReason, completedAt: new Date() } }); - } - continue; - } - if (delivery.application.status !== 'active' || !delivery.application.interfaceEnabled) { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'skipped', skipReason: '应用或投递能力已停用', completedAt: new Date() } }); - continue; - } - if (!delivery.payload || !['receipt', 'uplink'].includes(delivery.deliveryType)) { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'skipped', skipReason: '投递数据不完整', completedAt: new Date() } }); - continue; - } - const activeOther = await this.prisma.downstreamRequeueTaskItem.findFirst({ where: { deliveryId: item.deliveryId, id: { not: item.id }, status: { in: ['processing', 'waiting_ack', 'success'] } }, select: { id: true } }); - if (activeOther) { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'skipped', skipReason: '已被其他任务处理', completedAt: new Date() } }); - continue; - } - const result = await this.facade.requeueDownstreamDelivery(item.deliveryId) as { status?: string; lastError?: string | null }; - if (result?.status === 'awaiting_ack' || result?.status === 'delivered') { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: result.status === 'delivered' ? 'success' : 'waiting_ack', completedAt: result.status === 'delivered' ? new Date() : null } }); - consecutiveFailures = 0; - } else { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'failed', errorMessage: result?.lastError ?? 'Gateway未进入等待ACK状态', completedAt: new Date() } }); - consecutiveFailures += 1; - } - } catch (error) { - const message = error instanceof Error ? error.message : '后台重投失败'; - const skipReason = /已被其他操作处理|状态|等待客户端确认/.test(message) ? '执行前状态已变化' - : /payload|投递类型/.test(message) ? '投递数据不完整' - : /Submit|Msg_Id|Sequence/.test(message) ? '缺少原Submit映射,无法安全重投' - : null; - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: skipReason ? 'skipped' : 'failed', skipReason, errorMessage: skipReason ? null : message, completedAt: new Date() } }); - if (!skipReason) consecutiveFailures += 1; + const leaseOwner = randomUUID(); + const now = new Date(); + const lease = await this.prisma.downstreamRequeueTask.updateMany({ + where: { id: taskId, status: { in: ['queued', 'running'] }, OR: [{ scanLeaseUntil: null }, { scanLeaseUntil: { lt: now } }] }, + data: { scanLeaseOwner: leaseOwner, scanLeaseUntil: new Date(now.getTime() + SCAN_LEASE_MS) }, + }); + if (!lease.count) return; + try { + const task = await this.prisma.downstreamRequeueTask.findUnique({ where: { id: taskId } }); + if (!task || !['queued', 'running'].includes(task.status)) return; + // A process may die after the database claim but before the Gateway call. The lease makes that + // ambiguous window visible and recoverable; every recovered item is revalidated before replay. + await this.prisma.downstreamRequeueTaskItem.updateMany({ where: { taskId, status: 'processing', claimedAt: { lt: new Date(Date.now() - PROCESSING_LEASE_MS) } }, data: { status: 'queued', claimedAt: null, errorMessage: '执行进程中断,已回收并等待重新复核' } }); + let failures = await this.reconcileWaiting(taskId, jsonFailures(task.applicationFailures)); + const existingFailureEntry = Object.entries(failures).find(([, count]) => count >= task.consecutiveFailureLimit); + if (existingFailureEntry) { + await this.autoPause(task, existingFailureEntry[0], existingFailureEntry[1]); + await this.refreshTask(taskId); + return; } - await this.prisma.downstreamRequeueTask.update({ where: { id: taskId }, data: { consecutiveFailures } }); - if (consecutiveFailures >= task.consecutiveFailureLimit) { - // Stop before claiming another delivery: a customer or Gateway outage must not become a retry flood. - const pausedAt = new Date(); - await this.prisma.downstreamRequeueTask.update({ where: { id: taskId }, data: { status: 'paused', pausedAt, lastError: `连续失败达到安全阈值 ${task.consecutiveFailureLimit} 条,任务已自动暂停` } }); - await this.prisma.operationLog.create({ data: { action: 'gateway.downstream_requeue_task_auto_paused', resource: 'downstream_requeue_task', resourceId: taskId, detail: { taskNo: task.taskNo, consecutiveFailures, failureLimit: task.consecutiveFailureLimit } } }); - break; + await this.prisma.downstreamRequeueTask.update({ where: { id: taskId }, data: { status: 'running', startedAt: task.startedAt ?? new Date() } }); + const items = await this.prisma.downstreamRequeueTaskItem.findMany({ where: { taskId, status: 'queued' }, orderBy: { createdAt: 'asc' }, take: Math.min(200, task.ratePerSecond * 3), select: { id: true, deliveryId: true, applicationId: true } }); + for (const item of items) { + const latestTask = await this.prisma.downstreamRequeueTask.findUnique({ where: { id: taskId }, select: { status: true } }); + if (!latestTask || !['queued', 'running'].includes(latestTask.status)) break; + if (!(await this.consumeRate(item.applicationId, task.ratePerSecond))) continue; + const claimed = await this.prisma.downstreamRequeueTaskItem.updateMany({ where: { id: item.id, status: 'queued' }, data: { status: 'processing', claimedAt: new Date() } }); + if (!claimed.count) continue; + const outcome = await this.processItem(item.id, item.deliveryId); + if (outcome === 'success') failures[item.applicationId] = 0; + if (outcome === 'failed') failures[item.applicationId] = (failures[item.applicationId] ?? 0) + 1; + const maxFailures = Math.max(0, ...Object.values(failures)); + await this.prisma.downstreamRequeueTask.update({ where: { id: taskId }, data: { applicationFailures: failures as Prisma.InputJsonValue, consecutiveFailures: maxFailures } }); + if ((failures[item.applicationId] ?? 0) >= task.consecutiveFailureLimit) { + await this.autoPause(task, item.applicationId, failures[item.applicationId]); + break; + } } + failures = await this.reconcileWaiting(taskId, failures); + const maxFailures = Math.max(0, ...Object.values(failures)); + await this.prisma.downstreamRequeueTask.updateMany({ where: { id: taskId, status: { in: ['queued', 'running'] } }, data: { applicationFailures: failures as Prisma.InputJsonValue, consecutiveFailures: maxFailures } }); + const ackFailureEntry = Object.entries(failures).find(([, count]) => count >= task.consecutiveFailureLimit); + if (ackFailureEntry) await this.autoPause(task, ackFailureEntry[0], ackFailureEntry[1]); + await this.refreshTask(taskId); + } finally { + await this.prisma.downstreamRequeueTask.updateMany({ where: { id: taskId, scanLeaseOwner: leaseOwner }, data: { scanLeaseOwner: null, scanLeaseUntil: null } }); } - await this.reconcileWaiting(taskId); - await this.refreshTask(taskId); } - private async reconcileWaiting(taskId: string) { - const items = await this.prisma.downstreamRequeueTaskItem.findMany({ where: { taskId, status: { in: ['waiting_ack', 'waiting_external_ack'] } }, include: { delivery: { select: { status: true, ackResult: true, ackDeadlineAt: true, lastError: true } } }, take: 100 }); + private async processItem(itemId: string, deliveryId: string): Promise<'success' | 'failed' | 'waiting' | 'skipped'> { + try { + const delivery = await this.prisma.cmppDownstreamDelivery.findUnique({ + where: { id: deliveryId }, + include: { application: { select: { status: true, interfaceEnabled: true } } }, + }); + if (!delivery) return this.finishItem(itemId, 'skipped', '投递记录已不存在'); + if (!REPLAYABLE_STATUSES.includes(delivery.status)) { + if (delivery.status === 'awaiting_ack') { await this.prisma.downstreamRequeueTaskItem.update({ where: { id: itemId }, data: { status: 'waiting_external_ack', skipReason: null } }); return 'waiting'; } + return this.finishItem(itemId, 'skipped', delivery.status === 'delivered' ? '已被客户确认' : '执行前状态已变化'); + } + if (delivery.application.status !== 'active' || !delivery.application.interfaceEnabled) return this.finishItem(itemId, 'skipped', '应用或投递能力已停用'); + if (!delivery.payload || !['receipt', 'uplink'].includes(delivery.deliveryType)) return this.finishItem(itemId, 'skipped', '投递数据不完整'); + const activeOther = await this.prisma.downstreamRequeueTaskItem.findFirst({ where: { deliveryId, id: { not: itemId }, status: { in: ['processing', 'waiting_ack', 'success'] } }, select: { id: true } }); + if (activeOther) return this.finishItem(itemId, 'skipped', '已被其他任务处理'); + const connected = await this.prisma.cmppDownstreamConnection.count({ where: { applicationId: delivery.applicationId, status: 'connected' } }); + if (connected === 0) { await this.prisma.downstreamRequeueTaskItem.update({ where: { id: itemId }, data: { status: 'waiting_connection', claimedAt: null, errorMessage: '客户当前离线,等待连接恢复' } }); return 'waiting'; } + const result = await this.facade.requeueDownstreamDelivery(deliveryId) as { status?: string; lastError?: string | null }; + if (result?.status === 'delivered') { await this.prisma.downstreamRequeueTaskItem.update({ where: { id: itemId }, data: { status: 'success', completedAt: new Date() } }); return 'success'; } + if (result?.status === 'awaiting_ack') { await this.prisma.downstreamRequeueTaskItem.update({ where: { id: itemId }, data: { status: 'waiting_ack', completedAt: null } }); return 'waiting'; } + await this.prisma.downstreamRequeueTaskItem.update({ where: { id: itemId }, data: { status: 'failed', errorMessage: result?.lastError ?? 'Gateway未进入等待ACK状态', completedAt: new Date() } }); + return 'failed'; + } catch (error) { + const message = error instanceof Error ? error.message : '后台重投失败'; + const skipReason = /已被其他操作处理|状态|等待客户端确认/.test(message) ? '执行前状态已变化' + : /payload|投递类型/.test(message) ? '投递数据不完整' + : /Submit|Msg_Id|Sequence/.test(message) ? '缺少原Submit映射,无法安全重投' : null; + await this.prisma.downstreamRequeueTaskItem.update({ where: { id: itemId }, data: { status: skipReason ? 'skipped' : 'failed', skipReason, errorMessage: skipReason ? null : message, completedAt: new Date() } }); + return skipReason ? 'skipped' : 'failed'; + } + } + + private async finishItem(itemId: string, status: 'skipped', reason: string): Promise<'skipped'> { + await this.prisma.downstreamRequeueTaskItem.update({ where: { id: itemId }, data: { status, skipReason: reason, completedAt: new Date() } }); + return 'skipped'; + } + + private async consumeRate(applicationId: string, limit: number) { + const windowStartedAt = new Date(Math.floor(Date.now() / 1000) * 1000); + const rows = await this.prisma.$queryRaw>(Prisma.sql` + INSERT INTO "DownstreamRequeueRateWindow" ("id", "applicationId", "windowStartedAt", "consumed", "updatedAt") + VALUES (${randomUUID()}, ${applicationId}, ${windowStartedAt}, 1, NOW()) + ON CONFLICT ("applicationId", "windowStartedAt") DO UPDATE + SET "consumed" = "DownstreamRequeueRateWindow"."consumed" + 1, "updatedAt" = NOW() + WHERE "DownstreamRequeueRateWindow"."consumed" < ${limit} + RETURNING "consumed" + `); + return rows.length === 1; + } + + private async reconcileWaiting(taskId: string, currentFailures?: Record) { + const task = currentFailures ? null : await this.prisma.downstreamRequeueTask.findUnique({ where: { id: taskId }, select: { applicationFailures: true } }); + const failures = currentFailures ?? jsonFailures(task?.applicationFailures); + const connectionItems = await this.prisma.downstreamRequeueTaskItem.findMany({ where: { taskId, status: 'waiting_connection' }, select: { id: true, applicationId: true } }); + for (const item of connectionItems) { + const connected = await this.prisma.cmppDownstreamConnection.count({ where: { applicationId: item.applicationId, status: 'connected' } }); + if (connected > 0) await this.prisma.downstreamRequeueTaskItem.updateMany({ where: { id: item.id, status: 'waiting_connection' }, data: { status: 'queued', errorMessage: null, claimedAt: null } }); + } + const items = await this.prisma.downstreamRequeueTaskItem.findMany({ where: { taskId, status: { in: ['waiting_ack', 'waiting_external_ack'] } }, include: { delivery: { select: { status: true, ackResult: true, ackDeadlineAt: true, lastError: true } } }, take: 500 }); const now = new Date(); for (const item of items) { if (item.delivery.status === 'delivered' && item.delivery.ackResult === 0) { await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: item.status === 'waiting_external_ack' ? { status: 'skipped', skipReason: '已由其他投递链路完成', completedAt: now } : { status: 'success', completedAt: now } }); + if (item.status === 'waiting_ack') failures[item.applicationId] = 0; } else if (['failed', 'rejected', 'unconfirmed'].includes(item.delivery.status) || (item.delivery.ackDeadlineAt && item.delivery.ackDeadlineAt <= now)) { - await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: item.status === 'waiting_external_ack' ? { status: 'queued', skipReason: null, errorMessage: null, claimedAt: null } : { status: 'failed', errorMessage: item.delivery.lastError ?? '客户端ACK失败或超时', completedAt: now } }); + if (item.status === 'waiting_external_ack') { + await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'queued', skipReason: null, errorMessage: null, claimedAt: null } }); + } else { + await this.prisma.downstreamRequeueTaskItem.update({ where: { id: item.id }, data: { status: 'failed', errorMessage: item.delivery.lastError ?? '客户端ACK失败或超时', completedAt: now } }); + failures[item.applicationId] = (failures[item.applicationId] ?? 0) + 1; + } } } + return failures; + } + + private async autoPause(task: { id: string; taskNo: string; consecutiveFailureLimit: number }, applicationId: string, count: number) { + const pausedAt = new Date(); + const message = `应用 ${applicationId} 连续失败达到安全阈值 ${task.consecutiveFailureLimit} 条,任务已自动暂停`; + const updated = await this.prisma.downstreamRequeueTask.updateMany({ where: { id: task.id, status: { in: ['queued', 'running'] } }, data: { status: 'paused', pausedAt, lastError: message } }); + if (updated.count) await this.prisma.operationLog.create({ data: { action: 'gateway.downstream_requeue_task_auto_paused', resource: 'downstream_requeue_task', resourceId: task.id, detail: { taskNo: task.taskNo, applicationId, consecutiveFailures: count, failureLimit: task.consecutiveFailureLimit, pausedAt } } }); } private async refreshTask(taskId: string) { const groups = await this.prisma.downstreamRequeueTaskItem.groupBy({ by: ['status'], where: { taskId }, _count: { _all: true } }); const counts = new Map(groups.map((item) => [item.status, item._count._all])); const queued = counts.get('queued') ?? 0; - const active = (counts.get('processing') ?? 0) + (counts.get('waiting_ack') ?? 0) + (counts.get('waiting_external_ack') ?? 0); + const active = (counts.get('processing') ?? 0) + (counts.get('waiting_connection') ?? 0) + (counts.get('waiting_ack') ?? 0) + (counts.get('waiting_external_ack') ?? 0); const failed = counts.get('failed') ?? 0; const current = await this.prisma.downstreamRequeueTask.findUnique({ where: { id: taskId }, select: { status: true } }); const status = current?.status === 'paused' || current?.status === 'terminated' ? current.status : queued + active === 0 ? (failed > 0 ? 'partial_completed' : 'completed') : 'running'; diff --git a/docs/downstream-requeue-task-design-20260812.md b/docs/downstream-requeue-task-design-20260812.md new file mode 100644 index 0000000..bd05ec1 --- /dev/null +++ b/docs/downstream-requeue-task-design-20260812.md @@ -0,0 +1,306 @@ +# 下游投递后台批量重投任务设计与实现符合性审计 + +> 版本:V1.0
+> 需求确认日期:2026-08-12
+> 文档整理日期:2026-08-13
+> 适用页面:运营端 → 下游投递记录
+> 审计基线:当前工作区 `HEAD=67fee216162e638ba21004fcf87e7711facefb91`;本功能相关文件相对 HEAD 无未提交修改
+> 本文目的:还原 2026-08-12 已确认的设计口径,并将当前实现逐条映射到设计,不能以“已有代码”代替“符合设计”的结论。 + +## 1. 背景与目标 + +下游投递记录原有两种人工操作:单条重投、勾选当前页后批量重投。它们适合少量记录,但不适合事故期间处理跨页、跨应用的大批量积压。 + +新增“按当前筛选条件创建后台重投任务”的目标是: + +1. 后端按当前真实筛选条件确定范围,不受页面分页影响; +2. 创建时冻结任务快照,避免执行期间不断卷入新记录; +3. 后台限速执行,客户离线或 ACK 异常时不得形成重投洪峰; +4. 每条任务项可追踪、可恢复、可审计,成功确认后不得重复投递; +5. 运营人员能够预检、创建、暂停、继续、终止和查看完整结果。 + +本功能不替代单条重投和当前页勾选重投,三种入口应同时保留。 + +## 2. 已确认的第一版业务口径 + +以下七条是 2026-08-12 已写入正式需求文档和测试用例的确认口径: + +1. 保留单条重投、当前页勾选批量重投,新增“按筛选条件重投”;投递记录分页支持每页 10/25/50 条。 +2. 后台任务使用企业、应用、投递类型、状态、创建日期、关键词组成的筛选快照;页码和每页条数不属于任务范围。预检生成 `snapshotAt`,创建任务后产生的新记录不进入该任务。 +3. 第一版后台任务只允许 `pending`、`failed`、`unconfirmed`、`rejected`。不得批量重投客户端已确认的 `delivered`;处于 `awaiting_ack` 的记录不得并发重投。创建前展示真实命中数、可重投数、规则跳过数、状态分布,任务原因必填。 +4. 任务按企业应用分批执行,默认每个应用 10 条/秒。单条失败不阻断整批;连续失败达到 10 条,或 ACK 超时/拒绝达到安全阈值时,自动暂停。客户离线、已有链路等待 ACK 属于“等待”,不能记作“跳过”。 +5. “跳过”严格表示本任务没有调用 Gateway。第一版跳过原因包括:执行前状态变化、已被客户确认、已被其他任务处理、本任务已成功处理、不属于任务快照、应用或投递能力已停用、投递数据不完整、缺少原 Submit 映射。 +6. 任务支持列表、详情、暂停、继续、终止。终止只影响尚未发送的任务项。任务项以 `taskId + deliveryId` 幂等;执行前原子认领并复核当前状态;API 重启后继续执行;已获得成功 ACK 的任务项不得再次发送。 +7. 创建、暂停、继续、终止、自动暂停都必须写操作日志。任务必须使用真实 PostgreSQL、Gateway 和客户 ACK,不得使用 mock、静态数据或 localStorage。 + +## 3. 用户交互设计 + +### 3.1 投递记录筛选区 + +筛选条件应包含: + +| 条件 | 说明 | +| --- | --- | +| 企业 | 全部企业或指定企业 | +| 企业应用 | 受企业条件联动;全部应用或指定应用 | +| 投递类型 | 全部、状态回执、上行短信 | +| 状态 | 全部或单一状态 | +| 创建日期 | 开始日期、结束日期,按北京时间自然日 | +| 关键词 | 消息 ID、客户账号、手机号、最后错误、企业名称、应用名称 | + +“按筛选条件重投”使用上表条件,但明确排除页码、每页条数和当前页勾选状态。 + +### 3.2 创建任务流程 + +1. 用户设置筛选条件,点击“按筛选条件重投”。 +2. 后端在同一时点生成预检快照,返回:筛选条件、`snapshotAt`、筛选命中数、可重投数、规则跳过数、状态分布、涉及应用数、最早记录时间。 +3. 弹窗明确告知第一版允许和禁止的状态。 +4. 用户选择执行速度,填写不少于 5 个字的事故原因、工单号或处理说明。 +5. 用户确认后,后端必须重新按预检的筛选快照和 `snapshotAt` 物化任务项,而不是使用前端传入的记录 ID 列表。 +6. 创建成功后关闭弹窗,任务出现在任务列表,状态为“排队中”。 + +预检结果必须严格遵守当前筛选条件。例如当前状态选择“待投递”,命中数、可重投数和状态分布都只能基于“待投递”记录,不能擅自扩展为全部状态。 + +### 3.3 任务列表 + +任务列表至少展示:任务号、创建时间、企业/应用范围、任务原因、中文状态、总数、成功数、失败数、跳过数、等待数、创建人。 + +列表应支持分页和状态筛选,保证历史任务可查询,不能只展示固定数量的最近任务。 + +可执行操作: + +| 当前状态 | 可执行操作 | +| --- | --- | +| 排队中、执行中 | 查看详情、暂停、终止 | +| 已暂停 | 查看详情、继续、终止 | +| 已完成、部分完成、已终止 | 查看详情 | + +### 3.4 任务详情 + +任务详情分为三部分: + +1. 任务信息:筛选快照、快照时间、原因、创建人、速度、安全阈值、开始/暂停/完成时间; +2. 结果汇总:排队、处理中、等待连接、等待 ACK、成功、失败、跳过、未处理; +3. 完整任务项:支持分页及按结果状态、消息 ID、错误/跳过原因查询。 + +任务项必须使用中文状态和明确原因,不能只返回内部英文枚举,也不能只展示最近 50 项而使其余项目不可查询。 + +## 4. 后端数据设计 + +### 4.1 任务主表 + +`DownstreamRequeueTask` 保存: + +- 任务号、企业范围、应用范围; +- 完整筛选快照及 `snapshotAt`; +- 原因、执行速度、失败安全阈值; +- 任务状态和各结果计数; +- 创建人、开始时间、暂停时间、完成时间、最近错误; +- 创建时间、更新时间。 + +### 4.2 任务项表 + +`DownstreamRequeueTaskItem` 保存: + +- `taskId`、`deliveryId`、`applicationId`; +- 创建任务时的原状态; +- 当前任务项状态; +- 跳过原因、失败信息; +- 认领时间、完成时间、创建时间、更新时间。 + +数据库必须有 `taskId + deliveryId` 唯一约束。任务项不能级联删除原下游投递记录,历史下游投递和 ACK 证据必须保留。 + +## 5. 状态机设计 + +### 5.1 任务状态 + +| 状态 | 含义 | 进入条件 | +| --- | --- | --- | +| `queued` | 排队中 | 创建成功或人工继续 | +| `running` | 执行中 | 扫描器开始处理 | +| `paused` | 已暂停 | 人工暂停或触发安全阈值 | +| `completed` | 已完成 | 所有任务项完成且没有失败 | +| `partial_completed` | 部分完成 | 所有任务项结束但存在失败 | +| `terminated` | 已终止 | 人工终止;未认领项不再执行 | + +### 5.2 任务项状态 + +| 状态 | 是否调用过 Gateway | 说明 | +| --- | --- | --- | +| `queued` | 否 | 等待认领 | +| `processing` | 尚不确定 | 已原子认领,正在复核并准备调用 | +| `waiting_connection` | 否或尚未写出 | 客户离线,等待可用连接,不算失败、不算跳过 | +| `waiting_external_ack` | 否 | 其他链路已写出,等待其 ACK 结论 | +| `waiting_ack` | 是 | 本任务已写出,等待客户 ACK | +| `success` | 是 | 客户返回可关联原消息的成功 ACK | +| `failed` | 是或调用失败 | Gateway 调用失败、ACK 超时或 ACK 拒绝 | +| `skipped` | 否 | 复核后确认本任务不应调用 Gateway | +| `unprocessed` | 否 | 任务终止时尚未开始 | + +“写出成功”不能直接记为 `success`,只有客户有效 ACK 才是成功。 + +## 6. 执行与并发控制 + +1. 扫描器只处理 `queued/running` 任务。 +2. 多实例扫描时必须先原子认领任务项;同一任务项只能有一个执行者。 +3. 执行前重新读取下游投递及应用状态,再判断等待、跳过或调用 Gateway。 +4. 限速维度是“企业应用”,不是任务总量。一个跨 3 个应用的任务配置 10 条/秒时,每个应用各自最多 10 条/秒,并且各应用互不阻塞。 +5. 限速必须使用时间窗口或分布式令牌,不能依赖“定时器大约每秒执行一次 + 每轮取 N 条”,否则扫描重叠或执行耗时变化会突破或降低速率。 +6. 客户离线时任务项进入等待连接;客户恢复后继续,不应消耗连续失败阈值。 +7. 本任务写出后进入 `waiting_ack`;成功 ACK 转 `success`;ACK 超时、拒绝、无效 Msg_Id 转 `failed` 并计入安全阈值。 +8. 连续失败只能在真正成功 ACK 后清零,不能在“写出并开始等待 ACK”时提前清零。 +9. `processing` 必须有租约/超时恢复。API 进程中断后,超过认领租约的项目恢复为 `queued` 并重新复核,才能满足重启续跑。 +10. 暂停或终止与执行器并发时,执行器在认领下一项和调用 Gateway 前都必须复核任务状态。 + +## 7. 跳过、等待和失败判定 + +### 7.1 跳过 + +跳过意味着本任务没有调用 Gateway: + +| 原因 | 判定 | +| --- | --- | +| 执行前状态已变化 | 已不属于允许重投状态,且不是等待 ACK/已确认 | +| 已被客户确认 | 当前投递已是 `delivered` 且有效 ACK | +| 已被其他任务处理 | 其他任务已认领、等待 ACK 或成功 | +| 本任务已成功处理 | 同任务项已有成功结果,重复扫描不得再调用 | +| 不属于任务快照 | 创建时间晚于 `snapshotAt` 或不再满足冻结范围 | +| 应用或投递能力已停用 | 应用状态或接口能力不允许投递 | +| 投递数据不完整 | 缺少可重放 payload 或投递类型非法 | +| 缺少原 Submit 映射 | 无法形成可关联原短信的安全回执 | + +### 7.2 等待 + +- 客户离线、没有可用连接:`waiting_connection`; +- 其他链路已经写出且仍在 ACK 窗口:`waiting_external_ack`; +- 本任务已经写出:`waiting_ack`。 + +等待项不计入跳过数或失败数。等待结束后根据真实状态继续、成功或失败。 + +### 7.3 失败 + +- 调用 Gateway 发生非等待型错误; +- Gateway 明确拒绝或无法安全投递; +- 本任务写出后 ACK 超时; +- 客户 ACK 非零; +- 客户 ACK 无法关联原消息。 + +失败原因必须保留原始错误,同时归一为可统计的失败类别。 + +## 8. 安全阈值与自动暂停 + +默认安全策略: + +- 每应用默认 10 条/秒; +- 连续失败阈值默认 10 条; +- ACK 超时和 ACK 拒绝纳入失败阈值; +- 达到阈值时,在认领下一条之前将任务原子改为 `paused`; +- 写操作日志,记录任务号、失败类别、连续失败数、阈值和暂停时间; +- 页面展示明确的自动暂停原因; +- 人工继续后从待处理/等待项继续,不重放已成功项。 + +如果多个应用同时执行,连续失败计数至少应按应用隔离,避免一个客户应用故障暂停其他正常应用;第一版若选择整任务暂停,也必须在设计和页面中明确,且仍要保存触发暂停的应用。 + +## 9. API 设计 + +| 方法 | 路径 | 用途 | +| --- | --- | --- | +| POST | `/admin/operations/downstream-requeue-tasks/preview` | 按当前筛选条件生成预检快照 | +| POST | `/admin/operations/downstream-requeue-tasks` | 使用预检快照创建任务 | +| GET | `/admin/operations/downstream-requeue-tasks` | 分页查询任务列表,支持状态筛选 | +| GET | `/admin/operations/downstream-requeue-tasks/:id` | 查询任务汇总 | +| GET | `/admin/operations/downstream-requeue-tasks/:id/items` | 分页查询完整任务项及原因 | +| POST | `/admin/operations/downstream-requeue-tasks/:id/pause` | 暂停 | +| POST | `/admin/operations/downstream-requeue-tasks/:id/resume` | 继续 | +| POST | `/admin/operations/downstream-requeue-tasks/:id/terminate` | 终止 | + +创建接口不得仅信任前端回传的筛选条件和时间。建议预检生成短期有效、服务端签名的 `previewToken`,绑定筛选条件、`snapshotAt` 和操作人;创建时校验令牌,防止绕过页面篡改范围。 + +## 10. 验收用例 + +| 编号 | 验收点 | 预期 | +| --- | --- | --- | +| DRQ-001 | 筛选快照 | 企业、应用、类型、状态、日期、关键词均生效;分页无关;快照后新记录不进入 | +| DRQ-002 | 预检口径 | 命中、可重投、跳过和状态分布严格基于当前筛选条件 | +| DRQ-003 | 状态白名单 | 只物化 `pending/failed/unconfirmed/rejected`;拒绝 `delivered/awaiting_ack` | +| DRQ-004 | 每应用限速 | 多应用任务中每个应用独立达到配置速度,任意扫描重叠都不超速 | +| DRQ-005 | 客户离线 | 进入等待连接,不记失败或跳过,恢复连接后继续 | +| DRQ-006 | ACK 闭环 | 写出只进入等待;有效 ACK 成功;超时、拒绝、无效 Msg_Id 失败 | +| DRQ-007 | 自动暂停 | 立即失败与 ACK 失败均计入阈值;达到阈值自动暂停并审计 | +| DRQ-008 | 并发幂等 | 多实例同时扫描,同一投递最多调用一次 Gateway | +| DRQ-009 | 重启恢复 | 在 `processing`、等待连接、等待 ACK 三种阶段重启 API,任务均可继续且不重复成功项 | +| DRQ-010 | 人工控制 | 暂停、继续、终止与执行并发时状态正确;终止不撤回已写出项 | +| DRQ-011 | 完整可查 | 任务列表和任务项均可分页查询全部历史数据,不限最近 10/50 条 | +| DRQ-012 | 操作审计 | 创建、暂停、继续、终止、自动暂停均有操作人/触发源和完整上下文 | + +## 11. 当前实现符合性审计 + +当前实现的主要文件: + +- `api/src/send-chain/send-downstream-requeue-task.service.ts` +- `api/src/send-chain/send-downstream-state.service.ts` +- `api/src/operations/admin-operations.controller.ts` +- `api/prisma/schema.prisma` +- `src/apps/admin/AdminDownstreamDeliveriesPage.tsx` +- `api/src/send-chain/send-downstream-requeue-task.service.spec.ts` + +### 11.1 已符合或基本符合 + +| 设计项 | 结论 | 当前证据 | +| --- | --- | --- | +| 真实持久化 | 符合 | 已有任务表、任务项表和 migration,不使用前端本地状态代替任务 | +| 快照时间上限 | 基本符合 | 创建时按 `snapshotAt` 限制 `createdAt`,快照后记录不物化 | +| 后台状态白名单 | 符合 | `REPLAYABLE_STATUSES` 为 `pending/failed/unconfirmed/rejected` | +| 禁止后台任务处理已确认/等待 ACK 筛选 | 符合 | 创建接口明确拒绝 `delivered/awaiting_ack` | +| 原因校验 | 符合 | 少于 5 个字拒绝 | +| 任务项幂等 | 符合 | 数据库唯一约束 `taskId + deliveryId` | +| 执行前任务项认领 | 基本符合 | 通过 `status=queued` 的条件更新认领为 `processing` | +| 执行前复核 | 基本符合 | 重新查询投递、应用、payload 和其他任务状态 | +| 成功 ACK 不再发送 | 基本符合 | 已确认投递会跳过,其他任务 `success` 也会阻止调用 | +| 人工控制 | 基本符合 | 已有暂停、继续、终止接口和页面按钮 | +| 核心操作日志 | 基本符合 | 创建、暂停、继续、终止、自动暂停均写日志 | +| 投递记录分页 | 符合 | 页面支持每页 10/25/50 条并回到第一页 | + +### 11.2 明确偏差与缺口 + +| 优先级 | 偏差 | 当前实现 | 与设计冲突及风险 | +| --- | --- | --- | --- | +| P0 | API 重启后 `processing` 项无法恢复 | 只扫描 `queued`,没有 `processing` 认领租约或超时回收 | 进程在认领后中断会使任务项永久卡住,任务永久 `running`,不满足重启续跑 | +| P0 | ACK 失败不参与自动暂停阈值 | `reconcileWaiting`把 ACK 超时/拒绝改为 `failed`,但不增加`consecutiveFailures`;进入 `waiting_ack` 时反而立即清零 | 客户持续拒绝或 ACK 超时不会触发安全暂停,可能持续向故障客户重投 | +| P0 | 限速不是按应用,且可能被并发扫描突破 | 每个任务每轮最多取 `min(10, ratePerSecond)`;没有按`applicationId`分组,也没有分布式速率令牌;定时器可能重叠执行 | 不符合“每应用10条/秒”;配置20实际单轮最多10,多扫描器/多实例又可能超过配置 | +| P0 | 客户离线被记为任务项失败 | Gateway未写出时,底层通常返回`pending/failed`,任务层将非`awaiting_ack/delivered`结果直接记为`failed` | 不符合“客户离线属于等待”;会错误增加失败数并可能自动暂停 | +| P1 | 预检不严格遵守当前状态筛选 | 计算`matchedCount/statusCounts`时强制把状态改成`all` | 用户选择“待投递”时,命中数和状态分布仍可能包含其他状态;预检范围表达失真 | +| P1 | 页面缺少企业筛选 | 后端类型支持`tenantId`,但页面只有应用筛选,`currentTaskFilter`不传企业 | 未完整实现已确认的“企业 + 应用”筛选范围 | +| P1 | 任务列表不可完整查询 | 页面固定请求第1页、每页10条,没有任务分页和状态筛选 | 第11个以后历史任务在页面不可达,不满足任务列表要求 | +| P1 | 任务详情不可完整查询 | 详情接口固定返回最近50个任务项,没有任务项分页接口 | 大任务无法核对全部失败、跳过和等待项,难以验收和审计 | +| P1 | 扫描器缺少任务级/应用级分布式锁 | `setInterval`直接调用扫描,`runScan`可被重叠触发,多实例也会同时处理同一任务 | 虽有任务项条件认领可降低单项重复,但无法保证整体限速和连续失败统计一致 | +| P1 | 连续失败清零时点错误 | Gateway写出进入`waiting_ack`即把连续失败清零 | 写出不是业务成功;应等有效 ACK 后再清零 | +| P1 | 预检与创建没有不可篡改绑定 | 创建接口信任前端回传的`filter + snapshotAt`,没有预检令牌或服务端预检记录 | 可绕过页面修改筛选范围;虽仍受状态白名单和10万条上限保护,但不等于复用原预检结果 | +| P1 | 自动化覆盖远低于设计风险 | 专项仅4个测试:预检计数、参数拒绝、活动任务冲突、已确认跳过 | 未覆盖限速、多应用、离线等待、ACK阈值、重启恢复、并发扫描、完整分页及所有审计动作 | +| P2 | 状态和详情表达偏内部化 | 任务列表/详情直接展示英文状态;任务详情只展示汇总和最近项 | 运营人员不易区分等待连接、等待外部ACK、任务写出等待ACK等状态 | +| P2 | 终止后的未处理数未进入常规汇总 | 终止把`queued`改为`unprocessed`,但任务汇总只保存成功/失败/跳过/等待 | 任务进度分子可能小于总数,页面没有单独解释未处理数量 | +| P2 | 已确认跳过原因存在口径混用 | `waiting_external_ack`最终由其他链路成功后记为“跳过:已由其他投递链路完成” | 可以接受为“本任务未调用Gateway”,但详情必须明确这是外部链路成功,不应让用户误以为业务未处理 | + +### 11.3 综合结论 + +当前实现完成了数据库模型、基本预检、任务创建、任务项认领、Gateway调用、ACK结果回看和人工控制的骨架,但不能判定为严格按照 2026-08-12 方案完成。 + +尤其以下四项属于设计中的安全核心,当前存在实质缺口: + +1. `processing` 任务项缺少重启恢复; +2. ACK 超时/拒绝未进入自动暂停阈值; +3. 每应用限速没有实现,且扫描重叠可能突破限速; +4. 客户离线没有作为等待状态处理。 + +在上述 P0 修复并通过真实 PostgreSQL、Gateway、客户 ACK 链路测试前,不应把该后台任务认定为完整满足事故批量恢复方案。 + +## 12. 建议整改顺序 + +1. 先补 `processing` 租约恢复、扫描分布式锁和每应用速率令牌; +2. 重构任务项等待/失败状态,让离线、外部 ACK、本任务 ACK 三类等待分开; +3. 将 ACK 超时、拒绝、无效 ACK 纳入按应用安全阈值,并把清零时点改为有效 ACK; +4. 修正预检状态口径,增加企业筛选,并用服务端令牌绑定预检与创建; +5. 增加任务列表和任务项分页、中文状态及完整原因查询; +6. 补齐 DRQ-001~DRQ-012 自动化,并使用真实 PostgreSQL、Redis、Gateway 与本地客户连接完成专项验收。 + +整改应作为独立需求进行,不在未经授权的情况下直接修改或发布现有批量重投逻辑。 diff --git a/docs/first-version-development-requirements.md b/docs/first-version-development-requirements.md index 8dfe953..6535d83 100644 --- a/docs/first-version-development-requirements.md +++ b/docs/first-version-development-requirements.md @@ -2027,6 +2027,17 @@ - 通道选择控件必须为通用下拉可搜索控件,支持按通道名称或编码搜索;“全部通道”表示不按通道限制,点击“重置”必须同时清空通道组名称和通道条件并回到第一页。 # 下游投递后台重投任务(2026-08-12) +## 完整设计口径与安全整改(2026-08-13) + +1. 后台任务按企业、应用、投递类型、状态、北京时间创建日期和关键词冻结筛选快照,分页、每页条数和当前页勾选不属于范围。预检的命中数、可重投数、规则跳过数、状态分布、涉及应用和最早记录必须严格基于当前筛选条件,不得把单一状态擅自扩为全部状态。 +2. 预检返回短期有效且绑定当前操作人、筛选条件和 `snapshotAt` 的服务端签名凭证;创建接口只接受该凭证、原因、速度和安全阈值,不再信任前端重传的范围。创建时后端按签名快照重新物化真实 PostgreSQL 任务项。 +3. 第一版后台任务仅处理 `pending/failed/unconfirmed/rejected`;`delivered/awaiting_ack` 不得进入批量任务。跳过严格表示本任务没有调用 Gateway;客户离线进入 `waiting_connection`,其他链路等待 ACK 进入 `waiting_external_ack`,本任务写出后进入 `waiting_ack`,三类等待均不计失败或跳过。 +4. 限速以企业应用为维度,使用数据库原子秒级窗口在多实例、扫描重叠和执行耗时变化下保持每应用不超过配置速度。任务扫描使用数据库租约;`processing` 项使用认领租约,API 中断后过期恢复为 `queued` 并重新复核,已成功 ACK 的项目不得重放。 +5. 连续失败按应用隔离统计。Gateway 立即失败、ACK 超时、ACK 拒绝和无法安全关联均计入阈值;只有有效 ACK 成功才清零。达到阈值前原子暂停整个任务,记录触发应用、失败数、阈值和暂停时间,人工继续后从未完成项恢复。 +6. 任务列表支持状态筛选和真实分页,展示任务号、创建时间、企业/应用、原因、中文状态、总数、成功、失败、跳过、等待、创建人及进度。任务详情展示筛选快照、时间、安全参数、完整结果汇总和任务项分页;任务项可按中文结果、消息 ID、错误或跳过原因查询,不得仅返回最近 50 项。 +7. 暂停、继续和终止与执行器并发时,认领及调用 Gateway 前必须复核任务状态;终止只把尚未开始及等待连接的项目置为 `unprocessed`,不得撤回已写出内容。创建、暂停、继续、终止和自动暂停均写操作日志。 +8. 数据继续使用真实 PostgreSQL、Gateway、客户连接和 ACK;不得使用 mock、静态数据或 localStorage 代替任务状态。任务和任务项保留历史,下游投递及 ACK 证据不得级联删除。 + 1. 运营端“下游投递记录”必须同时保留单条重投、当前页勾选批量重投,并新增“按筛选条件重投”;分页支持每页 `10/25/50` 条,切换后回到第一页并重新查询真实后端。 2. 后台任务使用当前企业、应用、投递类型、状态、创建日期和关键词的后端筛选快照,分页不属于任务范围;任务创建时固定 `snapshotAt`,之后产生的记录不得被卷入。 3. 第一版只允许 `pending/failed/unconfirmed/rejected`,不支持批量重投客户端已确认的 `delivered`,`awaiting_ack` 不得并发重投。创建前必须真实预检命中、可重投、跳过和状态分布,原因必填。 diff --git a/docs/system-functional-test-cases.md b/docs/system-functional-test-cases.md index 2bc491f..93d668d 100644 --- a/docs/system-functional-test-cases.md +++ b/docs/system-functional-test-cases.md @@ -4596,6 +4596,20 @@ npm run verify:phase8 | TC-DOWNSTREAM-REQUEUE-TASK-007 | 任务控制 | 待执行/执行中任务可暂停、继续、终止;终止不撤回已写出消息。 | | TC-DOWNSTREAM-REQUEUE-TASK-008 | 审计 | 创建、暂停、继续、终止和自动暂停记录操作人、筛选快照、原因和结果。 | | TC-DOWNSTREAM-PAGE-SIZE-001 | 分页数量 | 可选 10/25/50;切换回第一页,后端返回对应条数,总数和筛选条件保持一致。 | + +## 下游投递后台重投任务安全整改专项用例(2026-08-13) + +| 编号 | 场景 | 预期 | +|---|---|---| +| TC-DOWNSTREAM-REQUEUE-TASK-009 | 单一状态严格预检 | 选择待投递后,命中数、状态分布和可重投数只基于 `pending`;企业、应用、类型、日期、关键词同时生效。 | +| TC-DOWNSTREAM-REQUEUE-TASK-010 | 预检签名绑定 | 篡改筛选、快照、操作人、签名或使用过期凭证均拒绝;合法凭证由后端按原快照物化任务项。 | +| TC-DOWNSTREAM-REQUEUE-TASK-011 | processing 租约恢复 | API 在认领后中断,租约过期后项目回到队列并重新复核;已成功项不重复调用 Gateway。 | +| TC-DOWNSTREAM-REQUEUE-TASK-012 | 每应用原子限速 | 多任务、多实例和重叠扫描同时执行时,每个应用每秒消耗不超过配置值;不同应用互不阻塞。 | +| TC-DOWNSTREAM-REQUEUE-TASK-013 | 离线等待与恢复 | 客户无 connected 连接时进入等待连接,不增加失败/跳过;连接恢复后回队列继续。 | +| TC-DOWNSTREAM-REQUEUE-TASK-014 | ACK 阈值与清零 | 写出进入等待 ACK 不清零;ACK 超时/拒绝按应用累加,达到阈值自动暂停并审计;有效 ACK 才清零。 | +| TC-DOWNSTREAM-REQUEUE-TASK-015 | 列表完整分页 | 任务列表支持状态和分页,第 11 条以后可访问,中文状态、创建人、原因、进度和各结果数准确。 | +| TC-DOWNSTREAM-REQUEUE-TASK-016 | 完整任务项查询 | 任务项支持分页、结果及关键词查询,等待连接/外部 ACK/本任务 ACK、跳过、失败和未处理均中文展示并保留原因。 | +| TC-DOWNSTREAM-REQUEUE-TASK-017 | 终止并发边界 | 终止后未认领和等待连接项置为未处理;处理中或已写出项不撤回;执行器不再认领新项。 | # 2026-08-13 HTTP 与 Gateway 报文容量专项用例 | 用例编号 | 优先级 | 验证内容 | 预期结果 | diff --git a/docs/testing-progress.md b/docs/testing-progress.md index 254eb91..7d2354c 100644 --- a/docs/testing-progress.md +++ b/docs/testing-progress.md @@ -3557,3 +3557,14 @@ git diff --check - 企业应用管理的状态、到达率、单价列宽均从130px缩至104px,字段和操作未隐藏。全局审计运营商数据展示后,监控、通道/通道组、应用路由、签名和报备、清退、短信审核、批次号码、发送记录及客户端发送详情统一复用`CarrierTag`;筛选选项、图表图例、导出和说明文字保留纯文本,三网通道展开为三个标签。 - API全量37套/460项、充值和HTTP容量专项14/14、前后端TypeScript、Vite 8.1.5生产构建、Gateway全量测试与`go vet`、4份Gateway队列契约、企业应用R11、短信记录R11、企业签名R4结构契约和`git diff --check`均通过;Vite仅保留既有约2.07MiB单chunk提示。企业签名R4哈希只按本次经验证的运营商标签JSX同步更新,未放宽模块边界。浏览器连接本地页面时既有登录会话已失效,未输入账号密码或验证码,因此页面级视觉验收未完成。 - 本轮与此前未提交的HTTP/Gateway容量修复一并进入待提交范围,未发送、补发或重投短信,未修改通道配置、余额、客户连接或生产数据;受保护缓存、`outputs/`和空文件`=`不删除、不提交、不归因。 + +# 2026-08-13 下游投递后台重投任务安全整改与页面优化(本地未提交、未发布) + +- 已将`docs/downstream-requeue-task-design-20260812.md`确认的完整设计口径合并进平台需求和系统测试文档;独立设计稿继续作为审计输入保留,不以现有实现反向改写设计结论。 +- 预检严格使用企业、应用、类型、状态、北京时间日期和关键词筛选,不再把单一状态扩为全部状态;预检生成绑定操作人、筛选和`snapshotAt`的15分钟服务端签名凭证,创建接口不再接受前端重传的范围,仅按签名快照在真实PostgreSQL重新物化任务项。 +- 新增任务扫描数据库租约、`processing`两分钟认领租约恢复和每应用秒级原子限速窗口;多实例/重叠扫描不能突破每应用配置速率。客户无connected连接时进入`waiting_connection`,连接恢复后回队列;写出进入`waiting_ack`不清零,只有有效ACK清零,ACK超时/拒绝按应用累计并达到阈值后自动暂停和审计。 +- 新增完整任务项分页接口,支持状态、消息ID、错误和跳过原因查询;终止将尚未开始和等待连接项置为`unprocessed`,不撤回处理中或已写出项。任务列表增加状态筛选和分页,页面补齐企业—应用联动、中文状态、创建人、原因、进度及结果数;创建弹窗重做命中/可重投/跳过/应用信息层级和安全边界,详情可查询全部历史项而非最近50条。 +- 新增向前兼容migration`20260813150000_harden_downstream_requeue_tasks`,只增加任务应用级失败计数、扫描租约和速率窗口表,不删除或改写历史投递/ACK。已在明确指向`localhost:5432/cmpp_platform`的本地真实PostgreSQL应用,当前本地87条migration且schema up to date;未连接或修改预生产数据库。 +- 后台任务专项7/7通过,覆盖严格预检、签名凭证、服务端物化、完整分页、processing恢复、离线等待和ACK自动暂停;API全量37套/463项、Prisma format/validate、API及前端TypeScript正式构建、Vite 8.1.5生产构建(2541 modules)均通过,仅保留既有约2.08MB单chunk提示。新增任务明细分页使`SendChainService`稳定门面方法由104增至105,对应R10结构契约已按真实新增接口精确同步并通过;`git diff --check`通过。 +- 浏览器优先接管现有会话后确认本地真实API与PostgreSQL服务可访问,但现有浏览器没有已登录页面,访问下游投递页被正常引导到运营端登录;本轮未读取、猜测或重置凭据,也未绕过登录,因此创建弹窗、任务列表和详情的登录后视觉验收留待具备既有会话时补充。 +- 本轮没有创建真实重投任务、调用真实Gateway重投、发送短信、修改通道配置、余额或客户连接;不提交、不推送、不部署。既有`api/tsconfig.build.tsbuildinfo`、`tsconfig.tsbuildinfo`、`outputs/`和空文件`=`继续保护。 diff --git a/src/api/admin/operations.api.ts b/src/api/admin/operations.api.ts index a082d3e..46226b4 100644 --- a/src/api/admin/operations.api.ts +++ b/src/api/admin/operations.api.ts @@ -1,5 +1,5 @@ import { request, requestBlob, requestForm, withQuery } from '../core/httpClient'; -import type { BatchRequeueResponse, BatchTaskMessagePage, DailyProfitReport, DailyQualityReport, DailyReconciliationReport, DashboardResponse, DownstreamDeliveryDashboard, DownstreamDeliveryRecord, DownstreamRecoveryStatusExportQuery, DownstreamRecoveryStatusResponse, DownstreamRequeueFilter, DownstreamRequeuePreview, DownstreamRequeueTask, GatewayDownstreamRecoveryStatus, GatewaySubmitException, GatewaySubmitExceptionResponse, OperationLogResponse, PagedResponse, PagedResult, PendingAuditCounts, ProfitReportSummary, ProtocolInteractionLogResponse, QualityReportSummary, ReceiptAnomalyResponse, ReconciliationReportSummary, SendQualityResponse, SignatureChannelQualityResponse, SmsBatchTask, SmsMessageRecord, SmsMessageSegmentAudit, SmsUplinkMessage, SystemLogExportResult } from '../types'; +import type { BatchRequeueResponse, BatchTaskMessagePage, DailyProfitReport, DailyQualityReport, DailyReconciliationReport, DashboardResponse, DownstreamDeliveryDashboard, DownstreamDeliveryRecord, DownstreamRecoveryStatusExportQuery, DownstreamRecoveryStatusResponse, DownstreamRequeueFilter, DownstreamRequeuePreview, DownstreamRequeueTask, DownstreamRequeueTaskItem, GatewayDownstreamRecoveryStatus, GatewaySubmitException, GatewaySubmitExceptionResponse, OperationLogResponse, PagedResponse, PagedResult, PendingAuditCounts, ProfitReportSummary, ProtocolInteractionLogResponse, QualityReportSummary, ReceiptAnomalyResponse, ReconciliationReportSummary, SendQualityResponse, SignatureChannelQualityResponse, SmsBatchTask, SmsMessageRecord, SmsMessageSegmentAudit, SmsUplinkMessage, SystemLogExportResult } from '../types'; // Read-heavy operations endpoints are isolated from configuration mutations. export const adminOperationsApi = { @@ -74,11 +74,13 @@ export const adminOperationsApi = { request('/admin/operations/downstream-deliveries/requeue', { method: 'POST', body: JSON.stringify({ ids }) }), previewDownstreamRequeueTask: (filter: DownstreamRequeueFilter) => request('/admin/operations/downstream-requeue-tasks/preview', { method: 'POST', body: JSON.stringify({ filter }) }), - createDownstreamRequeueTask: (body: { filter: DownstreamRequeueFilter; snapshotAt: string; reason: string; ratePerSecond: number; consecutiveFailureLimit: number }) => + createDownstreamRequeueTask: (body: { previewToken: string; reason: string; ratePerSecond: number; consecutiveFailureLimit: number }) => request('/admin/operations/downstream-requeue-tasks', { method: 'POST', body: JSON.stringify(body) }), listDownstreamRequeueTasks: (query: { status?: string; page?: number; pageSize?: number } = {}) => request>(withQuery('/admin/operations/downstream-requeue-tasks', query)), getDownstreamRequeueTask: (id: string) => request(`/admin/operations/downstream-requeue-tasks/${id}`), + listDownstreamRequeueTaskItems: (id: string, query: { status?: string; keyword?: string; page?: number; pageSize?: number } = {}) => + request>(withQuery(`/admin/operations/downstream-requeue-tasks/${id}/items`, query)), changeDownstreamRequeueTaskStatus: (id: string, action: 'pause' | 'resume' | 'terminate') => request(`/admin/operations/downstream-requeue-tasks/${id}/${action}`, { method: 'POST', body: JSON.stringify({}) }), }; diff --git a/src/api/types/operations.ts b/src/api/types/operations.ts index d313fa6..8a6eff7 100644 --- a/src/api/types/operations.ts +++ b/src/api/types/operations.ts @@ -495,6 +495,7 @@ export type DownstreamRequeueFilter = { export type DownstreamRequeuePreview = { snapshotAt: string; + previewToken: string; matchedCount: number; replayableCount: number; skippedCount: number; @@ -521,11 +522,23 @@ export type DownstreamRequeueTask = { createdAt: string; startedAt?: string | null; finishedAt?: string | null; + pausedAt?: string | null; tenant?: TenantOption | null; application?: EnterpriseApplication | null; createdBy?: { id: string; displayName: string; username: string } | null; itemCounts?: Record; - recentItems?: Array<{ id: string; status: string; skipReason?: string | null; errorMessage?: string | null; delivery: { messageId?: string | null; deliveryType: string; status: string; lastError?: string | null } }>; +}; + +export type DownstreamRequeueTaskItem = { + id: string; + status: string; + previousStatus: string; + skipReason?: string | null; + errorMessage?: string | null; + claimedAt?: string | null; + completedAt?: string | null; + updatedAt: string; + delivery: { messageId?: string | null; deliveryType: string; status: string; lastError?: string | null }; }; export type DownstreamDeliveryDashboard = { diff --git a/src/apps/admin/AdminDownstreamDeliveriesPage.tsx b/src/apps/admin/AdminDownstreamDeliveriesPage.tsx index 86f26eb..86b03d5 100644 --- a/src/apps/admin/AdminDownstreamDeliveriesPage.tsx +++ b/src/apps/admin/AdminDownstreamDeliveriesPage.tsx @@ -1,6 +1,6 @@ import { useCallback, useEffect, useMemo, useRef, useState } from 'react'; import { AlertTriangle, BarChart3, CheckCircle2, Eye, RefreshCw, Search, TimerReset } from 'lucide-react'; -import { adminApi, type DownstreamDeliveryDashboard, type DownstreamDeliveryRecord, type DownstreamRequeuePreview, type DownstreamRequeueTask, type EnterpriseApplication } from '@/api/adminApi'; +import { adminApi, type DownstreamDeliveryDashboard, type DownstreamDeliveryRecord, type DownstreamRequeuePreview, type DownstreamRequeueTask, type DownstreamRequeueTaskItem, type EnterpriseApplication, type TenantOption } from '@/api/adminApi'; import { Breadcrumb, Button, DateRangeInput, Input, Modal, Pagination, Select, Tag, type DateRangeValue } from '@/components/ui'; import { formatDateTime } from '@/utils/dateTime'; @@ -42,6 +42,24 @@ const attemptStatusLabel: Record = { failed: '投递失败', }; +const requeueTaskStatusLabel: Record = { + queued: '排队中', running: '执行中', paused: '已暂停', completed: '已完成', + partial_completed: '部分完成', terminated: '已终止', +}; + +const requeueItemStatusLabel: Record = { + queued: '排队中', processing: '处理中', waiting_connection: '等待连接', + waiting_external_ack: '等待其他链路确认', waiting_ack: '等待客户确认', success: '成功', + failed: '失败', skipped: '跳过', unprocessed: '未处理', +}; + +function requeueTone(status: string) { + if (status === 'completed' || status === 'success') return 'success' as const; + if (status === 'partial_completed' || status === 'failed') return 'danger' as const; + if (status === 'paused' || status === 'skipped' || status === 'unprocessed') return 'warning' as const; + return 'info' as const; +} + type RequeueTarget = | { kind: 'single'; record: DownstreamDeliveryRecord } | { kind: 'batch'; ids: string[] }; @@ -148,10 +166,12 @@ export function AdminDownstreamDeliveriesPage() { const [records, setRecords] = useState([]); const [dashboard, setDashboard] = useState(null); const [applications, setApplications] = useState([]); + const [tenants, setTenants] = useState([]); const [keyword, setKeyword] = useState(''); const [status, setStatus] = useState('all'); const [deliveryType, setDeliveryType] = useState('all'); const [applicationId, setApplicationId] = useState('all'); + const [tenantId, setTenantId] = useState('all'); const [page, setPage] = useState(1); const [pageSize, setPageSize] = useState(10); const [total, setTotal] = useState(0); @@ -170,28 +190,39 @@ export function AdminDownstreamDeliveriesPage() { const [taskRate, setTaskRate] = useState(10); const [taskCreateBusy, setTaskCreateBusy] = useState(false); const [requeueTasks, setRequeueTasks] = useState([]); + const [requeueTaskStatus, setRequeueTaskStatus] = useState('all'); + const [requeueTaskPage, setRequeueTaskPage] = useState(1); + const [requeueTaskTotal, setRequeueTaskTotal] = useState(0); const [selectedTask, setSelectedTask] = useState(null); + const [taskItems, setTaskItems] = useState([]); + const [taskItemStatus, setTaskItemStatus] = useState('all'); + const [taskItemKeyword, setTaskItemKeyword] = useState(''); + const [taskItemAppliedKeyword, setTaskItemAppliedKeyword] = useState(''); + const [taskItemPage, setTaskItemPage] = useState(1); + const [taskItemTotal, setTaskItemTotal] = useState(0); const currentTaskFilter = useCallback(() => ({ keyword: keyword || undefined, status, deliveryType, + tenantId, applicationId, createdAtFrom: dateRange.start, createdAtTo: dateRange.end, - }), [applicationId, dateRange.end, dateRange.start, deliveryType, keyword, status]); + }), [applicationId, dateRange.end, dateRange.start, deliveryType, keyword, status, tenantId]); const loadRequeueTasks = useCallback(() => { - adminApi.listDownstreamRequeueTasks({ page: 1, pageSize: 10 }) - .then((response) => setRequeueTasks(response.items)) - .catch(() => undefined); - }, []); + adminApi.listDownstreamRequeueTasks({ status: requeueTaskStatus, page: requeueTaskPage, pageSize: 10 }) + .then((response) => { setRequeueTasks(response.items); setRequeueTaskTotal(response.total); }) + .catch((failure: Error) => setError(failure.message || '后台重投任务加载失败')); + }, [requeueTaskPage, requeueTaskStatus]); const loadData = useCallback(() => { setLoading(true); Promise.all([ adminApi.getDownstreamDeliveryDashboard({ applicationId, + tenantId, deliveryType, createdAtFrom: dateRange.start, createdAtTo: dateRange.end, @@ -201,24 +232,27 @@ export function AdminDownstreamDeliveriesPage() { status, deliveryType, applicationId, + tenantId, page, pageSize, createdAtFrom: dateRange.start, createdAtTo: dateRange.end, }), adminApi.listEnterpriseApplications(), + adminApi.listTenants(), ]) - .then(([dashboardResponse, response, apps]) => { + .then(([dashboardResponse, response, apps, tenantOptions]) => { setDashboard(dashboardResponse); setRecords(response.items); setTotal(response.total); setApplications(apps); + setTenants(tenantOptions); setSelectedIds((current) => current.filter((id) => response.items.some((item) => item.id === id))); setError(''); }) .catch((failure: Error) => setError(failure.message || '下游投递记录加载失败')) .finally(() => setLoading(false)); - }, [applicationId, dateRange.end, dateRange.start, deliveryType, keyword, page, pageSize, status]); + }, [applicationId, dateRange.end, dateRange.start, deliveryType, keyword, page, pageSize, status, tenantId]); useEffect(() => { loadData(); @@ -247,7 +281,7 @@ export function AdminDownstreamDeliveriesPage() { if (!taskPreview || taskReason.trim().length < 5) return; setTaskCreateBusy(true); try { - await adminApi.createDownstreamRequeueTask({ filter: taskPreview.filter, snapshotAt: taskPreview.snapshotAt, reason: taskReason.trim(), ratePerSecond: taskRate, consecutiveFailureLimit: 10 }); + await adminApi.createDownstreamRequeueTask({ previewToken: taskPreview.previewToken, reason: taskReason.trim(), ratePerSecond: taskRate, consecutiveFailureLimit: 10 }); setTaskPreview(null); loadRequeueTasks(); } catch (failure) { @@ -257,6 +291,28 @@ export function AdminDownstreamDeliveriesPage() { } }; + const filteredApplications = useMemo( + () => applications.filter((item) => tenantId === 'all' || item.tenantId === tenantId), + [applications, tenantId], + ); + + const openTaskDetail = async (id: string) => { + try { + const task = await adminApi.getDownstreamRequeueTask(id); + setSelectedTask(task); + setTaskItemStatus('all'); setTaskItemKeyword(''); setTaskItemAppliedKeyword(''); setTaskItemPage(1); + } catch (failure) { setError(failure instanceof Error ? failure.message : '任务详情加载失败'); } + }; + + const loadTaskItems = useCallback(() => { + if (!selectedTask) return; + adminApi.listDownstreamRequeueTaskItems(selectedTask.id, { status: taskItemStatus, keyword: taskItemAppliedKeyword || undefined, page: taskItemPage, pageSize: 20 }) + .then((response) => { setTaskItems(response.items); setTaskItemTotal(response.total); }) + .catch((failure: Error) => setError(failure.message || '任务明细加载失败')); + }, [selectedTask, taskItemAppliedKeyword, taskItemPage, taskItemStatus]); + + useEffect(() => { loadTaskItems(); }, [loadTaskItems]); + const replayableStatuses = useMemo(() => new Set(['pending', 'delivered', 'failed', 'unconfirmed', 'rejected']), []); const selectableIds = useMemo( () => records.filter((item) => replayableStatuses.has(item.status)).map((item) => item.id), @@ -375,6 +431,12 @@ export function AdminDownstreamDeliveriesPage() {
setKeyword(event.target.value)} placeholder="请输入关键字" value={keyword} /> { setDateRange(value); setPage(1); }} value={dateRange} /> + ({ label: item.name, value: item.id })), + ...filteredApplications.map((item) => ({ label: item.name, value: item.id })), ]} value={applicationId} onChange={(event) => { @@ -424,6 +486,7 @@ export function AdminDownstreamDeliveriesPage() { setKeyword(''); setStatus('all'); setDeliveryType('all'); + setTenantId('all'); setApplicationId('all'); setDateRange(recentSevenDays()); setPage(1); @@ -615,18 +678,33 @@ export function AdminDownstreamDeliveriesPage() {
-
-

后台重投任务

大批量事故恢复按筛选快照分批执行,不受投递记录分页影响。

- +
+

后台重投任务

按筛选快照安全恢复,支持暂停、继续、终止和完整结果追踪。

+
+ setTaskRate(Number(event.target.value))} options={[{ label: '平稳(每应用10条/秒)', value: '10' }, { label: '快速(每应用20条/秒)', value: '20' }, { label: '低速(每应用5条/秒)', value: '5' }]} />