112 lines
2.3 KiB
JavaScript
112 lines
2.3 KiB
JavaScript
/**
|
|
* Tracks live ProtomuxRPC sessions for broadcast and cleanup.
|
|
*/
|
|
import logger from '../utils/logger.js'
|
|
|
|
const log = logger.child('peers')
|
|
|
|
export class PeerRegistry {
|
|
constructor() {
|
|
/** @type {Map<string, import('../rpc/session.js').PeerSession>} */
|
|
this.sessions = new Map()
|
|
/** @type {Set<(size: number, prevSize: number) => void>} */
|
|
this._listeners = new Set()
|
|
}
|
|
|
|
/**
|
|
* Subscribe to peer count changes (add / remove / clear).
|
|
* @param {(size: number, prevSize: number) => void} fn
|
|
* @returns {() => void} unsubscribe
|
|
*/
|
|
onChange(fn) {
|
|
if (typeof fn !== 'function') return () => {}
|
|
this._listeners.add(fn)
|
|
return () => {
|
|
this._listeners.delete(fn)
|
|
}
|
|
}
|
|
|
|
/**
|
|
* @param {number} prevSize
|
|
*/
|
|
_emitChange(prevSize) {
|
|
const size = this.sessions.size
|
|
if (size === prevSize) return
|
|
for (const fn of this._listeners) {
|
|
try {
|
|
fn(size, prevSize)
|
|
} catch (err) {
|
|
log.debug('onChange listener failed', { error: err?.message })
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* @param {import('../rpc/session.js').PeerSession} session
|
|
*/
|
|
add(session) {
|
|
const prev = this.sessions.size
|
|
this.sessions.set(session.id, session)
|
|
this._emitChange(prev)
|
|
}
|
|
|
|
/**
|
|
* @param {string} id
|
|
*/
|
|
remove(id) {
|
|
const prev = this.sessions.size
|
|
const had = this.sessions.delete(id)
|
|
if (had) this._emitChange(prev)
|
|
}
|
|
|
|
/**
|
|
* @param {string} id
|
|
*/
|
|
get(id) {
|
|
return this.sessions.get(id) || null
|
|
}
|
|
|
|
get size() {
|
|
return this.sessions.size
|
|
}
|
|
|
|
[Symbol.iterator]() {
|
|
return this.sessions.values()
|
|
}
|
|
|
|
/**
|
|
* Fire a push event on every open session.
|
|
* @param {string} method
|
|
* @param {unknown} payload
|
|
*/
|
|
broadcast(method, payload) {
|
|
if (this.sessions.size === 0) return
|
|
for (const session of this.sessions.values()) {
|
|
try {
|
|
session.push(method, payload)
|
|
} catch (err) {
|
|
log.debug('Broadcast failed', {
|
|
peerId: session.id.slice(0, 12),
|
|
method,
|
|
error: err.message,
|
|
})
|
|
}
|
|
}
|
|
}
|
|
|
|
clear() {
|
|
const prev = this.sessions.size
|
|
for (const session of this.sessions.values()) {
|
|
try {
|
|
session.destroy()
|
|
} catch {
|
|
// ignore
|
|
}
|
|
}
|
|
this.sessions.clear()
|
|
this._emitChange(prev)
|
|
}
|
|
}
|
|
|
|
export const peers = new PeerRegistry()
|