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 b3a1156..e53a45d 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,7 @@ 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 { 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 { SystemLogController } from '@/modules/admin/platform-config/system-log/system-log.controller'; import { SystemLogService } from '@/modules/admin/platform-config/system-log/system-log.service'; import { AdminTimezoneController } from '@/modules/admin/platform-config/timezone/admin-timezone.controller'; @@ -88,6 +89,7 @@ export const ADMIN_PLATFORM_CONFIG_PROVIDERS = [ NetworkManagementEventStreamService, NetworkDnsPodClient, NetworkDdnsService, + NetworkStunMessageSourceAdapter, NetworkAgentMqttService, ]; diff --git a/src/modules/admin/platform-config/network-management/network-ddns.service.ts b/src/modules/admin/platform-config/network-management/network-ddns.service.ts index 3711ab0..2aeccef 100644 --- a/src/modules/admin/platform-config/network-management/network-ddns.service.ts +++ b/src/modules/admin/platform-config/network-management/network-ddns.service.ts @@ -17,6 +17,7 @@ import { } from './network-dnspod.client'; import { NetworkManagementEventStreamService } from './network-management-event-stream.service'; import { NetworkPortForward } from './network-management.entity'; +import { classifyStunEndpointSource } from './network-source-eligibility'; import type { NetworkDdnsListQuery, NetworkDdnsRecordInput, @@ -757,21 +758,13 @@ export class NetworkDdnsService implements OnModuleInit, OnModuleDestroy { private portForwardSourceOption( mapping: NetworkPortForward, ): NetworkDdnsSourceOption { - let disabledReasonCode: null | string = null; - if (mapping.isDeleted || mapping.desiredPresence !== 'present') { - disabledReasonCode = 'SOURCE_DELETING'; - } else if (mapping.protocol !== 'udp') { - disabledReasonCode = 'UDP_REQUIRED'; - } else if (mapping.externalPort !== mapping.internalPort) { - disabledReasonCode = 'PORT_MISMATCH'; - } else if (!mapping.keeperDesiredEnabled) { - disabledReasonCode = 'KEEPER_DISABLED'; - } + const { disabledReasonCode, eligible } = + classifyStunEndpointSource(mapping); const leaseValid = isIP(mapping.currentPublicIpv4 || '') === 4 && !!mapping.currentValidUntil && new Date(mapping.currentValidUntil).getTime() > Date.now(); - const sourceUsable = disabledReasonCode === null && leaseValid; + const sourceUsable = eligible && leaseValid; return { currentAddress: sourceUsable ? mapping.currentPublicIpv4 || null : null, disabledReasonCode, @@ -996,13 +989,7 @@ export class NetworkDdnsService implements OnModuleInit, OnModuleDestroy { const mapping = await this.mappingRepository.findOne({ where: { id: portForwardId as string, isDeleted: false }, }); - if ( - !mapping || - mapping.desiredPresence !== 'present' || - mapping.protocol !== 'udp' || - mapping.externalPort !== mapping.internalPort || - !mapping.keeperDesiredEnabled - ) { + if (!mapping || !classifyStunEndpointSource(mapping).eligible) { throwVbenError( 'A 记录来源必须是已启用 Keeper 的同端口 UDP 转发', HttpStatus.BAD_REQUEST, diff --git a/src/modules/admin/platform-config/network-management/network-source-eligibility.ts b/src/modules/admin/platform-config/network-management/network-source-eligibility.ts new file mode 100644 index 0000000..912637e --- /dev/null +++ b/src/modules/admin/platform-config/network-management/network-source-eligibility.ts @@ -0,0 +1,46 @@ +import type { NetworkPortForward } from './network-management.entity'; + +/** Stable reason returned when a mapping cannot structurally act as a STUN endpoint source. */ +export type StunEndpointSourceDisabledReason = + | 'KEEPER_DISABLED' + | 'PORT_MISMATCH' + | 'SOURCE_DELETING' + | 'UDP_REQUIRED'; + +/** Structural, lease-independent STUN source classification shared by DDNS and message delivery. */ +export type StunEndpointSourceEligibility = { + disabledReasonCode: null | StunEndpointSourceDisabledReason; + eligible: boolean; +}; + +/** + * Classifies whether a mapping may own a UDP STUN endpoint, excluding its current lease. + * @param mapping - Persisted port-forward mapping including desired Keeper configuration. + * @returns Stable structural eligibility and existing DDNS-compatible reason codes. + */ +export function classifyStunEndpointSource( + mapping: Pick< + NetworkPortForward, + | 'desiredPresence' + | 'externalPort' + | 'internalPort' + | 'isDeleted' + | 'keeperDesiredEnabled' + | 'protocol' + >, +): StunEndpointSourceEligibility { + let disabledReasonCode: null | StunEndpointSourceDisabledReason = null; + if (mapping.isDeleted || mapping.desiredPresence !== 'present') { + disabledReasonCode = 'SOURCE_DELETING'; + } else if (mapping.protocol !== 'udp') { + disabledReasonCode = 'UDP_REQUIRED'; + } else if (mapping.externalPort !== mapping.internalPort) { + disabledReasonCode = 'PORT_MISMATCH'; + } else if (!mapping.keeperDesiredEnabled) { + disabledReasonCode = 'KEEPER_DISABLED'; + } + return { + disabledReasonCode, + eligible: disabledReasonCode === null, + }; +} diff --git a/src/modules/admin/platform-config/network-management/network-stun-message-source.adapter.ts b/src/modules/admin/platform-config/network-management/network-stun-message-source.adapter.ts new file mode 100644 index 0000000..975bc47 --- /dev/null +++ b/src/modules/admin/platform-config/network-management/network-stun-message-source.adapter.ts @@ -0,0 +1,533 @@ +import { isIP } from 'node:net'; +import { + Injectable, + type OnModuleDestroy, + type OnModuleInit, +} from '@nestjs/common'; +import { InjectRepository } from '@nestjs/typeorm'; +import { Repository } from 'typeorm'; +import { + SystemMessageContractError, + type SystemMessageDeliveryReadiness, + type SystemMessageScalar, + type SystemMessageSourceDefinition, +} from '@/modules/qqbot/core/contract/message-push/qqbot-message-push.types'; +import { SystemMessageSourceRegistry } from '@/modules/qqbot/core/application/message-push/system-message-source.registry'; +import { NetworkDdnsRecord } from './network-ddns.entity'; +import { NetworkPortForward } from './network-management.entity'; +import { classifyStunEndpointSource } from './network-source-eligibility'; + +const SOURCE_KEY = 'network.stun.mapping-port-changed'; +const SNOWFLAKE_ID_PATTERN = /^[1-9]\d{0,23}$/; +const RFC3339_PATTERN = + /^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(?:\.\d{1,9})?(?:Z|[+-]\d{2}:\d{2})$/; + +type StunSubscriptionConfig = { + ddnsRecordId: string; + portForwardId: string; +}; + +type StunEventPayload = { + changedAt: string; + currentPort: number; + portForwardId: string; + previousPort: number; + publicIpv4: string; +}; + +type ResolvedSubscription = { + config: StunSubscriptionConfig; + ddnsRecord: NetworkDdnsRecord; + mapping: NetworkPortForward; + sourceSummary: string; +}; + +/** + * Adapts Network-owned STUN mappings and DDNS records to the QQBot source contract. + * + * It deliberately owns no event staging or delivery queue state; every delivery is + * recomputed from the current Network records to prevent stale endpoint sends. + */ +@Injectable() +export class NetworkStunMessageSourceAdapter + implements OnModuleDestroy, OnModuleInit +{ + private registered = false; + + /** Immutable public schema for the first Network-owned system message source. */ + readonly definition: SystemMessageSourceDefinition = { + description: '当 UDP STUN 映射端口变更且 IPv4 DDNS 已同步时发送消息。', + displayName: 'STUN 映射端口变更', + sourceKey: SOURCE_KEY, + subscriptionFields: [ + { + key: 'portForwardId', + label: 'STUN 端口转发', + optionCollection: 'portForwards', + required: true, + type: 'select', + }, + { + dependsOn: 'portForwardId', + key: 'ddnsRecordId', + label: 'IPv4 DDNS 记录', + optionCollection: 'ddnsRecords', + required: true, + type: 'select', + }, + ], + variables: [ + { + description: '所选 DDNS 完整域名', + example: 'pal.kwitsukasa.top', + key: 'domain', + label: '域名', + type: 'string', + }, + { + description: '新的公网映射端口', + example: '38213', + key: 'port', + label: '端口', + type: 'number', + }, + { + description: '服务端组合的域名与端口', + example: 'pal.kwitsukasa.top:38213', + key: 'endpoint', + label: '访问端点', + type: 'string', + }, + { + description: '端口转发显示名称', + example: '帕鲁新世界', + key: 'mappingName', + label: '映射名称', + type: 'string', + }, + { + description: '变化前的公网映射端口', + example: '8213', + key: 'previousPort', + label: '原端口', + type: 'number', + }, + { + description: '事件对应的公网 IPv4', + example: '203.0.113.10', + key: 'publicIpv4', + label: '公网 IPv4', + type: 'string', + }, + { + description: '按上海时区格式化的变化时间', + example: '2026-07-23 20:30:00', + key: 'changedAt', + label: '变化时间', + type: 'string', + }, + ], + version: 1, + }; + + /** + * Creates the Network-owned source adapter. + * @param mappingRepository - Current port-forward and Keeper state reader. + * @param ddnsRepository - Current A-record DDNS binding reader. + * @param sourceRegistry - Core-owned process registry exported to Network. + */ + constructor( + @InjectRepository(NetworkPortForward) + private readonly mappingRepository: Repository, + @InjectRepository(NetworkDdnsRecord) + private readonly ddnsRepository: Repository, + private readonly sourceRegistry: SystemMessageSourceRegistry, + ) {} + + /** Registers this instance once after its owning Network module initializes. */ + onModuleInit(): void { + if (this.registered) return; + this.sourceRegistry.register(this); + this.registered = true; + } + + /** Unregisters only this instance during Network module teardown. */ + onModuleDestroy(): void { + if (!this.registered) return; + this.sourceRegistry.unregister(this.definition.sourceKey, this); + this.registered = false; + } + + /** + * Normalizes a submitted subscription to the two supported string identifiers. + * @param input - Untrusted admin configuration; unknown fields are intentionally discarded. + * @returns Canonical IDs, the mapping resource key, and a server-derived summary. + */ + async normalizeSubscriptionConfig(input: unknown): Promise<{ + canonicalConfig: Record; + resourceKey: string; + sourceSummary: string; + }> { + const resolved = await this.resolveSubscription(input); + return { + canonicalConfig: resolved.config, + resourceKey: resolved.config.portForwardId, + sourceSummary: resolved.sourceSummary, + }; + } + + /** + * Checks current subscription validity without writing subscription or Network state. + * @param config - Persisted or submitted source configuration. + * @returns Stable validity result with a safe source summary for management UI. + */ + async inspectSubscription(config: Record): Promise<{ + invalidReasonCode: null | string; + sourceSummary: string; + valid: boolean; + }> { + try { + const resolved = await this.resolveSubscription(config); + return { + invalidReasonCode: null, + sourceSummary: resolved.sourceSummary, + valid: true, + }; + } catch (error) { + return { + invalidReasonCode: messageSourceErrorCode(error), + sourceSummary: '未选择有效的 STUN 映射与 DDNS', + valid: false, + }; + } + } + + /** + * Lists current mapping and DDNS candidates with structural, lease-independent eligibility. + * @returns Locked `portForwards` and `ddnsRecords` option collections. + */ + async listSubscriptionOptions(): Promise> { + const [mappings, records] = await Promise.all([ + this.mappingRepository.find({ order: { id: 'ASC', name: 'ASC' } }), + this.ddnsRepository.find({ order: { id: 'ASC', name: 'ASC' } }), + ]); + const mappingsById = new Map( + mappings.map((mapping) => [String(mapping.id), mapping]), + ); + return { + ddnsRecords: records.map((record) => { + const mapping = record.portForwardId + ? mappingsById.get(String(record.portForwardId)) + : undefined; + const disabledReasonCode = ddnsOptionReason(record, mapping); + return { + disabledReasonCode, + eligible: disabledReasonCode === null, + fqdn: ddnsFqdn(record), + id: String(record.id), + name: record.name, + portForwardId: record.portForwardId + ? String(record.portForwardId) + : '', + }; + }), + portForwards: mappings.map((mapping) => { + const { disabledReasonCode, eligible } = + classifyStunEndpointSource(mapping); + return { + disabledReasonCode, + eligible, + externalPort: mapping.externalPort, + id: String(mapping.id), + internalPort: mapping.internalPort, + name: mapping.name, + protocol: mapping.protocol, + }; + }), + }; + } + + /** + * Rejects unsafe event data and retains only the source's permitted scalar fields. + * @param payload - Network event candidate before it is staged in the QQBot Outbox. + * @returns Canonical event variables with no client-controlled endpoint field. + */ + validateEventPayload( + payload: Record, + ): Record { + if (!isPlainRecord(payload)) { + throw new SystemMessageContractError('invalid_message_event_payload'); + } + const changedAt = normalizeRfc3339(payload.changedAt); + const currentPort = normalizePort(payload.currentPort); + const previousPort = normalizePort(payload.previousPort); + const portForwardId = normalizeSnowflakeId(payload.portForwardId); + if ( + typeof payload.publicIpv4 !== 'string' || + isIP(payload.publicIpv4) !== 4 || + currentPort === previousPort + ) { + throw new SystemMessageContractError('invalid_message_event_payload'); + } + return { + changedAt, + currentPort, + portForwardId, + previousPort, + publicIpv4: payload.publicIpv4, + }; + } + + /** + * Re-evaluates a frozen event against the current mapping and DDNS state before sending. + * @param input - Event payload and persisted subscription configuration. + * @returns Ready variables, a DDNS wait result, or terminal cancellation/supersession. + */ + async resolveDelivery(input: { + eventPayload: Record; + subscriptionConfig: Record; + }): Promise { + let resolved: ResolvedSubscription; + try { + resolved = await this.resolveSubscription(input.subscriptionConfig); + } catch (error) { + return { reasonCode: messageSourceErrorCode(error), status: 'cancelled' }; + } + let event: StunEventPayload; + try { + event = this.validateEventPayload(input.eventPayload) as StunEventPayload; + } catch (error) { + return { reasonCode: messageSourceErrorCode(error), status: 'cancelled' }; + } + if (event.portForwardId !== resolved.config.portForwardId) { + return { reasonCode: 'event_resource_mismatch', status: 'cancelled' }; + } + if (!hasCurrentEndpoint(resolved.mapping)) { + return { reasonCode: 'endpoint_unavailable', status: 'superseded' }; + } + if ( + resolved.mapping.currentPublicPort !== event.currentPort || + resolved.mapping.currentPublicIpv4 !== event.publicIpv4 + ) { + return { reasonCode: 'endpoint_changed', status: 'superseded' }; + } + const variables = deliveryVariables(resolved, event); + if ( + resolved.ddnsRecord.syncStatus !== 'synced' || + resolved.ddnsRecord.appliedAddress !== event.publicIpv4 + ) { + return { + reasonCode: 'ddns_not_synced', + status: 'waiting_ddns', + variables, + }; + } + return { reasonCode: null, status: 'ready', variables }; + } + + /** + * Resolves and structurally validates both persisted resources selected by a subscription. + * @param input - Unknown configuration from an Admin request or persisted subscription. + * @returns Current resource rows and canonical identity when their relationship is valid. + */ + private async resolveSubscription( + input: unknown, + ): Promise { + if (!isPlainRecord(input)) { + throw new SystemMessageContractError('invalid_message_source_config'); + } + let config: StunSubscriptionConfig; + try { + config = { + ddnsRecordId: normalizeSnowflakeId(input.ddnsRecordId), + portForwardId: normalizeSnowflakeId(input.portForwardId), + }; + } catch { + throw new SystemMessageContractError('invalid_message_source_config'); + } + const mapping = await this.mappingRepository.findOne({ + where: { id: config.portForwardId }, + }); + if (!mapping) { + throw new SystemMessageContractError('source_not_found'); + } + const sourceEligibility = classifyStunEndpointSource(mapping); + if (!sourceEligibility.eligible) { + throw new SystemMessageContractError( + sourceEligibility.disabledReasonCode as string, + ); + } + const ddnsRecord = await this.ddnsRepository.findOne({ + where: { id: config.ddnsRecordId }, + }); + const ddnsReason = ddnsOptionReason(ddnsRecord, mapping); + if (ddnsReason) { + throw new SystemMessageContractError(ddnsReason); + } + return { + config, + ddnsRecord: ddnsRecord as NetworkDdnsRecord, + mapping, + sourceSummary: `${mapping.name} · ${ddnsFqdn(ddnsRecord as NetworkDdnsRecord)}`, + }; + } +} + +/** + * Verifies that a value is a non-array object with string-addressable own values. + * @param value - Untrusted input at an adapter boundary. + * @returns Whether the value is safe to inspect as a plain record. + */ +function isPlainRecord(value: unknown): value is Record { + return !!value && typeof value === 'object' && !Array.isArray(value); +} + +/** + * Validates a Snowflake at JSON boundaries without converting it to an unsafe number. + * @param value - Candidate string identifier. + * @returns The unchanged decimal string. + */ +function normalizeSnowflakeId(value: unknown): string { + if (typeof value !== 'string' || !SNOWFLAKE_ID_PATTERN.test(value)) { + throw new SystemMessageContractError('invalid_message_source_config'); + } + return value; +} + +/** + * Validates an event port as an integer transport port. + * @param value - Candidate event scalar. + * @returns The safe port number. + */ +function normalizePort(value: unknown): number { + if ( + typeof value !== 'number' || + !Number.isInteger(value) || + value < 1 || + value > 65_535 + ) { + throw new SystemMessageContractError('invalid_message_event_payload'); + } + return value; +} + +/** + * Normalizes a timezone-qualified event timestamp to canonical ISO UTC text. + * @param value - Candidate RFC3339 timestamp. + * @returns Canonical ISO string safe for later Shanghai rendering. + */ +function normalizeRfc3339(value: unknown): string { + if (typeof value !== 'string' || !RFC3339_PATTERN.test(value)) { + throw new SystemMessageContractError('invalid_message_event_payload'); + } + const date = new Date(value); + if (Number.isNaN(date.getTime())) { + throw new SystemMessageContractError('invalid_message_event_payload'); + } + return date.toISOString(); +} + +/** + * Determines whether a persisted DDNS record is a usable structural subscription target. + * @param record - Candidate A-record binding, possibly absent or deleted. + * @param mapping - Mapping selected by the subscription, if currently present. + * @returns Null when structurally valid or a stable invalid-reason code. + */ +function ddnsOptionReason( + record: NetworkDdnsRecord | null | undefined, + mapping: NetworkPortForward | undefined, +): null | string { + if (!record) return 'ddns_not_found'; + if (record.isDeleted) return 'ddns_deleted'; + if (!record.enabled) return 'ddns_disabled'; + if (record.recordType !== 'A') return 'ddns_a_required'; + if (record.sourceType !== 'port_forward_ipv4') + return 'ddns_source_type_invalid'; + if (!mapping || String(record.portForwardId) !== String(mapping.id)) { + return 'ddns_mapping_mismatch'; + } + const source = classifyStunEndpointSource(mapping); + return source.disabledReasonCode; +} + +/** + * Forms the normalized fully-qualified hostname from persisted server-owned DDNS fields. + * @param record - Canonical DDNS binding. + * @returns Lowercase FQDN without a trailing dot. + */ +function ddnsFqdn(record: NetworkDdnsRecord): string { + const domain = record.domain.trim().toLowerCase().replace(/\.$/, ''); + const subDomain = record.subDomain.trim().toLowerCase().replace(/\.$/, ''); + return subDomain === '@' ? domain : `${subDomain}.${domain}`; +} + +/** + * Checks current Keeper endpoint readiness independently of structural subscription validity. + * @param mapping - Current mapping row after structural eligibility passed. + * @returns Whether a live public IPv4 endpoint is available at this instant. + */ +function hasCurrentEndpoint(mapping: NetworkPortForward): boolean { + return ( + isIP(mapping.currentPublicIpv4 || '') === 4 && + typeof mapping.currentPublicPort === 'number' && + Number.isInteger(mapping.currentPublicPort) && + mapping.currentPublicPort >= 1 && + mapping.currentPublicPort <= 65_535 && + !!mapping.currentValidUntil && + new Date(mapping.currentValidUntil).getTime() > Date.now() + ); +} + +/** + * Builds the whitelisted template variables from current Network records and frozen event data. + * @param resolved - Current structurally valid subscription resources. + * @param event - Validated immutable event payload. + * @returns Server-derived scalar values ready for formal template rendering. + */ +function deliveryVariables( + resolved: ResolvedSubscription, + event: StunEventPayload, +): Record { + const domain = ddnsFqdn(resolved.ddnsRecord); + return { + changedAt: formatShanghaiDateTime(event.changedAt), + domain, + endpoint: `${domain}:${event.currentPort}`, + mappingName: resolved.mapping.name, + port: event.currentPort, + previousPort: event.previousPort, + publicIpv4: event.publicIpv4, + }; +} + +/** + * Formats an event timestamp as the locked Asia/Shanghai template variable value. + * @param value - Valid RFC3339 event timestamp. + * @returns `YYYY-MM-DD HH:mm:ss` assembled without locale punctuation assumptions. + */ +function formatShanghaiDateTime(value: string): string { + const parts = new Intl.DateTimeFormat('en-CA', { + day: '2-digit', + hour: '2-digit', + hourCycle: 'h23', + minute: '2-digit', + month: '2-digit', + second: '2-digit', + timeZone: 'Asia/Shanghai', + year: 'numeric', + }).formatToParts(new Date(value)); + const part = (type: Intl.DateTimeFormatPartTypes): string => + parts.find((item) => item.type === type)?.value || ''; + return `${part('year')}-${part('month')}-${part('day')} ${part('hour')}:${part('minute')}:${part('second')}`; +} + +/** + * Converts adapter validation failures to one stable, non-sensitive reason code. + * @param error - Unknown caught adapter or repository error. + * @returns Domain code for expected contract failures or a generic source-unavailable code. + */ +function messageSourceErrorCode(error: unknown): string { + return error instanceof SystemMessageContractError + ? error.code + : 'message_source_unavailable'; +} diff --git a/src/modules/qqbot/core/application/message-push/system-message-source.registry.ts b/src/modules/qqbot/core/application/message-push/system-message-source.registry.ts new file mode 100644 index 0000000..7b31870 --- /dev/null +++ b/src/modules/qqbot/core/application/message-push/system-message-source.registry.ts @@ -0,0 +1,64 @@ +import { Injectable } from '@nestjs/common'; +import { + SystemMessageContractError, + type SystemMessageSourceAdapter, + type SystemMessageSourceDefinition, +} from '../../contract/message-push/qqbot-message-push.types'; + +/** + * Holds the process-local system message-source adapters registered by owning modules. + * + * Definitions are cloned on listing so management callers cannot mutate adapter state. + */ +@Injectable() +export class SystemMessageSourceRegistry { + private readonly adapters = new Map(); + + /** + * Registers one source adapter once for the current Nest process. + * @param adapter - Network, Core, or future module-owned source implementation. + * @throws {SystemMessageContractError} When its source key is already registered. + */ + register(adapter: SystemMessageSourceAdapter): void { + const key = adapter.definition.sourceKey; + if (this.adapters.has(key)) { + throw new SystemMessageContractError('duplicate_message_source'); + } + this.adapters.set(key, adapter); + } + + /** + * Removes only the exact adapter instance that registered the source key. + * @param sourceKey - Registered source identity. + * @param adapter - Instance performing module teardown. + */ + unregister(sourceKey: string, adapter: SystemMessageSourceAdapter): void { + if (this.adapters.get(sourceKey) === adapter) { + this.adapters.delete(sourceKey); + } + } + + /** + * Returns the adapter for a known source key. + * @param sourceKey - Stable source identity stored by subscriptions. + * @returns The registered adapter instance. + * @throws {SystemMessageContractError} When no source owns the requested key. + */ + get(sourceKey: string): SystemMessageSourceAdapter { + const adapter = this.adapters.get(sourceKey); + if (!adapter) { + throw new SystemMessageContractError('unknown_message_source'); + } + return adapter; + } + + /** + * Lists source definitions in stable source-key order without exposing internal objects. + * @returns Detached definition snapshots for read-only management use. + */ + list(): SystemMessageSourceDefinition[] { + return [...this.adapters.values()] + .map(({ definition }) => structuredClone(definition)) + .sort((left, right) => left.sourceKey.localeCompare(right.sourceKey)); + } +} diff --git a/src/modules/qqbot/core/qqbot-core.module.ts b/src/modules/qqbot/core/qqbot-core.module.ts index 22d620d..00db79f 100644 --- a/src/modules/qqbot/core/qqbot-core.module.ts +++ b/src/modules/qqbot/core/qqbot-core.module.ts @@ -47,6 +47,7 @@ import { QqbotMessagePublishBinding } from '@/modules/qqbot/core/infrastructure/ import { QqbotMessagePublishTarget } from '@/modules/qqbot/core/infrastructure/persistence/message-push/qqbot-message-publish-target.entity'; import { QqbotMessageSubscription } from '@/modules/qqbot/core/infrastructure/persistence/message-push/qqbot-message-subscription.entity'; import { QqbotMessageTemplate } from '@/modules/qqbot/core/infrastructure/persistence/message-push/qqbot-message-template.entity'; +import { SystemMessageSourceRegistry } from './application/message-push/system-message-source.registry'; export { QQBOT_CORE_DOMAIN_CONTRACT } from './contract/qqbot-core.contract'; @@ -82,6 +83,7 @@ export const QQBOT_CORE_CONTROLLERS = [ ]; export const QQBOT_CORE_PROVIDERS = [ + SystemMessageSourceRegistry, QqbotAccountService, QqbotBusService, QqbotCommandEngineService, @@ -102,6 +104,7 @@ export const QQBOT_CORE_PROVIDERS = [ ]; export const QQBOT_CORE_EXPORTS = [ + SystemMessageSourceRegistry, QqbotAccountService, QqbotConfigService, QqbotDashboardService, diff --git a/test/admin/network-management/network-stun-message-source.adapter.spec.ts b/test/admin/network-management/network-stun-message-source.adapter.spec.ts new file mode 100644 index 0000000..abb6550 --- /dev/null +++ b/test/admin/network-management/network-stun-message-source.adapter.spec.ts @@ -0,0 +1,228 @@ +import type { Repository } from 'typeorm'; +import { NetworkDdnsRecord } from '../../../src/modules/admin/platform-config/network-management/network-ddns.entity'; +import { NetworkPortForward } from '../../../src/modules/admin/platform-config/network-management/network-management.entity'; +import { NetworkStunMessageSourceAdapter } from '../../../src/modules/admin/platform-config/network-management/network-stun-message-source.adapter'; +import { SystemMessageSourceRegistry } from '../../../src/modules/qqbot/core/application/message-push/system-message-source.registry'; + +type Harness = { + adapter: NetworkStunMessageSourceAdapter; + ddns: NetworkDdnsRecord; + mapping: NetworkPortForward; + registry: SystemMessageSourceRegistry; +}; + +/** Creates a mutable in-memory repository harness for the adapter contract. */ +function createHarness(): Harness { + const mapping = Object.assign(new NetworkPortForward(), { + currentPublicIpv4: '203.0.113.10', + currentPublicPort: 38213, + currentValidUntil: new Date('2026-07-24T13:00:00.000Z'), + desiredPresence: 'present' as const, + externalPort: 8213, + id: '2041700000000000001', + internalPort: 8213, + isDeleted: false, + keeperDesiredEnabled: true, + name: '帕鲁新世界', + protocol: 'udp' as const, + }); + const ddns = Object.assign(new NetworkDdnsRecord(), { + appliedAddress: '203.0.113.10', + domain: 'kwitsukasa.top', + enabled: true, + id: '2041700000000000002', + isDeleted: false, + name: '帕鲁域名', + portForwardId: mapping.id, + recordType: 'A' as const, + sourceType: 'port_forward_ipv4' as const, + subDomain: 'pal', + syncStatus: 'synced' as const, + }); + const mappings = [mapping]; + const records = [ddns]; + const mappingRepository = { + find: jest.fn(async () => mappings), + findOne: jest.fn( + async ({ where }) => + mappings.find((item) => item.id === where.id) || null, + ), + } as unknown as Repository; + const recordRepository = { + find: jest.fn(async () => records), + findOne: jest.fn( + async ({ where }) => records.find((item) => item.id === where.id) || null, + ), + } as unknown as Repository; + const registry = new SystemMessageSourceRegistry(); + return { + adapter: new NetworkStunMessageSourceAdapter( + mappingRepository, + recordRepository, + registry, + ), + ddns, + mapping, + registry, + }; +} + +/** Returns a valid event payload, with optional untrusted-field overrides. */ +function eventPayload(overrides: Record = {}) { + return { + changedAt: '2026-07-24T12:30:00.000Z', + currentPort: 38213, + endpoint: 'attacker.example:1', + portForwardId: '2041700000000000001', + previousPort: 8213, + publicIpv4: '203.0.113.10', + ...overrides, + }; +} + +describe('NetworkStunMessageSourceAdapter', () => { + it('registers once and only unregisters its own source instance', () => { + const { adapter, registry } = createHarness(); + adapter.onModuleInit(); + adapter.onModuleInit(); + expect(registry.get(adapter.definition.sourceKey)).toBe(adapter); + adapter.onModuleDestroy(); + expect(() => registry.get(adapter.definition.sourceKey)).toThrow( + 'unknown_message_source', + ); + }); + + it('accepts only an enabled equal-port UDP Keeper and its linked enabled A record', async () => { + const { adapter } = createHarness(); + await expect( + adapter.normalizeSubscriptionConfig({ + ddnsRecordId: '2041700000000000002', + ignored: 'removed', + portForwardId: '2041700000000000001', + }), + ).resolves.toEqual({ + canonicalConfig: { + ddnsRecordId: '2041700000000000002', + portForwardId: '2041700000000000001', + }, + resourceKey: '2041700000000000001', + sourceSummary: '帕鲁新世界 · pal.kwitsukasa.top', + }); + }); + + it.each([ + ['tcp', (harness: Harness) => (harness.mapping.protocol = 'tcp')], + ['unequal ports', (harness: Harness) => (harness.mapping.internalPort = 1)], + [ + 'disabled keeper', + (harness: Harness) => (harness.mapping.keeperDesiredEnabled = false), + ], + [ + 'deleted mapping', + (harness: Harness) => (harness.mapping.isDeleted = true), + ], + ['disabled DDNS', (harness: Harness) => (harness.ddns.enabled = false)], + ['deleted DDNS', (harness: Harness) => (harness.ddns.isDeleted = true)], + ['non-A DDNS', (harness: Harness) => (harness.ddns.recordType = 'AAAA')], + [ + 'mismatched DDNS mapping', + (harness: Harness) => + (harness.ddns.portForwardId = '2041700000000000003'), + ], + ])('rejects %s subscriptions', async (_name, mutate) => { + const harness = createHarness(); + mutate(harness); + await expect( + harness.adapter.normalizeSubscriptionConfig({ + ddnsRecordId: harness.ddns.id, + portForwardId: harness.mapping.id, + }), + ).rejects.toThrow(); + }); + + it('rejects malformed Snowflake IDs and strips unknown event/config fields', async () => { + const { adapter } = createHarness(); + await expect( + adapter.normalizeSubscriptionConfig({ + ddnsRecordId: 2041700000000000002, + portForwardId: 'not-an-id', + }), + ).rejects.toThrow('invalid_message_source_config'); + expect(adapter.validateEventPayload(eventPayload())).toEqual({ + changedAt: '2026-07-24T12:30:00.000Z', + currentPort: 38213, + portForwardId: '2041700000000000001', + previousPort: 8213, + publicIpv4: '203.0.113.10', + }); + expect(() => + adapter.validateEventPayload(eventPayload({ currentPort: '38213' })), + ).toThrow('invalid_message_event_payload'); + }); + + it('returns ready variables derived from the server-owned DDNS FQDN and Shanghai time', async () => { + const { adapter } = createHarness(); + await expect( + adapter.resolveDelivery({ + eventPayload: eventPayload(), + subscriptionConfig: { + ddnsRecordId: '2041700000000000002', + portForwardId: '2041700000000000001', + }, + }), + ).resolves.toEqual({ + reasonCode: null, + status: 'ready', + variables: { + changedAt: '2026-07-24 20:30:00', + domain: 'pal.kwitsukasa.top', + endpoint: 'pal.kwitsukasa.top:38213', + mappingName: '帕鲁新世界', + port: 38213, + previousPort: 8213, + publicIpv4: '203.0.113.10', + }, + }); + }); + + it('waits for DDNS, supersedes a replaced endpoint, and cancels a changed relationship', async () => { + const waiting = createHarness(); + waiting.ddns.appliedAddress = null; + await expect( + waiting.adapter.resolveDelivery({ + eventPayload: eventPayload(), + subscriptionConfig: { + ddnsRecordId: waiting.ddns.id, + portForwardId: waiting.mapping.id, + }, + }), + ).resolves.toMatchObject({ + reasonCode: 'ddns_not_synced', + status: 'waiting_ddns', + }); + + const superseded = createHarness(); + superseded.mapping.currentPublicPort = 39000; + await expect( + superseded.adapter.resolveDelivery({ + eventPayload: eventPayload(), + subscriptionConfig: { + ddnsRecordId: superseded.ddns.id, + portForwardId: superseded.mapping.id, + }, + }), + ).resolves.toMatchObject({ status: 'superseded' }); + + const cancelled = createHarness(); + cancelled.ddns.enabled = false; + await expect( + cancelled.adapter.resolveDelivery({ + eventPayload: eventPayload(), + subscriptionConfig: { + ddnsRecordId: cancelled.ddns.id, + portForwardId: cancelled.mapping.id, + }, + }), + ).resolves.toMatchObject({ status: 'cancelled' }); + }); +}); diff --git a/test/modules/qqbot/message-push/system-message-source.registry.spec.ts b/test/modules/qqbot/message-push/system-message-source.registry.spec.ts new file mode 100644 index 0000000..3b18093 --- /dev/null +++ b/test/modules/qqbot/message-push/system-message-source.registry.spec.ts @@ -0,0 +1,60 @@ +import type { SystemMessageSourceAdapter } from '../../../../src/modules/qqbot/core/contract/message-push/qqbot-message-push.types'; +import { SystemMessageSourceRegistry } from '../../../../src/modules/qqbot/core/application/message-push/system-message-source.registry'; + +/** Creates a minimal adapter whose definition can be safely registered. */ +function createAdapter(sourceKey: string): SystemMessageSourceAdapter { + return { + definition: { + description: 'test', + displayName: sourceKey, + sourceKey, + subscriptionFields: [], + variables: [], + version: 1, + }, + inspectSubscription: jest.fn(), + listSubscriptionOptions: jest.fn(), + normalizeSubscriptionConfig: jest.fn(), + resolveDelivery: jest.fn(), + validateEventPayload: jest.fn(), + }; +} + +describe('SystemMessageSourceRegistry', () => { + it('rejects duplicate source registration and returns immutable definitions', () => { + const registry = new SystemMessageSourceRegistry(); + const adapter = createAdapter('network.stun.mapping-port-changed'); + registry.register(adapter); + + expect(() => registry.register(adapter)).toThrow( + 'duplicate_message_source', + ); + expect(registry.list()).toEqual([adapter.definition]); + expect(() => registry.get('missing')).toThrow('unknown_message_source'); + + const [definition] = registry.list(); + definition.displayName = 'mutated'; + expect( + registry.get(adapter.definition.sourceKey).definition.displayName, + ).toBe(adapter.definition.displayName); + }); + + it('sorts definitions and unregisters only the same adapter instance', () => { + const registry = new SystemMessageSourceRegistry(); + const first = createAdapter('z.source'); + const second = createAdapter('a.source'); + const replacement = createAdapter('z.source'); + registry.register(first); + registry.register(second); + + registry.unregister('z.source', replacement); + expect(registry.list().map((definition) => definition.sourceKey)).toEqual([ + 'a.source', + 'z.source', + ]); + registry.unregister('z.source', first); + expect(registry.list().map((definition) => definition.sourceKey)).toEqual([ + 'a.source', + ]); + }); +});