mirror of
https://github.com/nbd-wtf/nostr-tools.git
synced 2026-07-30 19:26:14 +00:00
fix reconnection relay leakage and double onclose calls.
This commit is contained in:
+26
-6
@@ -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()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@nostr/tools",
|
||||
"version": "2.24.0",
|
||||
"version": "2.24.1",
|
||||
"exports": {
|
||||
".": "./index.ts",
|
||||
"./core": "./core.ts",
|
||||
|
||||
+1
-1
@@ -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",
|
||||
|
||||
+35
-40
@@ -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<void>((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<void>((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<void>((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<void>((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)
|
||||
|
||||
+95
-8
@@ -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<void>((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<void>((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<string, any>
|
||||
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<void>((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()
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user