diff --git a/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts b/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts index 415bde0..ee80152 100644 --- a/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts +++ b/src/modules/admin/platform-config/network-management/network-agent-mqtt.service.ts @@ -652,6 +652,7 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { throw new NetworkMessageValidationError('Invalid event revision'); } const mapping = await manager.getRepository(NetworkPortForward).findOne({ + lock: { mode: 'pessimistic_write' }, where: { id: event.mappingId }, }); if (!mapping) { @@ -662,7 +663,8 @@ export class NetworkAgentMqttService implements OnModuleInit, OnModuleDestroy { return false; } const previousHistory = await repository.findOne({ - order: { id: 'DESC', occurredAt: 'DESC' }, + lock: { mode: 'pessimistic_read' }, + order: { occurredAt: 'DESC', id: 'DESC' }, where: { mappingId: event.mappingId }, }); const history = repository.create({ diff --git a/test/admin/network-management/network-agent-mqtt.service.spec.ts b/test/admin/network-management/network-agent-mqtt.service.spec.ts index f00187b..ca2d7a5 100644 --- a/test/admin/network-management/network-agent-mqtt.service.spec.ts +++ b/test/admin/network-management/network-agent-mqtt.service.spec.ts @@ -39,6 +39,29 @@ type MqttHarness = { transactionCalls: () => number; }; +type Deferred = { + promise: Promise; + resolve: (value: T | PromiseLike) => void; +}; + +type ConcurrentEndpointHarness = { + historyStageEntered: Promise; + histories: NetworkEndpointHistory[]; + releaseFirstStage: () => void; + service: NetworkAgentMqttService; + stageCalls: SystemMessageEventInput[]; + stagedEvents: SystemMessageEventInput[]; +}; + +/** Creates one externally controlled promise for deterministic transaction interleaving. */ +function createDeferred(): Deferred { + let resolve: (value: T | PromiseLike) => void = () => undefined; + const promise = new Promise((nextResolve) => { + resolve = nextResolve; + }); + return { promise, resolve }; +} + /** Creates a fake MQTT client plus in-memory TypeORM state for bridge tests. */ function createHarness(): MqttHarness { const state = Object.assign(new NetworkAgentState(), { @@ -93,9 +116,16 @@ function createHarness(): MqttHarness { if (!order) return matches[0] || null; return ( [...matches].sort((left, right) => { - const occurredAt = - right.occurredAt.getTime() - left.occurredAt.getTime(); - return occurredAt || String(right.id).localeCompare(String(left.id)); + for (const [key, direction] of Object.entries(order)) { + const leftValue = left[key as keyof NetworkEndpointHistory]; + const rightValue = right[key as keyof NetworkEndpointHistory]; + const compare = + leftValue instanceof Date && rightValue instanceof Date + ? leftValue.getTime() - rightValue.getTime() + : String(leftValue).localeCompare(String(rightValue)); + if (compare !== 0) return direction === 'DESC' ? -compare : compare; + } + return 0; })[0] || null ); }); @@ -207,6 +237,169 @@ function createHarness(): MqttHarness { }; } +/** + * Creates isolated transaction views with a mapping-row lock and delayed first Outbox stage. + * The committed arrays change only after each transaction callback returns successfully. + */ +function createConcurrentEndpointHarness(): ConcurrentEndpointHarness { + const state = Object.assign(new NetworkAgentState(), { + agentId: 'nas-main', + desiredRevision: '7', + }); + const mapping = Object.assign(new NetworkPortForward(), { + id: '100', + isDeleted: false, + }); + const histories = [endpointHistory({ publicPort: 8213 })]; + const stageCalls: SystemMessageEventInput[] = []; + const stagedEvents: SystemMessageEventInput[] = []; + const firstStageEntered = createDeferred(); + const firstStageReleased = createDeferred(); + const transactions = new WeakMap< + object, + { + pendingHistories: NetworkEndpointHistory[]; + pendingStagedEvents: SystemMessageEventInput[]; + releaseMappingLock?: () => void; + } + >(); + let mappingLockTail = Promise.resolve(); + + /** Acquires the fake exclusive mapping-row lock in transaction commit order. */ + async function acquireMappingLock(manager: object): Promise { + const predecessor = mappingLockTail; + const released = createDeferred(); + mappingLockTail = released.promise; + await predecessor; + const transaction = transactions.get(manager); + if (transaction) { + transaction.releaseMappingLock = () => released.resolve(undefined); + } + } + + /** Sorts histories according to the precise TypeORM order-property insertion order. */ + function findNewestHistory( + rows: readonly NetworkEndpointHistory[], + order: Record, + ): NetworkEndpointHistory | null { + return ( + [...rows].sort((left, right) => { + for (const [key, direction] of Object.entries(order)) { + const leftValue = left[key as keyof NetworkEndpointHistory]; + const rightValue = right[key as keyof NetworkEndpointHistory]; + const compare = + leftValue instanceof Date && rightValue instanceof Date + ? leftValue.getTime() - rightValue.getTime() + : String(leftValue).localeCompare(String(rightValue)); + if (compare !== 0) return direction === 'DESC' ? -compare : compare; + } + return 0; + })[0] || null + ); + } + + const stager = { + /** Stages into the caller transaction and pauses only the first event before commit. */ + stage: jest.fn( + async (manager: EntityManager, input: SystemMessageEventInput) => { + stageCalls.push(input); + if (input.eventId === 'endpoint-event-2') { + firstStageEntered.resolve(undefined); + await firstStageReleased.promise; + } + const transaction = transactions.get(manager); + if (!transaction) + throw new Error('missing transaction-local Outbox state'); + transaction.pendingStagedEvents.push(input); + return 'accepted' as const; + }, + ), + } as jest.Mocked; + const dataSource = { + transaction: async (work) => { + const pendingHistories: NetworkEndpointHistory[] = []; + const pendingStagedEvents: SystemMessageEventInput[] = []; + const manager = { + getRepository: (entity) => { + if (entity === NetworkAgentState) { + return { + findOne: async () => state, + } as unknown as Repository; + } + if (entity === NetworkPortForward) { + return { + findOne: async ({ lock, where }) => { + if (where.id !== mapping.id) return null; + if (lock?.mode === 'pessimistic_write') { + await acquireMappingLock(manager); + } + return mapping; + }, + } as unknown as Repository; + } + if (entity === NetworkEndpointHistory) { + return { + create: (input) => + Object.assign(new NetworkEndpointHistory(), input), + findOne: async ({ order, where }) => { + const visible = [...histories, ...pendingHistories]; + if (where.eventId) { + return ( + visible.find((item) => item.eventId === where.eventId) || + null + ); + } + return findNewestHistory( + visible.filter((item) => item.mappingId === where.mappingId), + order, + ); + }, + save: async (history: NetworkEndpointHistory) => { + history.id ||= String( + histories.length + pendingHistories.length + 1, + ); + pendingHistories.push(history); + return history; + }, + } as unknown as Repository; + } + throw new Error('unexpected repository'); + }, + } as unknown as EntityManager; + transactions.set(manager, { pendingHistories, pendingStagedEvents }); + try { + const result = await work(manager); + histories.push(...pendingHistories); + stagedEvents.push(...pendingStagedEvents); + return result; + } finally { + transactions.get(manager)?.releaseMappingLock?.(); + } + }, + } as unknown as DataSource; + const configService = { + get: (key) => + ({ + NETWORK_AGENT_ID: 'nas-main', + NETWORK_AGENT_MQTT_RETRY_MS: '60000', + })[key], + } as ConfigService; + const service = new NetworkAgentMqttService( + configService, + dataSource, + { publishCommitted: jest.fn() } as never, + stager, + ); + return { + historyStageEntered: firstStageEntered.promise, + histories, + releaseFirstStage: () => firstStageReleased.resolve(undefined), + service, + stageCalls, + stagedEvents, + }; +} + /** Builds one endpoint event payload with a valid STUN endpoint transition shape. */ function endpointEvent(overrides: Record = {}): Buffer { return Buffer.from( @@ -930,15 +1123,15 @@ describe('NetworkAgentMqttService', () => { it('orders prior history by occurredAt then id before deciding the previous port', async () => { const harness = createHarness(); harness.histories.push( - endpointHistory({ id: '2', publicPort: 8213 }), endpointHistory({ - eventId: 'endpoint-event-older', - id: '1', - publicPort: 38213, + id: '999', + occurredAt: new KtDateTime('2026-07-22T01:02:03.000Z'), + publicPort: 8213, }), endpointHistory({ - eventId: 'endpoint-event-newest', - id: '3', + eventId: 'endpoint-event-newer', + id: '1', + occurredAt: new KtDateTime('2026-07-22T01:02:04.000Z'), publicPort: 45000, }), ); @@ -949,7 +1142,8 @@ describe('NetworkAgentMqttService', () => { ); expect(harness.historyFindOne).toHaveBeenCalledWith({ - order: { id: 'DESC', occurredAt: 'DESC' }, + lock: { mode: 'pessimistic_read' }, + order: { occurredAt: 'DESC', id: 'DESC' }, where: { mappingId: '100' }, }); expect(harness.stagedEvents[0]?.payload).toMatchObject({ @@ -957,6 +1151,65 @@ describe('NetworkAgentMqttService', () => { }); }); + it('uses history ID as the deterministic tie-breaker for equal occurrence times', async () => { + const harness = createHarness(); + const occurredAt = new KtDateTime('2026-07-22T01:02:04.000Z'); + harness.histories.push( + endpointHistory({ id: '1', occurredAt, publicPort: 8213 }), + endpointHistory({ id: '2', occurredAt, publicPort: 45000 }), + ); + + await harness.service.consumeMessage( + 'kt/network/v1/agents/nas-main/events', + endpointEvent(), + ); + + expect(harness.stagedEvents[0]?.payload).toMatchObject({ + previousPort: 45000, + }); + }); + + it('serializes concurrent changes and stages the second transition from the first committed port', async () => { + const harness = createConcurrentEndpointHarness(); + const topic = 'kt/network/v1/agents/nas-main/events'; + const first = harness.service.consumeMessage(topic, endpointEvent()); + await harness.historyStageEntered; + + const second = harness.service.consumeMessage( + topic, + endpointEvent({ + endpoint: { + observedAt: '2026-07-22T01:02:06.000Z', + publicIpv4: '8.8.4.4', + publicPort: 45000, + validUntil: '2026-07-22T01:04:06.000Z', + }, + eventId: 'endpoint-event-3', + occurredAt: '2026-07-22T01:02:07.000Z', + }), + ); + await flushPromises(); + + expect(harness.stageCalls).toHaveLength(1); + expect(harness.stagedEvents).toHaveLength(0); + harness.releaseFirstStage(); + await Promise.all([first, second]); + + expect(harness.histories).toHaveLength(3); + expect(harness.stagedEvents).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + eventId: 'endpoint-event-2', + payload: expect.objectContaining({ previousPort: 8213 }), + }), + expect.objectContaining({ + eventId: 'endpoint-event-3', + payload: expect.objectContaining({ previousPort: 38213 }), + }), + ]), + ); + }); + it('does not save or stage a duplicate endpoint event ID', async () => { const harness = createHarness(); harness.histories.push(endpointHistory({ eventId: 'endpoint-event-2' }));