From f9fd81cd800557dd733de700fcfc3cc3f8507a0d Mon Sep 17 00:00:00 2001 From: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> Date: Tue, 15 Sep 2026 11:05:53 +0200 Subject: [PATCH 1/6] Add bounded snapshot presence using existing session and socket owners Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- AGENTS.md | 21 ++ dev/relay-broker-api.test.mjs | 46 +++ dev/relay-broker.mjs | 173 ++++++---- docs/presence.md | 52 +++ src/bundled/profiles/ProfilePanel.tsx | 2 + src/features/communities/service.ts | 9 +- src/features/messages/MessageRow.tsx | 7 + src/features/presence/activity.test.ts | 76 +++++ src/features/presence/activity.ts | 80 +++++ src/features/presence/presence.test.ts | 326 ++++++++++++++++++ src/features/presence/presence.ts | 290 ++++++++++++++++ src/features/presence/react.tsx | 42 +++ src/features/relay/broker-live.ts | 18 + src/features/relay/host-admission.test.ts | 31 +- src/features/relay/host-admission.ts | 8 +- src/features/relay/http-admission.ts | 57 +++- src/features/relay/live.test.ts | 87 +++++ src/features/relay/live.ts | 114 +++++++ src/features/relay/service.ts | 3 + src/features/relay/session.ts | 15 + src/features/relay/transport.test.ts | 75 ++++ src/features/relay/transport.ts | 88 +++++ tests/browser/channel-opening.spec.mjs | 83 +++++ tests/browser/fixture.mjs | 52 +++ tests/browser/policy-relay.mjs | 62 ++++ tests/browser/presence.spec.mjs | 396 ++++++++++++++++++++++ 26 files changed, 2147 insertions(+), 66 deletions(-) create mode 100644 docs/presence.md create mode 100644 src/features/presence/activity.test.ts create mode 100644 src/features/presence/activity.ts create mode 100644 src/features/presence/presence.test.ts create mode 100644 src/features/presence/presence.ts create mode 100644 src/features/presence/react.tsx create mode 100644 tests/browser/presence.spec.mjs diff --git a/AGENTS.md b/AGENTS.md index a760a3f5..af22d8c3 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -28,6 +28,27 @@ complete alternative. Require a current caller or explicit approval for adapter parity. Review necessity separately from correctness; passing tests do not justify scope growth. Split at real ownership boundaries, not by deleting safety coverage. +For non-trivial work, make that standard operational: + +- Before coding, publish a short scope checkpoint: required behavior, non-goals, + existing owners/platform support to reuse, expected files, and a rough production + diff budget (separate from tests/docs). A small fix needs only a sentence, not a + design ceremony. +- Use one implementation owner per end-to-end change. Reviewers challenge necessity + as well as correctness; delegate bounded evidence/review, not competing rewrites. + Review the first working slice before expanding the design, without blocking + ordinary human UI feedback on a full validation cycle. +- Justify each new abstraction, lifecycle owner, timer, retry policy, or shared + contract expansion against a current requirement. If the implementation materially + exceeds the checkpoint, stop adding machinery and show the smallest alternative + and any behavior tradeoff before continuing. Do not silently weaken agreed behavior. +- Assess the combined feature diff, including stacked PRs. Passing tests, splitting + PRs, or already-invested work do not establish proportionality. Preserve required + regression coverage; do not game the budget by deleting tests or compressing code. + Keep speculative hardening and unrelated failures outside the task. +- Close with one verified end-to-end result and explicit remaining gaps, not a chain + of green intermediate repairs presented as completion. + Keep files cohesive and group modules and tests by owner. Treat size as a review signal, not a quota. Extract stable boundaries only when they simplify the requested change. diff --git a/dev/relay-broker-api.test.mjs b/dev/relay-broker-api.test.mjs index c6818119..61d13f06 100644 --- a/dev/relay-broker-api.test.mjs +++ b/dev/relay-broker-api.test.mjs @@ -494,3 +494,49 @@ test("both real sign and publish routes admit direct replies but reject arbitrar await h.close(); } }); + +test("held optional snapshot body leaves ordinary broker capacity free and start credit untouched", async () => { + let release; + const h = await harness((call) => { + if (call.body?.[0]?.kinds?.includes(20001)) + return new Response( + new ReadableStream({ + start(controller) { + release = () => { + controller.enqueue(new TextEncoder().encode("[]")); + controller.close(); + }; + }, + }), + ); + return Response.json([]); + }); + const presence = [{ kinds: [20001], authors: [h.event.pubkey], limit: 1 }]; + try { + const snapshot = h.post("presence-snapshot", presence); + await vi.waitFor(() => expect(release).toBeTypeOf("function")); + const duplicate = await h.post("presence-snapshot", presence); + expect(duplicate.status).toBe(204); + expect((await h.post("query", filters)).status).toBe(200); + expect(h.calls).toHaveLength(2); + expect(h.calls[1].at - h.calls[0].at).toBeLessThan(400); + release(); + release = undefined; + expect(await (await snapshot).json()).toEqual([]); + expect( + (await h.post("presence-snapshot", [{ ...presence[0], authors: [] }])) + .status, + ).toBe(400); + expect( + ( + await h.post("presence-snapshot", [ + { ...presence[0], authors: Array(257).fill(h.event.pubkey) }, + ]) + ).status, + ).toBe(400); + expect(h.calls).toHaveLength(2); + } finally { + release?.(); + await h.close(); + } +}); diff --git a/dev/relay-broker.mjs b/dev/relay-broker.mjs index 4df3a188..11971022 100644 --- a/dev/relay-broker.mjs +++ b/dev/relay-broker.mjs @@ -39,6 +39,8 @@ import { ApiPaused, ApiCapacity, apiFailure, + presenceFilter, + presenceText, } from "../src/features/relay/http-admission.ts"; import { execFileSync } from "node:child_process"; import { createHash, randomBytes } from "node:crypto"; @@ -331,6 +333,7 @@ export function relayBrokerPlugin({ }; const stats = { queries: 0, errors: 0, media: 0, connects: 0 }; let inflight = 0; + let presenceFlight = false; let sidebarUploads = 0; let libraryRead; const streams = new Map(); @@ -528,6 +531,7 @@ export function relayBrokerPlugin({ readState: true, agentLibrary: true, live: true, + presence: true, agentActivity: true, }); if ( @@ -535,9 +539,11 @@ export function relayBrokerPlugin({ "/api/relay/stream-retry", "/api/relay/stream-priority", "/api/relay/stream-observer", + "/api/relay/stream-presence", ].includes(route) && req.method === "POST" ) { + const publishingPresence = route === "/api/relay/stream-presence"; const prioritizing = route === "/api/relay/stream-priority"; const observing = route === "/api/relay/stream-observer"; let raw = ""; @@ -546,10 +552,15 @@ export function relayBrokerPlugin({ if (Buffer.byteLength(raw) > (prioritizing ? 9000 : 256)) return json(res, 413, { error: "Live control too large" }); } - let streamId, priority, observer; + let streamId, priority, observer, status; try { const body = JSON.parse(raw); streamId = body.streamId; + if (publishingPresence) { + status = body.status; + if (status !== "online" && status !== "away") + throw new Error("Invalid presence"); + } if (observing) observer = observerGeneration(body.observer); if (prioritizing) { liveChannels(body.channels); @@ -566,10 +577,27 @@ export function relayBrokerPlugin({ ) return json(res, 400, { error: "Invalid live control" }); const stream = streams.get(streamId); + if (publishingPresence && (!stream || stream.relay !== relay)) + return json(res, 200, { accepted: false }); if (!stream || stream.relay !== relay) return json(res, 404, { error: "Live stream no longer available", }); + if (publishingPresence) { + const cancel = new AbortController(); + const abort = () => cancel.abort(); + res.once("close", abort); + try { + const accepted = await stream.traffic.publishPresence( + status, + cancel.signal, + ); + if (!res.destroyed) return json(res, 200, { accepted }); + } finally { + res.off("close", abort); + } + return; + } if (prioritizing) stream.traffic.prioritize(priority); else if (observing) stream.traffic.observe(observer); else stream.traffic.retry(); @@ -726,6 +754,7 @@ export function relayBrokerPlugin({ if ( ![ "/api/relay/query", + "/api/relay/presence-snapshot", "/api/relay/sign", "/api/relay/publish", "/api/relay/read-state-sign", @@ -738,10 +767,11 @@ export function relayBrokerPlugin({ req.method !== "POST" ) return json(res, 404, { error: "Unknown broker route" }); + const presence = route === "/api/relay/presence-snapshot"; let raw = ""; for await (const part of req) { raw += part; - if (raw.length > 65536) + if (Buffer.byteLength(raw) > (presence ? 20 * 1024 : 65536)) return json(res, 413, { error: "Filter body too large" }); } let filters; @@ -750,6 +780,8 @@ export function relayBrokerPlugin({ } catch { return json(res, 400, { error: "Filter body is not JSON" }); } + if (presence && !presenceFilter(filters)) + return json(res, 400, { error: "Invalid presence filter" }); const profile = route === "/api/relay/profile"; const claim = route === "/api/relay/claim"; const policy = route === "/api/relay/accept-policy"; @@ -891,14 +923,24 @@ export function relayBrokerPlugin({ : policy ? "/api/invites/accept-policy" : "/query"; - if (inflight >= MAX_INFLIGHT) + const lane = admissions(relay, viewer).api; + let releasePresence; + if (presence) { + releasePresence = + !presenceFlight && !inflight ? lane.tryPresence() : undefined; + if (!releasePresence) { + res.writeHead(204); + return res.end(); + } + presenceFlight = true; + } + if (!presence && inflight >= MAX_INFLIGHT) return json(res, 429, { error: "Query concurrency limit", sent: false, }); - inflight++; + if (!presence) inflight++; try { - const lane = admissions(relay, viewer).api; const body = JSON.stringify(filters); // A browser that gave up (the client's ten-second deadline) must also release // this upstream request, or hung requests exhaust the inflight budget. @@ -911,65 +953,69 @@ export function relayBrokerPlugin({ try { const requestSignal = AbortSignal.any([ cancel.signal, - AbortSignal.timeout(UPSTREAM_TIMEOUT_MS), + AbortSignal.timeout(presence ? 10000 : UPSTREAM_TIMEOUT_MS), ]); - response = await admittedApiRequest( - lane, - () => { - // Auth freshness and network timings begin at dispatch, not queue entry. - requestSignal.throwIfAborted(); - timings.push( - `admission;dur=${(performance.now() - admissionStart).toFixed(2)}`, - ); - const authStart = performance.now(); - const auth = finalizeEvent( - { - kind: 27235, - created_at: Math.floor(Date.now() / 1000), - content: "", - tags: [ - ["u", `${relay}${upstreamPath}`], - ["method", "POST"], - [ - "payload", - createHash("sha256").update(body).digest("hex"), - ], - ["nonce", randomBytes(16).toString("hex")], + const request = () => { + // Auth freshness and network timings begin at dispatch, not queue entry. + requestSignal.throwIfAborted(); + timings.push( + `admission;dur=${(performance.now() - admissionStart).toFixed(2)}`, + ); + const authStart = performance.now(); + const auth = finalizeEvent( + { + kind: 27235, + created_at: Math.floor(Date.now() / 1000), + content: "", + tags: [ + ["u", `${relay}${upstreamPath}`], + ["method", "POST"], + [ + "payload", + createHash("sha256").update(body).digest("hex"), ], - }, - key, - ); + ["nonce", randomBytes(16).toString("hex")], + ], + }, + key, + ); + timings.push( + `auth;dur=${(performance.now() - authStart).toFixed(2)}`, + ); + connectsBefore = upstream.connects(); + upstreamStart = performance.now(); + return fetchUpstream(`${relay}${upstreamPath}`, { + method: "POST", + headers: { + "Content-Type": "application/json", + Authorization: + "Nostr " + + Buffer.from(JSON.stringify(auth)).toString("base64"), + }, + body, + redirect: "error", + signal: requestSignal, + }).then((response) => { timings.push( - `auth;dur=${(performance.now() - authStart).toFixed(2)}`, + `ttfb;dur=${(performance.now() - upstreamStart).toFixed(2)}`, ); - connectsBefore = upstream.connects(); - upstreamStart = performance.now(); - return fetchUpstream(`${relay}${upstreamPath}`, { - method: "POST", - headers: { - "Content-Type": "application/json", - Authorization: - "Nostr " + - Buffer.from(JSON.stringify(auth)).toString("base64"), - }, - body, - redirect: "error", - signal: requestSignal, - }).then((response) => { - timings.push( - `ttfb;dur=${(performance.now() - upstreamStart).toFixed(2)}`, - ); - return response; - }); - }, - requestSignal, - route === "/api/relay/query" && - req.headers["x-buzz-read-priority"] === "background" - ? "background" - : "foreground", - ); - const text = - snapshot && response.ok + return response; + }); + }; + response = presence + ? await request() + : await admittedApiRequest( + lane, + request, + requestSignal, + route === "/api/relay/query" && + req.headers["x-buzz-read-priority"] === "background" + ? "background" + : "foreground", + ); + const text = presence + ? await presenceText(response) + : snapshot && response.ok ? await readSnapshotText(response) : await response.text(); // The relay's own service time separates server work from network time. @@ -994,6 +1040,8 @@ export function relayBrokerPlugin({ } catch { failure = apiFailure(response.status, undefined); } + if (presence && failure.quota === "api") + lane.pause(failure.retryAfterMs); return json(res, response.status, failure); } if (profile) { @@ -1015,7 +1063,10 @@ export function relayBrokerPlugin({ res.off("close", release); } } finally { - inflight--; + if (presence) { + presenceFlight = false; + releasePresence(); + } else inflight--; } } catch (error) { if (res.destroyed) return; // The browser gave up first; nothing to answer. diff --git a/docs/presence.md b/docs/presence.md new file mode 100644 index 00000000..6279f9b8 --- /dev/null +++ b/docs/presence.md @@ -0,0 +1,52 @@ +# Periodically refreshed community presence + +Message and thread bylines and profiles show Online, Away, Offline or Unknown +with text and distinct symbols, not color alone. This is snapshot presence, not +an immediate live-status stream. + +## Ownership and bounds + +- One app input source derives Away after ten minutes without Buzz input. Focus + loss and changing communities alone do not mean Away. Same-origin windows share + recent input through BroadcastChannel; no machine-idle or cross-device claim. +- Each retained connected session owns one volatile presence directory. Mounted + rows (including timeline overscan and offscreen thread replies) demand authors. + At most 256 unique authors are selected; profiles take priority, then existing + selections and stable acquisition order. Overflow stays Unknown with a tooltip. +- Initial/new demand coalesces for 100ms behind a five-second start gate. Successful + views refresh after 60–65 seconds. Evidence expires 75 seconds after request start. + Empty demand makes no request; removed authors lose their evidence. Hidden views, + disconnect, access/cache invalidation and disposal invalidate observations. +- One bounded complete snapshot is validated before any status changes. The + configured relay must sign each unique requested subject; only a successful + complete response can make omitted subjects Offline. This is read-time evidence, + not the relay's remaining lease. Invalid or failed responses mean Unknown. +- The development broker permits one optional HTTP flight, a principal-wide + five-second start gate, 256 subjects, a 20 KiB request and 1 MiB response, and + a ten-second request lifetime. Optional work never uses ordinary read slots or + dispatch credit. Local busy skips retry after 5–6 seconds; network failures wait + 60–65 seconds, honoring longer server cooldowns. +- Renewal uses the existing authenticated socket, with a Web Lock per scope/viewer + serializing same-origin windows. Publication is lossy and bounded; it never + enters the durable outbox, replays missed ticks, or publishes Offline on close. + Hidden observation does not stop connected-community renewal. Ordinary setup and + actual shared relay cooldowns still take priority. + +Sustained ordinary traffic may starve optional reads and renewal. Unknown can last +indefinitely when busy/unavailable. These limits bound client work; they do not +promise zero CPU/network/backend cost, instant transitions, or a delivery SLA. +There are no presence REQs, added sockets, relay changes, or direct-adapter parity. + +## Validation + +Owner tests live with `features/presence`, relay transport/admission and the broker. +`tests/browser/presence.spec.mjs` exercises the built app and production broker with +modeled upstream and ephemeral keys: held snapshots versus chat, shared row demand, +300 distinct thread authors, Unknown during held replacement reads, and real +same-origin Web Lock handoff. `channel-opening.spec.mjs` includes matched-thread +measurement support. See [browser measurement limits](browser-testing.md). + +These fixtures do not certify deployed relay capacity, native-window behavior, +attended account use, or cross-device availability. Full browser activity/idle +transition integration and the complete changed-call-site mutation audit remain +separate validation work; activity derivation has controlled owner tests. diff --git a/src/bundled/profiles/ProfilePanel.tsx b/src/bundled/profiles/ProfilePanel.tsx index 4259775a..496234ef 100644 --- a/src/bundled/profiles/ProfilePanel.tsx +++ b/src/bundled/profiles/ProfilePanel.tsx @@ -1,3 +1,4 @@ +import { PresenceIndicator } from "../../features/presence/react"; import { useEffect, useMemo, @@ -102,6 +103,7 @@ function ProfileDetails({ />

{name}

+ {profile?.about &&

{profile.about}

} {context?.canOpen(activity) && (
diff --git a/src/features/communities/service.ts b/src/features/communities/service.ts index ee7951ca..8a9dfac5 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"; @@ -31,6 +32,7 @@ export function createCommunities(ctx: Context, live: boolean) { let unresolvedSelection: string | null = null; let disposed = false; const controller = new AbortController(); + const presenceActivity = createPresenceActivity(); const listeners = new Set<() => void>(); const relayListeners = new Set<() => void>(); const sessions = new Map(); @@ -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), + presenceActivity, ); sessions.set(id, session); session.subscribe(() => { @@ -185,6 +189,7 @@ export function createCommunities(ctx: Context, live: boolean) { ctx.effect(() => () => { disposed = true; controller.abort(); + presenceActivity.dispose(); listeners.clear(); relayListeners.clear(); return Promise.all(scopes.map((scope) => scope.fiber.dispose())); diff --git a/src/features/messages/MessageRow.tsx b/src/features/messages/MessageRow.tsx index 323a9580..4df88f5f 100644 --- a/src/features/messages/MessageRow.tsx +++ b/src/features/messages/MessageRow.tsx @@ -1,4 +1,5 @@ import { memo, useCallback, useSyncExternalStore } from "react"; +import { PresenceIndicator } from "../presence/react"; import type { UnreadCapability } from "../relay/unread"; import { profileTarget } from "../profiles/target"; import { InlineText } from "../conversation/InlineText"; @@ -119,6 +120,12 @@ export const MessageRow = memo(function MessageRow({
{name} + {session && ( + + )}