From f6ef9c9f2855ba8e4e01b941c923509492600707 Mon Sep 17 00:00:00 2001 From: sunlei Date: Sun, 26 Jul 2026 19:17:31 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=A2=9E=E5=8A=A0=E9=80=BB=E8=BE=91?= =?UTF-8?q?=E7=AB=AF=E5=8F=A3=E8=BD=AC=E5=8F=91=E7=BB=84=E6=8E=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../admin-platform-config.module.ts | 4 + .../network-management.service.ts | 152 +-- .../network-port-forward-group.controller.ts | 174 +++ .../network-port-forward-group.dto.ts | 126 ++ .../network-port-forward-group.service.ts | 1028 +++++++++++++++++ .../network-management.controller.spec.ts | 24 + .../network-management.service.spec.ts | 101 +- ...work-port-forward-group.controller.spec.ts | 242 ++++ ...network-port-forward-group.service.spec.ts | 491 ++++++++ 9 files changed, 2194 insertions(+), 148 deletions(-) create mode 100644 src/modules/admin/platform-config/network-management/network-port-forward-group.controller.ts create mode 100644 src/modules/admin/platform-config/network-management/network-port-forward-group.dto.ts create mode 100644 src/modules/admin/platform-config/network-management/network-port-forward-group.service.ts create mode 100644 test/admin/network-management/network-port-forward-group.controller.spec.ts create mode 100644 test/admin/network-management/network-port-forward-group.service.spec.ts diff --git a/src/modules/admin/platform-config/admin-platform-config.module.ts b/src/modules/admin/platform-config/admin-platform-config.module.ts index f846025..1112212 100644 --- a/src/modules/admin/platform-config/admin-platform-config.module.ts +++ b/src/modules/admin/platform-config/admin-platform-config.module.ts @@ -18,6 +18,8 @@ import { NetworkManagementController } from '@/modules/admin/platform-config/net import { NetworkManagementEventStreamService } from '@/modules/admin/platform-config/network-management/network-management-event-stream.service'; import { NetworkPortForward } from '@/modules/admin/platform-config/network-management/network-management.entity'; import { NetworkPortForwardGroup } from '@/modules/admin/platform-config/network-management/network-port-forward-group.entity'; +import { NetworkPortForwardGroupController } from '@/modules/admin/platform-config/network-management/network-port-forward-group.controller'; +import { NetworkPortForwardGroupService } from '@/modules/admin/platform-config/network-management/network-port-forward-group.service'; import { NetworkManagementService } from '@/modules/admin/platform-config/network-management/network-management.service'; import { NetworkStunMessageSourceAdapter } from '@/modules/admin/platform-config/network-management/network-stun-message-source.adapter'; import { NetworkTcpReleasePolicyService } from '@/modules/admin/platform-config/network-management/network-tcp-release-policy.service'; @@ -55,6 +57,7 @@ export const ADMIN_PLATFORM_CONFIG_DIRECT_CONTROLLERS = [ AdminTimezoneController, EnvironmentDashboardController, NetworkManagementController, + NetworkPortForwardGroupController, ]; export const ADMIN_PLATFORM_CONFIG_IMPORTED_CONTROLLERS = [ @@ -88,6 +91,7 @@ export const ADMIN_PLATFORM_CONFIG_PROVIDERS = [ EnvironmentEventMaterializer, EnvironmentEventStreamService, NetworkManagementService, + NetworkPortForwardGroupService, NetworkManagementEventStreamService, NetworkDnsPodClient, NetworkDdnsService, diff --git a/src/modules/admin/platform-config/network-management/network-management.service.ts b/src/modules/admin/platform-config/network-management/network-management.service.ts index 6ade128..f8eac23 100644 --- a/src/modules/admin/platform-config/network-management/network-management.service.ts +++ b/src/modules/admin/platform-config/network-management/network-management.service.ts @@ -14,10 +14,7 @@ import { NetworkPortForwardUpdateDto, } from './network-management.dto'; import { NetworkPortForward } from './network-management.entity'; -import { - isIpv4Address, - portForwardActiveKey, -} from './network-management.types'; +import { isIpv4Address } from './network-management.types'; import { NetworkTcpReleasePolicyError, NetworkTcpReleasePolicyService, @@ -25,6 +22,7 @@ import { type TcpReleaseMutation, type TcpReleaseState, } from './network-tcp-release-policy.service'; +import { NetworkPortForwardGroupService } from './network-port-forward-group.service'; const DEFAULT_AGENT_ID = 'nas-main'; const DEFAULT_TARGET_IPV4 = '192.168.31.224'; @@ -42,6 +40,7 @@ export class NetworkManagementService { private readonly configService: ConfigService, private readonly mqttService: NetworkAgentMqttService, private readonly tcpReleasePolicy: NetworkTcpReleasePolicyService, + private readonly groupService: NetworkPortForwardGroupService, ) {} async list(query: NetworkPortForwardListQueryDto = {}) { @@ -79,135 +78,15 @@ export class NetworkManagementService { } async create(input: NetworkPortForwardCreateDto) { - this.assertReleaseMutation({ - after: this.releaseState({ - externalPort: input.externalPort, - internalPort: input.internalPort, - natmapDesiredEnabled: false, - protocol: input.protocol, - }), - kind: 'create', - }); - const saved = await this.dataSource.transaction(async (manager) => { - const state = await this.lockAgentState(manager); - const repository = manager.getRepository(NetworkPortForward); - if ((await repository.count({ where: { isDeleted: false } })) >= 64) { - throwVbenError('端口转发记录已达到 64 条上限', HttpStatus.CONFLICT); - } - const activeKey = portForwardActiveKey( - input.protocol, - input.externalPort, - ); - if (await repository.findOne({ where: { activeKey } })) { - throwVbenError('同协议外部端口已存在', HttpStatus.CONFLICT); - } - const mapping = repository.create({ - activeKey, - currentObservedAt: null, - currentPublicIpv4: null, - currentPublicPort: null, - currentValidUntil: null, - desiredPresence: 'present', - externalPort: input.externalPort, - internalPort: input.internalPort, - isDeleted: false, - keeperDesiredEnabled: false, - keeperStatus: 'disabled', - lastErrorCode: null, - lastErrorMessage: null, - name: this.normalizeName(input.name), - probeRequestId: null, - protocol: input.protocol, - remark: input.remark?.trim() || null, - reportedRevision: '0', - syncStatus: 'pending', - targetIpv4: state.targetIpv4, - }); - this.advanceRevision(state, mapping); - try { - await repository.save(mapping); - await manager.getRepository(NetworkAgentState).save(state); - } catch (error) { - if (this.isDuplicateKeyError(error)) { - throwVbenError('同协议外部端口已存在', HttpStatus.CONFLICT); - } - throw error; - } - return mapping; - }); - this.notifyDesiredChanged(); - return this.serialize(saved); + return this.groupService.createV1(input); } async update(id: string, input: NetworkPortForwardUpdateDto) { - if (Object.keys(input).length === 0) { - throwVbenError('至少提供一个修改字段', HttpStatus.BAD_REQUEST); - } - return this.mutate(id, async (mapping, manager) => { - if (mapping.desiredPresence === 'absent') { - throwVbenError('端口转发正在删除', HttpStatus.CONFLICT); - } - const protocol = input.protocol || mapping.protocol; - const externalPort = input.externalPort || mapping.externalPort; - const internalPort = input.internalPort || mapping.internalPort; - this.assertReleaseMutation({ - after: this.releaseState({ - externalPort, - internalPort, - natmapDesiredEnabled: mapping.natmapDesiredEnabled, - protocol, - }), - before: this.releaseState(mapping), - kind: 'update', - }); - if ( - mapping.keeperDesiredEnabled && - (protocol !== 'udp' || externalPort !== internalPort) - ) { - throwVbenError( - '请先停用 Keeper 再修改协议或同源端口', - HttpStatus.BAD_REQUEST, - ); - } - const activeKey = portForwardActiveKey(protocol, externalPort); - if (activeKey !== mapping.activeKey) { - const conflict = await manager - .getRepository(NetworkPortForward) - .findOne({ where: { activeKey } }); - if (conflict && conflict.id !== mapping.id) { - throwVbenError('同协议外部端口已存在', HttpStatus.CONFLICT); - } - } - if (input.name !== undefined) { - mapping.name = this.normalizeName(input.name); - } - mapping.remark = - input.remark === undefined - ? mapping.remark - : input.remark.trim() || null; - mapping.protocol = protocol; - mapping.externalPort = externalPort; - mapping.internalPort = internalPort; - mapping.activeKey = activeKey; - mapping.syncStatus = 'pending'; - }); + return this.groupService.updateV1(id, input); } async remove(id: string) { - return this.mutate(id, async (mapping) => { - if (mapping.desiredPresence === 'absent') { - throwVbenError('端口转发正在删除', HttpStatus.CONFLICT); - } - this.assertReleaseMutation({ - current: this.releaseState(mapping), - kind: 'delete', - }); - mapping.desiredPresence = 'absent'; - mapping.keeperDesiredEnabled = false; - mapping.probeRequestId = null; - mapping.syncStatus = 'deleting'; - this.withdrawCurrentEndpoint(mapping); - }); + return this.groupService.removeV1(id); } async retry(id: string) { @@ -389,10 +268,7 @@ export class NetworkManagementService { this.tcpReleasePolicy.assertMutationAllowed(mutation); } catch (error) { if (error instanceof NetworkTcpReleasePolicyError) { - throwVbenError( - '当前发布模式不允许该 TCP 操作', - HttpStatus.CONFLICT, - ); + throwVbenError('当前发布模式不允许该 TCP 操作', HttpStatus.CONFLICT); } throw error; } @@ -454,6 +330,7 @@ export class NetworkManagementService { new Date(mapping.currentValidUntil).getTime() > Date.now(); return { id: String(mapping.id), + groupId: String(mapping.groupId), name: mapping.name, remark: mapping.remark || null, protocol: mapping.protocol, @@ -462,12 +339,14 @@ export class NetworkManagementService { targetIpv4: mapping.targetIpv4, desiredPresence: mapping.desiredPresence, keeperDesiredEnabled: mapping.keeperDesiredEnabled, + natmapDesiredEnabled: mapping.natmapDesiredEnabled, probeRequestId: mapping.probeRequestId || null, desiredRevision: String(mapping.desiredRevision), desiredIssuedAt: mapping.desiredIssuedAt, reportedRevision: String(mapping.reportedRevision), syncStatus: mapping.syncStatus, keeperStatus: mapping.keeperStatus, + natmapStatus: mapping.natmapStatus, currentPublicIpv4: leaseValid ? mapping.currentPublicIpv4 : null, currentPublicPort: leaseValid ? mapping.currentPublicPort : null, currentPublicEndpoint: leaseValid @@ -535,17 +414,6 @@ export class NetworkManagementService { return target; } - private normalizeName(value: string): string { - const normalized = value.trim(); - if (!normalized || Buffer.byteLength(normalized, 'utf8') > 128) { - throwVbenError( - '规则名称超出 Agent UTF-8 长度限制', - HttpStatus.BAD_REQUEST, - ); - } - return normalized; - } - private isDuplicateKeyError(error: unknown): boolean { if (!error || typeof error !== 'object') return false; const record = error as { code?: unknown; errno?: unknown }; diff --git a/src/modules/admin/platform-config/network-management/network-port-forward-group.controller.ts b/src/modules/admin/platform-config/network-management/network-port-forward-group.controller.ts new file mode 100644 index 0000000..955e7ac --- /dev/null +++ b/src/modules/admin/platform-config/network-management/network-port-forward-group.controller.ts @@ -0,0 +1,174 @@ +import { + Body, + Controller, + Delete, + Get, + HttpCode, + HttpStatus, + Param, + Post, + Put, + Query, + Res, + UseGuards, + UsePipes, + ValidationPipe, +} from '@nestjs/common'; +import { ApiOperation, ApiTags } from '@nestjs/swagger'; +import type { Response } from 'express'; +import { vbenPage, vbenSuccess } from '@/common'; +import { AdminSuperGuard } from '@/modules/admin/identity/auth/admin-super.guard'; +import { JwtAuthGuard } from '@/modules/admin/identity/auth/jwt-auth.guard'; +import { NetworkEndpointHistoryQueryDto } from './network-management.dto'; +import { + NetworkPortForwardGroupChannelParamsDto, + NetworkPortForwardGroupCreateDto, + NetworkPortForwardGroupListQueryDto, + NetworkPortForwardGroupParamsDto, + NetworkPortForwardGroupUpdateDto, +} from './network-port-forward-group.dto'; +import { NetworkPortForwardGroupService } from './network-port-forward-group.service'; + +@ApiTags('Admin - 网络逻辑端口转发组') +@Controller('system/network/port-forward-group') +@UseGuards(JwtAuthGuard, AdminSuperGuard) +@UsePipes( + new ValidationPipe({ + forbidNonWhitelisted: true, + transform: true, + whitelist: true, + }), +) +export class NetworkPortForwardGroupController { + constructor(private readonly service: NetworkPortForwardGroupService) {} + + @Get('list') + @ApiOperation({ summary: '分页查询逻辑端口转发组' }) + async list( + @Query() query: NetworkPortForwardGroupListQueryDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + const page = await this.service.list(query); + return vbenPage(page.items, page.total); + } + + @Post() + @ApiOperation({ summary: '新增逻辑端口转发组' }) + async create( + @Body() body: NetworkPortForwardGroupCreateDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.create(body)); + } + + @Put(':groupId') + @ApiOperation({ summary: '修改逻辑端口转发组' }) + async update( + @Param() params: NetworkPortForwardGroupParamsDto, + @Body() body: NetworkPortForwardGroupUpdateDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.update(params.groupId, body)); + } + + @Delete(':groupId') + @ApiOperation({ summary: '删除逻辑端口转发组' }) + async remove( + @Param() params: NetworkPortForwardGroupParamsDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.remove(params.groupId)); + } + + @Post(':groupId/channels/:protocol/retry') + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: '重试协议通道同步' }) + async retry( + @Param() params: NetworkPortForwardGroupChannelParamsDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess( + await this.service.retry(params.groupId, params.protocol), + ); + } + + @Get(':groupId/channels/:protocol/endpoint-history') + @ApiOperation({ summary: '查询协议通道公网端点历史' }) + async endpointHistory( + @Param() params: NetworkPortForwardGroupChannelParamsDto, + @Query() query: NetworkEndpointHistoryQueryDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + const page = await this.service.endpointHistory( + params.groupId, + params.protocol, + query, + ); + return vbenPage(page.items, page.total); + } + + @Post(':groupId/channels/tcp/natmap/enable') + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: '启用 TCP NATMap' }) + async enableNatmap( + @Param() params: NetworkPortForwardGroupParamsDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.enableNatmap(params.groupId)); + } + + @Post(':groupId/channels/tcp/natmap/disable') + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: '停用 TCP NATMap' }) + async disableNatmap( + @Param() params: NetworkPortForwardGroupParamsDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.disableNatmap(params.groupId)); + } + + @Post(':groupId/channels/udp/keeper/enable') + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: '启用 UDP STUN Keeper' }) + async enableKeeper( + @Param() params: NetworkPortForwardGroupParamsDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.enableKeeper(params.groupId)); + } + + @Post(':groupId/channels/udp/keeper/disable') + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: '停用 UDP STUN Keeper' }) + async disableKeeper( + @Param() params: NetworkPortForwardGroupParamsDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.disableKeeper(params.groupId)); + } + + @Post(':groupId/channels/udp/keeper/probe') + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: '立即刷新 UDP 公网端点' }) + async probe( + @Param() params: NetworkPortForwardGroupParamsDto, + @Res({ passthrough: true }) response: Response, + ) { + this.noStore(response); + return vbenSuccess(await this.service.probe(params.groupId)); + } + + private noStore(response: Response): void { + response.setHeader('Cache-Control', 'no-store'); + } +} diff --git a/src/modules/admin/platform-config/network-management/network-port-forward-group.dto.ts b/src/modules/admin/platform-config/network-management/network-port-forward-group.dto.ts new file mode 100644 index 0000000..a4b8bbb --- /dev/null +++ b/src/modules/admin/platform-config/network-management/network-port-forward-group.dto.ts @@ -0,0 +1,126 @@ +import { Type } from 'class-transformer'; +import { + IsIn, + IsInt, + IsOptional, + IsString, + Length, + Matches, + Max, + MaxLength, + Min, + ValidateIf, +} from 'class-validator'; +import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger'; +import type { TcpProtocolMode } from './network-tcp-release-policy.service'; + +const DECIMAL_ID_PATTERN = /^\d{1,24}$/; + +export class NetworkPortForwardGroupCreateDto { + @ApiProperty({ maxLength: 100 }) + @IsString() + @Length(1, 100) + @Matches(/\S/, { message: 'name must contain a non-whitespace character' }) + name: string; + + @ApiPropertyOptional({ maxLength: 500 }) + @ValidateIf(isProvided) + @IsString() + @MaxLength(500) + remark?: string; + + @ApiProperty({ enum: ['tcp', 'tcp_udp', 'udp'] }) + @IsIn(['tcp', 'tcp_udp', 'udp']) + protocolMode: TcpProtocolMode; + + @ApiProperty({ maximum: 65535, minimum: 1 }) + @IsInt() + @Min(1) + @Max(65535) + externalPort: number; + + @ApiProperty({ maximum: 65535, minimum: 1 }) + @IsInt() + @Min(1) + @Max(65535) + internalPort: number; +} + +export class NetworkPortForwardGroupUpdateDto { + @ApiPropertyOptional({ maxLength: 100 }) + @ValidateIf(isProvided) + @IsString() + @Length(1, 100) + @Matches(/\S/, { message: 'name must contain a non-whitespace character' }) + name?: string; + + @ApiPropertyOptional({ maxLength: 500, nullable: true }) + @ValidateIf(isProvided) + @IsString() + @MaxLength(500) + remark?: string; + + @ApiPropertyOptional({ enum: ['tcp', 'tcp_udp', 'udp'] }) + @ValidateIf(isProvided) + @IsIn(['tcp', 'tcp_udp', 'udp']) + protocolMode?: TcpProtocolMode; + + @ApiPropertyOptional({ maximum: 65535, minimum: 1 }) + @ValidateIf(isProvided) + @IsInt() + @Min(1) + @Max(65535) + externalPort?: number; + + @ApiPropertyOptional({ maximum: 65535, minimum: 1 }) + @ValidateIf(isProvided) + @IsInt() + @Min(1) + @Max(65535) + internalPort?: number; +} + +export class NetworkPortForwardGroupListQueryDto { + @ApiPropertyOptional({ minimum: 1 }) + @IsOptional() + @Type(() => Number) + @IsInt() + @Min(1) + pageNo?: number; + + @ApiPropertyOptional({ maximum: 100, minimum: 1 }) + @IsOptional() + @Type(() => Number) + @IsInt() + @Min(1) + @Max(100) + pageSize?: number; + + @ApiPropertyOptional({ maxLength: 100 }) + @IsOptional() + @IsString() + @MaxLength(100) + name?: string; + + @ApiPropertyOptional({ enum: ['tcp', 'tcp_udp', 'udp'] }) + @IsOptional() + @IsIn(['tcp', 'tcp_udp', 'udp']) + protocolMode?: TcpProtocolMode; +} + +export class NetworkPortForwardGroupParamsDto { + @ApiProperty() + @IsString() + @Matches(DECIMAL_ID_PATTERN) + groupId: string; +} + +export class NetworkPortForwardGroupChannelParamsDto extends NetworkPortForwardGroupParamsDto { + @ApiProperty({ enum: ['tcp', 'udp'] }) + @IsIn(['tcp', 'udp']) + protocol: 'tcp' | 'udp'; +} + +function isProvided(_object: object, value: unknown): boolean { + return value !== undefined; +} diff --git a/src/modules/admin/platform-config/network-management/network-port-forward-group.service.ts b/src/modules/admin/platform-config/network-management/network-port-forward-group.service.ts new file mode 100644 index 0000000..f208d9a --- /dev/null +++ b/src/modules/admin/platform-config/network-management/network-port-forward-group.service.ts @@ -0,0 +1,1028 @@ +import { randomUUID } from 'node:crypto'; +import { HttpStatus, Injectable } from '@nestjs/common'; +import { ConfigService } from '@nestjs/config'; +import { InjectRepository } from '@nestjs/typeorm'; +import { DataSource, type EntityManager, Repository } from 'typeorm'; +import { KtDateTime, throwVbenError } from '@/common'; +import { NetworkAgentMqttService } from './network-agent-mqtt.service'; +import { NetworkAgentState } from './network-agent-state.entity'; +import { NETWORK_AGENT_V2_MAX_CHANNELS } from './network-agent-v2.types'; +import { NetworkEndpointHistory } from './network-endpoint-history.entity'; +import type { + NetworkEndpointHistoryQueryDto, + NetworkPortForwardCreateDto, + NetworkPortForwardUpdateDto, +} from './network-management.dto'; +import { NetworkPortForward } from './network-management.entity'; +import { + isIpv4Address, + portForwardActiveKey, + type PortForwardProtocol, +} from './network-management.types'; +import type { + NetworkPortForwardGroupCreateDto, + NetworkPortForwardGroupListQueryDto, + NetworkPortForwardGroupUpdateDto, +} from './network-port-forward-group.dto'; +import { NetworkPortForwardGroup } from './network-port-forward-group.entity'; +import { + NetworkTcpReleasePolicyError, + NetworkTcpReleasePolicyService, + type TcpProtocolMode, + type TcpReleaseMutation, + type TcpReleaseState, +} from './network-tcp-release-policy.service'; + +const DEFAULT_AGENT_ID = 'nas-main'; +const DEFAULT_TARGET_IPV4 = '192.168.31.224'; + +type GroupTransactionResult = { + channels: NetworkPortForward[]; + changed: boolean; + group: NetworkPortForwardGroup; +}; + +@Injectable() +export class NetworkPortForwardGroupService { + constructor( + @InjectRepository(NetworkPortForwardGroup) + private readonly groupRepository: Repository, + @InjectRepository(NetworkPortForward) + private readonly mappingRepository: Repository, + @InjectRepository(NetworkEndpointHistory) + private readonly historyRepository: Repository, + private readonly dataSource: DataSource, + private readonly configService: ConfigService, + private readonly mqttService: NetworkAgentMqttService, + private readonly tcpReleasePolicy: NetworkTcpReleasePolicyService, + ) {} + + async list(query: NetworkPortForwardGroupListQueryDto = {}) { + const pageNo = query.pageNo || 1; + const pageSize = query.pageSize || 20; + const builder = this.groupRepository + .createQueryBuilder('group') + .where('group.isDeleted = :isDeleted', { isDeleted: false }); + if (query.name) { + builder.andWhere('group.name LIKE :name', { + name: `%${query.name.trim()}%`, + }); + } + if (query.protocolMode) { + builder.andWhere('group.protocolMode = :protocolMode', { + protocolMode: query.protocolMode, + }); + } + const [groups, total] = await builder + .orderBy('group.createTime', 'DESC') + .skip((pageNo - 1) * pageSize) + .take(pageSize) + .getManyAndCount(); + const items = await Promise.all( + groups.map(async (group) => + this.serializeGroup( + group, + await this.mappingRepository.find({ + order: { protocol: 'ASC' }, + where: { groupId: group.id, isDeleted: false }, + }), + ), + ), + ); + return { items, total }; + } + + async create(input: NetworkPortForwardGroupCreateDto) { + const result = await this.createInternal(input); + return this.serializeGroup(result.group, result.channels); + } + + async createV1(input: NetworkPortForwardCreateDto) { + const result = await this.createInternal({ + externalPort: input.externalPort, + internalPort: input.internalPort, + name: input.name, + protocolMode: input.protocol, + remark: input.remark, + }); + const channel = result.channels.find( + (item) => item.protocol === input.protocol, + ); + if (!channel) { + throwVbenError('端口转发通道创建失败', HttpStatus.INTERNAL_SERVER_ERROR); + } + return this.serializeChannel(channel); + } + + async update(groupId: string, input: NetworkPortForwardGroupUpdateDto) { + this.assertId(groupId, '逻辑组'); + this.assertUpdateInput(input); + const result = await this.dataSource.transaction(async (manager) => { + const state = await this.lockAgentState(manager); + const group = await this.findLockedGroup(manager, groupId); + const channels = await this.findChannels(manager, groupId); + return this.updateInTransaction(manager, state, group, channels, input); + }); + this.notifyDesiredChanged(); + return this.serializeGroup(result.group, result.channels); + } + + async updateV1(channelId: string, input: NetworkPortForwardUpdateDto) { + this.assertId(channelId, '端口转发'); + this.assertUpdateInput(input); + const result = await this.dataSource.transaction(async (manager) => { + const state = await this.lockAgentState(manager); + const repository = manager.getRepository(NetworkPortForward); + const requested = await repository.findOne({ + lock: { mode: 'pessimistic_write' }, + where: { id: channelId, isDeleted: false }, + }); + if (!requested) { + throwVbenError('端口转发不存在', HttpStatus.NOT_FOUND); + } + const group = await this.findLockedGroup(manager, requested.groupId); + const channels = await this.findChannels(manager, group.id); + this.assertV1SingleChannel(channels); + return this.updateInTransaction(manager, state, group, channels, { + externalPort: input.externalPort, + internalPort: input.internalPort, + name: input.name, + protocolMode: input.protocol, + remark: input.remark, + }); + }); + this.notifyDesiredChanged(); + const protocol = input.protocol || result.channels[0]?.protocol; + const channel = result.channels.find( + (item) => + !item.isDeleted && + item.desiredPresence === 'present' && + item.protocol === protocol, + ); + if (!channel) { + throwVbenError('端口转发通道更新失败', HttpStatus.INTERNAL_SERVER_ERROR); + } + return this.serializeChannel(channel); + } + + async remove(groupId: string) { + this.assertId(groupId, '逻辑组'); + const result = await this.dataSource.transaction(async (manager) => { + const state = await this.lockAgentState(manager); + const group = await this.findLockedGroup(manager, groupId); + const channels = await this.findChannels(manager, groupId); + return this.removeInTransaction(manager, state, group, channels); + }); + this.notifyDesiredChanged(); + return this.serializeGroup(result.group, result.channels); + } + + async removeV1(channelId: string) { + this.assertId(channelId, '端口转发'); + const result = await this.dataSource.transaction(async (manager) => { + const state = await this.lockAgentState(manager); + const repository = manager.getRepository(NetworkPortForward); + const requested = await repository.findOne({ + lock: { mode: 'pessimistic_write' }, + where: { id: channelId, isDeleted: false }, + }); + if (!requested) { + throwVbenError('端口转发不存在', HttpStatus.NOT_FOUND); + } + const group = await this.findLockedGroup(manager, requested.groupId); + const channels = await this.findChannels(manager, group.id); + this.assertV1SingleChannel(channels); + return this.removeInTransaction(manager, state, group, channels); + }); + this.notifyDesiredChanged(); + return this.serializeChannel(result.channels[0]); + } + + async retry(groupId: string, protocol: PortForwardProtocol) { + return this.mutateChannel(groupId, protocol, async (group, channel) => { + if (channel.desiredPresence !== 'present') { + throwVbenError('删除中的通道不能重试', HttpStatus.CONFLICT); + } + if (protocol === 'tcp') { + this.assertReleaseMutation({ + current: this.releaseState(group, channel.natmapDesiredEnabled), + kind: 'retry', + }); + } + channel.lastErrorCode = null; + channel.lastErrorMessage = null; + if (protocol === 'tcp') { + channel.natmapLastErrorCode = null; + channel.natmapLastErrorMessage = null; + } else { + channel.keeperLastErrorCode = null; + channel.keeperLastErrorMessage = null; + } + channel.syncStatus = 'pending'; + return true; + }); + } + + async enableNatmap(groupId: string) { + return this.mutateChannel( + groupId, + 'tcp', + async (group, channel, channels) => { + if (channel.natmapDesiredEnabled) return false; + this.assertMechanismTransitionAllowed(channels); + this.assertReleaseMutation({ + current: this.releaseState(group, true), + kind: 'natmap-enable', + }); + channel.natmapDesiredEnabled = true; + channel.syncStatus = 'pending'; + return true; + }, + true, + ); + } + + async disableNatmap(groupId: string) { + return this.mutateChannel( + groupId, + 'tcp', + async (group, channel, channels) => { + if (!channel.natmapDesiredEnabled) return false; + this.assertMechanismTransitionAllowed(channels); + this.assertReleaseMutation({ + after: this.releaseState(group, false), + before: this.releaseState(group, true), + kind: 'natmap-disable', + }); + channel.natmapDesiredEnabled = false; + channel.syncStatus = 'pending'; + this.withdrawCurrentEndpoint(channel); + return true; + }, + true, + ); + } + + async enableKeeper(groupId: string) { + return this.mutateChannel( + groupId, + 'udp', + async (_, channel, channels) => { + this.assertKeeperPorts(channel); + if (channel.keeperDesiredEnabled) return false; + this.assertMechanismTransitionAllowed(channels); + channel.keeperDesiredEnabled = true; + channel.probeRequestId = randomUUID(); + channel.syncStatus = 'pending'; + return true; + }, + true, + ); + } + + async disableKeeper(groupId: string) { + return this.mutateChannel( + groupId, + 'udp', + async (_, channel, channels) => { + this.assertKeeperPorts(channel); + if (!channel.keeperDesiredEnabled) return false; + this.assertMechanismTransitionAllowed(channels); + channel.keeperDesiredEnabled = false; + channel.probeRequestId = null; + channel.syncStatus = 'pending'; + this.withdrawCurrentEndpoint(channel); + return true; + }, + true, + ); + } + + async probe(groupId: string) { + return this.mutateChannel( + groupId, + 'udp', + async (_, channel, channels) => { + this.assertMechanismTransitionAllowed(channels); + this.assertKeeperPorts(channel); + if (!channel.keeperDesiredEnabled) { + throwVbenError('请先启用 UDP Keeper', HttpStatus.BAD_REQUEST); + } + channel.probeRequestId = randomUUID(); + channel.syncStatus = 'pending'; + return true; + }, + true, + ); + } + + async endpointHistory( + groupId: string, + protocol: PortForwardProtocol, + query: NetworkEndpointHistoryQueryDto = {}, + ) { + this.assertId(groupId, '逻辑组'); + const group = await this.groupRepository.findOne({ + where: { id: groupId, isDeleted: false }, + }); + if (!group) throwVbenError('逻辑端口转发组不存在', HttpStatus.NOT_FOUND); + const channel = await this.mappingRepository.findOne({ + where: { groupId, isDeleted: false, protocol }, + }); + if (!channel) throwVbenError('协议通道不存在', HttpStatus.NOT_FOUND); + const pageNo = query.pageNo || 1; + const pageSize = query.pageSize || 20; + const mechanism = protocol === 'tcp' ? 'tcp_natmap' : 'udp_stun'; + const [items, total] = await this.historyRepository.findAndCount({ + order: { occurredAt: 'DESC' }, + skip: (pageNo - 1) * pageSize, + take: pageSize, + where: { mappingId: channel.id, mechanism }, + }); + return { items: items.map((item) => this.serializeHistory(item)), total }; + } + + private async createInternal( + input: NetworkPortForwardGroupCreateDto, + ): Promise { + const name = this.normalizeName(input.name); + this.assertReleaseMutation({ + after: { + externalPort: input.externalPort, + internalPort: input.internalPort, + natmapDesiredEnabled: false, + protocolMode: input.protocolMode, + }, + kind: 'create', + }); + const result = await this.dataSource.transaction(async (manager) => { + const state = await this.lockAgentState(manager); + const mappingRepository = manager.getRepository(NetworkPortForward); + const groupRepository = manager.getRepository(NetworkPortForwardGroup); + const protocols = this.protocols(input.protocolMode); + const count = await mappingRepository.count({ + where: { isDeleted: false }, + }); + if (count + protocols.length > NETWORK_AGENT_V2_MAX_CHANNELS) { + throwVbenError('端口转发通道已达到 64 条上限', HttpStatus.CONFLICT); + } + for (const protocol of protocols) { + await this.assertActiveKeyAvailable( + mappingRepository, + protocol, + input.externalPort, + ); + } + const group = groupRepository.create({ + externalPort: input.externalPort, + internalPort: input.internalPort, + isDeleted: false, + name, + protocolMode: input.protocolMode, + remark: input.remark?.trim() || null, + targetIpv4: state.targetIpv4, + }); + await groupRepository.save(group); + const channels = protocols.map((protocol) => + this.createChannel(mappingRepository, group, protocol), + ); + const revision = this.advanceGlobalRevision(state); + this.assignRevision(channels, revision, state.desiredIssuedAt); + try { + await mappingRepository.save(channels); + await manager.getRepository(NetworkAgentState).save(state); + } catch (error) { + this.rethrowDuplicate(error); + } + return { changed: true, channels, group }; + }); + this.notifyDesiredChanged(); + return result; + } + + private async updateInTransaction( + manager: EntityManager, + state: NetworkAgentState, + group: NetworkPortForwardGroup, + channels: NetworkPortForward[], + input: NetworkPortForwardGroupUpdateDto, + ): Promise { + const protocolMode = + input.protocolMode || (group.protocolMode as TcpProtocolMode); + const externalPort = input.externalPort || group.externalPort; + const internalPort = input.internalPort || group.internalPort; + const structuralChange = + protocolMode !== group.protocolMode || + externalPort !== group.externalPort || + internalPort !== group.internalPort; + if (structuralChange) this.assertStructuralEditAllowed(channels); + if (channels.every((channel) => channel.desiredPresence === 'absent')) { + throwVbenError('逻辑组正在删除', HttpStatus.CONFLICT); + } + const tcp = channels.find((channel) => channel.protocol === 'tcp'); + const mutation: TcpReleaseMutation = + group.protocolMode === 'tcp_udp' && protocolMode === 'udp' + ? { + after: { + externalPort, + internalPort, + natmapDesiredEnabled: false, + protocolMode, + }, + before: this.releaseState( + group, + tcp?.natmapDesiredEnabled || false, + ), + kind: 'protocol-shrink', + } + : { + after: { + externalPort, + internalPort, + natmapDesiredEnabled: + protocolMode === 'udp' + ? false + : tcp?.natmapDesiredEnabled || false, + protocolMode, + }, + before: this.releaseState( + group, + tcp?.natmapDesiredEnabled || false, + ), + kind: 'update', + }; + this.assertReleaseMutation(mutation); + + const oldProtocols = new Set( + this.protocols(group.protocolMode as TcpProtocolMode), + ); + const nextProtocols = new Set(this.protocols(protocolMode)); + const changedChannels = new Set(); + const mappingRepository = manager.getRepository(NetworkPortForward); + const additions = [...nextProtocols].filter( + (protocol) => !oldProtocols.has(protocol), + ); + if (additions.length) { + const count = await mappingRepository.count({ + where: { isDeleted: false }, + }); + if (count + additions.length > NETWORK_AGENT_V2_MAX_CHANNELS) { + throwVbenError('端口转发通道已达到 64 条上限', HttpStatus.CONFLICT); + } + } + for (const protocol of additions) { + if ( + channels.some( + (channel) => + channel.protocol === protocol && + channel.desiredPresence === 'absent', + ) + ) { + throwVbenError('协议通道正在删除,不能重新添加', HttpStatus.CONFLICT); + } + await this.assertActiveKeyAvailable( + mappingRepository, + protocol, + externalPort, + ); + const channel = this.createChannel(mappingRepository, group, protocol); + channels.push(channel); + changedChannels.add(channel); + } + for (const channel of channels) { + if (!nextProtocols.has(channel.protocol)) { + channel.desiredPresence = 'absent'; + channel.keeperDesiredEnabled = false; + channel.natmapDesiredEnabled = false; + channel.probeRequestId = null; + channel.syncStatus = 'deleting'; + this.withdrawCurrentEndpoint(channel); + changedChannels.add(channel); + } + } + + const name = + input.name === undefined ? group.name : this.normalizeName(input.name); + const authorityPayloadChanged = + name !== group.name || + externalPort !== group.externalPort || + internalPort !== group.internalPort; + group.name = name; + group.remark = + input.remark === undefined ? group.remark : input.remark.trim() || null; + group.externalPort = externalPort; + group.internalPort = internalPort; + group.protocolMode = protocolMode; + for (const channel of channels) { + if (authorityPayloadChanged) changedChannels.add(channel); + channel.name = group.name; + channel.remark = group.remark; + channel.externalPort = group.externalPort; + channel.internalPort = group.internalPort; + channel.targetIpv4 = group.targetIpv4; + if (channel.desiredPresence === 'present') { + channel.activeKey = portForwardActiveKey( + channel.protocol, + group.externalPort, + ); + } + } + const revision = this.advanceGlobalRevision(state); + this.assignRevision([...changedChannels], revision, state.desiredIssuedAt); + try { + await manager.getRepository(NetworkPortForwardGroup).save(group); + await mappingRepository.save(channels); + await manager.getRepository(NetworkAgentState).save(state); + } catch (error) { + this.rethrowDuplicate(error); + } + return { changed: true, channels, group }; + } + + private async removeInTransaction( + manager: EntityManager, + state: NetworkAgentState, + group: NetworkPortForwardGroup, + channels: NetworkPortForward[], + ): Promise { + if ( + channels.length === 0 || + channels.some((channel) => channel.desiredPresence === 'absent') + ) { + throwVbenError('逻辑组正在删除或协调中', HttpStatus.CONFLICT); + } + const tcp = channels.find((channel) => channel.protocol === 'tcp'); + this.assertReleaseMutation({ + current: this.releaseState(group, tcp?.natmapDesiredEnabled || false), + kind: 'delete', + }); + for (const channel of channels) { + channel.desiredPresence = 'absent'; + channel.keeperDesiredEnabled = false; + channel.natmapDesiredEnabled = false; + channel.probeRequestId = null; + channel.syncStatus = 'deleting'; + this.withdrawCurrentEndpoint(channel); + } + const revision = this.advanceGlobalRevision(state); + this.assignRevision(channels, revision, state.desiredIssuedAt); + await manager.getRepository(NetworkPortForward).save(channels); + await manager.getRepository(NetworkAgentState).save(state); + return { changed: true, channels, group }; + } + + private async mutateChannel( + groupId: string, + protocol: PortForwardProtocol, + change: ( + group: NetworkPortForwardGroup, + channel: NetworkPortForward, + channels: NetworkPortForward[], + ) => Promise, + invalidProtocolIsBadRequest = false, + ) { + this.assertId(groupId, '逻辑组'); + const result = await this.dataSource.transaction(async (manager) => { + const state = await this.lockAgentState(manager); + const group = await this.findLockedGroup(manager, groupId); + const channels = await this.findChannels(manager, groupId); + const channel = channels.find((item) => item.protocol === protocol); + if (!channel) { + throwVbenError( + `逻辑组不包含 ${protocol.toUpperCase()} 协议通道`, + invalidProtocolIsBadRequest + ? HttpStatus.BAD_REQUEST + : HttpStatus.NOT_FOUND, + ); + } + if (channel.desiredPresence !== 'present') { + throwVbenError('协议通道正在删除', HttpStatus.CONFLICT); + } + const changed = await change(group, channel, channels); + if (!changed) return { changed, channel }; + const revision = this.advanceGlobalRevision(state); + this.assignRevision([channel], revision, state.desiredIssuedAt); + await manager.getRepository(NetworkPortForward).save(channel); + await manager.getRepository(NetworkAgentState).save(state); + return { changed, channel }; + }); + if (result.changed) this.notifyDesiredChanged(); + return this.serializeChannel(result.channel); + } + + private async findLockedGroup( + manager: EntityManager, + groupId: string, + ): Promise { + const group = await manager.getRepository(NetworkPortForwardGroup).findOne({ + lock: { mode: 'pessimistic_write' }, + where: { id: groupId, isDeleted: false }, + }); + if (!group) throwVbenError('逻辑端口转发组不存在', HttpStatus.NOT_FOUND); + return group; + } + + private async findChannels( + manager: EntityManager, + groupId: string, + ): Promise { + return manager.getRepository(NetworkPortForward).find({ + order: { protocol: 'ASC' }, + where: { groupId, isDeleted: false }, + }); + } + + private createChannel( + repository: Repository, + group: NetworkPortForwardGroup, + protocol: PortForwardProtocol, + ): NetworkPortForward { + return repository.create({ + activeGroupProtocolKey: `${group.id}:${protocol}`, + activeKey: portForwardActiveKey(protocol, group.externalPort), + currentObservedAt: null, + currentPublicIpv4: null, + currentPublicPort: null, + currentValidUntil: null, + desiredPresence: 'present', + externalPort: group.externalPort, + groupId: group.id, + internalPort: group.internalPort, + isDeleted: false, + keeperDesiredEnabled: false, + keeperStatus: 'disabled', + lastErrorCode: null, + lastErrorMessage: null, + name: group.name, + natmapDesiredEnabled: false, + natmapStatus: 'disabled', + probeRequestId: null, + protocol, + remark: group.remark, + reportedRevision: '0', + syncStatus: 'pending', + targetIpv4: group.targetIpv4, + }); + } + + private assertStructuralEditAllowed(channels: NetworkPortForward[]): void { + if ( + channels.some( + (channel) => + channel.desiredPresence !== 'present' || + channel.syncStatus !== 'synced', + ) + ) { + throwVbenError('逻辑组正在删除或协调中', HttpStatus.CONFLICT); + } + if ( + channels.some( + (channel) => + channel.keeperDesiredEnabled || + channel.natmapDesiredEnabled || + channel.keeperStatus !== 'disabled' || + channel.natmapStatus !== 'disabled', + ) + ) { + throwVbenError( + '请先停用 Keeper 和 NATMap 再修改端口或协议', + HttpStatus.BAD_REQUEST, + ); + } + } + + private assertMechanismTransitionAllowed( + channels: NetworkPortForward[], + ): void { + if ( + channels.some( + (channel) => + channel.desiredPresence !== 'present' || + channel.syncStatus !== 'synced', + ) + ) { + throwVbenError('逻辑组正在删除或协调中', HttpStatus.CONFLICT); + } + } + + private assertV1SingleChannel(channels: NetworkPortForward[]): void { + if ( + channels.length !== 1 || + channels[0].desiredPresence !== 'present' || + channels[0].syncStatus === 'deleting' || + channels[0].syncStatus === 'pending' || + channels[0].syncStatus === 'syncing' + ) { + throwVbenError( + '多协议或协调中的逻辑组请使用新版管理接口', + HttpStatus.CONFLICT, + ); + } + } + + private assertKeeperPorts(channel: NetworkPortForward): void { + if (channel.externalPort !== channel.internalPort) { + throwVbenError( + 'UDP Keeper 要求外部端口与内部端口一致', + HttpStatus.BAD_REQUEST, + ); + } + } + + private async assertActiveKeyAvailable( + repository: Repository, + protocol: PortForwardProtocol, + externalPort: number, + currentId?: string, + ): Promise { + const conflict = await repository.findOne({ + where: { activeKey: portForwardActiveKey(protocol, externalPort) }, + }); + if (conflict && conflict.id !== currentId) { + throwVbenError('同协议外部端口已存在', HttpStatus.CONFLICT); + } + } + + private advanceGlobalRevision(state: NetworkAgentState): string { + const revision = (BigInt(state.desiredRevision) + 1n).toString(); + state.desiredRevision = revision; + state.desiredIssuedAt = new KtDateTime(); + return revision; + } + + private assignRevision( + channels: NetworkPortForward[], + revision: string, + issuedAt: KtDateTime, + ): void { + for (const channel of channels) { + channel.desiredRevision = revision; + channel.desiredIssuedAt = issuedAt; + } + } + + private async lockAgentState( + manager: EntityManager, + ): Promise { + const repository = manager.getRepository(NetworkAgentState); + const agentId = this.agentId(); + await repository + .createQueryBuilder() + .insert() + .into(NetworkAgentState) + .values({ + agentId, + appliedRevision: '0', + desiredIssuedAt: new KtDateTime(), + desiredRevision: '0', + online: false, + publishedRevision: '0', + targetIpv4: this.targetIpv4(), + }) + .orIgnore() + .execute(); + const state = await repository.findOne({ + lock: { mode: 'pessimistic_write' }, + where: { agentId }, + }); + if (!state) { + throwVbenError( + 'Agent 状态行初始化失败', + HttpStatus.INTERNAL_SERVER_ERROR, + ); + } + if (state.targetIpv4 !== this.targetIpv4()) { + throwVbenError( + 'Agent 目标 IPv4 与服务端固定配置不一致', + HttpStatus.CONFLICT, + ); + } + return state; + } + + private releaseState( + group: Pick< + NetworkPortForwardGroup, + 'externalPort' | 'internalPort' | 'protocolMode' + >, + natmapDesiredEnabled: boolean, + ): TcpReleaseState { + return { + externalPort: group.externalPort, + internalPort: group.internalPort, + natmapDesiredEnabled, + protocolMode: group.protocolMode as TcpProtocolMode, + }; + } + + private assertReleaseMutation(mutation: TcpReleaseMutation): void { + try { + this.tcpReleasePolicy.assertMutationAllowed(mutation); + } catch (error) { + if (error instanceof NetworkTcpReleasePolicyError) { + throwVbenError('当前发布模式不允许该 TCP 操作', HttpStatus.CONFLICT); + } + throw error; + } + } + + private protocols(mode: TcpProtocolMode): PortForwardProtocol[] { + if (mode === 'tcp_udp') return ['tcp', 'udp']; + return [mode]; + } + + private appliedProtocolMode( + channels: NetworkPortForward[], + ): TcpProtocolMode | null { + const applied = channels + .filter( + (channel) => + channel.desiredPresence === 'present' && + channel.syncStatus === 'synced' && + this.revisionCaughtUp( + channel.reportedRevision, + channel.desiredRevision, + ), + ) + .map((channel) => channel.protocol); + const tcp = applied.includes('tcp'); + const udp = applied.includes('udp'); + if (tcp && udp) return 'tcp_udp'; + if (tcp) return 'tcp'; + if (udp) return 'udp'; + return null; + } + + private revisionCaughtUp(reported: string, desired: string): boolean { + try { + return BigInt(reported) >= BigInt(desired); + } catch { + return false; + } + } + + private serializeGroup( + group: NetworkPortForwardGroup, + channels: NetworkPortForward[], + ) { + const tcp = channels.find( + (channel) => !channel.isDeleted && channel.protocol === 'tcp', + ); + const udp = channels.find( + (channel) => !channel.isDeleted && channel.protocol === 'udp', + ); + return { + id: String(group.id), + name: group.name, + remark: group.remark || null, + externalPort: group.externalPort, + internalPort: group.internalPort, + protocolMode: group.protocolMode, + appliedProtocolMode: this.appliedProtocolMode(channels), + targetIpv4: group.targetIpv4, + channels: { + tcp: tcp ? this.serializeChannel(tcp) : null, + udp: udp ? this.serializeChannel(udp) : null, + }, + isDeleted: group.isDeleted, + createTime: group.createTime, + updateTime: group.updateTime, + }; + } + + private serializeChannel(channel: NetworkPortForward) { + const leaseValid = + !!channel.currentPublicIpv4 && + !!channel.currentPublicPort && + !!channel.currentValidUntil && + new Date(channel.currentValidUntil).getTime() > Date.now(); + return { + id: String(channel.id), + groupId: String(channel.groupId), + name: channel.name, + remark: channel.remark || null, + protocol: channel.protocol, + externalPort: channel.externalPort, + internalPort: channel.internalPort, + targetIpv4: channel.targetIpv4, + desiredPresence: channel.desiredPresence, + keeperDesiredEnabled: channel.keeperDesiredEnabled, + natmapDesiredEnabled: channel.natmapDesiredEnabled, + probeRequestId: channel.probeRequestId || null, + desiredRevision: String(channel.desiredRevision), + desiredIssuedAt: channel.desiredIssuedAt, + reportedRevision: String(channel.reportedRevision), + syncStatus: channel.syncStatus, + keeperStatus: channel.keeperStatus, + natmapStatus: channel.natmapStatus, + currentPublicIpv4: leaseValid ? channel.currentPublicIpv4 : null, + currentPublicPort: leaseValid ? channel.currentPublicPort : null, + currentPublicEndpoint: leaseValid + ? `${channel.currentPublicIpv4}:${channel.currentPublicPort}` + : null, + currentObservedAt: leaseValid ? channel.currentObservedAt : null, + currentValidatedAt: leaseValid ? channel.currentValidatedAt : null, + currentValidUntil: leaseValid ? channel.currentValidUntil : null, + lastObservedIpv4: channel.lastObservedIpv4 || null, + lastObservedPort: channel.lastObservedPort || null, + lastObservedAt: channel.lastObservedAt || null, + lastErrorCode: channel.lastErrorCode || null, + lastErrorMessage: channel.lastErrorMessage || null, + keeperLastErrorCode: channel.keeperLastErrorCode || null, + keeperLastErrorMessage: channel.keeperLastErrorMessage || null, + natmapLastErrorCode: channel.natmapLastErrorCode || null, + natmapLastErrorMessage: channel.natmapLastErrorMessage || null, + isDeleted: channel.isDeleted, + createTime: channel.createTime, + updateTime: channel.updateTime, + }; + } + + private serializeHistory(history: NetworkEndpointHistory) { + return { + id: String(history.id), + eventId: history.eventId, + eventType: history.eventType, + mechanism: history.mechanism, + firstObservedAt: history.firstObservedAt, + lastObservedAt: history.lastObservedAt, + occurredAt: history.occurredAt, + portForwardId: String(history.mappingId), + publicIpv4: history.publicIpv4 || null, + publicPort: history.publicPort || null, + withdrawalReason: history.reason || null, + createTime: history.createTime, + }; + } + + private withdrawCurrentEndpoint(channel: NetworkPortForward): void { + channel.currentPublicIpv4 = null; + channel.currentPublicPort = null; + channel.currentObservedAt = null; + channel.currentValidatedAt = null; + channel.currentValidUntil = null; + } + + private assertUpdateInput(input: object): void { + if (Object.keys(input).length === 0) { + throwVbenError('至少提供一个修改字段', HttpStatus.BAD_REQUEST); + } + } + + private assertId(id: string, label: string): void { + if (!/^\d{1,24}$/.test(id)) { + throwVbenError(`${label} ID 无效`, HttpStatus.BAD_REQUEST); + } + } + + private normalizeName(value: string): string { + const normalized = value.trim(); + if (!normalized || Buffer.byteLength(normalized, 'utf8') > 128) { + throwVbenError( + '规则名称超出 Agent UTF-8 长度限制', + HttpStatus.BAD_REQUEST, + ); + } + return normalized; + } + + private agentId(): string { + return ( + this.configService.get('NETWORK_AGENT_ID') || DEFAULT_AGENT_ID + ); + } + + private targetIpv4(): string { + const target = + this.configService.get('NETWORK_AGENT_TARGET_IPV4') || + DEFAULT_TARGET_IPV4; + if (!isIpv4Address(target)) { + throwVbenError( + 'NETWORK_AGENT_TARGET_IPV4 配置无效', + HttpStatus.INTERNAL_SERVER_ERROR, + ); + } + return target; + } + + private rethrowDuplicate(error: unknown): never { + if (this.isDuplicateKeyError(error)) { + throwVbenError('同协议外部端口或组内协议已存在', HttpStatus.CONFLICT); + } + throw error; + } + + private isDuplicateKeyError(error: unknown): boolean { + if (!error || typeof error !== 'object') return false; + const record = error as { code?: unknown; errno?: unknown }; + return record.code === 'ER_DUP_ENTRY' || record.errno === 1062; + } + + private notifyDesiredChanged(): void { + try { + this.mqttService.requestDesiredPublish(); + } catch { + // Desired state is durable; the periodic publisher will retry independently. + } + } +} diff --git a/test/admin/network-management/network-management.controller.spec.ts b/test/admin/network-management/network-management.controller.spec.ts index cecf76b..1ca3261 100644 --- a/test/admin/network-management/network-management.controller.spec.ts +++ b/test/admin/network-management/network-management.controller.spec.ts @@ -159,6 +159,30 @@ describe('NetworkManagementController', () => { ); }); + it('keeps v1 responses flat and forwards only the channel ID', async () => { + service.update.mockResolvedValue({ + groupId: '200', + id: '100', + protocol: 'udp', + }); + + const response = await request(apiUrl) + .put('/system/network/port-forward/100') + .send({ remark: 'legacy' }) + .expect(200); + + expect(response.body.data).toMatchObject({ + groupId: '200', + id: '100', + protocol: 'udp', + }); + expect(response.body.data.channels).toBeUndefined(); + expect(service.update).toHaveBeenCalledWith( + '100', + expect.objectContaining({ remark: 'legacy' }), + ); + }); + it('exposes asynchronous actions, history, and agent status routes', async () => { service.retry.mockResolvedValue({ id: '100' }); service.enableKeeper.mockResolvedValue({ id: '100' }); diff --git a/test/admin/network-management/network-management.service.spec.ts b/test/admin/network-management/network-management.service.spec.ts index 09c13a0..01bfaaf 100644 --- a/test/admin/network-management/network-management.service.spec.ts +++ b/test/admin/network-management/network-management.service.spec.ts @@ -7,11 +7,14 @@ import { NetworkAgentState } from '../../../src/modules/admin/platform-config/ne import { NetworkEndpointHistory } from '../../../src/modules/admin/platform-config/network-management/network-endpoint-history.entity'; import { NetworkPortForward } from '../../../src/modules/admin/platform-config/network-management/network-management.entity'; import { NetworkManagementService } from '../../../src/modules/admin/platform-config/network-management/network-management.service'; +import { NetworkPortForwardGroup } from '../../../src/modules/admin/platform-config/network-management/network-port-forward-group.entity'; +import { NetworkPortForwardGroupService } from '../../../src/modules/admin/platform-config/network-management/network-port-forward-group.service'; import { NetworkTcpReleasePolicyService } from '../../../src/modules/admin/platform-config/network-management/network-tcp-release-policy.service'; type Harness = { bootstrapOrder: string[]; bootstrapExecute: jest.Mock; + groups: NetworkPortForwardGroup[]; histories: NetworkEndpointHistory[]; mappings: NetworkPortForward[]; mqtt: jest.Mocked>; @@ -29,6 +32,24 @@ function createHarness( releaseConfig: ReleaseConfig = {}, ): Harness { const mappings = initialMappings; + const groups = Array.from( + new Set(mappings.map((mapping) => mapping.groupId)), + ).map((groupId) => { + const channels = mappings.filter((mapping) => mapping.groupId === groupId); + const protocols = new Set(channels.map((mapping) => mapping.protocol)); + const first = channels[0]; + return Object.assign(new NetworkPortForwardGroup(), { + externalPort: first.externalPort, + id: groupId, + internalPort: first.internalPort, + isDeleted: false, + name: first.name, + protocolMode: + protocols.size === 2 ? 'tcp_udp' : protocols.has('tcp') ? 'tcp' : 'udp', + remark: first.remark || null, + targetIpv4: first.targetIpv4, + }); + }); const histories: NetworkEndpointHistory[] = []; const bootstrapOrder: string[] = []; const bootstrapExecute = jest.fn().mockImplementation(async () => { @@ -49,17 +70,41 @@ function createHarness( create: (input) => Object.assign(new NetworkPortForward(), { id: '100' }, input), createQueryBuilder: () => createListBuilder(mappings), + find: async ({ where }) => + mappings.filter((mapping) => + Object.entries(where).every(([key, value]) => mapping[key] === value), + ), findOne: async ({ where }) => mappings.find((mapping) => Object.entries(where).every(([key, value]) => mapping[key] === value), ) || null, - save: async (mapping) => { - const index = mappings.findIndex((item) => item.id === mapping.id); - if (index >= 0) mappings[index] = mapping; - else mappings.push(mapping); - return mapping; + save: async (value) => { + const values = Array.isArray(value) ? value : [value]; + for (const mapping of values) { + const index = mappings.findIndex((item) => item.id === mapping.id); + if (index >= 0) mappings[index] = mapping; + else mappings.push(mapping); + } + return value; }, } as unknown as Repository; + const groupRepository = { + create: (input) => + Object.assign(new NetworkPortForwardGroup(), { id: '200' }, input), + findOne: async ({ where }) => + groups.find((group) => + Object.entries(where).every(([key, value]) => group[key] === value), + ) || null, + save: async (value) => { + const values = Array.isArray(value) ? value : [value]; + for (const group of values) { + const index = groups.findIndex((item) => item.id === group.id); + if (index >= 0) groups[index] = group; + else groups.push(group); + } + return value; + }, + } as unknown as Repository; const stateRepository = { create: (input) => Object.assign(new NetworkAgentState(), input), createQueryBuilder: () => { @@ -84,6 +129,7 @@ function createHarness( const manager = { getRepository: (entity) => { if (entity === NetworkPortForward) return mappingRepository; + if (entity === NetworkPortForwardGroup) return groupRepository; if (entity === NetworkAgentState) return stateRepository; if (entity === NetworkEndpointHistory) return historyRepository; throw new Error('unexpected repository'); @@ -112,10 +158,20 @@ function createHarness( configService, mqtt as unknown as NetworkAgentMqttService, new NetworkTcpReleasePolicyService(configService), + new NetworkPortForwardGroupService( + groupRepository, + mappingRepository, + historyRepository, + dataSource, + configService, + mqtt as unknown as NetworkAgentMqttService, + new NetworkTcpReleasePolicyService(configService), + ), ); return { bootstrapExecute, bootstrapOrder, + groups, histories, mappings, mqtt, @@ -129,6 +185,7 @@ function createMapping( ): NetworkPortForward { return Object.assign(new NetworkPortForward(), { activeKey: 'udp:9000', + activeGroupProtocolKey: '200:udp', currentObservedAt: null, currentPublicIpv4: null, currentPublicPort: null, @@ -137,12 +194,15 @@ function createMapping( desiredPresence: 'present', desiredRevision: '3', externalPort: 9000, + groupId: '200', id: '100', internalPort: 9000, isDeleted: false, keeperDesiredEnabled: false, keeperStatus: 'disabled', name: 'rule', + natmapDesiredEnabled: false, + natmapStatus: 'disabled', protocol: 'udp', reportedRevision: '0', syncStatus: 'synced', @@ -173,7 +233,10 @@ function createListBuilder(mappings: NetworkPortForward[]) { } return builder; }, - getManyAndCount: async () => [rows.slice(offset, offset + limit), rows.length], + getManyAndCount: async () => [ + rows.slice(offset, offset + limit), + rows.length, + ], orderBy: () => builder, skip: (value: number) => { offset = value; @@ -212,8 +275,10 @@ describe('NetworkManagementService', () => { }); expect(harness.state.desiredRevision).toBe('1'); expect(harness.mappings[0]).toMatchObject({ + activeGroupProtocolKey: expect.any(String), activeKey: 'udp:9000', desiredRevision: '1', + groupId: expect.any(String), }); expect(harness.mappings[0].desiredIssuedAt.toISOString()).toBe( harness.state.desiredIssuedAt.toISOString(), @@ -222,6 +287,29 @@ describe('NetworkManagementService', () => { expect(harness.bootstrapExecute).toHaveBeenCalledTimes(1); }); + it('keeps v1 channel-ID update and delete from mutating a multi-channel group', async () => { + const udp = createMapping({ + activeGroupProtocolKey: '200:udp', + groupId: '200', + }); + const tcp = createMapping({ + activeGroupProtocolKey: '200:tcp', + activeKey: 'tcp:9000', + groupId: '200', + id: '101', + protocol: 'tcp', + }); + const harness = createHarness([udp, tcp], { mode: 'on' }); + + await expect( + harness.service.update('100', { remark: 'legacy update' }), + ).rejects.toMatchObject({ status: 409 }); + await expect(harness.service.remove('100')).rejects.toMatchObject({ + status: 409, + }); + expect(harness.state.desiredRevision).toBe('3'); + }); + it('uses insert-ignore bootstrap before the singleton pessimistic lock', async () => { const harness = createHarness(); @@ -550,6 +638,7 @@ describe('NetworkManagementService', () => { await expect(harness.service.disableKeeper('100')).resolves.toMatchObject({ keeperDesiredEnabled: false, }); + mapping.syncStatus = 'synced'; await expect(harness.service.remove('100')).resolves.toMatchObject({ desiredPresence: 'absent', }); diff --git a/test/admin/network-management/network-port-forward-group.controller.spec.ts b/test/admin/network-management/network-port-forward-group.controller.spec.ts new file mode 100644 index 0000000..8c5f2b1 --- /dev/null +++ b/test/admin/network-management/network-port-forward-group.controller.spec.ts @@ -0,0 +1,242 @@ +import type { + CanActivate, + ExecutionContext, + INestApplication, +} from '@nestjs/common'; +import { HttpException, HttpStatus } from '@nestjs/common'; +import { Test } from '@nestjs/testing'; +import * as request from 'supertest'; +import { AdminSuperGuard } from '../../../src/modules/admin/identity/auth/admin-super.guard'; +import { JwtAuthGuard } from '../../../src/modules/admin/identity/auth/jwt-auth.guard'; +import { NetworkPortForwardGroupController } from '../../../src/modules/admin/platform-config/network-management/network-port-forward-group.controller'; +import { NetworkPortForwardGroupService } from '../../../src/modules/admin/platform-config/network-management/network-port-forward-group.service'; + +describe('NetworkPortForwardGroupController', () => { + let app: INestApplication; + let apiUrl: string; + const service = { + create: jest.fn(), + disableKeeper: jest.fn(), + disableNatmap: jest.fn(), + enableKeeper: jest.fn(), + enableNatmap: jest.fn(), + endpointHistory: jest.fn(), + list: jest.fn(), + probe: jest.fn(), + remove: jest.fn(), + retry: jest.fn(), + update: jest.fn(), + }; + const authGuard: CanActivate = { + canActivate(context: ExecutionContext) { + const requestValue = context.switchToHttp().getRequest(); + const roleCode = requestValue.headers['x-test-role'] || 'super'; + requestValue.adminUser = { + roles: [{ isDeleted: false, roleCode, status: 1 }], + }; + return true; + }, + }; + + beforeAll(async () => { + const moduleRef = await Test.createTestingModule({ + controllers: [NetworkPortForwardGroupController], + providers: [ + AdminSuperGuard, + { provide: NetworkPortForwardGroupService, useValue: service }, + ], + }) + .overrideGuard(JwtAuthGuard) + .useValue(authGuard) + .compile(); + + app = moduleRef.createNestApplication(); + await app.listen(0, '127.0.0.1'); + apiUrl = await app.getUrl(); + }); + + beforeEach(() => { + jest.clearAllMocks(); + service.list.mockResolvedValue({ items: [], total: 0 }); + service.create.mockResolvedValue({ + appliedProtocolMode: null, + channels: { tcp: { id: '100' }, udp: { id: '101' } }, + id: '90071992547409930', + protocolMode: 'tcp_udp', + }); + service.update.mockResolvedValue({ id: '90071992547409930' }); + service.remove.mockResolvedValue({ id: '90071992547409930' }); + service.retry.mockResolvedValue({ id: '100' }); + service.enableNatmap.mockResolvedValue({ id: '100' }); + service.disableNatmap.mockResolvedValue({ id: '100' }); + service.enableKeeper.mockResolvedValue({ id: '101' }); + service.disableKeeper.mockResolvedValue({ id: '101' }); + service.probe.mockResolvedValue({ id: '101' }); + service.endpointHistory.mockResolvedValue({ + items: [{ mechanism: 'tcp_natmap', portForwardId: '100' }], + total: 1, + }); + }); + + afterAll(async () => { + await app?.close(); + }); + + it('serves exact v2 CRUD routes with strict DTO validation and string IDs over real HTTP', async () => { + await request(apiUrl) + .get( + '/system/network/port-forward-group/list?pageNo=1&pageSize=20&protocolMode=tcp_udp', + ) + .expect(200) + .expect('Cache-Control', 'no-store'); + await request(apiUrl) + .post('/system/network/port-forward-group') + .send({ + externalPort: 9000, + internalPort: 9000, + name: 'dual', + protocolMode: 'icmp', + }) + .expect(400); + await request(apiUrl) + .post('/system/network/port-forward-group') + .send({ + externalPort: 9000, + internalPort: 9000, + name: 'dual', + protocolMode: 'tcp_udp', + routerPassword: 'forbidden', + }) + .expect(400); + + const created = await request(apiUrl) + .post('/system/network/port-forward-group') + .send({ + externalPort: 9000, + internalPort: 9000, + name: 'dual', + protocolMode: 'tcp_udp', + }) + .expect(201) + .expect('Cache-Control', 'no-store'); + expect(created.body.data.id).toBe('90071992547409930'); + expect(typeof created.body.data.id).toBe('string'); + await request(apiUrl) + .put('/system/network/port-forward-group/90071992547409930') + .send({ remark: 'updated' }) + .expect(200) + .expect('Cache-Control', 'no-store'); + await request(apiUrl) + .delete('/system/network/port-forward-group/90071992547409930') + .expect(200) + .expect('Cache-Control', 'no-store'); + expect(service.update).toHaveBeenCalledWith( + '90071992547409930', + expect.objectContaining({ remark: 'updated' }), + ); + }); + + it('serves every protocol-scoped action and history route over real HTTP', async () => { + const groupId = '90071992547409930'; + await request(apiUrl) + .post(`/system/network/port-forward-group/${groupId}/channels/tcp/retry`) + .expect(200); + await request(apiUrl) + .post(`/system/network/port-forward-group/${groupId}/channels/udp/retry`) + .expect(200); + await request(apiUrl) + .post( + `/system/network/port-forward-group/${groupId}/channels/tcp/natmap/enable`, + ) + .expect(200); + await request(apiUrl) + .post( + `/system/network/port-forward-group/${groupId}/channels/tcp/natmap/disable`, + ) + .expect(200); + await request(apiUrl) + .post( + `/system/network/port-forward-group/${groupId}/channels/udp/keeper/enable`, + ) + .expect(200); + await request(apiUrl) + .post( + `/system/network/port-forward-group/${groupId}/channels/udp/keeper/disable`, + ) + .expect(200); + await request(apiUrl) + .post( + `/system/network/port-forward-group/${groupId}/channels/udp/keeper/probe`, + ) + .expect(200); + const tcpHistory = await request(apiUrl) + .get( + `/system/network/port-forward-group/${groupId}/channels/tcp/endpoint-history`, + ) + .expect(200) + .expect('Cache-Control', 'no-store'); + const udpHistory = await request(apiUrl) + .get( + `/system/network/port-forward-group/${groupId}/channels/udp/endpoint-history`, + ) + .expect(200); + expect(tcpHistory.body.data.items[0].mechanism).toBe('tcp_natmap'); + expect(service.retry).toHaveBeenNthCalledWith(1, groupId, 'tcp'); + expect(service.retry).toHaveBeenNthCalledWith(2, groupId, 'udp'); + expect(service.endpointHistory).toHaveBeenNthCalledWith( + 1, + groupId, + 'tcp', + expect.any(Object), + ); + expect(service.endpointHistory).toHaveBeenNthCalledWith( + 2, + groupId, + 'udp', + expect.any(Object), + ); + expect(udpHistory.body.data.total).toBe(1); + }); + + it('enforces super-admin permission and maps service 400/409 errors to Chinese Vben msg', async () => { + await request(apiUrl) + .get('/system/network/port-forward-group/list') + .set('X-Test-Role', 'member') + .expect(403); + + service.enableNatmap.mockRejectedValueOnce( + new HttpException( + { code: HttpStatus.BAD_REQUEST, msg: 'TCP NATMap 状态不允许启用' }, + HttpStatus.BAD_REQUEST, + ), + ); + const badRequest = await request(apiUrl) + .post('/system/network/port-forward-group/200/channels/tcp/natmap/enable') + .expect(400); + expect(badRequest.body.msg).toMatch(/TCP NATMap/); + + service.update.mockRejectedValueOnce( + new HttpException( + { code: HttpStatus.CONFLICT, msg: '逻辑组正在协调中' }, + HttpStatus.CONFLICT, + ), + ); + const conflict = await request(apiUrl) + .put('/system/network/port-forward-group/200') + .send({ protocolMode: 'udp' }) + .expect(409); + expect(conflict.body.msg).toMatch(/协调/); + }); + + it('rejects invalid group IDs and protocols before calling the service', async () => { + await request(apiUrl) + .post( + '/system/network/port-forward-group/not-a-number/channels/tcp/retry', + ) + .expect(400); + await request(apiUrl) + .post('/system/network/port-forward-group/200/channels/icmp/retry') + .expect(400); + expect(service.retry).not.toHaveBeenCalled(); + }); +}); diff --git a/test/admin/network-management/network-port-forward-group.service.spec.ts b/test/admin/network-management/network-port-forward-group.service.spec.ts new file mode 100644 index 0000000..c00eb19 --- /dev/null +++ b/test/admin/network-management/network-port-forward-group.service.spec.ts @@ -0,0 +1,491 @@ +import { HttpException } from '@nestjs/common'; +import type { ConfigService } from '@nestjs/config'; +import type { DataSource, EntityManager, Repository } from 'typeorm'; +import { KtDateTime } from '../../../src/common'; +import type { NetworkAgentMqttService } from '../../../src/modules/admin/platform-config/network-management/network-agent-mqtt.service'; +import { NetworkAgentState } from '../../../src/modules/admin/platform-config/network-management/network-agent-state.entity'; +import { NetworkEndpointHistory } from '../../../src/modules/admin/platform-config/network-management/network-endpoint-history.entity'; +import { NetworkPortForward } from '../../../src/modules/admin/platform-config/network-management/network-management.entity'; +import { NetworkPortForwardGroup } from '../../../src/modules/admin/platform-config/network-management/network-port-forward-group.entity'; +import { NetworkPortForwardGroupService } from '../../../src/modules/admin/platform-config/network-management/network-port-forward-group.service'; +import { NetworkTcpReleasePolicyService } from '../../../src/modules/admin/platform-config/network-management/network-tcp-release-policy.service'; + +type Harness = { + groups: NetworkPortForwardGroup[]; + histories: NetworkEndpointHistory[]; + mappings: NetworkPortForward[]; + mqtt: jest.Mocked>; + service: NetworkPortForwardGroupService; + state: NetworkAgentState; +}; + +function createHarness( + initialGroups: NetworkPortForwardGroup[] = [], + initialMappings: NetworkPortForward[] = [], +): Harness { + const groups = initialGroups; + const mappings = initialMappings; + const histories: NetworkEndpointHistory[] = []; + let nextGroupId = 200; + let nextMappingId = 100; + const state = Object.assign(new NetworkAgentState(), { + agentId: 'nas-main', + appliedRevision: '0', + desiredIssuedAt: new KtDateTime('2026-07-26T00:00:00.000Z'), + desiredRevision: initialMappings.length ? '3' : '0', + online: false, + publishedRevision: '0', + targetIpv4: '192.168.31.224', + }); + const groupRepository = { + count: async () => groups.filter((group) => !group.isDeleted).length, + create: (input) => Object.assign(new NetworkPortForwardGroup(), input), + createQueryBuilder: () => createGroupListBuilder(groups), + findOne: async ({ where }) => findOne(groups, where), + save: async (value) => { + const values = Array.isArray(value) ? value : [value]; + for (const group of values) { + if (!group.id) group.id = String(nextGroupId++); + const index = groups.findIndex((item) => item.id === group.id); + if (index >= 0) groups[index] = group; + else groups.push(group); + } + return value; + }, + } as unknown as Repository; + const mappingRepository = { + count: async () => mappings.filter((mapping) => !mapping.isDeleted).length, + create: (input) => Object.assign(new NetworkPortForward(), input), + find: async ({ where }) => + mappings.filter((mapping) => matches(mapping, where)), + findOne: async ({ where }) => findOne(mappings, where), + save: async (value) => { + const values = Array.isArray(value) ? value : [value]; + for (const mapping of values) { + if (!mapping.id) mapping.id = String(nextMappingId++); + const duplicate = mappings.find( + (item) => + item.id !== mapping.id && + ((mapping.activeKey && item.activeKey === mapping.activeKey) || + (mapping.activeGroupProtocolKey && + item.activeGroupProtocolKey === + mapping.activeGroupProtocolKey)), + ); + if (duplicate) { + const error = new Error('duplicate') as Error & { code: string }; + error.code = 'ER_DUP_ENTRY'; + throw error; + } + const index = mappings.findIndex((item) => item.id === mapping.id); + if (index >= 0) mappings[index] = mapping; + else mappings.push(mapping); + } + return value; + }, + } as unknown as Repository; + const historyRepository = { + findAndCount: jest.fn(async ({ where }) => { + const rows = histories.filter((history) => matches(history, where)); + return [rows, rows.length]; + }), + } as unknown as Repository; + const stateRepository = { + createQueryBuilder: () => { + const builder = { + execute: async () => ({ identifiers: [] }), + insert: () => builder, + into: () => builder, + orIgnore: () => builder, + values: () => builder, + }; + return builder; + }, + findOne: async () => state, + save: async (value) => Object.assign(state, value), + } as unknown as Repository; + const manager = { + getRepository: (entity) => { + if (entity === NetworkPortForwardGroup) return groupRepository; + if (entity === NetworkPortForward) return mappingRepository; + if (entity === NetworkEndpointHistory) return historyRepository; + if (entity === NetworkAgentState) return stateRepository; + throw new Error('unexpected repository'); + }, + } as unknown as EntityManager; + const dataSource = { + transaction: async (work) => { + const groupSnapshot = groups.map((group) => ({ ...group })); + const mappingSnapshot = mappings.map((mapping) => ({ ...mapping })); + const stateSnapshot = { ...state }; + try { + return await work(manager); + } catch (error) { + groups.splice( + 0, + groups.length, + ...groupSnapshot.map((group) => + Object.assign(new NetworkPortForwardGroup(), group), + ), + ); + mappings.splice( + 0, + mappings.length, + ...mappingSnapshot.map((mapping) => + Object.assign(new NetworkPortForward(), mapping), + ), + ); + Object.assign(state, stateSnapshot); + throw error; + } + }, + } as unknown as DataSource; + const configService = { + get: (key) => + ({ + NETWORK_AGENT_ID: 'nas-main', + NETWORK_AGENT_TARGET_IPV4: '192.168.31.224', + NETWORK_TCP_NATMAP_RELEASE_MODE: 'on', + })[key], + } as ConfigService; + const mqtt = { + requestDesiredPublish: jest.fn(), + } as jest.Mocked>; + const service = new NetworkPortForwardGroupService( + groupRepository, + mappingRepository, + historyRepository, + dataSource, + configService, + mqtt as unknown as NetworkAgentMqttService, + new NetworkTcpReleasePolicyService(configService), + ); + return { groups, histories, mappings, mqtt, service, state }; +} + +function createGroup( + patch: Partial = {}, +): NetworkPortForwardGroup { + return Object.assign(new NetworkPortForwardGroup(), { + externalPort: 9000, + id: '200', + internalPort: 9000, + isDeleted: false, + name: 'group', + protocolMode: 'tcp_udp', + remark: null, + targetIpv4: '192.168.31.224', + ...patch, + }); +} + +function createMapping( + patch: Partial = {}, +): NetworkPortForward { + const protocol = patch.protocol || 'udp'; + const id = patch.id || (protocol === 'tcp' ? '101' : '100'); + return Object.assign(new NetworkPortForward(), { + activeGroupProtocolKey: `200:${protocol}`, + activeKey: `${protocol}:9000`, + currentObservedAt: null, + currentPublicIpv4: null, + currentPublicPort: null, + currentValidUntil: null, + desiredIssuedAt: new KtDateTime('2026-07-26T00:00:00.000Z'), + desiredPresence: 'present', + desiredRevision: '3', + externalPort: 9000, + groupId: '200', + id, + internalPort: 9000, + isDeleted: false, + keeperDesiredEnabled: false, + keeperStatus: 'disabled', + name: 'group', + natmapDesiredEnabled: false, + natmapStatus: 'disabled', + probeRequestId: null, + protocol, + remark: null, + reportedRevision: '3', + syncStatus: 'synced', + targetIpv4: '192.168.31.224', + ...patch, + }); +} + +function matches(value: object, where: Record): boolean { + return Object.entries(where).every( + ([key, expected]) => value[key] === expected, + ); +} + +function findOne( + values: T[], + where: Record, +): T | null { + return values.find((value) => matches(value, where)) || null; +} + +function createGroupListBuilder(groups: NetworkPortForwardGroup[]) { + let rows = groups.filter((group) => !group.isDeleted); + let offset = 0; + let limit = rows.length; + const builder = { + andWhere: (clause: string, values: Record) => { + if (clause.includes('group.name LIKE')) { + const name = String(values.name).replaceAll('%', ''); + rows = rows.filter((group) => group.name.includes(name)); + } else if (clause.includes('group.protocolMode =')) { + rows = rows.filter( + (group) => group.protocolMode === values.protocolMode, + ); + } + return builder; + }, + getManyAndCount: async () => [ + rows.slice(offset, offset + limit), + rows.length, + ], + orderBy: () => builder, + skip: (value: number) => { + offset = value; + return builder; + }, + take: (value: number) => { + limit = value; + return builder; + }, + where: () => builder, + }; + return builder; +} + +function errorStatus(error: unknown): number { + return error instanceof HttpException ? error.getStatus() : 0; +} + +describe('NetworkPortForwardGroupService', () => { + it('creates tcp_udp atomically with one group authority and one global revision', async () => { + const harness = createHarness(); + + await expect( + harness.service.create({ + externalPort: 9000, + internalPort: 9000, + name: ' Dual Rule ', + protocolMode: 'tcp_udp', + remark: 'both', + }), + ).resolves.toMatchObject({ + appliedProtocolMode: null, + channels: { + tcp: { desiredRevision: '1', id: expect.any(String) }, + udp: { desiredRevision: '1', id: expect.any(String) }, + }, + id: expect.any(String), + name: 'Dual Rule', + protocolMode: 'tcp_udp', + }); + expect(harness.groups).toHaveLength(1); + expect(harness.mappings).toHaveLength(2); + expect(harness.state.desiredRevision).toBe('1'); + expect(harness.mappings).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + activeGroupProtocolKey: '200:tcp', + activeKey: 'tcp:9000', + groupId: '200', + name: 'Dual Rule', + }), + expect.objectContaining({ + activeGroupProtocolKey: '200:udp', + activeKey: 'udp:9000', + groupId: '200', + name: 'Dual Rule', + }), + ]), + ); + expect(harness.mqtt.requestDesiredPublish).toHaveBeenCalledTimes(1); + }); + + it('projects group authority while preserving channel-specific revisions and derives appliedProtocolMode', async () => { + const group = createGroup(); + const tcp = createMapping({ protocol: 'tcp' }); + const udp = createMapping({ protocol: 'udp' }); + const harness = createHarness([group], [tcp, udp]); + + await expect(harness.service.list()).resolves.toMatchObject({ + items: [{ appliedProtocolMode: 'tcp_udp' }], + }); + await harness.service.enableNatmap('200'); + expect(tcp.desiredRevision).toBe('4'); + expect(udp.desiredRevision).toBe('3'); + await harness.service.enableNatmap('200'); + expect(harness.state.desiredRevision).toBe('4'); + expect(harness.mqtt.requestDesiredPublish).toHaveBeenCalledTimes(1); + + tcp.reportedRevision = '4'; + await harness.service.update('200', { name: 'Projected' }); + expect(harness.state.desiredRevision).toBe('5'); + expect(tcp).toMatchObject({ desiredRevision: '5', name: 'Projected' }); + expect(udp).toMatchObject({ desiredRevision: '5', name: 'Projected' }); + expect(group.name).toBe('Projected'); + await expect(harness.service.list()).resolves.toMatchObject({ + items: [{ appliedProtocolMode: null }], + }); + }); + + it('keeps removed protocols as tombstones and never resurrects a prior channel ID', async () => { + const group = createGroup(); + const tcp = createMapping({ protocol: 'tcp' }); + const udp = createMapping({ protocol: 'udp' }); + const harness = createHarness([group], [tcp, udp]); + + await harness.service.update('200', { protocolMode: 'udp' }); + expect(tcp).toMatchObject({ + activeGroupProtocolKey: '200:tcp', + activeKey: 'tcp:9000', + desiredPresence: 'absent', + desiredRevision: '4', + isDeleted: false, + syncStatus: 'deleting', + }); + expect(udp.desiredRevision).toBe('3'); + expect(group.isDeleted).toBe(false); + + await harness.service + .update('200', { protocolMode: 'tcp_udp' }) + .catch((error) => expect(errorStatus(error)).toBe(409)); + expect(harness.mappings).toHaveLength(2); + + const persistedTcp = harness.mappings.find( + (mapping) => mapping.protocol === 'tcp', + )!; + const persistedUdp = harness.mappings.find( + (mapping) => mapping.protocol === 'udp', + )!; + persistedTcp.isDeleted = true; + persistedTcp.activeKey = null; + persistedTcp.activeGroupProtocolKey = null; + persistedUdp.syncStatus = 'synced'; + await harness.service.update('200', { protocolMode: 'tcp_udp' }); + const replacement = harness.mappings.find( + (mapping) => !mapping.isDeleted && mapping.protocol === 'tcp', + ); + expect(replacement?.id).not.toBe('101'); + expect(replacement).toMatchObject({ + activeGroupProtocolKey: '200:tcp', + activeKey: 'tcp:9000', + desiredPresence: 'present', + isDeleted: false, + }); + }); + + it('marks every channel absent on group delete without soft-deleting channels or the group', async () => { + const group = createGroup(); + const tcp = createMapping({ + currentPublicIpv4: '8.8.8.8', + currentPublicPort: 45000, + natmapDesiredEnabled: true, + protocol: 'tcp', + }); + const udp = createMapping({ + currentPublicIpv4: '1.1.1.1', + currentPublicPort: 45001, + keeperDesiredEnabled: true, + protocol: 'udp', + }); + const harness = createHarness([group], [tcp, udp]); + + await harness.service.remove('200'); + expect(group.isDeleted).toBe(false); + expect(harness.mappings).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + desiredPresence: 'absent', + isDeleted: false, + natmapDesiredEnabled: false, + }), + expect.objectContaining({ + desiredPresence: 'absent', + isDeleted: false, + keeperDesiredEnabled: false, + }), + ]), + ); + expect(harness.mappings.every((mapping) => mapping.activeKey)).toBe(true); + tcp.reportedRevision = tcp.desiredRevision; + udp.reportedRevision = udp.desiredRevision; + await harness.service.list(); + expect(group.isDeleted).toBe(false); + expect(harness.mappings.every((mapping) => !mapping.isDeleted)).toBe(true); + }); + + it('fails an atomic dual-channel create without partial state at the 64-channel limit', async () => { + const existing = Array.from({ length: 63 }, (_, index) => + createMapping({ + activeGroupProtocolKey: `${300 + index}:udp`, + activeKey: `udp:${10000 + index}`, + externalPort: 10000 + index, + groupId: String(300 + index), + id: String(1000 + index), + internalPort: 10000 + index, + }), + ); + const harness = createHarness([], existing); + + await harness.service + .create({ + externalPort: 9000, + internalPort: 9000, + name: 'too many', + protocolMode: 'tcp_udp', + }) + .catch((error) => expect(errorStatus(error)).toBe(409)); + expect(harness.groups).toHaveLength(0); + expect(harness.mappings).toHaveLength(63); + expect(harness.state.desiredRevision).toBe('3'); + expect(harness.mqtt.requestDesiredPublish).not.toHaveBeenCalled(); + }); + + it('keeps UDP Keeper switches idempotent and filters history by channel mechanism', async () => { + const group = createGroup({ protocolMode: 'udp' }); + const udp = createMapping({ protocol: 'udp' }); + const harness = createHarness([group], [udp]); + harness.histories.push( + Object.assign(new NetworkEndpointHistory(), { + eventId: 'udp-event', + id: '300', + mappingId: '100', + mechanism: 'udp_stun', + }), + Object.assign(new NetworkEndpointHistory(), { + eventId: 'wrong-mechanism', + id: '301', + mappingId: '100', + mechanism: 'tcp_natmap', + }), + ); + + await harness.service.enableKeeper('200'); + const enableRevision = harness.state.desiredRevision; + await harness.service.enableKeeper('200'); + expect(harness.state.desiredRevision).toBe(enableRevision); + udp.syncStatus = 'synced'; + await harness.service.disableKeeper('200'); + const disableRevision = harness.state.desiredRevision; + await harness.service.disableKeeper('200'); + expect(harness.state.desiredRevision).toBe(disableRevision); + expect(harness.mqtt.requestDesiredPublish).toHaveBeenCalledTimes(2); + + await expect( + harness.service.endpointHistory('200', 'udp'), + ).resolves.toMatchObject({ + items: [{ eventId: 'udp-event', mechanism: 'udp_stun' }], + total: 1, + }); + await harness.service + .enableNatmap('200') + .catch((error) => expect(errorStatus(error)).toBe(400)); + }); +});