import type { C2RMessage, EventVerifier, Filter, NostrEvent } from "@nostr-fetch/kernel/nostr"; import { FilterMatcher, generateSubId, parseR2CMessage } from "@nostr-fetch/kernel/nostr"; import { type WSCloseEvent, type WSMessageEvent, type WebSocketMin, type WebSocketMinCtor, WebSocketReadyState, } from "@nostr-fetch/kernel/webSocket"; type Callback = E extends void ? () => void : (ev: E) => void; export interface Relay { url: string; wsReadyState: number; connect(): Promise; close(): void; prepareSub(filters: Filter[], options: SubscriptionOptions): Subscription; on(type: E, cb: RelayEventCbTypes[E]): void; off(type: E, cb: RelayEventCbTypes[E]): void; } export type RelayOptions = { connectTimeoutMs: number; webSocketConstructor: WebSocketMinCtor; }; export const initRelay = (relayUrl: string, options: RelayOptions): Relay => { return new RelayImpl(relayUrl, options); }; export type RelayConnectCb = Callback; export type RelayDisconnectCb = Callback; export type RelayNoticeCb = Callback; export type RelayErrorCb = Callback; export type RelayEventCbTypes = { connect: RelayConnectCb; disconnect: RelayDisconnectCb; notice: RelayNoticeCb; error: RelayErrorCb; }; export type RelayEventTypes = keyof RelayEventCbTypes; type RelayListenersTable = { [E in RelayEventTypes]: Set; }; class RelayImpl implements Relay { #relayUrl: string; #ws: WebSocketMin | undefined; #options: Required; #listeners: RelayListenersTable = { connect: new Set(), disconnect: new Set(), notice: new Set(), error: new Set(), }; #subscriptions: Map = new Map(); #msgQueue: string[] = []; #handleMsgsInterval: ReturnType | undefined; constructor(relayUrl: string, options: RelayOptions) { this.#relayUrl = relayUrl; this.#options = options; } public get url(): string { return this.#relayUrl; } public get wsReadyState(): number { return this.#ws?.readyState ?? WebSocketReadyState.CONNECTING; } private async forwardToSubAsync( subId: string, forwardFn: (sub: RelaySubscription) => Promise, ) { const targSub = this.#subscriptions.get(subId); if (targSub !== undefined) { await forwardFn(targSub); } } private async forwardToSub(subId: string, forwardFn: (sub: RelaySubscription) => void) { const targSub = this.#subscriptions.get(subId); if (targSub !== undefined) { forwardFn(targSub); } } private async handleMsgs() { if (this.#msgQueue.length === 0) { clearInterval(this.#handleMsgsInterval); this.#handleMsgsInterval = undefined; return; } const dispatchStartedAt = performance.now(); while (this.#msgQueue.length > 0 && performance.now() - dispatchStartedAt < 5.0) { const rawMsg = this.#msgQueue.shift() as string; const parsed = parseR2CMessage(rawMsg); if (parsed === undefined) { continue; } switch (parsed[0]) { case "EVENT": { const [, subId, ev] = parsed; await this.forwardToSubAsync(subId, (sub) => sub._forwardEvent(ev)); break; } case "EOSE": { const [, subId] = parsed; this.forwardToSub(subId, (sub) => sub._forwardEose()); break; } case "CLOSED": { const [, subId, msg] = parsed; this.forwardToSub(subId, (sub) => sub._forwardClosed(msg)); break; } case "NOTICE": { const [, notice] = parsed; this.#listeners.notice.forEach((cb) => cb(notice)); break; } } } } public async connect(): Promise { return new Promise((resolve, reject) => { let isTimedout = false; const timeout = setTimeout(() => { isTimedout = true; reject(Error(`attempt to connect to the relay '${this.#relayUrl}' timed out`)); }, this.#options.connectTimeoutMs); const onErrorBeforeOpen = () => { reject(Error("WebSocket error")); clearTimeout(timeout); }; const ws = new this.#options.webSocketConstructor(this.#relayUrl); ws.addEventListener("open", () => { if (!isTimedout) { this.#listeners.connect.forEach((cb) => cb()); this.#ws = ws; // set error listeners after the connection opened successfully ws.removeEventListener("error", onErrorBeforeOpen); ws.addEventListener("error", () => { this.#listeners.error.forEach((cb) => cb()); }); resolve(this); clearTimeout(timeout); } }); // error listeners are *not* activated while attempt to connect is in progress ws.addEventListener("error", onErrorBeforeOpen); ws.addEventListener("close", (e: WSCloseEvent) => { const reducted = { code: e.code, reason: e.reason, wasClean: e.wasClean, }; this.#listeners.disconnect.forEach((cb) => cb(reducted)); }); ws.addEventListener("message", (e: WSMessageEvent) => { this.#msgQueue.push(e.data); if (this.#handleMsgsInterval === undefined) { this.#handleMsgsInterval = setInterval(() => this.handleMsgs(), 0); } }); }); } public close() { if (this.#ws !== undefined) { this.#ws.close(); } } public prepareSub(filters: Filter[], options: SubscriptionOptions): Subscription { const subId = options.subId ?? generateSubId(); const sub = new RelaySubscription(this, subId, filters, options); this.#subscriptions.set(subId, sub); return sub; } _removeSub(subId: string) { this.#subscriptions.delete(subId); } public on(type: E, cb: RelayEventCbTypes[E]) { this.#listeners[type].add(cb); } public off(type: E, cb: RelayEventCbTypes[E]) { this.#listeners[type].delete(cb); } _sendC2RMessage(msg: C2RMessage) { if (this.#ws === undefined || this.#ws.readyState !== WebSocketReadyState.OPEN) { throw Error("not connected to the relay"); } this.#ws.send(JSON.stringify(msg)); } } type EoseEventPayload = { aborted: boolean; }; export type SubEventCb = Callback; export type SubEoseCb = Callback; export type SubClosedCb = Callback; export type SubEventCbTypes = { event: SubEventCb; eose: SubEoseCb; closed: SubClosedCb; }; export type SubEventTypes = keyof SubEventCbTypes; export interface Subscription { subId: string; req(): void; close(): void; on(type: E, cb: SubEventCbTypes[E]): void; off(type: E, cb: SubEventCbTypes[E]): void; } export interface SubscriptionOptions { subId?: string; eventVerifier: EventVerifier; skipVerification: boolean; skipFilterMatching: boolean; abortSubBeforeEoseTimeoutMs: number; } type SubListenersTable = { [E in SubEventTypes]: Set; }; class RelaySubscription implements Subscription { #relay: RelayImpl; #subId: string; #filters: Filter[]; #filterMatcher: FilterMatcher; #options: SubscriptionOptions; #listeners: SubListenersTable = { event: new Set(), eose: new Set(), closed: new Set(), }; #abortSubTimer: ReturnType | undefined; constructor(relay: RelayImpl, subId: string, filters: Filter[], options: SubscriptionOptions) { this.#relay = relay; this.#subId = subId; this.#filters = filters; this.#filterMatcher = new FilterMatcher(filters); this.#options = options; } public get subId(): string { return this.#subId; } public req() { this.#relay._sendC2RMessage(["REQ", this.#subId, ...this.#filters]); this.#resetAbortSubTimer(); } public close() { this.#clearListeners(); this.#relay._removeSub(this.#subId); this.#relay._sendC2RMessage(["CLOSE", this.#subId]); } public on(type: E, cb: SubEventCbTypes[E]) { this.#listeners[type].add(cb); } public off(type: E, cb: SubEventCbTypes[E]) { this.#listeners[type].delete(cb); } #clearListeners() { for (const s of Object.values(this.#listeners)) { s.clear(); } if (this.#abortSubTimer !== undefined) { clearTimeout(this.#abortSubTimer); } } #resetAbortSubTimer() { if (this.#abortSubTimer !== undefined) { clearTimeout(this.#abortSubTimer); this.#abortSubTimer = undefined; } this.#abortSubTimer = setTimeout(() => { this.#listeners.eose.forEach((cb) => cb({ aborted: true })); }, this.#options.abortSubBeforeEoseTimeoutMs); } async _forwardEvent(ev: NostrEvent) { this.#resetAbortSubTimer(); if (!this.#options.skipVerification && !(await this.#options.eventVerifier(ev))) { return; } if (!this.#options.skipFilterMatching && !this.#filterMatcher.match(ev)) { return; } this.#listeners.event.forEach((cb) => cb(ev)); } _forwardEose() { if (this.#abortSubTimer !== undefined) { clearTimeout(this.#abortSubTimer); } this.#listeners.eose.forEach((cb) => cb({ aborted: false })); } _forwardClosed(msg: string) { this.#listeners.closed.forEach((cb) => cb(msg)); // subscription has been closed by the relay -> clean up things this.#clearListeners(); this.#relay._removeSub(this.#subId); } }