From fed4a5561398a93dfae0663089b7196e4d6e534f Mon Sep 17 00:00:00 2001 From: fiatjaf Date: Tue, 21 Jul 2026 08:50:32 -0300 Subject: [PATCH] fix reconnection relay leakage and double onclose calls. --- abstract-relay.ts | 32 +++++++++++--- jsr.json | 2 +- package.json | 2 +- pool.test.ts | 75 ++++++++++++++++----------------- relay.test.ts | 103 ++++++++++++++++++++++++++++++++++++++++++---- 5 files changed, 158 insertions(+), 56 deletions(-) diff --git a/abstract-relay.ts b/abstract-relay.ts index 81cf690..44f19dc 100644 --- a/abstract-relay.ts +++ b/abstract-relay.ts @@ -133,6 +133,17 @@ export class AbstractRelay { } private handleHardClose(reason: string) { + // browsers and the node `ws` library fire BOTH onerror and onclose on a + // failed/abnormal socket. detach the handlers on the first invocation so the + // error->close pair collapses into a single terminal action -- otherwise + // onclose would fire twice (terminal paths) and reconnect() would be + // scheduled twice (reconnect paths, skipping a backoff slot). + if (this.ws) { + this.ws.onopen = null + this.ws.onerror = null + this.ws.onclose = null + } + if (this.pingIntervalHandle) { clearInterval(this.pingIntervalHandle) this.pingIntervalHandle = undefined @@ -164,12 +175,11 @@ export class AbstractRelay { connectionTimeoutHandle = setTimeout(() => { reject('connection timed out') this.connectionPromise = undefined - // Only give up on the initial connect; a slow reconnect + // only give up on the initial connect; a slow reconnect // should fall through to the next backoff slot. if (this.reconnectAttempts === 0) { this.skipReconnection = true } - this.onclose?.() this.handleHardClose('relay connection timed out') }, opts.timeout) } @@ -220,13 +230,12 @@ export class AbstractRelay { clearTimeout(connectionTimeoutHandle) reject('connection failed') this.connectionPromise = undefined - // Only give up on the initial connect. A failed reconnect + // only give up on the initial connect. a failed reconnect // attempt must fall through so the next backoff slot fires; // otherwise one failed retry tears down every subscription. if (this.reconnectAttempts === 0) { this.skipReconnection = true } - this.onclose?.() this.handleHardClose('relay connection failed') } @@ -435,11 +444,22 @@ export class AbstractRelay { } this.closeAllSubscriptions('relay connection closed by us') this._connected = false + this.connectionPromise = undefined this.idleSince = undefined this.clearIdleTimeout() this.onclose?.() - if (this.ws?.readyState === this._WebSocket.OPEN) { - this.ws?.close() + if (this.ws) { + // detach the handlers before closing so the resulting ws.onclose doesn't + // route into handleHardClose and fire onclose / closeAllSubscriptions a + // second time -- this user-initiated close() is the single owner of that. + this.ws.onopen = null + this.ws.onerror = null + this.ws.onclose = null + // close the socket unless it's already closing/closed; this also aborts a + // still-CONNECTING socket so it can't open and linger untracked. + if (this.ws.readyState !== this._WebSocket.CLOSING && this.ws.readyState !== this._WebSocket.CLOSED) { + this.ws.close() + } } } diff --git a/jsr.json b/jsr.json index eb74188..bce3fc3 100644 --- a/jsr.json +++ b/jsr.json @@ -1,6 +1,6 @@ { "name": "@nostr/tools", - "version": "2.24.0", + "version": "2.24.1", "exports": { ".": "./index.ts", "./core": "./core.ts", diff --git a/package.json b/package.json index 43dd80c..9c389a9 100644 --- a/package.json +++ b/package.json @@ -1,7 +1,7 @@ { "type": "module", "name": "nostr-tools", - "version": "2.24.0", + "version": "2.24.1", "description": "Tools for making a Nostr client.", "repository": { "type": "git", diff --git a/pool.test.ts b/pool.test.ts index 76c244e..8d4ff3d 100644 --- a/pool.test.ts +++ b/pool.test.ts @@ -255,10 +255,8 @@ test('ping-pong timeout in pool', async () => { test('reconnect on disconnect in pool', async () => { const mockRelay = mockRelays[0] - pool = new SimplePool({ enablePing: true, enableReconnect: true }) + pool = new SimplePool({ enableReconnect: true }) const relay = await pool.ensureRelay(mockRelay.url) - relay.pingTimeout = 50 - relay.pingFrequency = 50 relay.resubscribeBackoff = [50, 100] let closes = 0 @@ -268,51 +266,46 @@ test('reconnect on disconnect in pool', async () => { expect(relay.connected).toBeTrue() - // wait for the first ping to succeed - await new Promise(resolve => setTimeout(resolve, 75)) - expect(closes).toBe(0) + // drop the live socket, which schedules a reconnect (but must not fire onclose) + ;(relay as any).ws?.close() - // now make it unresponsive - mockRelay.unresponsive = true - - // wait for the second ping to fail, which will trigger a close - await new Promise(resolve => { + // wait for the connection to drop + await new Promise((resolve, reject) => { + const deadline = setTimeout(() => reject(new Error('relay never disconnected')), 2000) const interval = setInterval(() => { - if (closes > 0) { + if (!relay.connected) { + clearTimeout(deadline) clearInterval(interval) - resolve(null) + resolve() } }, 10) }) - expect(closes).toBe(1) expect(relay.connected).toBeFalse() + // a transient drop that is going to reconnect must NOT fire onclose + expect(closes).toBe(0) - // now make it responsive again - mockRelay.unresponsive = false - - // wait for reconnect - await new Promise(resolve => { + // wait for reconnect (the mock relay server is still running) + await new Promise((resolve, reject) => { + const deadline = setTimeout(() => reject(new Error('relay never reconnected')), 2000) const interval = setInterval(() => { if (relay.connected) { + clearTimeout(deadline) clearInterval(interval) - resolve(null) + resolve() } }, 10) }) expect(relay.connected).toBeTrue() - expect(closes).toBe(1) + expect(closes).toBe(0) }) test('reconnect with filter update in pool', async () => { const mockRelay = mockRelays[0] pool = new SimplePool({ - enablePing: true, enableReconnect: true, }) const relay = await pool.ensureRelay(mockRelay.url) - relay.pingTimeout = 50 - relay.pingFrequency = 50 relay.resubscribeBackoff = [50, 100] let closes = 0 @@ -325,40 +318,42 @@ test('reconnect with filter update in pool', async () => { const sub = relay.subscribe([{ kinds: [1], since: 0 }], { onevent: () => {} }) expect(sub.filters[0].since).toBe(0) - // wait for the first ping to succeed - await new Promise(resolve => setTimeout(resolve, 75)) + // wait for events to arrive so lastEmitted gets set (used to bump `since` on reconnect) + await new Promise(resolve => setTimeout(resolve, 50)) expect(closes).toBe(0) - // now make it unresponsive - mockRelay.unresponsive = true + // drop the live socket, which schedules a reconnect (but must not fire onclose) + ;(relay as any).ws?.close() - // wait for the second ping to fail, which will trigger a close - await new Promise(resolve => { + // wait for the connection to drop + await new Promise((resolve, reject) => { + const deadline = setTimeout(() => reject(new Error('relay never disconnected')), 2000) const interval = setInterval(() => { - if (closes > 0) { + if (!relay.connected) { + clearTimeout(deadline) clearInterval(interval) - resolve(null) + resolve() } }, 10) }) - expect(closes).toBe(1) expect(relay.connected).toBeFalse() + // a transient drop that is going to reconnect must NOT fire onclose + expect(closes).toBe(0) - // now make it responsive again - mockRelay.unresponsive = false - - // wait for reconnect - await new Promise(resolve => { + // wait for reconnect (the mock relay server is still running) + await new Promise((resolve, reject) => { + const deadline = setTimeout(() => reject(new Error('relay never reconnected')), 2000) const interval = setInterval(() => { if (relay.connected) { + clearTimeout(deadline) clearInterval(interval) - resolve(null) + resolve() } }, 10) }) expect(relay.connected).toBeTrue() - expect(closes).toBe(1) + expect(closes).toBe(0) // check if filter was updated expect(sub.filters[0].since).toBeGreaterThan(1) diff --git a/relay.test.ts b/relay.test.ts index a6bf83f..22a6b0c 100644 --- a/relay.test.ts +++ b/relay.test.ts @@ -3,6 +3,7 @@ import { Server } from 'mock-socket' import { finalizeEvent, generateSecretKey, getPublicKey } from './pure.ts' import { NostrEvent } from './pure.ts' import { Relay, useWebSocketImplementation } from './relay.ts' +import { AbstractSimplePool } from './abstract-pool.ts' import { MockRelay, MockWebSocketClient } from './test-helpers.ts' useWebSocketImplementation(MockWebSocketClient) @@ -309,33 +310,43 @@ test('reconnect on disconnect', async () => { // now make it unresponsive mockRelay.unresponsive = true - // wait for the second ping to fail, which will trigger a close - await new Promise(resolve => { + // wait for the second ping to fail, which will drop the connection (but schedule a reconnect) + let sawDisconnect = false + await new Promise((resolve, reject) => { + const deadline = setTimeout(() => reject(new Error('relay never disconnected')), 2000) const interval = setInterval(() => { - if (closes > 0) { + if (!relay.connected) { + sawDisconnect = true + clearTimeout(deadline) clearInterval(interval) - resolve(null) + resolve() } }, 10) }) - expect(closes).toBe(1) + expect(sawDisconnect).toBeTrue() expect(relay.connected).toBeFalse() + // a transient drop that is going to reconnect must NOT fire onclose + expect(closes).toBe(0) // now make it responsive again mockRelay.unresponsive = false // wait for reconnect - await new Promise(resolve => { + await new Promise((resolve, reject) => { + const deadline = setTimeout(() => reject(new Error('relay never reconnected')), 2000) const interval = setInterval(() => { if (relay.connected) { + clearTimeout(deadline) clearInterval(interval) - resolve(null) + resolve() } }, 10) }) expect(relay.connected).toBeTrue() - expect(closes).toBe(1) // should not have closed again + expect(closes).toBe(0) // reconnecting must never fire onclose + + relay.close() }) test('reconnect survives a failed reconnect attempt and recovers when the relay returns', async () => { @@ -437,3 +448,79 @@ test('oninvalidevent is called for events that do not match subscription filters relay._onmessage({ data: JSON.stringify(['EVENT', sub.id, event]) } as MessageEvent) }) + +test('a failing-then-succeeding reconnect never orphans or duplicates the relay in the pool map', async () => { + const url = 'wss://reconnect.leak.test/1' + let phase: 'up' | 'down' = 'up' + + class FlakyWS extends EventTarget { + static OPEN = 1 + static CLOSED = 3 + readyState = 0 + onopen: any + onclose: any + onerror: any + onmessage: any + constructor(public url: string) { + super() + setTimeout(() => { + if (phase === 'up') { + this.readyState = 1 + this.onopen?.() + } else { + this.onerror?.(new Event('error')) + } + }, 5) + } + send() {} + close() { + this.readyState = 3 + this.onclose?.({}) + } + } + + const pool = new AbstractSimplePool({ + verifyEvent: () => true, + enableReconnect: true, + websocketImplementation: FlakyWS as any, + maxWaitForConnection: 3000, + }) + const relay = await pool.ensureRelay(url) + relay.resubscribeBackoff = [30, 30, 30, 30] + relay.subscribe([{ kinds: [1] }], { onevent: () => {} }) + expect(relay.openSubs.size).toBe(1) + + const map = (pool as any).relays as Map + expect(map.get(url)).toBe(relay) + + // go down and drop the live socket -> schedules a reconnect whose attempt will error via onerror + phase = 'down' + ;(relay as any).ws.close() + + // let at least one reconnect attempt FAIL via onerror (the previously-leaky path) + await new Promise(r => setTimeout(r, 150)) + expect(relay.connected).toBeFalse() + // INVARIANT: still the SAME tracked relay, no orphan, no duplicate + expect(map.get(url)).toBe(relay) + expect(map.size).toBe(1) + expect(relay.openSubs.size).toBe(1) + + // bring it back; next backoff slot reconnects successfully + phase = 'up' + await new Promise((resolve, reject) => { + const deadline = setTimeout(() => reject(new Error('never reconnected')), 1500) + const interval = setInterval(() => { + if (relay.connected) { + clearTimeout(deadline) + clearInterval(interval) + resolve() + } + }, 10) + }) + expect(relay.connected).toBeTrue() + expect(map.get(url)).toBe(relay) + expect(map.size).toBe(1) + expect(relay.openSubs.size).toBe(1) + + relay.close() +})