// The base (plaintext) SyncTransport: one WebSocket to the relay, multiplexing // every joined document plus the invitation inbox. Portable across Deno // (runtime/tests) and the browser. Auth is a caller-supplied token provider // (dev HMAC now; atproto service-auth later) — the transport never mints it. import { type AcceptResult, type ClientFrame, type EphemeralMessage, fromB64, type Invitation, type JoinOptions, type LogMessage, type Member, type ServerFrame, type Session, type SyncManifest, type SyncTransport, toB64, } from './protocol.ts'; interface Pending { resolve: (frame: ServerFrame) => void; reject: (err: Error) => void; } class RelaySession implements Session { me!: Member; #roster: Member[] = []; #msg = new Set<(m: LogMessage) => void>(); #eph = new Set<(m: EphemeralMessage) => void>(); #ros = new Set<(r: Member[]) => void>(); constructor(readonly docId: string, private conn: RelaySyncTransport) {} _init(me: Member, roster: Member[]) { this.me = me; this.#roster = roster; } _message(m: LogMessage) { for (const h of this.#msg) h(m); } _ephemeral(m: EphemeralMessage) { for (const h of this.#eph) h(m); } _roster(r: Member[]) { this.#roster = r; for (const h of this.#ros) h(r); } roster() { return this.#roster.slice(); } append(bytes: Uint8Array): Promise { return this.conn._append(this.docId, bytes); } ephemeral(bytes: Uint8Array): void { this.conn._send({ t: 'ephemeral', docId: this.docId, b64: toB64(bytes) }); } snapshot(bytes: Uint8Array, atSeq: number): Promise { this.conn._send({ t: 'snapshot', docId: this.docId, atSeq, b64: toB64(bytes) }); return Promise.resolve(); } onMessage(h: (m: LogMessage) => void) { this.#msg.add(h); return () => this.#msg.delete(h); } onEphemeral(h: (m: EphemeralMessage) => void) { this.#eph.add(h); return () => this.#eph.delete(h); } onRoster(h: (r: Member[]) => void) { this.#ros.add(h); return () => this.#ros.delete(h); } leave() { this.conn._send({ t: 'leave', docId: this.docId }); this.conn._dropSession(this.docId); } } export interface RelayClientOptions { url: string; // ws://host/ws did: string; handle?: string; /** Returns a fresh bearer token (dev HMAC or atproto service-auth JWT). */ token: () => Promise | string; /** For creating documents: the role manifest to register. */ manifestFor?: (docId: string) => SyncManifest | undefined; /** The socket closed (relay restart, network drop). Fires however the connection ends, including our own close(). */ onClose?: () => void; } export class RelaySyncTransport implements SyncTransport { #ws?: WebSocket; #ready?: Promise; #closed = false; #seq = 0; #pending = new Map(); #sessions = new Map(); #joinWaiters = new Map(); #inboxHandlers = new Set<(inv: Invitation) => void>(); #readyHandlers = new Set<(token: string, docId: string, inviteeDid: string) => void>(); #inboxBuffer: Invitation[] = []; #subbedInbox = false; constructor(private opts: RelayClientOptions) {} #connect(): Promise { if (this.#ready) return this.#ready; this.#ready = (async () => { const token = await this.opts.token(); // close() may have been called while the token was being minted; a // socket opened now would outlive its owner as a zombie connection. if (this.#closed) throw new Error('transport closed'); const url = `${this.opts.url}?token=${encodeURIComponent(token)}` + (this.opts.handle ? `&handle=${encodeURIComponent(this.opts.handle)}` : ''); const ws = new WebSocket(url); this.#ws = ws; ws.addEventListener('message', (ev) => { if (typeof ev.data === 'string') this.#onFrame(JSON.parse(ev.data)); }); await new Promise((res, rej) => { // A blackholed host can leave a raw WebSocket in CONNECTING for a // minute or more; surface "unreachable" while it still means something. const deadline = setTimeout(() => { rej(new Error('relay connection timed out')); try { ws.close(); } catch { /* already closing */ } }, 15_000); ws.addEventListener('open', () => { clearTimeout(deadline); if (this.#closed) { ws.close(); rej(new Error('transport closed')); return; } res(); }); ws.addEventListener('error', () => { clearTimeout(deadline); rej(new Error('relay connection failed')); }); ws.addEventListener('close', () => { clearTimeout(deadline); for (const p of this.#pending.values()) p.reject(new Error('connection closed')); this.#pending.clear(); for (const w of this.#joinWaiters.values()) w.reject(new Error('connection closed')); this.#joinWaiters.clear(); this.opts.onClose?.(); }); }); })(); return this.#ready; } #onFrame(f: ServerFrame) { switch (f.t) { case 'joined': { const s = this.#sessions.get(f.docId); if (s) s._init(f.me, f.roster); this.#joinWaiters.get(f.docId)?.resolve(f); this.#joinWaiters.delete(f.docId); break; } case 'msg': this.#sessions.get(f.docId)?._message({ seq: f.seq, author: f.author, role: f.role, bytes: fromB64(f.b64) }); break; case 'eph': this.#sessions.get(f.docId)?._ephemeral({ author: f.author, role: f.role, bytes: fromB64(f.b64) }); break; case 'roster': this.#sessions.get(f.docId)?._roster(f.roster); break; case 'ack': case 'invited': case 'accepted': case 'kp': case 'invite-info': if (f.id) { this.#pending.get(f.id)?.resolve(f); this.#pending.delete(f.id); } break; case 'inbox': if (this.#inboxHandlers.size) for (const h of this.#inboxHandlers) h(f.invitation); else this.#inboxBuffer.push(f.invitation); break; case 'invite-ready': for (const h of this.#readyHandlers) h(f.token, f.docId, f.inviteeDid); break; case 'err': { const err = new Error(`${f.code}: ${f.message}`); // Scope the failure to what actually failed: a request by id, a join // by docId. Only an unattributable error takes down every join wait // (better than letting them hang). if (f.id) { this.#pending.get(f.id)?.reject(err); this.#pending.delete(f.id); break; } if (f.docId) { this.#joinWaiters.get(f.docId)?.reject(err); this.#joinWaiters.delete(f.docId); break; } for (const w of this.#joinWaiters.values()) w.reject(err); this.#joinWaiters.clear(); break; } } } /** Send or say why not: a silent no-op on a dead socket turns every caller into a promise that never settles ("Invite" spinning forever). */ _send(frame: ClientFrame) { const ws = this.#ws; if (!ws || ws.readyState !== WebSocket.OPEN) throw new Error('the relay connection is down'); ws.send(JSON.stringify(frame)); } #request(build: (id: string) => ClientFrame): Promise { const id = `r${++this.#seq}`; return new Promise((resolve, reject) => { this.#pending.set(id, { resolve, reject }); try { this._send(build(id)); } catch (err) { this.#pending.delete(id); reject(err); } }); } async _append(docId: string, bytes: Uint8Array): Promise { const f = await this.#request((id) => ({ t: 'append', docId, id, b64: toB64(bytes) })); return f.t === 'ack' ? (f.seq ?? 0) : 0; } _dropSession(docId: string) { this.#sessions.delete(docId); } async join(docId: string, opts?: JoinOptions): Promise { await this.#connect(); const session = new RelaySession(docId, this); this.#sessions.set(docId, session); const manifest = this.opts.manifestFor?.(docId); await new Promise((resolve, reject) => { this.#joinWaiters.set(docId, { resolve: () => resolve(), reject }); this._send({ t: 'join', docId, sinceSeq: opts?.sinceSeq, ...(manifest ? { manifest } : {}) } as ClientFrame); }); return session; } async invite( docId: string, inviteeDid: string, role: string, extra?: { title?: string; payload?: string; pending?: boolean }, ): Promise { await this.#connect(); const f = await this.#request((id) => ({ t: 'invite', docId, inviteeDid, role, id, ...extra })); if (f.t !== 'invited') throw new Error('invite failed'); return f.invitation; } /** Supply the Welcome for an invitation issued before the invitee had keys. */ async completeInvite(token: string, payload: string): Promise { await this.#connect(); await this.#request((id) => ({ t: 'complete-invite', token, payload, id })); } /** Fires when a pending invitation this user issued becomes completable. */ onInviteReady(handler: (token: string, docId: string, inviteeDid: string) => void): () => void { this.#readyHandlers.add(handler); return () => this.#readyHandlers.delete(handler); } async *inbox(): AsyncIterable { await this.#connect(); if (!this.#subbedInbox) { this.#subbedInbox = true; this._send({ t: 'sub-inbox' }); } const queue: Invitation[] = [...this.#inboxBuffer]; this.#inboxBuffer = []; let notify: (() => void) | undefined; const handler = (inv: Invitation) => { queue.push(inv); notify?.(); }; this.#inboxHandlers.add(handler); try { while (true) { if (queue.length) { yield queue.shift()!; continue; } await new Promise((r) => (notify = r)); } } finally { this.#inboxHandlers.delete(handler); } } async accept(token: string): Promise { await this.#connect(); const f = await this.#request((id) => ({ t: 'accept', token, id })); if (f.t !== 'accepted') throw new Error('accept failed'); return { docId: f.result.docId, role: f.result.role, snapshot: f.result.snapshotB64 ? fromB64(f.result.snapshotB64) : null, atSeq: f.result.atSeq, }; } fetchSnapshot(): Promise<{ bytes: Uint8Array; atSeq: number } | null> { // Snapshots arrive via accept() in this transport; direct fetch is unused. return Promise.resolve(null); } /** Publish this DID's MLS KeyPackage to the relay's directory. */ async publishKeyPackage(bytes: Uint8Array): Promise { await this.#connect(); this._send({ t: 'pub-kp', b64: toB64(bytes) }); } /** Fetch a member's published KeyPackage, if any. */ async getKeyPackage(did: string): Promise { await this.#connect(); const f = await this.#request((id) => ({ t: 'get-kp', did, id })); return f.t === 'kp' && f.b64 ? fromB64(f.b64) : null; } /** Fetch an invitation by token (deep-link accepts; invitee-only). */ async getInvite(token: string): Promise { await this.#connect(); const f = await this.#request((id) => ({ t: 'get-invite', token, id })); return f.t === 'invite-info' ? f.invitation : null; } /** Whether the relay connection is currently up. */ get connected(): boolean { return this.#ws?.readyState === WebSocket.OPEN; } close() { this.#closed = true; this.#ws?.close(); } }