From 525231ea1606904cdd18e6164ff7a5b89106c622 Mon Sep 17 00:00:00 2001 From: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 09:47:15 -0600 Subject: [PATCH 01/14] feat(relay): add bounded same-socket presence transport Signed-off-by: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> --- dev/relay-broker-live.test.mjs | 239 ++++++++++- dev/relay-broker.mjs | 71 +++- src/features/relay/broker-live.test.ts | 86 +++- src/features/relay/broker-live.ts | 134 +++++- src/features/relay/live.test.ts | 7 +- src/features/relay/live.ts | 88 ++-- src/features/relay/presence-contract.ts | 51 +++ src/features/relay/presence-live.test.ts | 517 +++++++++++++++++++++++ src/features/relay/presence-live.ts | 308 ++++++++++++++ 9 files changed, 1460 insertions(+), 41 deletions(-) create mode 100644 src/features/relay/presence-contract.ts create mode 100644 src/features/relay/presence-live.test.ts create mode 100644 src/features/relay/presence-live.ts diff --git a/dev/relay-broker-live.test.mjs b/dev/relay-broker-live.test.mjs index 471ecf80..3bbc7750 100644 --- a/dev/relay-broker-live.test.mjs +++ b/dev/relay-broker-live.test.mjs @@ -2,7 +2,7 @@ import { fixtureRelayUrl, fixtureAliases } from "../tests/relay-config.ts"; import { createServer } from "node:http"; import { setTimeout as delay } from "node:timers/promises"; import { test, expect, vi } from "vitest"; -import { getPublicKey } from "nostr-tools"; +import { finalizeEvent, getPublicKey } from "nostr-tools"; import { relayBrokerPlugin } from "./relay-broker.mjs"; import { connectBrokerTransport } from "../src/features/relay/transport.ts"; @@ -16,6 +16,7 @@ async function harness( key[31] = 1; const requests = []; const sockets = []; + const frames = []; let handler; const server = createServer((req, res) => handler?.(req, res)); const plugin = relayBrokerPlugin({ @@ -29,6 +30,7 @@ async function harness( readyState: 1, send(text) { const [kind, id, filter] = JSON.parse(text); + frames.push({ kind, id, filter, socket, at: performance.now() }); if (kind === "AUTH") queueMicrotask(() => this.receive(["OK", id.id, true])); if (kind !== "REQ") return; @@ -66,6 +68,8 @@ async function harness( const controllers = []; return { sockets, + frames, + key, requests, base, async post(channels, origin = base) { @@ -373,3 +377,236 @@ test("priority control cannot allocate interests or bypass owner, origin, commun await h.close(); } }); + +test("actual browser/broker presence controls preserve socket and healthy routes; only a correlated WS receipt resolves publication", async () => { + const h = await harness(); + const nativeFetch = globalThis.fetch; + let traffic; + try { + const fetcher = vi.fn((input, init) => + nativeFetch(input, { + ...init, + headers: { + ...init?.headers, + ...(init?.method === "POST" ? { Origin: h.base } : {}), + }, + }), + ); + vi.stubGlobal("fetch", fetcher); + const callbacks = { + receive: vi.fn(), + state: vi.fn(), + established: vi.fn(), + denied: vi.fn(), + presence: vi.fn(), + presenceState: vi.fn(), + }; + const transport = await connectBrokerTransport(h.base); + traffic = transport.subscribe(callbacks); + traffic.update(["a"]); + await until(() => callbacks.established.mock.calls.length === 3); + const sockets = h.sockets.length; + const streams = fetcher.mock.calls.filter(([url]) => + String(url).endsWith("/stream"), + ).length; + const author = getPublicKey(h.key); + for (let i = 0; i < 1000; i++) traffic.presence.update([author]); + await until( + () => callbacks.presenceState.mock.lastCall?.[0].status === "ready", + ); + const route = h.frames.find( + (f) => f.kind === "REQ" && f.filter.kinds[0] === 20001, + ); + expect(route.filter).toEqual({ + kinds: [20001], + authors: [author], + limit: 0, + }); + const event = finalizeEvent( + { + kind: 20001, + content: "online", + created_at: Math.floor(Date.now() / 1000), + tags: [], + }, + h.key, + ); + await route.socket.receive(["EVENT", route.id, event]); + await until(() => callbacks.presence.mock.calls.length === 1); + expect(callbacks.receive).not.toHaveBeenCalled(); + expect(callbacks.established).toHaveBeenCalledTimes(3); + // Coalesced A -> B -> A must return to Ready even though the server union never changed. + traffic.presence.update(["f".repeat(64)]); + traffic.presence.update([author]); + expect(callbacks.presenceState.mock.lastCall?.[0].status).toBe("pending"); + await until( + () => callbacks.presenceState.mock.lastCall?.[0].status === "ready", + ); + const completed = vi.fn(); + const operation = traffic.presence + .publish("away", new AbortController().signal) + .then(completed); + await until(() => h.frames.some((f) => f.kind === "EVENT")); + const frame = h.frames.find((f) => f.kind === "EVENT"); + expect(frame.id).toMatchObject({ + kind: 20001, + content: "away", + tags: [], + pubkey: author, + }); + await frame.socket.receive(["OK", "wrong-id", true]); + await delay(20); + expect(completed).not.toHaveBeenCalled(); + await frame.socket.receive(["OK", frame.id.id, true]); + await operation; + expect(completed).toHaveBeenCalledOnce(); + traffic.presence.update([]); + await until(() => + h.frames.some((f) => f.kind === "CLOSE" && f.id === route.id), + ); + expect(h.sockets).toHaveLength(sockets); + expect( + fetcher.mock.calls.filter(([url]) => String(url).endsWith("/stream")), + ).toHaveLength(streams); + expect(h.requests.filter((r) => r.filter.kinds[0] !== 20001)).toHaveLength( + 3, + ); + expect( + fetcher.mock.calls.filter(([url]) => + String(url).endsWith("/stream-presence"), + ), + ).toHaveLength(3); + expect( + fetcher.mock.calls.some(([url]) => + /\/(events|publish|sign)$/.test(String(url)), + ), + ).toBe(false); + } finally { + traffic?.dispose(); + vi.unstubAllGlobals(); + await h.close(); + } +}); + +test("presence controls enforce origin, owner, community, shape and body bounds without extra sockets/signing", async () => { + const h = await harness(); + const post = (path, body, origin = h.base) => + fetch(`${h.base}${path}`, { + method: "POST", + headers: { Origin: origin, "Content-Type": "application/json" }, + body: JSON.stringify(body), + }); + try { + const stream = await h.post([]); + const streamId = stream.response.headers.get("x-buzz-live-id"); + const author = getPublicKey(h.key); + await until(() => h.requests.length === 2); + for (const [path, body, code] of [ + ["stream-presence", { streamId, authors: ["bad"] }, 400], + ["stream-presence", { streamId, authors: Array(257).fill(author) }, 400], + ["stream-presence", { streamId, authors: ["x".repeat(18001)] }, 413], + ["stream-presence-publish", { streamId, status: "offline" }, 400], + [ + "stream-presence-publish", + { streamId, status: { status: "online" } }, + 400, + ], + ["stream-presence-publish", { streamId, status: "x".repeat(300) }, 413], + [ + "stream-presence-publish", + { streamId: "f".repeat(32), status: "online" }, + 404, + ], + ]) + expect((await post(`/api/relay/${path}`, body)).status).toBe(code); + for (const path of ["stream-presence", "stream-presence-publish"]) { + const body = { streamId, authors: [author], status: "online" }; + expect( + (await post(`/api/relay/${path}`, body, "https://wrong.invalid")) + .status, + ).toBe(403); + expect((await post(`/api/relay/secondary/${path}`, body)).status).toBe( + 404, + ); + } + expect(h.sockets).toHaveLength(1); + expect(h.frames.some((f) => f.kind === "EVENT")).toBe(false); + expect(h.requests).toHaveLength(2); + stream.abort(); + await until(() => h.sockets[0].readyState === 3); + expect( + ( + await post("/api/relay/stream-presence", { + streamId, + authors: [author], + }) + ).status, + ).toBe(404); + expect( + ( + await post("/api/relay/stream-presence-publish", { + streamId, + status: "online", + }) + ).status, + ).toBe(404); + } finally { + await h.close(); + } +}); + +test("broker publication cancellation frees its owner and quota rejection remains unconfirmed with no retry", async () => { + const h = await harness(); + const nativeFetch = globalThis.fetch; + let traffic; + try { + vi.stubGlobal("fetch", (input, init) => + nativeFetch(input, { + ...init, + headers: { + ...init?.headers, + ...(init?.method === "POST" ? { Origin: h.base } : {}), + }, + }), + ); + const transport = await connectBrokerTransport(h.base); + let ready = 0; + traffic = transport.subscribe({ + receive() {}, + state() {}, + established() { + ready++; + }, + denied() {}, + }); + await until(() => ready === 2); + const controller = new AbortController(); + const operation = traffic.presence.publish("online", controller.signal); + const failure = expect(operation).rejects.toThrow(); + await until(() => h.frames.some((f) => f.kind === "EVENT")); + controller.abort(); + await failure; + await delay(50); // Let HTTP close cancellation reach the server owner. + const second = traffic.presence.publish( + "away", + new AbortController().signal, + ); + const rejected = expect(second).rejects.toThrow("unconfirmed"); + await until(() => h.frames.filter((f) => f.kind === "EVENT").length === 2); + const frame = h.frames.filter((f) => f.kind === "EVENT")[1]; + await frame.socket.receive([ + "OK", + frame.id.id, + false, + "rate-limited: quota exceeded; retry in 0s", + ]); + await rejected; + await delay(1100); + expect(h.frames.filter((f) => f.kind === "EVENT")).toHaveLength(2); + expect(h.sockets).toHaveLength(1); + } finally { + traffic?.dispose(); + vi.unstubAllGlobals(); + await h.close(); + } +}); diff --git a/dev/relay-broker.mjs b/dev/relay-broker.mjs index 9ec90c45..3a68e091 100644 --- a/dev/relay-broker.mjs +++ b/dev/relay-broker.mjs @@ -15,6 +15,10 @@ import { SIDEBAR_UPLOAD_MS, SIDEBAR_UPLOAD_SLOTS, } from "./sidebar-preferences.mjs"; +import { + presenceAuthors, + presenceStatus, +} from "../src/features/relay/presence-contract.ts"; import { createHostAdmission } from "../src/features/relay/host-admission.ts"; import { relayKlipySearchPath } from "../src/features/relay/gifs.ts"; // Dev-only relay broker. Holds the local Buzz identity in this Node process and signs NIP-98 reads @@ -513,22 +517,33 @@ export function relayBrokerPlugin({ live: true, }); if ( - ["/api/relay/stream-retry", "/api/relay/stream-priority"].includes( - route, - ) && + [ + "/api/relay/stream-retry", + "/api/relay/stream-priority", + "/api/relay/stream-presence", + "/api/relay/stream-presence-publish", + ].includes(route) && req.method === "POST" ) { const prioritizing = route === "/api/relay/stream-priority"; + const observing = route === "/api/relay/stream-presence"; + const publishingPresence = + route === "/api/relay/stream-presence-publish"; let raw = ""; for await (const part of req) { raw += part; - if (Buffer.byteLength(raw) > (prioritizing ? 9000 : 256)) + if ( + Buffer.byteLength(raw) > + (observing ? 18000 : prioritizing ? 9000 : 256) + ) return json(res, 413, { error: "Live control too large" }); } - let streamId, priority; + let streamId, priority, authors, status; try { const body = JSON.parse(raw); streamId = body.streamId; + if (observing) authors = presenceAuthors(body.authors); + if (publishingPresence) status = presenceStatus(body.status); if (prioritizing) { liveChannels(body.channels); if (body.channels.length > 64) @@ -548,7 +563,32 @@ export function relayBrokerPlugin({ return json(res, 404, { error: "Live stream no longer available", }); - if (prioritizing) stream.traffic.prioritize(priority); + if (publishingPresence) { + const controller = new AbortController(); + const abort = () => controller.abort(); + res.once("close", abort); + try { + await stream.traffic.presence.publish( + status, + controller.signal, + ); + if (!res.destroyed) return json(res, 200, { accepted: true }); + } catch { + if (!res.destroyed) + return json(res, 503, { + error: "Presence publication unconfirmed", + }); + } finally { + res.off("close", abort); + } + return; + } + if (observing) { + stream.traffic.presence.update(authors); + // Reassert state on the ordered SSE lane even when the union is unchanged + // (e.g. A -> B -> A coalesced in the browser, or retry after a lost response). + stream.presence(); + } else if (prioritizing) stream.traffic.prioritize(priority); else stream.traffic.retry(); return json(res, 200, { accepted: true }); } @@ -559,10 +599,11 @@ export function relayBrokerPlugin({ if (Buffer.byteLength(raw) > 150000) return json(res, 413, { error: "Live interests too large" }); } - let channels, priority; + let channels, priority, authors; try { const body = JSON.parse(raw); channels = liveChannels(body.channels); + authors = presenceAuthors(body.authors ?? []); liveChannels(body.priority ?? []); if (body.priority?.length > 64) throw new Error("Priority capacity reached"); @@ -598,6 +639,7 @@ export function relayBrokerPlugin({ `${kind ? `event: ${kind}\n` : ""}data: ${JSON.stringify(value)}\n\n`, ); }; + let presenceState = { status: "idle", authors: [] }; const traffic = subscribeRelayTraffic( relay.replace(/^http/, "ws"), async (event) => finalizeEvent(event, key), @@ -606,6 +648,13 @@ export function relayBrokerPlugin({ receive: (events) => { for (const event of events) write("", event); }, + presence: (events) => { + for (const event of events) write("presence", event); + }, + presenceState: (state) => { + presenceState = state; + write("presence-state", state); + }, state: (state) => write("state", state), established: (channelId) => write("established", { channelId }), denied: (channelId, reason) => @@ -617,6 +666,7 @@ export function relayBrokerPlugin({ principal.streams++; traffic.prioritize(priority); traffic.update(channels); + traffic.presence.update(authors); const keepAlive = setInterval( () => res.write(": keepalive\n\n"), 15000, @@ -631,7 +681,12 @@ export function relayBrokerPlugin({ streams.delete(streamId); res.destroy(); }; - streams.set(streamId, { relay, traffic, close }); + streams.set(streamId, { + relay, + traffic, + close, + presence: () => write("presence-state", presenceState), + }); res.once("close", close); if (res.destroyed) close(); return; diff --git a/src/features/relay/broker-live.test.ts b/src/features/relay/broker-live.test.ts index 5fd2440f..58031ef6 100644 --- a/src/features/relay/broker-live.test.ts +++ b/src/features/relay/broker-live.test.ts @@ -38,7 +38,7 @@ function fixture() { live: true, }), ); - if (url.endsWith("/stream-retry")) { + if (url.endsWith("/stream-retry") || url.endsWith("/stream-presence")) { const d = deferred(); controls.push(d); signals.push(init.signal as AbortSignal); @@ -81,6 +81,13 @@ function fixture() { snapshots, accept, publish, + presence(index: number, state: unknown) { + required(bodyControllers[index]).enqueue( + new TextEncoder().encode( + `event: presence-state\ndata: ${JSON.stringify(state)}\n\n`, + ), + ); + }, callbacks: { state(s: unknown) { snapshots.push(s); @@ -161,3 +168,80 @@ for (const finish of ["replacement", "dispose"] as const) owner.dispose(); } }); + +it("presence controls coalesce continuous demand with a fixed deadline and stale author state cannot replace current demand", async () => { + vi.useFakeTimers(); + const f = fixture(); + const states = vi.fn(); + const t = await connectBrokerTransport(); + const owner = required(t.subscribe)({ + ...f.callbacks, + presenceState: states, + }); + try { + f.accept(0); + await tick(); + for (let i = 1; i <= 10; i++) { + owner.presence?.update([i.toString(16).padStart(64, "0")]); + await vi.advanceTimersByTimeAsync(50); + } + expect(f.controls).toHaveLength(1); + const latest = "a".padStart(64, "0"); + f.presence(0, { status: "ready", authors: ["1".padStart(64, "0")] }); + await tick(); + expect(states.mock.lastCall?.[0]).toEqual({ + status: "pending", + authors: [latest], + }); + required(f.controls[0]).resolve(new Response(null, { status: 200 })); + await tick(); + await vi.advanceTimersByTimeAsync(600); + expect(f.controls).toHaveLength(2); + f.presence(0, { status: "ready", authors: [latest] }); + await tick(); + expect(states.mock.lastCall?.[0]).toEqual({ + status: "ready", + authors: [latest], + }); + required(f.controls[1]).resolve(new Response(null, { status: 200 })); + await tick(); + } finally { + owner.dispose(); + } +}); + +it("stream replacement clears a queued presence timer without poisoning later controls; late failed controls are fenced", async () => { + vi.useFakeTimers(); + const f = fixture(); + const states = vi.fn(); + const t = await connectBrokerTransport(); + const owner = required(t.subscribe)({ + ...f.callbacks, + presenceState: states, + }); + try { + f.accept(0); + await tick(); + owner.presence?.update(["a".repeat(64)]); + owner.update(["a"]); // replacement while the first control is still queued + f.accept(1); + await tick(); + owner.presence?.update(["b".repeat(64)]); + await vi.advanceTimersByTimeAsync(100); + expect(f.controls).toHaveLength(1); + owner.update(["b"]); + f.accept(2); + await tick(); + const before = states.mock.calls.length; + required(f.controls[0]).resolve(new Response(null, { status: 503 })); + await tick(); + expect(states).toHaveBeenCalledTimes(before); + owner.presence?.update(["c".repeat(64)]); + await vi.advanceTimersByTimeAsync(1000); + expect(f.controls).toHaveLength(2); + required(f.controls[1]).resolve(new Response(null, { status: 200 })); + await tick(); + } finally { + owner.dispose(); + } +}); diff --git a/src/features/relay/broker-live.ts b/src/features/relay/broker-live.ts index cac64c77..037eeef3 100644 --- a/src/features/relay/broker-live.ts +++ b/src/features/relay/broker-live.ts @@ -1,4 +1,10 @@ import { eventDto } from "./events"; +import { + presenceAuthors, + presenceStatus, + presenceState, + type PresenceStatus, +} from "./presence-contract"; import { liveChannels, type LiveCallbacks, @@ -18,6 +24,11 @@ export function subscribeBrokerTraffic( attempts = 0; let channels: string[] = []; let priority: string[] = []; + let authors: string[] = []; + let presencePending = false; + let presenceTimer: ReturnType | undefined; + let presenceDue = 0; + let publishing = false; let priorityPending = false; let controller: AbortController | undefined; let streamId: string | undefined; @@ -38,6 +49,9 @@ export function subscribeBrokerTraffic( streamId = undefined; controlPending = false; priorityPending = false; + presencePending = false; + clearTimeout(presenceTimer); + presenceTimer = undefined; receiving = true; controller?.abort(); clearTimeout(retryTimer); @@ -56,13 +70,14 @@ export function subscribeBrokerTraffic( }; pulse(); const startingPriority = JSON.stringify(priority); + const startingPresence = JSON.stringify(authors); void (async () => { try { const response = await fetch(`${endpoint}/stream`, { method: "POST", credentials: "same-origin", headers: { "Content-Type": "application/json" }, - body: JSON.stringify({ channels, priority }), + body: JSON.stringify({ channels, priority, authors }), signal: owned.signal, }); if (!valid()) return; @@ -85,6 +100,7 @@ export function subscribeBrokerTraffic( throw new Error("Invalid live broker control identity"); streamId = identity ?? undefined; if (startingPriority !== JSON.stringify(priority)) sendPriority(); + if (startingPresence !== JSON.stringify(authors)) schedulePresence(); const reader = response.body.getReader(); const decoder = new TextDecoder(); let buffer = ""; @@ -113,8 +129,19 @@ export function subscribeBrokerTraffic( if (!lines.length) continue; // Keepalives carry no data. const data: unknown = JSON.parse(lines.join("\n")); if (!valid()) return; - if (kind === "message") callbacks.receive([eventDto(data)]); - else if (kind === "state") { + if (kind === "message") { + const event = eventDto(data); + if (event.kind !== 20001) callbacks.receive([event]); + } else if (kind === "presence") { + const event = eventDto(data); + if (event.kind === 20001 && authors.includes(event.pubkey)) + callbacks.presence?.([event]); + } else if (kind === "presence-state") { + const state = presenceState(data); + // SSE already in transit can describe an older control's author set. + if (JSON.stringify(state.authors) === JSON.stringify(authors)) + callbacks.presenceState?.(state); + } else if (kind === "state") { const snapshot = liveSnapshot(data); publish(snapshot); } else if (kind === "established") { @@ -154,6 +181,7 @@ export function subscribeBrokerTraffic( retryTimer = setTimeout(start, 500 * 2 ** attempts++); } finally { if (current === generation) { + owned.abort(); streamId = undefined; receiving = false; clearTimeout(heartbeat); @@ -190,8 +218,106 @@ export function subscribeBrokerTraffic( if (sent !== JSON.stringify(priority)) sendPriority(); }); } + function schedulePresence() { + if (closed || presencePending || !streamId || presenceTimer) return; + // First dirty update fixes the deadline; scrolling only replaces authors. + presenceTimer = setTimeout( + () => { + presenceTimer = undefined; + sendPresence(); + }, + Math.max(100, presenceDue - performance.now()), + ); + } + function sendPresence() { + if (closed || !streamId || presencePending) return; + const current = generation; + const sent = JSON.stringify(authors); + presencePending = true; + presenceDue = performance.now() + 1000; + void fetch(`${endpoint}/stream-presence`, { + method: "POST", + credentials: "same-origin", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ streamId, authors }), + signal: AbortSignal.any([ + controller?.signal ?? new AbortController().signal, + AbortSignal.timeout(5000), + ]), + }) + .then((response) => { + if (!closed && current === generation && !response.ok) + throw new Error(`Presence control failed (${response.status})`); + }) + .catch((error) => { + if (!closed && current === generation) + callbacks.presenceState?.({ + status: authors.length ? "error" : "idle", + authors: [...authors], + error: String(error), + }); + }) + .finally(() => { + if (current !== generation) return; + presencePending = false; + if (sent !== JSON.stringify(authors)) schedulePresence(); + }); + } start(); return { + presence: { + update(input) { + const next = presenceAuthors(input); + if (closed || JSON.stringify(next) === JSON.stringify(authors)) return; + authors = next; + callbacks.presenceState?.({ + status: authors.length ? "pending" : "idle", + authors: [...authors], + }); + schedulePresence(); + }, + async publish(status: PresenceStatus, signal: AbortSignal) { + presenceStatus(status); + signal.throwIfAborted(); + if (closed || !streamId) throw new Error("Presence stream unavailable"); + if (publishing) + throw new Error("Presence publication already in flight"); + const current = generation; + publishing = true; + try { + const response = await fetch(`${endpoint}/stream-presence-publish`, { + method: "POST", + credentials: "same-origin", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ streamId, status }), + signal: AbortSignal.any([ + signal, + controller?.signal ?? new AbortController().signal, + AbortSignal.timeout(11000), + ]), + }); + signal.throwIfAborted(); + if (closed || current !== generation) + throw new Error("Presence stream replaced; outcome unknown"); + if (!response.ok) + throw new Error( + `Presence publication unconfirmed (${response.status})`, + ); + const receipt: unknown = await response.json(); + signal.throwIfAborted(); + if ( + closed || + current !== generation || + !receipt || + typeof receipt !== "object" || + (receipt as { accepted?: unknown }).accepted !== true + ) + throw new Error("Invalid presence receipt or replaced stream"); + } finally { + publishing = false; + } + }, + }, prioritize(input) { liveChannels(input); const next = [...new Set(input)].slice(0, 64); @@ -214,6 +340,7 @@ export function subscribeBrokerTraffic( start(); return; } + schedulePresence(); // Retry a failed author control too, not just server-side routes. const current = generation; controlPending = true; void fetch(`${endpoint}/stream-retry`, { @@ -250,6 +377,7 @@ export function subscribeBrokerTraffic( controller?.abort(); clearTimeout(retryTimer); clearTimeout(heartbeat); + clearTimeout(presenceTimer); }, }; } diff --git a/src/features/relay/live.test.ts b/src/features/relay/live.test.ts index ff1520df..1ada3376 100644 --- a/src/features/relay/live.test.ts +++ b/src/features/relay/live.test.ts @@ -1,5 +1,6 @@ import { assert, afterEach, expect, it, vi } from "vitest"; import { + LIVE_CHANNEL_CAPACITY, createLiveAdmission, liveChannels, subscribeRelayTraffic, @@ -185,19 +186,19 @@ it("bounds interests and exposes every omitted ID without exceeding 1024 subscri await h.first.auth(); await vi.advanceTimersByTimeAsync(750); let index = 0; - while (index < 1024) { + while (index < LIVE_CHANNEL_CAPACITY + 2) { if (index >= h.first.requests().length) await vi.advanceTimersByTimeAsync(250); const request = h.first.requests()[index++]; assert.exists(request); await h.first.receive(["EOSE", request[1]]); } - expect(h.first.requests()).toHaveLength(1024); + expect(h.first.requests()).toHaveLength(LIVE_CHANNEL_CAPACITY + 2); expect( h.callbacks.state.mock.lastCall?.[0].routes .filter((r) => r.status === "limited") .map((r) => r.channelId), - ).toEqual(ids.slice(1022)); + ).toEqual(ids.slice(LIVE_CHANNEL_CAPACITY)); expect(() => liveChannels([...ids, "excess"])).toThrow(); for (const invalid of [[""], ["a b"], ["x".repeat(129)], [9], {}]) expect(() => liveChannels(invalid)).toThrow(); diff --git a/src/features/relay/live.ts b/src/features/relay/live.ts index e4a6c2ab..151f9c55 100644 --- a/src/features/relay/live.ts +++ b/src/features/relay/live.ts @@ -2,7 +2,10 @@ import type { EventTemplate, VerifiedEvent } from "nostr-tools"; import { eventDto } from "./events.ts"; import { EMOJI_SET } from "./emoji.ts"; -export const LIVE_CHANNEL_CAPACITY = 1022; // Reserve two of the relay's 1024 slots. +import { createPresenceLive } from "./presence-live.ts"; +import type { PresenceCapability, PresenceState } from "./presence-contract.ts"; + +export const LIVE_CHANNEL_CAPACITY = 1020; // Two globals plus confirmed/candidate presence slots. export const LIVE_REPLAY_LIMIT = 500; const SETUP_CONCURRENCY = 4; const REQUEST_INTERVAL_MS = 250; // 4 starts/s leaves room below the reference 10/s quota. @@ -12,9 +15,19 @@ const MAX_QUOTA_RETRIES = 3; export function createLiveAdmission() { let next = 0; let cooldown = 0; + let nextPresence = 0, + nextPublish = 0; return { delay: () => Math.max(0, next - performance.now(), cooldown - performance.now()), + presenceDelay: () => Math.max(0, nextPresence - performance.now()), + publishDelay: () => Math.max(0, nextPublish - performance.now()), + takePresence() { + nextPresence = performance.now() + 1000; + }, + takePublish() { + nextPublish = performance.now() + 1000; + }, take() { next = performance.now() + REQUEST_INTERVAL_MS; }, @@ -40,11 +53,14 @@ export type LiveSnapshot = Readonly<{ }>; export type LiveCallbacks = { receive(events: readonly VerifiedEvent[]): void; + presence?(events: readonly VerifiedEvent[]): void; + presenceState?(state: PresenceState): void; state(snapshot: LiveSnapshot): void; established(channelId?: string): void; denied(channelId: string, reason: string): void; }; export type LiveSubscription = { + presence?: PresenceCapability; update(channels: readonly string[]): void; /** Host demand only: reorder existing pending routes, never grant new interests. */ prioritize?(channels: readonly string[]): void; @@ -127,6 +143,34 @@ export function subscribeRelayTraffic( const send = (value: unknown) => { if (!closed && socket?.readyState === 1) socket.send(JSON.stringify(value)); }; + function cooldown(reason: string): boolean | undefined { + if (!reason.startsWith("rate-limited:")) return undefined; + const hint = /^rate-limited: quota exceeded; retry in (\d+)s$/.exec(reason); + const seconds = hint ? Number(hint[1]) : 5; + const supported = Number.isSafeInteger(seconds) && seconds <= 60; + admission.pause( + Number.isSafeInteger(seconds) && seconds <= 86400 ? seconds : 86400, + ); + if (!supported) { + for (const queued of routes.values()) + if (queued.status === "pending" && !queued.wire) { + queued.status = "error"; + queued.error = "Unsupported live cooldown; automatic setup stopped"; + } + notify(); + } + return supported; + } + const presence = createPresenceLive({ + sign, + viewer, + callbacks, + admission, + send, + connected: () => !closed && authenticated, + wake: pump, + cooldown, + }); function remove(route: Route) { clearTimeout(route.deadline); if (route.wire) { @@ -191,33 +235,15 @@ export function subscribeRelayTraffic( delete route.wire; route.status = "error"; route.error = reason; - if (reason.startsWith("rate-limited:")) { - const hint = /^rate-limited: quota exceeded; retry in (\d+)s$/.exec( - reason, - ); - const seconds = hint ? Number(hint[1]) : 5; - if (!Number.isSafeInteger(seconds) || seconds > 60) { + if (cooldown(reason) === true) { + if (++route.quotaRetries <= MAX_QUOTA_RETRIES) route.status = "pending"; + else { for (const queued of routes.values()) if (queued.status === "pending" && !queued.wire) { queued.status = "error"; - queued.error = "Unsupported live cooldown; automatic setup stopped"; + queued.error = + "Live request cooldown retries exhausted; retry available"; } - // Conservative shared pause survives replacement; never overflow a timer. - admission.pause( - Number.isSafeInteger(seconds) && seconds <= 86400 ? seconds : 86400, - ); - } else { - admission.pause(seconds); - if (++route.quotaRetries <= MAX_QUOTA_RETRIES) route.status = "pending"; - else { - // Stop the unsent queue too: rejection must never drain it into an exhausted budget. - for (const queued of routes.values()) - if (queued.status === "pending" && !queued.wire) { - queued.status = "error"; - queued.error = - "Live request cooldown retries exhausted; retry available"; - } - } } } notify(); @@ -231,6 +257,13 @@ export function subscribeRelayTraffic( let active = [...routes.values()].filter( (route) => route.wire && route.status === "pending", ).length; + const foreground = [...routes.values()].some( + (route) => + route.status === "pending" && + (!route.channelId || priority.includes(route.channelId)), + ); + presence.dispatch(foreground || active >= SETUP_CONCURRENCY); + if (presence.pending()) active++; const rank = (route: Route) => !route.channelId ? -2 @@ -297,6 +330,7 @@ export function subscribeRelayTraffic( function clearSocket() { generation++; authenticated = false; + presence.reset("Presence socket disconnected; outcome unknown"); clearTimeout(dispatchTimer); clearTimeout(deadline); for (const route of routes.values()) clearTimeout(route.deadline); @@ -404,6 +438,7 @@ export function subscribeRelayTraffic( notify(); return; } + presence.message(data); const route = typeof data[1] === "string" ? wires.get(data[1]) : undefined; if (!authenticated || !route) return; @@ -416,7 +451,7 @@ export function subscribeRelayTraffic( return; } if (route.status === "pending") route.count++; - callbacks.receive([incoming]); + if (incoming.kind !== 20001) callbacks.receive([incoming]); } else if (data[0] === "EOSE" && route.status === "pending") { clearTimeout(route.deadline); route.status = "live"; @@ -440,6 +475,7 @@ export function subscribeRelayTraffic( } connect(); return { + presence: presence.capability, prioritize(input) { liveChannels(input); // Same bounded ID validation, but preserve demand order. priority = [...new Set(input)].slice(0, 64); @@ -456,6 +492,7 @@ export function subscribeRelayTraffic( clearTimeout(retryTimer); attempts = 0; if (authenticated) { + presence.retry(); for (const route of routes.values()) { if (route.status !== "error") continue; route.status = "pending"; @@ -468,6 +505,7 @@ export function subscribeRelayTraffic( }, dispose() { if (closed) return; + presence.dispose(); closed = true; clearTimeout(retryTimer); clearSocket(); diff --git a/src/features/relay/presence-contract.ts b/src/features/relay/presence-contract.ts new file mode 100644 index 00000000..e4b4ebec --- /dev/null +++ b/src/features/relay/presence-contract.ts @@ -0,0 +1,51 @@ +/** Ephemeral kind 20001, never a durable outbox operation. + * Publish bare online/away content and no tags. Reads also recognize offline and + * legacy JSON {"status":"online"|"away"|"offline"}; other values are unknown. + * Live subject = verified author. Only HTTP snapshots use relay-signed p subjects. + */ +export type PresenceStatus = "online" | "away"; +export type PresenceState = Readonly<{ + status: "idle" | "pending" | "ready" | "error"; + /** Confirmed authors when ready; desired authors otherwise. EOSE is not a seed. */ + authors: readonly string[]; + error?: string; +}>; +export type PresenceCapability = { + update(authors: readonly string[]): void; + /** Resolves only for matching WS OK. Failure/abort is not evidence of Offline. + * One in flight; callers coalesce status changes. No transport heartbeat replay. */ + publish(status: PresenceStatus, signal: AbortSignal): Promise; +}; +export const PRESENCE_AUTHOR_CAPACITY = 256; +export function presenceAuthors(input: unknown): string[] { + if ( + !Array.isArray(input) || + input.length > PRESENCE_AUTHOR_CAPACITY || + input.some((id) => typeof id !== "string" || !/^[0-9a-f]{64}$/.test(id)) + ) + throw new Error("Invalid presence authors (maximum 256 full public keys)"); + return [...new Set(input as string[])].sort(); +} +export function presenceStatus(input: unknown): PresenceStatus { + if (input !== "online" && input !== "away") + throw new Error("Invalid presence publication status"); + return input; +} +export function presenceState(input: unknown): PresenceState { + if (!input || typeof input !== "object") + throw new Error("Invalid presence state"); + const value = input as PresenceState; + if ( + !["idle", "pending", "ready", "error"].includes(value.status) || + (value.error !== undefined && typeof value.error !== "string") + ) + throw new Error("Invalid presence state"); + const authors = presenceAuthors(value.authors); + if ((value.status === "idle") !== (authors.length === 0)) + throw new Error("Invalid presence state authors"); + return Object.freeze({ + status: value.status, + authors: Object.freeze(authors), + ...(value.error ? { error: value.error } : {}), + }); +} diff --git a/src/features/relay/presence-live.test.ts b/src/features/relay/presence-live.test.ts new file mode 100644 index 00000000..31a63d50 --- /dev/null +++ b/src/features/relay/presence-live.test.ts @@ -0,0 +1,517 @@ +import { afterEach, assert, expect, it, vi } from "vitest"; +import type { EventTemplate, VerifiedEvent } from "nostr-tools"; +import { + createLiveAdmission, + LIVE_CHANNEL_CAPACITY, + subscribeRelayTraffic, + type LiveCallbacks, +} from "./live"; +import { connectSignedTransport } from "./transport"; +import { presenceAuthors } from "./presence-contract"; +import { keypair, signed } from "./testing"; + +class Socket { + readyState = 1; + sent: unknown[][] = []; + times: number[] = []; + onmessage?: (event: { data: string }) => Promise; + onclose?: () => void; + send(text: string) { + this.sent.push(JSON.parse(text)); + this.times.push(performance.now()); + } + close() { + this.readyState = 3; + this.onclose?.(); + } + async receive(frame: unknown[]) { + await this.onmessage?.({ data: JSON.stringify(frame) }); + } + requests(presence = false) { + return this.sent.filter( + (f) => + f[0] === "REQ" && + ((f[2] as { kinds: number[] }).kinds[0] === 20001) === presence, + ) as [ + string, + string, + { kinds: number[]; authors?: string[]; limit: number; "#h"?: string[] }, + ][]; + } + events() { + return this.sent + .filter((f) => f[0] === "EVENT") + .map((f) => f[1] as VerifiedEvent); + } + async auth() { + await this.receive(["AUTH", "test"]); + const event = this.sent.find((f) => f[0] === "AUTH")?.[1] as VerifiedEvent; + await this.receive(["OK", event.id, true]); + } + async globals() { + await this.auth(); + for (let i = 0; i < 2; i++) { + const request = this.requests()[i]; + assert.exists(request); + await this.receive(["EOSE", request[1]]); + await vi.advanceTimersByTimeAsync(250); + } + } +} +const author = (n: number) => n.toString(16).padStart(64, "0"); +const signal = () => new AbortController().signal; +afterEach(() => { + vi.useRealTimers(); + vi.unstubAllGlobals(); +}); +function setup( + signOverride?: (event: EventTemplate) => Promise, + admission = createLiveAdmission(), +) { + vi.useFakeTimers(); + const key = keypair(); + const sockets: Socket[] = []; + const callbacks = { + receive: vi.fn(), + established: vi.fn(), + denied: vi.fn(), + state: vi.fn(), + presence: vi.fn(), + presenceState: vi.fn(), + } satisfies LiveCallbacks; + const sign = vi.fn( + signOverride ?? (async (event: EventTemplate) => signed(key, event)), + ); + const owner = subscribeRelayTraffic( + "wss://test.invalid", + sign, + key.pubkey, + callbacks, + () => { + const socket = new Socket(); + sockets.push(socket); + return socket as unknown as WebSocket; + }, + admission, + ); + const socket = sockets[0]; + assert.exists(socket); + const presence = owner.presence; + assert.exists(presence); + return { key, sockets, socket, owner, presence, callbacks, sign, admission }; +} + +it("zero demand opens no presence route; production update is bounded, explicit and separate from ordinary establishment/receive", async () => { + const h = setup(); + try { + await h.socket.globals(); + await vi.advanceTimersByTimeAsync(60000); + expect(h.socket.requests(true)).toHaveLength(0); + expect(h.sign).toHaveBeenCalledTimes(1); // AUTH, no implicit heartbeat. + const person = keypair(); + for (let i = 0; i < 1000; i++) h.presence.update([person.pubkey]); + await vi.advanceTimersByTimeAsync(100); + const route = h.socket.requests(true)[0]; + assert.exists(route); + expect(route[2]).toEqual({ + kinds: [20001], + authors: [person.pubkey], + limit: 0, + }); + const event = signed(person, { kind: 20001, content: "online", tags: [] }); + await h.socket.receive(["EVENT", route[1], event]); + await h.socket.receive(["EOSE", route[1]]); + expect(h.callbacks.presence).toHaveBeenCalledWith([event]); + expect(h.callbacks.presenceState).toHaveBeenLastCalledWith({ + status: "ready", + authors: [person.pubkey], + }); + expect(h.callbacks.established).toHaveBeenCalledTimes(2); + expect(h.callbacks.receive).not.toHaveBeenCalled(); + // Misrouted ephemeral events can never enter generic acceptance either. + await h.socket.receive(["EVENT", h.socket.requests()[0]?.[1], event]); + expect(h.callbacks.receive).not.toHaveBeenCalled(); + h.presence.update([]); + await h.socket.receive(["EVENT", route[1], event]); + expect(h.callbacks.presence).toHaveBeenCalledTimes(1); + expect(h.callbacks.presenceState).toHaveBeenLastCalledWith({ + status: "idle", + authors: [], + }); + for (const value of [ + [""], + ["a"], + [author(1).toUpperCase().replace("1", "A")], + Array(257).fill(author(1)), + {}, + ]) + expect(() => presenceAuthors(value)).toThrow(); + } finally { + h.owner.dispose(); + } + expect(vi.getTimerCount()).toBe(0); +}); + +it("continuous scrolling makes progress at <=1 start/sec, retains only confirmed+candidate and fences obsolete EOSE", async () => { + const h = setup(); + try { + await h.socket.globals(); + h.presence.update([author(1)]); + await vi.advanceTimersByTimeAsync(100); + const first = h.socket.requests(true)[0]; + assert.exists(first); + await h.socket.receive(["EOSE", first[1]]); + for (let i = 2; i <= 50; i++) { + h.presence.update([author(i)]); + await vi.advanceTimersByTimeAsync(50); + } + expect(h.socket.requests(true)).toHaveLength(2); // One candidate, not 49 cancellations. + const candidate = h.socket.requests(true)[1]; + assert.exists(candidate); + expect(h.socket.sent).not.toContainEqual(["CLOSE", first[1]]); + await h.socket.receive(["EOSE", candidate[1]]); + expect(h.socket.sent).toContainEqual(["CLOSE", candidate[1]]); + expect(h.socket.sent).not.toContainEqual(["CLOSE", first[1]]); + await vi.advanceTimersByTimeAsync(100); + const latest = h.socket.requests(true)[2]; + assert.exists(latest); + expect(latest[2].authors).toEqual([author(50)]); + await h.socket.receive(["EOSE", latest[1]]); + expect(h.socket.sent).toContainEqual(["CLOSE", first[1]]); + const count = h.callbacks.presenceState.mock.calls.length; + await h.socket.receive(["EOSE", candidate[1]]); + expect(h.callbacks.presenceState).toHaveBeenCalledTimes(count); + let active = 0, + max = 0; + for (const frame of h.socket.sent) { + if (frame[0] === "REQ" && String(frame[1]).startsWith("presence-")) + max = Math.max(max, ++active); + if (frame[0] === "CLOSE" && String(frame[1]).startsWith("presence-")) + active--; + } + const starts = h.socket.sent.flatMap((frame, i) => + frame[0] === "REQ" && String(frame[1]).startsWith("presence-") + ? [h.socket.times[i] as number] + : [], + ); + for (let i = 1; i < starts.length; i++) + expect( + (starts[i] as number) - (starts[i - 1] as number), + ).toBeGreaterThanOrEqual(1000); + expect(max).toBe(2); + expect(active).toBe(1); + } finally { + h.owner.dispose(); + } +}); + +it("failed candidates preserve confirmed routes and retries stay capped even as desired authors change", async () => { + const h = setup(); + try { + await h.socket.globals(); + h.presence.update([author(1)]); + await vi.advanceTimersByTimeAsync(100); + const first = h.socket.requests(true)[0]; + assert.exists(first); + await h.socket.receive(["EOSE", first[1]]); + h.presence.update([author(2)]); + for (let i = 0; i < 4; i++) { + await vi.advanceTimersByTimeAsync(8000); + const route = h.socket.requests(true).at(-1); + assert.exists(route); + expect(route[1]).not.toBe(first[1]); + await h.socket.receive([ + "CLOSED", + route[1], + "temporary: presence unavailable", + ]); + h.presence.update([author(i + 3)]); + } + const count = h.socket.requests(true).length; + await vi.advanceTimersByTimeAsync(120000); + expect(h.socket.requests(true)).toHaveLength(count); + expect(count).toBe(5); + expect(h.socket.sent).not.toContainEqual(["CLOSE", first[1]]); + h.owner.retry(); + await vi.advanceTimersByTimeAsync(100); + expect(h.socket.requests(true)).toHaveLength(count + 1); + } finally { + h.owner.dispose(); + } +}); + +it("channel opening wins over presence and all routes fit the reserved 1024 slots", async () => { + const h = setup(); + try { + const ids = Array.from({ length: 1024 }, (_, i) => `c-${i}`); + h.owner.update(ids); + h.owner.prioritize?.(["c-999"]); + h.presence.update([author(1)]); + await h.socket.auth(); + for (let i = 0; i < 3; i++) { + const request = h.socket.requests()[i]; + assert.exists(request); + await h.socket.receive(["EOSE", request[1]]); + await vi.advanceTimersByTimeAsync(250); + } + expect(h.socket.requests()[2]?.[2]["#h"]).toEqual(["c-999"]); + const p = h.socket.requests(true)[0]; + assert.exists(p); + await h.socket.receive(["EOSE", p[1]]); + let index = 3; + while (index < LIVE_CHANNEL_CAPACITY + 2) { + if (index >= h.socket.requests().length) + await vi.advanceTimersByTimeAsync(250); + const request = h.socket.requests()[index++]; + assert.exists(request); + await h.socket.receive(["EOSE", request[1]]); + } + h.presence.update([author(2)]); + await vi.advanceTimersByTimeAsync(1000); + expect(h.socket.requests().length + h.socket.requests(true).length).toBe( + 1024, + ); + expect( + h.callbacks.state.mock.lastCall?.[0].routes.filter( + (r: { status: string }) => r.status === "limited", + ), + ).toHaveLength(4); + } finally { + h.owner.dispose(); + } +}); + +it("publication signs only after authenticated foreground admission; matching OK alone resolves and no echo/read is needed", async () => { + const h = setup(); + try { + await expect(h.presence.publish("online", signal())).rejects.toThrow( + "authenticated", + ); + await h.socket.auth(); + const operation = h.presence.publish("away", signal()); + const done = vi.fn(); + void operation.then(done); + await vi.advanceTimersByTimeAsync(250); + expect(h.sign).toHaveBeenCalledTimes(1); // globals await EOSE + for (const route of h.socket.requests()) + await h.socket.receive(["EOSE", route[1]]); + await vi.advanceTimersByTimeAsync(250); + const event = h.socket.events()[0]; + assert.exists(event); + expect(event).toMatchObject({ + kind: 20001, + content: "away", + tags: [], + pubkey: h.key.pubkey, + }); + await h.socket.receive(["OK", "wrong", true]); + expect(done).not.toHaveBeenCalled(); + await expect(h.presence.publish("online", signal())).rejects.toThrow( + "in flight", + ); + await h.socket.receive(["OK", event.id, true]); + await operation; + expect(done).toHaveBeenCalledOnce(); + expect(h.socket.events()).toHaveLength(1); + await vi.advanceTimersByTimeAsync(180000); + expect(h.socket.events()).toHaveLength(1); + } finally { + h.owner.dispose(); + } +}); + +it("publication rejection shares cooldown with channel setup and cannot create automatic retries", async () => { + const h = setup(); + try { + await h.socket.globals(); + const promise = h.presence.publish("online", signal()); + const failure = expect(promise).rejects.toThrow("rate-limited"); + await vi.advanceTimersByTimeAsync(0); + const event = h.socket.events()[0]; + assert.exists(event); + await h.socket.receive([ + "OK", + event.id, + false, + "rate-limited: quota exceeded; retry in 2s", + ]); + await failure; + h.owner.update(["a"]); + await vi.advanceTimersByTimeAsync(2999); + expect(h.socket.requests()).toHaveLength(2); + await vi.advanceTimersByTimeAsync(1); + expect(h.socket.requests()).toHaveLength(3); + expect(h.socket.events()).toHaveLength(1); + } finally { + h.owner.dispose(); + } +}); + +it.each(["abort", "disconnect", "dispose", "timeout"])( + "%s fences a pending signer and leaves no heartbeat replay", + async (finish) => { + let release!: (value: VerifiedEvent) => void; + const key = keypair(); + let template!: EventTemplate; + // Match the viewer returned by setup by replacing only the publication signing phase. + const h = setup(); + h.sign.mockImplementation(async (event) => { + if (event.kind !== 20001) return signed(h.key, event); + template = event; + return new Promise((resolve) => { + release = resolve; + }); + }); + try { + await h.socket.globals(); + const controller = new AbortController(); + const promise = h.presence.publish("online", controller.signal); + const failed = expect(promise).rejects.toThrow(); + await vi.advanceTimersByTimeAsync(0); + if (finish === "abort") controller.abort(); + if (finish === "disconnect") h.socket.close(); + if (finish === "dispose") h.owner.dispose(); + if (finish === "timeout") await vi.advanceTimersByTimeAsync(10000); + await failed; + release(signed(key, template)); + await vi.advanceTimersByTimeAsync(0); + expect(h.socket.events()).toHaveLength(0); + } finally { + h.owner.dispose(); + } + expect(vi.getTimerCount()).toBe(0); + }, +); + +it("a signing delay cannot bypass a newly learned shared cooldown", async () => { + const h = setup(); + let release!: (value: VerifiedEvent) => void; + let template!: EventTemplate; + try { + await h.socket.globals(); + h.sign.mockImplementation((event) => { + template = event; + return new Promise((resolve) => { + release = resolve; + }); + }); + const promise = h.presence.publish("online", signal()); + h.admission.pause(1); + release(signed(h.key, template)); + await vi.advanceTimersByTimeAsync(1999); + expect(h.socket.events()).toHaveLength(0); + await vi.advanceTimersByTimeAsync(1); + const event = h.socket.events()[0]; + assert.exists(event); + await h.socket.receive(["OK", event.id, true]); + await promise; + } finally { + h.owner.dispose(); + } +}); + +it("the signed transport exposes the actual same-socket presence capability, never HTTP publication", async () => { + vi.useFakeTimers(); + const key = keypair(); + const sockets: Socket[] = []; + vi.stubGlobal( + "WebSocket", + class extends Socket { + constructor() { + super(); + sockets.push(this); + } + }, + ); + const fetcher = vi.fn(); + vi.stubGlobal("fetch", fetcher); + const transport = await connectSignedTransport( + { + getPublicKey: async () => key.pubkey, + signEvent: async (event) => signed(key, event), + }, + "https://presence-fixture.invalid", + key.pubkey, + ); + const owner = transport.subscribe?.({ + receive() {}, + established() {}, + state() {}, + denied() {}, + }); + assert.exists(owner); + try { + const socket = sockets[0]; + assert.exists(socket); + await socket.globals(); + const operation = owner.presence?.publish("online", signal()); + await vi.advanceTimersByTimeAsync(0); + const event = socket.events()[0]; + assert.exists(event); + await socket.receive(["OK", event.id, true]); + await operation; + expect(sockets).toHaveLength(1); + expect(fetcher).not.toHaveBeenCalled(); + } finally { + owner.dispose(); + } +}); + +it("presence CLOSED honors shared cooldown and setup timeout advances only the latest candidate", async () => { + const h = setup(); + try { + await h.socket.globals(); + h.presence.update([author(1)]); + await vi.advanceTimersByTimeAsync(100); + const first = h.socket.requests(true)[0]; + assert.exists(first); + await h.socket.receive([ + "CLOSED", + first[1], + "rate-limited: quota exceeded; retry in 2s", + ]); + h.owner.update(["a"]); + h.owner.prioritize?.(["a"]); + h.presence.update([author(2)]); + await vi.advanceTimersByTimeAsync(2999); + expect(h.socket.requests()).toHaveLength(2); + expect(h.socket.requests(true)).toHaveLength(1); + await vi.advanceTimersByTimeAsync(1); + const channel = h.socket.requests()[2]; + assert.exists(channel); + await h.socket.receive(["EOSE", channel[1]]); + await vi.advanceTimersByTimeAsync(250); + const candidate = h.socket.requests(true)[1]; + assert.exists(candidate); + expect(candidate[2].authors).toEqual([author(2)]); + h.presence.update([author(3)]); + await vi.advanceTimersByTimeAsync(12000); + const replacement = h.socket.requests(true)[2]; + assert.exists(replacement); + expect(replacement[2].authors).toEqual([author(3)]); + expect(h.socket.sent).toContainEqual(["CLOSE", candidate[1]]); + expect(h.callbacks.established).toHaveBeenCalledTimes(3); + } finally { + h.owner.dispose(); + } +}); + +it("lost WS OK is an unknown outcome, never an automatic resend or acceptance of late receipts", async () => { + const h = setup(); + try { + await h.socket.globals(); + const operation = h.presence.publish("online", signal()); + const failed = expect(operation).rejects.toThrow("unknown"); + await vi.advanceTimersByTimeAsync(0); + const event = h.socket.events()[0]; + assert.exists(event); + await vi.advanceTimersByTimeAsync(10000); + await failed; + await h.socket.receive(["OK", event.id, true]); + await vi.advanceTimersByTimeAsync(120000); + expect(h.socket.events()).toHaveLength(1); + expect(h.callbacks.receive).not.toHaveBeenCalled(); + } finally { + h.owner.dispose(); + } +}); diff --git a/src/features/relay/presence-live.ts b/src/features/relay/presence-live.ts new file mode 100644 index 00000000..158b1c7d --- /dev/null +++ b/src/features/relay/presence-live.ts @@ -0,0 +1,308 @@ +import type { EventTemplate, VerifiedEvent } from "nostr-tools"; +import { eventDto } from "./events.ts"; +import type { LiveAdmission, LiveCallbacks } from "./live.ts"; +import { + presenceAuthors, + presenceStatus, + type PresenceState, + type PresenceStatus, +} from "./presence-contract.ts"; + +type Route = { + wire: string; + authors: string[]; + deadline?: ReturnType; +}; +type Publication = { + status: PresenceStatus; + event?: VerifiedEvent; + signing: boolean; + sent: boolean; + finish(error?: unknown): void; +}; +const same = (a: readonly string[], b: readonly string[]) => + a.length === b.length && a.every((id, i) => id === b[i]); +/** Socket-owned ephemeral lane. No reconnect, snapshots, renewal clock or event history. + * The ordinary route pump calls dispatch only after foreground setup has yielded. */ +export function createPresenceLive(host: { + sign(event: EventTemplate): Promise; + viewer: string; + callbacks: LiveCallbacks; + admission: LiveAdmission; + connected(): boolean; + send(frame: unknown): void; + wake(): void; + cooldown(reason: string): boolean | undefined; +}) { + let closed = false, + serial = 0, + failures = 0; + let desired: string[] = []; + let confirmed: Route | undefined, candidate: Route | undefined; + let timer: ReturnType | undefined; + let due = 0; + let error: string | undefined; + let lastState = ""; + let publication: Publication | undefined; + function notify() { + if (closed) return; + const ready = confirmed && same(confirmed.authors, desired); + const state: PresenceState = Object.freeze({ + status: !desired.length + ? "idle" + : ready + ? "ready" + : error + ? "error" + : "pending", + authors: Object.freeze([...desired]), + ...(!ready && desired.length && error ? { error } : {}), + }); + const key = JSON.stringify(state); + if (key === lastState) return; + lastState = key; + host.callbacks.presenceState?.(state); + } + function remove(route: Route | undefined) { + if (!route) return; + // Fence first: even a reentrant transport cannot deliver after CLOSE. + if (confirmed === route) confirmed = undefined; + if (candidate === route) candidate = undefined; + clearTimeout(route.deadline); + host.send(["CLOSE", route.wire]); + } + function schedule() { + if (!due) due = performance.now() + 100; + notify(); + host.wake(); + } + function failed(route: Route, reason: string) { + remove(route); + error = reason; + const retryable = host.cooldown(reason); + failures = retryable === false ? 4 : failures + 1; + due = performance.now() + 1000 * 2 ** Math.min(failures - 1, 3); + notify(); + host.wake(); + } + function dispatch(blocked: boolean) { + clearTimeout(timer); + if (closed || !host.connected() || blocked) return; + const needsRoute = + desired.length > 0 && + !candidate && + failures <= 3 && + (!confirmed || !same(confirmed.authors, desired)); + const pendingWrite = + publication && !publication.sent && !publication.signing; + if (!needsRoute && !pendingWrite) return; + const routeDelay = needsRoute + ? Math.max(0, due - performance.now(), host.admission.presenceDelay()) + : Infinity; + const writeDelay = pendingWrite ? host.admission.publishDelay() : Infinity; + const delay = Math.max( + host.admission.delay(), + Math.min(routeDelay, writeDelay), + ); + if (delay > 0) { + timer = setTimeout(host.wake, delay); + return; + } + if (needsRoute && routeDelay === 0) { + const route: Route = { + wire: `presence-${++serial}`, + authors: [...desired], + }; + candidate = route; + due = 0; + host.admission.take(); + host.admission.takePresence(); + route.deadline = setTimeout(() => { + if (candidate === route) + failed(route, "Presence setup timed out; retry available"); + }, 10000); + host.send([ + "REQ", + route.wire, + { kinds: [20001], authors: route.authors, limit: 0 }, + ]); + host.wake(); + return; + } + const operation = publication; + if (!operation || operation.sent || operation.signing) return; + if (!operation.event) { + operation.signing = true; + const template = { + kind: 20001, + content: operation.status, + tags: [], + created_at: Math.floor(Date.now() / 1000), + }; + void (async () => { + const raw = await host.sign(template); + if (publication !== operation || closed) return; + const event = eventDto(raw); + if ( + event.pubkey !== host.viewer || + event.kind !== template.kind || + event.content !== template.content || + event.created_at !== template.created_at || + event.tags.length + ) + throw new Error("Presence signer changed the publication"); + operation.event = event; + operation.signing = false; + host.wake(); // Recheck foreground work and cooldown learned during signing. + })().catch((error) => operation.finish(error)); + return; + } + operation.sent = true; + host.admission.take(); + host.admission.takePublish(); + host.send(["EVENT", operation.event]); + } + return { + capability: { + update(input: readonly string[]) { + const next = presenceAuthors(input); + if (closed || same(next, desired)) return; + desired = next; + if (failures <= 3) error = undefined; + if (!desired.length) { + remove(candidate); + remove(confirmed); + failures = 0; + due = 0; + } else if (confirmed && same(confirmed.authors, desired)) + remove(candidate); + schedule(); + }, + publish(status: PresenceStatus, signal: AbortSignal): Promise { + return new Promise((resolve, reject) => { + presenceStatus(status); + signal.throwIfAborted(); + if (closed || !host.connected()) + throw new Error("Presence socket is not authenticated"); + if (publication) + throw new Error("Presence publication already in flight"); + const deadline = setTimeout( + () => + operation.finish( + new Error( + "Presence publication outcome unknown (deadline exceeded)", + ), + ), + 10000, + ); + const abort = () => + operation.finish( + signal.reason ?? + new DOMException("Presence cancelled", "AbortError"), + ); + const operation: Publication = { + status, + signing: false, + sent: false, + finish(error) { + if (publication !== operation) return; + publication = undefined; + clearTimeout(deadline); + signal.removeEventListener("abort", abort); + if (error === undefined) resolve(); + else reject(error); + }, + }; + publication = operation; + signal.addEventListener("abort", abort, { once: true }); + host.wake(); + }); + }, + }, + dispatch, + pending: () => candidate !== undefined, + message(data: unknown[]) { + if (closed || !host.connected()) return; + if ( + data[0] === "OK" && + publication?.sent && + data[1] === publication.event?.id && + typeof data[2] === "boolean" + ) { + const reason = + typeof data[3] === "string" + ? data[3].slice(0, 512) + : "Presence publication rejected"; + if (!data[2]) host.cooldown(reason); + publication.finish(data[2] ? undefined : new Error(reason)); + host.wake(); + return; + } + const route = + candidate?.wire === data[1] + ? candidate + : confirmed?.wire === data[1] + ? confirmed + : undefined; + if (!route) return; + if (data[0] === "EVENT") { + let event: VerifiedEvent; + try { + event = eventDto(data[2]); + } catch { + failed(route, "Invalid presence signature"); + return; + } + if ( + event.kind === 20001 && + route.authors.includes(event.pubkey) && + desired.includes(event.pubkey) + ) + host.callbacks.presence?.([event]); + } else if (data[0] === "EOSE" && candidate === route) { + clearTimeout(route.deadline); + if (!same(route.authors, desired)) { + remove(route); + schedule(); + return; + } + remove(confirmed); + confirmed = route; + candidate = undefined; + failures = 0; + error = undefined; + notify(); + host.wake(); + } else if (data[0] === "CLOSED") { + failed( + route, + typeof data[2] === "string" + ? data[2].slice(0, 512) + : "Presence subscription closed", + ); + } + }, + retry() { + failures = 0; + error = undefined; + schedule(); + }, + reset(reason: string) { + clearTimeout(timer); + // Socket generation owns the wires. No CLOSE on a dead/replaced socket. + clearTimeout(candidate?.deadline); + clearTimeout(confirmed?.deadline); + candidate = confirmed = undefined; + publication?.finish(new Error(reason)); + error = reason; + notify(); + }, + dispose() { + closed = true; + clearTimeout(timer); + remove(candidate); + remove(confirmed); + publication?.finish(new Error("Presence owner disposed")); + }, + }; +} From 0bcab083819cfc8e0abd8558476d2d31b4276cf1 Mon Sep 17 00:00:00 2001 From: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 09:48:20 -0600 Subject: [PATCH 02/14] feat(presence): add session directory, activity renewal and visible UI demand Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- src/bundled/profiles/ProfilePanel.tsx | 6 + src/features/communities/service.ts | 9 +- src/features/messages/ChannelTimeline.tsx | 3 + src/features/messages/MessageRow.tsx | 9 +- src/features/messages/ThreadPanel.tsx | 4 + src/features/presence/Presence.module.css | 31 ++ src/features/presence/activity.test.ts | 39 +++ src/features/presence/activity.ts | 88 +++++ src/features/presence/directory.test.ts | 228 ++++++++++++ src/features/presence/directory.ts | 365 ++++++++++++++++++++ src/features/presence/publisher.test.ts | 187 ++++++++++ src/features/presence/publisher.ts | 151 ++++++++ src/features/presence/react.tsx | 147 ++++++++ src/features/relay/presence-session.test.ts | 141 ++++++++ src/features/relay/service.ts | 3 + src/features/relay/session.ts | 58 +++- tests/browser/presence.spec.mjs | 100 ++++++ tests/fixtures/presence.html | 1 + tests/fixtures/presence.tsx | 143 ++++++++ 19 files changed, 1709 insertions(+), 4 deletions(-) create mode 100644 src/features/presence/Presence.module.css create mode 100644 src/features/presence/activity.test.ts create mode 100644 src/features/presence/activity.ts create mode 100644 src/features/presence/directory.test.ts create mode 100644 src/features/presence/directory.ts create mode 100644 src/features/presence/publisher.test.ts create mode 100644 src/features/presence/publisher.ts create mode 100644 src/features/presence/react.tsx create mode 100644 src/features/relay/presence-session.test.ts create mode 100644 tests/browser/presence.spec.mjs create mode 100644 tests/fixtures/presence.html create mode 100644 tests/fixtures/presence.tsx diff --git a/src/bundled/profiles/ProfilePanel.tsx b/src/bundled/profiles/ProfilePanel.tsx index 584bdd31..dd39b238 100644 --- a/src/bundled/profiles/ProfilePanel.tsx +++ b/src/bundled/profiles/ProfilePanel.tsx @@ -1,3 +1,7 @@ +import { + PresenceIndicator, + usePresenceDemand, +} from "../../features/presence/react"; import { useEffect, useMemo, @@ -40,6 +44,7 @@ function ProfileDetails({ session: RelaySession; pubkey: string; }) { + usePresenceDemand(session.presence, pubkey); const selection = useMemo( () => selectProfiles(session.profiles, [pubkey]), [session.profiles, pubkey], @@ -95,6 +100,7 @@ function ProfileDetails({ size="large" />

{name}

+ {profile?.about &&

{profile.about}

}
diff --git a/src/features/communities/service.ts b/src/features/communities/service.ts index ee7951ca..f7b8acd9 100644 --- a/src/features/communities/service.ts +++ b/src/features/communities/service.ts @@ -1,4 +1,5 @@ // FOUNDATION: Client identity and membership selection outlive community query sessions. +import { createPresenceActivity } from "../presence/activity"; import { Context } from "@deepseek-ai/cordis"; import { provideRelay, type RelayData } from "../relay/service"; import { connectBrokerTransport } from "../relay/transport"; @@ -22,6 +23,7 @@ const empty = (): Saved => ({ selected: null, }); export function createCommunities(ctx: Context, live: boolean) { + const activity = live ? createPresenceActivity() : undefined; let state: ClientSnapshot = { ...empty(), status: live ? "loading" : "unavailable", @@ -72,8 +74,10 @@ export function createCommunities(ctx: Context, live: boolean) { const acquire = (id: string) => { let session = sessions.get(id); if (!session) { - session = provideRelay(newScope(), (signal) => - connectBrokerTransport("", signal, id), + session = provideRelay( + newScope(), + (signal) => connectBrokerTransport("", signal, id), + activity?.activity, ); sessions.set(id, session); session.subscribe(() => { @@ -185,6 +189,7 @@ export function createCommunities(ctx: Context, live: boolean) { ctx.effect(() => () => { disposed = true; controller.abort(); + activity?.dispose(); listeners.clear(); relayListeners.clear(); return Promise.all(scopes.map((scope) => scope.fiber.dispose())); diff --git a/src/features/messages/ChannelTimeline.tsx b/src/features/messages/ChannelTimeline.tsx index 5f6497fe..bc873219 100644 --- a/src/features/messages/ChannelTimeline.tsx +++ b/src/features/messages/ChannelTimeline.tsx @@ -1,4 +1,5 @@ // biome-ignore-all lint/a11y/noNoninteractiveTabindex: The history region must support keyboard scrolling. +import { usePresenceSurface } from "../presence/react"; import type { ConversationExtensions } from "../conversation/contracts"; import type { RelaySession } from "../relay/session"; import { useCallback, useLayoutEffect, useMemo, useRef, useState } from "react"; @@ -105,6 +106,7 @@ function Timeline({ [rows, profiles], ); const scroller = useRef(null); + usePresenceSurface(queries.presence, scroller); const handle = useRef(null); const [size, setSize] = useState({ width: 0, height: 0 }); const width = size.width; @@ -358,6 +360,7 @@ function Timeline({ key={row.id} row={row} unread={queries.unread} + presence={queries.presence} extensions={extensions} profile={profiles.get(row.authorId)} participantProfiles={profiles} diff --git a/src/features/messages/MessageRow.tsx b/src/features/messages/MessageRow.tsx index 7e5e2e53..5fd0e531 100644 --- a/src/features/messages/MessageRow.tsx +++ b/src/features/messages/MessageRow.tsx @@ -1,3 +1,5 @@ +import { PresenceIndicator } from "../presence/react"; +import type { PresenceQueries } from "../presence/directory"; import { memo, useCallback, useSyncExternalStore } from "react"; import type { UnreadCapability } from "../relay/unread"; import { profileTarget } from "../profiles/target"; @@ -12,6 +14,7 @@ import { usesLargeEmojiPresentation } from "./emoji-size"; export type MessageRowProps = { row: ChannelMessage; + presence?: PresenceQueries | undefined; unread?: UnreadCapability | undefined; extensions?: ConversationExtensions | undefined; profile: Profile | undefined; @@ -26,6 +29,7 @@ export type MessageRowProps = { export const MessageRow = memo(function MessageRow({ row, + presence, unread, extensions, profile, @@ -57,7 +61,7 @@ export const MessageRow = memo(function MessageRow({ const AvatarTag = clickable ? "button" : "div"; const emojiOnly = usesLargeEmojiPresentation(row.content, row.emoji); return ( -
+
{day && (
@@ -92,6 +96,9 @@ export const MessageRow = memo(function MessageRow({
{name} + {presence && ( + + )}