Compare commits

...
3 Commits
6 changed files with 165 additions and 60 deletions
+32 -10
View File
@@ -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()
}
}
}
@@ -650,10 +670,12 @@ export class Subscription {
this.relay.openSubs.delete(this.id)
// compute idleness state
this.relay.ongoingOperations--
if (this.relay.ongoingOperations === 0) {
this.relay.idleSince = Date.now()
this.relay.scheduleIdleClose()
if (!this.id.startsWith('<forced-ping>')) {
this.relay.ongoingOperations--
if (this.relay.ongoingOperations === 0) {
this.relay.idleSince = Date.now()
this.relay.scheduleIdleClose()
}
}
this.onclose?.(reason)
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@nostr/tools",
"version": "2.24.0",
"version": "2.24.1",
"exports": {
".": "./index.ts",
"./core": "./core.ts",
+1 -1
View File
@@ -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
View File
@@ -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
View File
@@ -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()
})
+1
View File
@@ -128,6 +128,7 @@ export function mergeReverseSortedLists(list1: NostrEvent[], list2: NostrEvent[]
* Checks if a string is a 64-char lowercase hex string as most Nostr ids and pubkeys.
*/
export function isHex32(input: string): boolean {
if (input.length !== 64) return false
for (let i = 0; i < 64; i++) {
let cc = input.charCodeAt(i)
if (isNaN(cc) || cc < 48 || cc > 102 || (cc > 57 && cc < 97)) {