diff --git a/dev/relay-broker-fixture.test.mjs b/dev/relay-broker-fixture.test.mjs
index 044703b8..81379bdf 100644
--- a/dev/relay-broker-fixture.test.mjs
+++ b/dev/relay-broker-fixture.test.mjs
@@ -177,7 +177,7 @@ test("actual fixture fails an unclassified 503 even with no console event", asyn
fixture(async (app) => {
await publication(app, false);
}),
- ).rejects.toThrow("Unclassified presence publication 503");
+ ).rejects.toThrow("Unclassified presence publication 404/503");
});
test("same endpoint disposal console cannot hide a separate unclassified response", async () => {
@@ -186,7 +186,7 @@ test("same endpoint disposal console cannot hide a separate unclassified respons
console503(page, await publication(app, true));
await publication(app, false);
}),
- ).rejects.toThrow("Unclassified presence publication 503");
+ ).rejects.toThrow("Unclassified presence publication 404/503");
});
test("one classified response permits one console diagnostic", async () => {
@@ -232,10 +232,197 @@ test.each([
req.emit("end");
res.end(body);
expect(() => evidence.assertPublications()).toThrow(
- "Unclassified presence publication 503",
+ "Unclassified presence publication 404/503",
);
});
+async function retiredPublication(app, community = "primary") {
+ const headers = { Origin: app.origin, "Content-Type": "application/json" };
+ const controller = new AbortController();
+ let streamId;
+ try {
+ const response = await fetch(`${app.origin}/api/relay/primary/stream`, {
+ method: "POST",
+ headers,
+ body: JSON.stringify({ channels: [] }),
+ signal: controller.signal,
+ });
+ expect(response.status).toBe(200);
+ streamId = response.headers.get("x-buzz-live-id");
+ } finally {
+ controller.abort();
+ }
+ await until(() =>
+ app.report.brokerRequests.some(
+ (record) => record.url.endsWith("/stream") && record.close,
+ ),
+ );
+ const endpoint = `${app.origin}/api/relay/${community}/stream-presence-publish`;
+ const response = await fetch(endpoint, {
+ method: "POST",
+ headers,
+ body: JSON.stringify({ streamId, status: "online" }),
+ signal: AbortSignal.timeout(3000),
+ });
+ expect(response.status).toBe(404);
+ expect(await response.json()).toEqual({
+ error: "Live stream no longer available",
+ });
+ expect(app.report.presencePublications).toEqual([]);
+ return { endpoint, streamId };
+}
+
+const console404 = (page, url) =>
+ page.emit("console", {
+ type: () => "error",
+ text: () =>
+ "Failed to load resource: the server responded with a status of 404 (Not Found)",
+ location: () => ({ url }),
+ });
+
+test.each([false, true])(
+ "actual fixture accounts an already-retired publication 404 (console %s)",
+ async (console) => {
+ await fixture(async (app, page) => {
+ const { endpoint, streamId } = await retiredPublication(app);
+ expect(app.report.presencePublicationResponses).toEqual([
+ expect.objectContaining({
+ url: endpoint,
+ streamId,
+ status: 404,
+ disposed: false,
+ retired: true,
+ body: { error: "Live stream no longer available" },
+ }),
+ ]);
+ if (console) console404(page, endpoint);
+ });
+ },
+);
+
+test("retirement in another community cannot classify a real publication 404", async () => {
+ await expect(
+ fixture(async (app) => {
+ await retiredPublication(app, "secondary");
+ }),
+ ).rejects.toThrow("Unclassified presence publication 404/503");
+});
+
+test.each([404, 503])(
+ "retired 404 console cannot hide an unclassified %s at the same endpoint",
+ async (status) => {
+ await expect(
+ fixture(async (app, page) => {
+ const { endpoint } = await retiredPublication(app);
+ console404(page, endpoint);
+ if (status === 503) await publication(app, false);
+ else {
+ const response = await fetch(endpoint, {
+ method: "POST",
+ headers: { Origin: app.origin, "Content-Type": "application/json" },
+ body: JSON.stringify({
+ streamId: "0".repeat(32),
+ status: "online",
+ }),
+ });
+ expect(response.status).toBe(404);
+ await response.text();
+ }
+ }),
+ ).rejects.toThrow("Unclassified presence publication 404/503");
+ },
+);
+
+test.each(["duplicate", "wrong status", "wrong endpoint"])(
+ "retired publication accounting rejects %s console evidence",
+ async (failure) => {
+ await expect(
+ fixture(async (app, page) => {
+ const { endpoint } = await retiredPublication(app);
+ if (failure === "duplicate") {
+ console404(page, endpoint);
+ console404(page, endpoint);
+ } else if (failure === "wrong status") console503(page, endpoint);
+ else console404(page, `${endpoint}/other`);
+ }),
+ ).rejects.toThrow();
+ },
+);
+
+// The real fixture tests above prove production wiring. These controlled response
+// boundaries cover impossible/malformed evidence without altering broker behavior.
+test.each([
+ { name: "current stream", retirement: "never" },
+ { name: "unknown stream", requestId: "b".repeat(32) },
+ { name: "malformed stream ID", requestId: "bad", streamId: "bad" },
+ { name: "different relay", streamPath: "/api/relay/secondary/stream" },
+ { name: "retirement during upload", retirement: "after arrival" },
+ { name: "retirement after response", retirement: "after response" },
+ { name: "missing body", body: undefined },
+ { name: "malformed JSON", body: "not-json" },
+ { name: "wrong error", body: JSON.stringify({ error: "other" }) },
+ {
+ name: "extra field",
+ body: JSON.stringify({
+ error: "Live stream no longer available",
+ accepted: true,
+ }),
+ },
+ { name: "truncated response", truncated: true },
+])("publication 404 fails closed for $name", (options) => {
+ const report = { brokerRequests: [] };
+ const evidence = brokerEvidence(report, new Set());
+ const streamId = options.streamId ?? "a".repeat(32);
+ const request = (url) => {
+ const req = new EventEmitter();
+ req.url = url;
+ req.headers = { host: "127.0.0.1:1234" };
+ return req;
+ };
+ const response = () => {
+ const res = new EventEmitter();
+ res.statusCode = 200;
+ res.writeHead = () => {};
+ res.end = () => {};
+ res.getHeader = () => undefined;
+ return res;
+ };
+ const stream = response();
+ evidence.middleware(
+ request(options.streamPath ?? "/api/relay/primary/stream"),
+ stream,
+ () => {},
+ );
+ stream.writeHead(200, { "X-Buzz-Live-ID": streamId });
+ const retirement = options.retirement ?? "before arrival";
+ if (retirement === "before arrival") stream.emit("close");
+ const req = request("/api/relay/primary/stream-presence-publish");
+ const res = response();
+ evidence.middleware(req, res, () => {});
+ if (retirement === "after arrival") stream.emit("close");
+ req.emit("data", JSON.stringify({ streamId: options.requestId ?? streamId }));
+ req.emit("end");
+ res.statusCode = 404;
+ if (options.truncated) res.emit("close");
+ else
+ res.end(
+ Object.hasOwn(options, "body")
+ ? options.body
+ : JSON.stringify({ error: "Live stream no longer available" }),
+ );
+ if (retirement === "after response") stream.emit("close");
+ expect(report.presencePublicationResponses).toHaveLength(1);
+ expect(() => evidence.assertPublications()).toThrow(
+ "Unclassified presence publication 404/503",
+ );
+ expect(
+ evidence.consoleFilter()(
+ "Failed to load resource: the server responded with a status of 404 (Not Found)",
+ "http://127.0.0.1:1234/api/relay/primary/stream-presence-publish",
+ ),
+ ).toBe(false);
+});
+
test("passive completion evidence preserves Server-Timing and distinguishes an unfinished close", async () => {
const report = { brokerRequests: [] };
const evidence = brokerEvidence(report, new Set());
diff --git a/docs/browser-testing.md b/docs/browser-testing.md
index 0ad69e29..c7e91c2c 100644
--- a/docs/browser-testing.md
+++ b/docs/browser-testing.md
@@ -40,7 +40,8 @@ bin/pnpm test:browser tests/browser/layout.spec.mjs --no-deps
bin/pnpm test:browser --no-deps --workers=1
```
-The default local gate runs `channel-opening.spec.mjs` and `scroll.spec.mjs` first,
+The default local gate runs `channel-opening.spec.mjs`, `scroll.spec.mjs`,
+`presence-contention.spec.mjs`, and `presence-control.spec.mjs` first,
one browser/worker at a time, through the `chromium-measurements` →
`webkit-measurements` dependency chain. Only then may functional journeys run
with two workers. This preserves timing/heap samples without unrelated browser
@@ -198,6 +199,37 @@ The separate `channel-opening.test.ts` exercises catch-up ownership and terminal
retry states through the production session. A held-response reproducer establishes
a failure mechanism; it does not on its own identify a live incident's cause.
+## Presence contention controls
+
+`presence-contention.spec.mjs` forces a real signed presence conflict through the
+session directory, verified reader, and production broker, then holds the upstream
+snapshot. Immediately afterward it submits the actual composer or opens a cold
+channel. The foreground host request must arrive within **200ms** of snapshot
+start (inside the old 500ms residual pacing window); broker admission must be
+under **100ms**, with host-arrival-to-upstream under **150ms**. These generous
+regression ceilings detect the inherited pacing interval, not universal zero-cost
+service. Send keeps the snapshot pending; cold navigation can abort it but must
+not inherit its consumed credit. Optimistic message text is not a send receipt.
+
+`presence-control.spec.mjs` repeats both journeys with only the directory and
+publisher disabled in a test build. It preserves normal reader, broker, signer,
+outbox, cache and connection behavior and asserts that no presence traffic occurs.
+Both files run serially in each measurement engine. Evidence separates runner-clock
+host/upstream timestamps, observed browser request events, browser resource timing,
+Server-Timing admission/auth/network, and the exported production read/write profiler.
+These clocks must not be subtracted across domains.
+
+The opt-in policy relay numerically enforces audited reference defaults: shared
+API **300/min**, WS REQ/EVENT **50/5s**, plus **60 EVENT/min**. Counters are shared
+across sockets for the same community/viewer; HTTP publications share the API
+counter with snapshots. `dev/policy-relay.test.mjs` runs once under Vitest (no browser) and deliberately exceeds each budget
+and requires correlated rejection, community isolation, and window expiry. This
+positive control prevents an empty refusal list from masquerading as enforcement.
+The short browser journeys do not saturate the client's maximum envelope or prove
+capacity isolation; colocated transport tests cover those contracts. Other devices,
+lower deployed quotas, CPU contention, SQL/Redis cost, and native GUI remain outside
+this offline model.
+
## DM label recovery
`dm-labels.spec.mjs` builds the actual page and uses the production broker with
diff --git a/docs/presence.md b/docs/presence.md
new file mode 100644
index 00000000..bdb0c8b3
--- /dev/null
+++ b/docs/presence.md
@@ -0,0 +1,170 @@
+# Presence
+
+`session.presence` is a volatile current-value directory, separate from message
+retention, generic event reconciliation, and the durable outbox. Its values are
+`online`, `away`, `offline`, and `unknown`. Unknown is not evidence of being offline.
+
+Message bylines distinguish status without color: Online is a filled circle, Away
+is a filled square, Offline is a hollow circle, and Unknown has a dotted outline.
+The profile panel also displays the status text; all indicators have an accessible label.
+
+## Ownership and traffic
+
+The session requires explicit transport `presence: true` and live support before
+starting observation or publishing. The production broker supplies both snapshot
+admission and shared-socket controls. Ordinary direct-signed reads remain supported,
+but that alternate adapter does not enable partial presence through its generic
+subscription. Unsupported transports stay Unknown and start no publisher/renewals.
+
+- A timeline or thread owns one demand handle for its viewport plus a 160px margin;
+ the profile panel owns its selected author. Avatar indicators only subscribe to
+ per-author values. Repeated authors share work. Without IntersectionObserver,
+ demand falls back to the bounded rendered window, not all retained history.
+- A session accepts at most 256 unique demanded authors and 64 surface handles.
+ Over-cap demand is rejected, not broadened. Empty or hidden demand removes its
+ observer route and snapshot work; it does not stop the availability publisher.
+- Presence shares the existing authenticated socket. Author interests have a
+ one-second minimum REQ interval, at most two overlapping presence routes, and
+ share the total 1,024 subscription-slot ceiling with normal routes. Presence
+ updates do not restart the message stream and yield to foreground work. Two
+ presence slots remain reserved; channel capacity is 1,020, or 1,019 when the
+ independent Agent Activity observer is enabled.
+- One background snapshot owner uses the existing verified reader. Reads have a
+ five-second cooldown after completion or cancellation, not enqueue time: shared
+ reader/broker queue delays must not compress successive actual reads. Repeated
+ changes coalesce without moving an already scheduled deadline. Stable demand
+ uses one 60–65 second backstop; failures wait at least 60 seconds and honor longer
+ relay retry advice. These are per-session budgets, not a fleet-wide rate limit.
+- Same-status renewals do not trigger a read or UI notification per heartbeat.
+ Conflicting evidence becomes unknown and requests one rate-bounded confirmation.
+ No separate socket, per-avatar polling, heartbeat journal, or durable replay.
+
+## Foreground isolation
+
+Optional presence does not spend ordinary dispatch credit or occupy ordinary
+capacity. Limits are additive, not a claim that the old total concurrency is
+unchanged:
+
+| Boundary | Ordinary work | Optional presence |
+| --- | --- | --- |
+| Host/principal HTTP starts | 500ms minimum | 5s minimum snapshot interval |
+| Session reader | 3 active, at most 1 background; 128 distinct pending | 1 snapshot |
+| Development broker | 6 process-wide inflight | 1 process-wide snapshot |
+| Host/principal WS starts | 250ms minimum REQ interval | 1s REQ / 5s EVENT minimum |
+| Live setup per stream | 4 pending ordinary routes | 1 candidate presence route |
+
+The shared host principal is keyed by relay origin and viewer, not by stream.
+Foreground wins ready HTTP ties; presence setup/publication yields to foreground
+setup. Two presence wires (confirmed plus candidate) remain inside the existing
+1,024 subscription ceiling. Ordinary HTTP preparation and queue budgets stay 128;
+one separate optional preparation lease spans signing, fetch/body and verification.
+A cancelled consumer cannot free unresolved underlying optional work. The analogous
+WS publication lease stays owned until signing settles, fences late completion,
+and never makes ordinary signing wait behind the optional lease.
+
+Snapshot classification requires exactly one filter with only `kinds`, `authors`
+and `limit`: kind `[20001]`, 1–256 unique lowercase full 64-hex authors, and limit
+exactly their count. Priority headers do not grant optional admission, and reader
+normalization cannot upgrade malformed requests into that class.
+
+Server API cooldown remains shared by ordinary and optional HTTP work; WS cooldown
+remains shared across ordinary/presence work and streams. HTTP and WS quota families
+remain independent. Presence must not bypass a real shared quota rejection.
+
+The combined-budget test drives eight production broker transports over real local
+HTTP for a full minute. The existing modeled relay charges all attempts across
+callers and enforces first-call-anchored quota windows. Assertions allow at most
+134 API charges/min, 27 combined REQ/EVENT charges per 5s, and 13 presence EVENTs/min
+including setup and drain; lower throughput bounds prevent a vacuous pass. Exact
+500ms/250ms/1s/5s start clocks are tested deterministically at their admission owners.
+This bounds the modeled client-owned traffic, not requests from other devices or hosts. It is below the inspected reference defaults
+(API300/min, WS50/5s, Messages60/min), not a deployed configuration guarantee. The
+normal publisher still renews only every 60–65s, not every 5s. No local scheduler can
+promise zero CPU, signer-provider, network, SQL/Redis cost or end-user latency.
+
+## Authority and lifecycle
+
+Live presence subjects come from the verified event author. Only a verified
+relay-authored snapshot can use its single requested `p` tag as subject. A
+successful bounded snapshot's omission means no current entry for that subject;
+failed or obsolete reads do not establish offline. Snapshot `created_at` is the
+relay's synthesis time, not heartbeat age or remaining lease duration.
+
+Per-author lifetime/revision and session generation guards fence conflicts,
+removal/re-add, visibility changes, access/cache clearing, and disposal. They do
+not create a globally ordered snapshot/stream protocol. Exact expiry and perfectly
+current status cannot be inferred from the existing wire contract.
+
+Presence-only setup retries are bounded. After four failed attempts, unchanged
+nonempty demand can remain Unknown across an automatic socket reconnect. One
+explicit Live Retry resets that exhaustion whether authenticated or disconnected;
+recovery still requires fresh authentication/EOSE and respects shared cooldowns.
+Replacing or losing the browser broker stream invalidates its separate presence
+readiness; ordinary connected frames cannot restore it.
+
+## Activity and publishing
+
+One app-level activity detector measures **Buzz input**, not operating-system
+idle. Ten minutes without input produces Away (sampled every 30 seconds).
+Window blur or switching communities is not Away. Each connected community/viewer
+publisher sends its current status after startup, on status transitions, and
+roughly every 60–65 seconds. It keeps one in-flight write and the latest desired
+status; missed renewals are not replayed. Acceptance requires the matching socket
+`OK`, not merely handing bytes to the broker.
+
+If stream disposal wins an in-flight publication's settlement, the development
+broker preserves that typed cause as `code: "presence_owner_disposed"` on its
+existing **503 / Presence publication unconfirmed** response. This describes why
+confirmation stopped, not whether an EVENT was sent or accepted. Earlier rejection,
+deadline, socket reset, caller abort, or local failure cannot be relabeled by later
+disposal. No new wait, retry, round trip, or successful browser outcome is added.
+
+Browser fixtures record each publication 404/503 at the host response boundary,
+including when browser cancellation hides its body or console diagnostic. A 503
+is classified only by the explicit disposal code with the unconfirmed body. A 404
+requires the exact missing-stream rejection and evidence that this same relay/stream
+was already retired when the request arrived, not later during teardown. Neither
+classification means delivery. All other recorded 404/503s fail independently of
+console output. Each classified response permits at most one diagnostic matching
+its endpoint and status, and all responses remain in the evidence. Host request
+completion records distinguish `finish` from `close` and retain status/Server-Timing;
+neither closing a response nor handing bytes to the host establishes relay delivery.
+
+Optional same-origin Web Locks coordinate one publisher per community/viewer;
+BroadcastChannel shares recent local input. Unsupported hosts can publish once per
+window. A frozen browser leader can miss renewals; neither mechanism coordinates
+other devices. Connected retained communities continue renewing after navigation.
+Ordinary teardown never publishes identity-wide Offline on another device's behalf.
+
+## Backend limits
+
+This client changes no relay storage or aggregation semantics. The audited relay
+still performs a community lifecycle SQL lookup for a WS heartbeat, access queries
+for HTTP snapshots, and SQL/pool work for `REQ limit:0`. Snapshots also require Redis
+reads and relay signatures. Author scoping reduces delivered traffic/client work;
+it does not remove candidate subscription matching or all backend work.
+
+The relay aggregates by identity with last-arrival-wins status, not device leases.
+Pod-local final disconnect can clear another device's entry without offline fanout.
+Removing these costs and ambiguities requires separate server work with preserved
+authorization/lifecycle fencing; this implementation does not claim zero SQL,
+zero polling, exact expiry, or distributed-device correctness.
+
+## Validation shape
+
+Owner tests live in `src/features/presence/` and relay presence tests. Browser
+journeys are `tests/browser/presence.spec.mjs` (observer lifecycle and real same-origin
+lock/activity handoff) and `presence-integration.spec.mjs` (production conversation,
+broker, snapshots, conflict repair, and route teardown). The paired
+`presence-contention.spec.mjs` / `presence-control.spec.mjs` measure real composer
+send and cold-open admission with a held snapshot versus equivalent no-presence
+owners. Colocated reader, broker and live tests separately hold capacity/signing
+work; `dev/relay-broker-fixture.test.mjs` exercises combined eight-caller budgets; browser
+journeys alone do not establish those bounds. `dev/policy-relay.test.mjs` proves that
+the numerical quota fixture rejects deliberate overload. Channel-opening measurements
+retain their existing warm-switch budget.
+
+These browser tests use isolated identities and modeled relay policy. Passing
+native compilation/tests is not native GUI acceptance; none of these establishes
+deployed SQL/Redis load or production latency. See [the contribution workflow](contributing.md)
+for batch gates and [browser testing](browser-testing.md) for measurement limits.
diff --git a/docs/relay-queries.md b/docs/relay-queries.md
index 8dce9f23..34643b54 100644
--- a/docs/relay-queries.md
+++ b/docs/relay-queries.md
@@ -62,6 +62,12 @@ locally authored event is **not proof of relay acceptance**. Signature-verified
membership, bounds and persistence. Domain folds can consume local payloads, but
must not let them manufacture relay-authored authority.
+## Presence
+
+`session.presence` exposes volatile per-author status and surface-owned demand,
+using the same authenticated socket and verified background reader without message
+retention or outbox replay. See [presence ownership, traffic budgets and limits](presence.md).
+
## Community emoji
`session.emoji` owns the current community's kind-30030 `d=buzz:custom-emoji`
diff --git a/src/app/services.test.ts b/src/app/services.test.ts
index ff070391..cfdbcca2 100644
--- a/src/app/services.test.ts
+++ b/src/app/services.test.ts
@@ -121,8 +121,17 @@ async function openCommunities() {
function expectHostStopped() {
expect(signals.every((signal) => signal.aborted)).toBe(true);
for (const stream of streams) expect(stream.close).toHaveBeenCalledTimes(1);
- expect(document.addEventListener).toHaveBeenCalledTimes(3);
- expect(document.removeEventListener).toHaveBeenCalledTimes(3);
+ // Three retained service visibility listeners plus one shared activity detector.
+ const added = vi.mocked(document.addEventListener).mock.calls;
+ const removed = vi.mocked(document.removeEventListener).mock.calls;
+ expect(added).toHaveLength(8);
+ expect(removed).toHaveLength(added.length);
+ for (const [type, listener] of added)
+ expect(removed.some(([name, fn]) => name === type && fn === listener)).toBe(
+ true,
+ );
+ for (const type of ["pointerdown", "pointermove", "keydown", "wheel"])
+ expect(added.filter(([name]) => name === type)).toHaveLength(1);
expect(services.pages.snapshot()).toHaveLength(0);
}
diff --git a/src/bundled/profiles/ProfilePanel.tsx b/src/bundled/profiles/ProfilePanel.tsx
index 4259775a..ac5a5083 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,
@@ -45,6 +49,7 @@ function ProfileDetails({
pubkey: string;
context: PanelProps["context"];
}) {
+ usePresenceDemand(session.presence, pubkey);
const selection = useMemo(
() => selectProfiles(session.profiles, [pubkey]),
[session.profiles, pubkey],
@@ -101,6 +106,7 @@ function ProfileDetails({
size="large"
/>
{name}
+
{profile?.about && {profile.about}
}
{context?.canOpen(activity) && (
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.test.tsx b/src/features/messages/ChannelTimeline.test.tsx
index db09202f..fbbbb8e4 100644
--- a/src/features/messages/ChannelTimeline.test.tsx
+++ b/src/features/messages/ChannelTimeline.test.tsx
@@ -92,6 +92,12 @@ vi.mock("react", async (original) => ({
}
},
}));
+// These shallow geometry tests do not implement browser observer APIs.
+// Presence's real observer/render lifecycle runs in tests/browser/presence.spec.mjs.
+vi.mock("../presence/react", async (original) => ({
+ ...(await original()),
+ usePresenceSurface: vi.fn(),
+}));
vi.mock("../relay/react", () => ({
useRowProfiles: () => new Map(),
}));
diff --git a/src/features/messages/ChannelTimeline.tsx b/src/features/messages/ChannelTimeline.tsx
index 9db2483c..522230ac 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 { MembershipRow } from "./MembershipRow";
import { membershipRows } from "./membership-rows";
import type { ConversationExtensions } from "../conversation/contracts";
@@ -115,6 +116,7 @@ function Timeline({
const [focusedMessageId, setFocusedMessageId] = useState();
const focusedIndex = rows.findIndex((row) => row.id === focusedMessageId);
const scroller = useRef(null);
+ usePresenceSurface(queries.presence, scroller);
const handle = useRef(null);
const [size, setSize] = useState({ width: 0, height: 0 });
const width = size.width;
@@ -488,6 +490,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 9cc74e38..58083fc7 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";
@@ -13,6 +15,7 @@ import { usesLargeEmojiPresentation } from "./emoji-size";
export type MessageRowProps = {
row: ChannelMessage;
+ presence?: PresenceQueries | undefined;
unread?: UnreadCapability | undefined;
extensions?: ConversationExtensions | undefined;
profile: Profile | undefined;
@@ -27,6 +30,7 @@ export type MessageRowProps = {
export const MessageRow = memo(function MessageRow({
row,
+ presence,
unread,
extensions,
profile,
@@ -58,7 +62,7 @@ export const MessageRow = memo(function MessageRow({
const AvatarTag = clickable ? "button" : "div";
const emojiOnly = usesLargeEmojiPresentation(row.content, row.emoji);
return (
-
+
{day && (
@@ -93,6 +97,9 @@ export const MessageRow = memo(function MessageRow({
{name}
+ {presence && (
+
+ )}
{new Date(row.createdAt * 1000).toLocaleTimeString(undefined, {
hour: "numeric",
diff --git a/src/features/messages/ThreadPanel.test.tsx b/src/features/messages/ThreadPanel.test.tsx
index 7642f9ba..12adaf0f 100644
--- a/src/features/messages/ThreadPanel.test.tsx
+++ b/src/features/messages/ThreadPanel.test.tsx
@@ -71,6 +71,12 @@ vi.mock("react", async (original) => ({
useSyncExternalStore: (_subscribe: unknown, snapshot: () => unknown) =>
snapshot(),
}));
+// These shallow geometry tests do not implement browser observer APIs.
+// Presence's real observer/render lifecycle runs in tests/browser/presence.spec.mjs.
+vi.mock("../presence/react", async (original) => ({
+ ...(await original()),
+ usePresenceSurface: vi.fn(),
+}));
vi.mock("../relay/react", () => {
const profiles = new Map();
return { useRowProfiles: () => profiles };
diff --git a/src/features/messages/ThreadPanel.tsx b/src/features/messages/ThreadPanel.tsx
index 4bd07991..2c73f737 100644
--- a/src/features/messages/ThreadPanel.tsx
+++ b/src/features/messages/ThreadPanel.tsx
@@ -1,4 +1,5 @@
// biome-ignore-all lint/a11y/noNoninteractiveTabindex: The thread region supports keyboard scrolling and Escape.
+import { usePresenceSurface } from "../presence/react";
import {
useCallback,
useEffect,
@@ -191,6 +192,7 @@ function ThreadMessages({
}, [session.profiles, authors]);
const profiles = useRowProfiles(session.profiles, rows);
const scroller = useRef(null);
+ usePresenceSurface(session.presence, scroller);
const positioned = useRef(false);
const follow = useRef(true);
const targetAnchor = useRef(undefined);
@@ -313,6 +315,7 @@ function ThreadMessages({
{
+ vi.useFakeTimers();
+ vi.setSystemTime(1000000);
+});
+afterEach(() => {
+ vi.useRealTimers();
+ vi.unstubAllGlobals();
+});
+it("shares one detector, keeps raw input out of subscribers, and blur alone never means away", async () => {
+ const doc = new EventTarget() as EventTarget & { visibilityState: string };
+ doc.visibilityState = "visible";
+ vi.stubGlobal("document", doc);
+ const owner = createPresenceActivity();
+ const listener = vi.fn();
+ owner.activity.subscribe(listener);
+ for (let i = 0; i < 1000; i++) doc.dispatchEvent(new Event("pointermove"));
+ expect(listener).not.toHaveBeenCalled();
+ doc.visibilityState = "hidden";
+ doc.dispatchEvent(new Event("visibilitychange"));
+ expect(owner.activity.snapshot()).toEqual({
+ status: "online",
+ visible: false,
+ });
+ await vi.advanceTimersByTimeAsync(600000);
+ expect(owner.activity.snapshot().status).toBe("away");
+ doc.visibilityState = "visible";
+ doc.dispatchEvent(new Event("visibilitychange"));
+ expect(owner.activity.snapshot().status).toBe("away");
+ doc.dispatchEvent(new Event("keydown"));
+ expect(owner.activity.snapshot().status).toBe("online");
+ const count = listener.mock.calls.length;
+ owner.dispose();
+ await vi.advanceTimersByTimeAsync(600000);
+ doc.dispatchEvent(new Event("pointermove"));
+ expect(listener).toHaveBeenCalledTimes(count);
+});
diff --git a/src/features/presence/activity.ts b/src/features/presence/activity.ts
new file mode 100644
index 00000000..dd748533
--- /dev/null
+++ b/src/features/presence/activity.ts
@@ -0,0 +1,88 @@
+export type ActivitySnapshot = Readonly<{
+ status: "online" | "away";
+ visible: boolean;
+}>;
+export type PresenceActivity = {
+ snapshot(): ActivitySnapshot;
+ subscribe(listener: () => void): () => void;
+};
+const IDLE_MS = 10 * 60 * 1000;
+/** One app-owned detector. Web fallback measures Buzz input, not machine-wide idle.
+ * Same-origin windows exchange activity only; no identity, message or status history. */
+export function createPresenceActivity() {
+ let closed = false;
+ let lastInput = Date.now();
+ let lastBroadcast = 0;
+ const doc = typeof document === "undefined" ? undefined : document;
+ const page = typeof window === "undefined" ? undefined : window;
+ let state: ActivitySnapshot = Object.freeze({
+ status: "online",
+ visible: doc?.visibilityState !== "hidden",
+ });
+ const listeners = new Set<() => void>();
+ let channel: BroadcastChannel | undefined;
+ try {
+ if (page && typeof BroadcastChannel !== "undefined")
+ channel = new BroadcastChannel("buzz-presence-activity-v1");
+ } catch {
+ /* Local input remains usable when cross-window messaging is unavailable. */
+ }
+ function sample() {
+ if (closed) return;
+ const status = Date.now() - lastInput >= IDLE_MS ? "away" : "online";
+ const visible = doc?.visibilityState !== "hidden";
+ if (status === state.status && visible === state.visible) return;
+ state = Object.freeze({ status, visible });
+ for (const listener of listeners) listener();
+ }
+ function input() {
+ const now = Date.now();
+ if (closed || doc?.visibilityState === "hidden") return;
+ // Raw activity remains local, outside React; network publication observes derived state only.
+ lastInput = now;
+ if (now - lastBroadcast >= 1000) {
+ lastBroadcast = now;
+ channel?.postMessage(now);
+ sample();
+ }
+ }
+ if (channel)
+ channel.onmessage = ({ data }: MessageEvent) => {
+ const now = Date.now();
+ if (
+ typeof data !== "number" ||
+ !Number.isSafeInteger(data) ||
+ data > now + 1000 ||
+ data < now - IDLE_MS
+ )
+ return;
+ lastInput = Math.max(lastInput, Math.min(data, now));
+ sample();
+ };
+ const events = ["pointerdown", "pointermove", "keydown", "wheel"];
+ for (const event of events)
+ doc?.addEventListener(event, input, { passive: true });
+ doc?.addEventListener("visibilitychange", sample);
+ page?.addEventListener("pageshow", sample);
+ const timer = setInterval(sample, 30000);
+ return {
+ activity: {
+ snapshot: () => state,
+ subscribe(listener: () => void) {
+ listeners.add(listener);
+ return () => {
+ listeners.delete(listener);
+ };
+ },
+ } satisfies PresenceActivity,
+ dispose() {
+ closed = true;
+ clearInterval(timer);
+ channel?.close();
+ listeners.clear();
+ for (const event of events) doc?.removeEventListener(event, input);
+ doc?.removeEventListener("visibilitychange", sample);
+ page?.removeEventListener("pageshow", sample);
+ },
+ };
+}
diff --git a/src/features/presence/directory.test.ts b/src/features/presence/directory.test.ts
new file mode 100644
index 00000000..0ef67c3b
--- /dev/null
+++ b/src/features/presence/directory.test.ts
@@ -0,0 +1,261 @@
+import { afterEach, beforeEach, expect, it, vi } from "vitest";
+import { createPresenceDirectory, parsePresenceStatus } from "./directory";
+import { keypair, signed } from "../relay/testing";
+import type { RelayEvent } from "../relay/events";
+import type { RelayReader } from "../relay/reader";
+
+beforeEach(() => {
+ vi.useFakeTimers();
+ vi.setSystemTime(1000000);
+});
+afterEach(() => vi.useRealTimers());
+function setup() {
+ const relay = keypair();
+ const user = keypair();
+ const updates = vi.fn();
+ const reads: {
+ resolve(events: RelayEvent[]): void;
+ reject(error: Error): void;
+ signal?: AbortSignal;
+ }[] = [];
+ const read = vi.fn(
+ (_filters, options) =>
+ new Promise((resolve, reject) =>
+ reads.push({
+ resolve,
+ reject,
+ ...(options?.signal ? { signal: options.signal } : {}),
+ }),
+ ),
+ );
+ const owner = createPresenceDirectory({
+ reader: { read },
+ relayAuthor: relay.pubkey,
+ updateInterests: updates,
+ supported: true,
+ random: () => 0,
+ });
+ const demand = owner.queries.demand();
+ owner.connection(true);
+ const mount = (authors = [user.pubkey]) => {
+ demand.update(authors);
+ owner.route({ status: "ready", authors });
+ };
+ const snapshot = (status = "online", who = user.pubkey) =>
+ signed(relay, { kind: 20001, content: status, tags: [["p", who]] });
+ const live = (status = "online", time = 1) =>
+ signed(user, { kind: 20001, content: status, tags: [], created_at: time });
+ return {
+ owner,
+ demand,
+ user,
+ relay,
+ reads,
+ read,
+ updates,
+ mount,
+ snapshot,
+ live,
+ };
+}
+it("one surface deduplicates 1000 avatars, shares a background snapshot and stops on empty demand", async () => {
+ const s = setup();
+ const authors = Array.from({ length: 20 }, (_, i) =>
+ i.toString(16).padStart(64, "0"),
+ );
+ s.mount(Array.from({ length: 1000 }, (_, i) => authors[i % 20] ?? ""));
+ await vi.advanceTimersByTimeAsync(100);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ expect(s.read.mock.calls[0]?.[0]).toEqual([
+ { kinds: [20001], authors, limit: 20 },
+ ]);
+ expect(s.read.mock.calls[0]?.[1]?.priority).toBe("background");
+ expect(s.owner.diagnostics()).toMatchObject({ authors: 20, handles: 1 });
+ s.reads[0]?.resolve([]);
+ await vi.advanceTimersByTimeAsync(0);
+ s.demand.dispose();
+ expect(s.updates).toHaveBeenLastCalledWith([]);
+ await vi.advanceTimersByTimeAsync(180000);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ s.owner.dispose();
+});
+it("same-status renewals do not notify or poll per heartbeat; the backstop is shared", async () => {
+ const s = setup();
+ s.mount();
+ const listener = vi.fn();
+ s.owner.queries.subscribe(s.user.pubkey, listener);
+ await vi.advanceTimersByTimeAsync(100);
+ s.reads[0]?.resolve([s.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(listener).toHaveBeenCalledTimes(1);
+ for (let i = 0; i < 100; i++) s.owner.receive([s.live("online", i)]);
+ await vi.advanceTimersByTimeAsync(5000);
+ expect(listener).toHaveBeenCalledTimes(1);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ await vi.advanceTimersByTimeAsync(55200);
+ expect(s.read).toHaveBeenCalledTimes(2);
+ expect(listener).toHaveBeenCalledTimes(1);
+ s.reads[1]?.resolve([s.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(listener).toHaveBeenCalledTimes(1);
+ s.owner.dispose();
+});
+it("offline snapshot then delayed heartbeat stays unknown until one rate-bounded confirmation", async () => {
+ const s = setup();
+ s.mount();
+ await vi.advanceTimersByTimeAsync(100);
+ s.reads[0]?.resolve([]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("offline");
+ for (let i = 0; i < 100; i++) s.owner.receive([s.live("online", i)]);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ await vi.advanceTimersByTimeAsync(4999);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ await vi.advanceTimersByTimeAsync(1);
+ expect(s.read).toHaveBeenCalledTimes(2);
+ s.reads[1]?.resolve([]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("offline");
+ s.owner.dispose();
+});
+it("queue delay cannot compress the snapshot cooldown after completion", async () => {
+ const s = setup();
+ s.mount();
+ await vi.advanceTimersByTimeAsync(100);
+ // The shared reader/broker may hold this background read behind foreground work.
+ await vi.advanceTimersByTimeAsync(7000);
+ s.reads[0]?.resolve([s.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ s.owner.receive([s.live("away")]);
+ await vi.advanceTimersByTimeAsync(4999);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ await vi.advanceTimersByTimeAsync(1);
+ expect(s.read).toHaveBeenCalledTimes(2);
+ s.owner.dispose();
+});
+it("cancelling a delayed read preserves cooldown through demand replacement and late completion", async () => {
+ const s = setup();
+ s.mount();
+ await vi.advanceTimersByTimeAsync(7100);
+ s.demand.update([]);
+ expect(s.reads[0]?.signal?.aborted).toBe(true);
+ s.mount();
+ await vi.advanceTimersByTimeAsync(1000);
+ s.reads[0]?.resolve([s.snapshot()]);
+ await vi.advanceTimersByTimeAsync(3999);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ await vi.advanceTimersByTimeAsync(1);
+ expect(s.read).toHaveBeenCalledTimes(2);
+ s.owner.dispose();
+});
+it("conflict during snapshot, clock skew and malicious live p tags cannot supply authority", async () => {
+ const s = setup();
+ s.mount();
+ await vi.advanceTimersByTimeAsync(100);
+ s.owner.receive([s.live("away", 9000000000)]);
+ s.reads[0]?.resolve([s.snapshot("online")]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ const attacker = keypair();
+ s.owner.receive([
+ signed(attacker, {
+ kind: 20001,
+ content: "online",
+ tags: [["p", s.user.pubkey]],
+ }),
+ ]);
+ expect(s.owner.diagnostics().received).toBe(1);
+ await vi.advanceTimersByTimeAsync(5000);
+ s.reads[1]?.resolve([s.snapshot("away")]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("away");
+ s.owner.dispose();
+});
+it("failed or malformed snapshots are not offline and retry only on the slow budget", async () => {
+ const s = setup();
+ s.mount();
+ await vi.advanceTimersByTimeAsync(100);
+ s.reads[0]?.resolve([s.live()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ expect(s.owner.diagnostics().error).toContain("Invalid relay");
+ await vi.advanceTimersByTimeAsync(59000);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ await vi.advanceTimersByTimeAsync(1000);
+ expect(s.read).toHaveBeenCalledTimes(2);
+ s.reads[1]?.reject(new Error("Redis unavailable"));
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ s.owner.dispose();
+});
+it("removed/re-added demand, hidden state, access clear and disposal fence late reads", async () => {
+ const s = setup();
+ s.mount();
+ await vi.advanceTimersByTimeAsync(100);
+ s.demand.update([]);
+ s.mount();
+ s.reads[0]?.resolve([s.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ await vi.advanceTimersByTimeAsync(5000);
+ s.owner.visibility(false);
+ expect(s.updates).toHaveBeenLastCalledWith([]);
+ expect(s.reads[1]?.signal?.aborted).toBe(true);
+ s.reads[1]?.resolve([s.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ s.owner.visibility(true);
+ s.owner.route({ status: "ready", authors: [s.user.pubkey] });
+ await vi.advanceTimersByTimeAsync(5000);
+ s.owner.clear();
+ s.reads[2]?.resolve([s.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(s.owner.queries.get(s.user.pubkey)).toBe("unknown");
+ s.owner.dispose();
+ await vi.advanceTimersByTimeAsync(180000);
+ expect(s.owner.diagnostics()).toMatchObject({
+ authors: 0,
+ handles: 0,
+ pending: false,
+ });
+});
+it("continuous changed demand still flushes and never creates concurrent reads", async () => {
+ const s = setup();
+ s.mount();
+ for (let i = 0; i < 20; i++) {
+ const authors = [s.user.pubkey, i.toString(16).padStart(64, "0")];
+ s.demand.update(authors);
+ s.owner.route({ status: "ready", authors });
+ await vi.advanceTimersByTimeAsync(50);
+ }
+ expect(s.read).toHaveBeenCalledTimes(1);
+ s.reads[0]?.resolve([]);
+ await vi.advanceTimersByTimeAsync(0);
+ await vi.advanceTimersByTimeAsync(4999);
+ expect(s.read).toHaveBeenCalledTimes(1);
+ await vi.advanceTimersByTimeAsync(1);
+ expect(s.read).toHaveBeenCalledTimes(2);
+ s.owner.dispose();
+});
+it("caps authors and surface handles without a broad fallback", () => {
+ const s = setup();
+ const tooMany = Array.from({ length: 257 }, (_, i) =>
+ i.toString(16).padStart(64, "0"),
+ );
+ expect(s.demand.update(tooMany)).toBe(false);
+ expect(s.updates).toHaveBeenLastCalledWith([]);
+ for (let i = 0; i < 100; i++) s.owner.queries.demand();
+ expect(s.owner.diagnostics().handles).toBe(64);
+ s.owner.dispose();
+});
+it("accepts bare statuses and legacy JSON, not arbitrary values", () => {
+ for (const status of ["online", "away", "offline"]) {
+ expect(parsePresenceStatus({ content: status })).toBe(status);
+ expect(parsePresenceStatus({ content: JSON.stringify({ status }) })).toBe(
+ status,
+ );
+ }
+ expect(parsePresenceStatus({ content: "banana" })).toBeUndefined();
+ expect(parsePresenceStatus({ content: "null" })).toBeUndefined();
+});
diff --git a/src/features/presence/directory.ts b/src/features/presence/directory.ts
new file mode 100644
index 00000000..1bb7c4ab
--- /dev/null
+++ b/src/features/presence/directory.ts
@@ -0,0 +1,368 @@
+import type { RelayEvent } from "../relay/events";
+import type { RelayReader } from "../relay/reader";
+import { ReadError } from "../relay/errors";
+
+export type PresenceStatus = "online" | "away" | "offline" | "unknown";
+export type PresenceDemand = {
+ update(authors: readonly string[]): boolean;
+ dispose(): void;
+};
+export type PresenceQueries = {
+ get(author: string): PresenceStatus;
+ subscribe(author: string, listener: () => void): () => void;
+ demand(): PresenceDemand;
+};
+type Entry = {
+ status: PresenceStatus;
+ confirmed?: PresenceStatus;
+ live?: PresenceStatus;
+ revision: number;
+ seen: string[];
+ dirty: boolean;
+};
+const AUTHOR_LIMIT = 256;
+const HANDLE_LIMIT = 64;
+const MIN_READ_MS = 5000;
+const BACKSTOP_MS = 60000;
+const validAuthor = (value: string) => /^[0-9a-f]{64}$/.test(value);
+
+/** Volatile, bounded current values. Never feeds the event journal or durable outbox. */
+export function createPresenceDirectory({
+ reader,
+ relayAuthor,
+ updateInterests,
+ supported,
+ notify = (listener) => listener(),
+ now = () => Date.now(),
+ random = Math.random,
+}: {
+ reader: RelayReader;
+ relayAuthor: string;
+ updateInterests(authors: readonly string[]): void;
+ supported: boolean;
+ notify?: (listener: () => void) => void;
+ now?: () => number;
+ random?: () => number;
+}) {
+ let closed = false;
+ let visible = true;
+ let connected = false;
+ let ready = new Set();
+ let generation = 0;
+ let readAfter = 0;
+ let retryAt = 0;
+ let timer: ReturnType | undefined;
+ let backstop: ReturnType | undefined;
+ let pending: AbortController | undefined;
+ let error: string | undefined;
+ const entries = new Map();
+ const handles = new Map();
+ const listeners = new Map void>>();
+ const counters = { reads: 0, received: 0, notifications: 0, limited: 0 };
+ const emit = (author: string) => {
+ for (const listener of listeners.get(author) ?? []) {
+ counters.notifications++;
+ notify(listener);
+ }
+ };
+ const set = (author: string, entry: Entry, status: PresenceStatus) => {
+ if (entry.status === status) return;
+ entry.status = status;
+ emit(author);
+ };
+ const eligible = () =>
+ !closed && supported && visible && connected && entries.size > 0;
+ function stopRead() {
+ generation++;
+ if (pending) readAfter = now() + MIN_READ_MS;
+ pending?.abort();
+ pending = undefined;
+ clearTimeout(timer);
+ clearTimeout(backstop);
+ timer = backstop = undefined;
+ }
+ function dirty() {
+ for (const [author, entry] of entries) {
+ entry.dirty = true;
+ set(author, entry, "unknown");
+ }
+ }
+ function schedule() {
+ if (
+ !eligible() ||
+ pending ||
+ timer ||
+ ![...entries].some(([key, value]) => value.dirty && ready.has(key))
+ )
+ return;
+ // Fixed deadline: later input changes the pending set, never postpones this flush.
+ timer = setTimeout(
+ () => {
+ timer = undefined;
+ void read();
+ },
+ Math.max(100, readAfter - now(), retryAt - now()),
+ );
+ }
+ function periodic() {
+ clearTimeout(backstop);
+ if (!eligible()) return;
+ backstop = setTimeout(
+ () => {
+ backstop = undefined;
+ for (const entry of entries.values()) entry.dirty = true;
+ schedule();
+ },
+ BACKSTOP_MS + random() * 5000,
+ );
+ }
+ async function read() {
+ if (!eligible() || pending) return;
+ const requested = new Map([...entries].filter(([key]) => ready.has(key)));
+ if (!requested.size) return;
+ const revisions = new Map(
+ [...requested].map(([key, value]) => [key, value.revision]),
+ );
+ const current = generation;
+ const owned = new AbortController();
+ pending = owned;
+ counters.reads++;
+ const valid = () =>
+ !closed && !owned.signal.aborted && generation === current;
+ try {
+ const events = await reader.read(
+ [
+ {
+ kinds: [20001],
+ authors: [...requested.keys()],
+ limit: requested.size,
+ },
+ ],
+ { signal: owned.signal, priority: "background", fresh: true },
+ );
+ if (!valid()) return;
+ const values = new Map();
+ if (events.length > requested.size)
+ throw new Error("Presence snapshot exceeds requested subjects");
+ for (const event of events) {
+ const subjects = event.tags.filter(([name]) => name === "p");
+ const author = subjects[0]?.[1];
+ const status = parsePresenceStatus(event);
+ if (
+ event.kind !== 20001 ||
+ event.pubkey !== relayAuthor ||
+ subjects.length !== 1 ||
+ subjects[0]?.length !== 2 ||
+ !author ||
+ !requested.has(author) ||
+ values.has(author) ||
+ !status
+ )
+ throw new Error("Invalid relay presence snapshot");
+ values.set(author, status);
+ }
+ error = undefined;
+ retryAt = 0;
+ for (const [author, entry] of requested) {
+ if (entries.get(author) !== entry) continue; // Removed/re-added is another lifetime.
+ const status = values.get(author) ?? "offline";
+ if (entry.revision !== revisions.get(author) && entry.live !== status) {
+ entry.dirty = true;
+ set(author, entry, "unknown");
+ } else {
+ entry.confirmed = status;
+ entry.dirty = false;
+ set(author, entry, status);
+ }
+ }
+ } catch (failure) {
+ if (!valid()) return;
+ error =
+ failure instanceof Error ? failure.message : "Presence unavailable";
+ retryAt =
+ now() +
+ Math.max(
+ BACKSTOP_MS,
+ failure instanceof ReadError ? (failure.retryAfterMs ?? 0) : 0,
+ );
+ dirty();
+ } finally {
+ if (pending === owned) {
+ // Start the cooldown after the whole shared reader/broker operation, not
+ // enqueue time: admission delay must not compress actual relay reads.
+ readAfter = now() + MIN_READ_MS;
+ pending = undefined;
+ periodic();
+ schedule();
+ }
+ }
+ }
+ function interests() {
+ const wanted = new Set([...handles.values()].flat());
+ for (const [author] of entries)
+ if (!wanted.has(author)) {
+ entries.delete(author);
+ emit(author);
+ }
+ for (const author of wanted)
+ if (!entries.has(author))
+ entries.set(author, {
+ status: "unknown",
+ revision: 0,
+ seen: [],
+ dirty: true,
+ });
+ if (!entries.size) stopRead();
+ if (!closed && supported)
+ updateInterests(visible ? [...entries.keys()].sort() : []);
+ schedule();
+ }
+ const queries: PresenceQueries = Object.freeze({
+ get: (author: string) => entries.get(author)?.status ?? "unknown",
+ subscribe(author: string, listener: () => void) {
+ if (closed) return () => {};
+ let selected = listeners.get(author);
+ if (!selected) {
+ selected = new Set();
+ listeners.set(author, selected);
+ }
+ selected.add(listener);
+ return () => {
+ selected.delete(listener);
+ if (!selected.size && listeners.get(author) === selected)
+ listeners.delete(author);
+ };
+ },
+ demand() {
+ const token = {};
+ let disposed = closed || handles.size >= HANDLE_LIMIT;
+ if (disposed) counters.limited++;
+ else handles.set(token, []);
+ return {
+ update(input) {
+ if (disposed || closed) return false;
+ const authors = [...new Set(input)];
+ const union = new Set(
+ [...handles].flatMap(([key, value]) =>
+ key === token ? [] : [...value],
+ ),
+ );
+ for (const author of authors) union.add(author);
+ const allowed =
+ authors.length <= AUTHOR_LIMIT &&
+ authors.every(validAuthor) &&
+ union.size <= AUTHOR_LIMIT;
+ if (!allowed) counters.limited++;
+ handles.set(token, allowed ? authors : []);
+ interests();
+ return allowed;
+ },
+ dispose() {
+ if (disposed) return;
+ disposed = true;
+ handles.delete(token);
+ interests();
+ },
+ };
+ },
+ });
+ return {
+ queries,
+ receive(events: readonly RelayEvent[]) {
+ if (!eligible()) return;
+ for (const event of events) {
+ const entry = entries.get(event.pubkey);
+ const status = parsePresenceStatus(event);
+ if (
+ event.kind !== 20001 ||
+ !entry ||
+ !status ||
+ entry.seen.includes(event.id)
+ )
+ continue;
+ counters.received++;
+ entry.seen.push(event.id);
+ if (entry.seen.length > 8) entry.seen.shift();
+ entry.live = status;
+ entry.revision++;
+ if (entry.confirmed !== status || entry.status === "unknown") {
+ entry.dirty = true;
+ set(event.pubkey, entry, "unknown");
+ }
+ }
+ schedule();
+ },
+ route(state: {
+ status: string;
+ authors: readonly string[];
+ error?: string | undefined;
+ }) {
+ if (closed) return;
+ ready = new Set(state.status === "ready" ? state.authors : []);
+ if (state.status === "error") {
+ error = state.error ?? "Presence subscription unavailable";
+ stopRead();
+ dirty();
+ }
+ schedule();
+ },
+ connection(active: boolean) {
+ if (closed || connected === active) return;
+ connected = active;
+ stopRead();
+ ready.clear();
+ dirty();
+ },
+ visibility(active: boolean) {
+ if (closed || visible === active) return;
+ visible = active;
+ stopRead();
+ ready.clear();
+ dirty();
+ interests();
+ },
+ clear() {
+ if (closed) return;
+ stopRead();
+ for (const entry of entries.values()) {
+ delete entry.confirmed;
+ delete entry.live;
+ entry.seen = [];
+ }
+ dirty();
+ schedule();
+ },
+ dispose() {
+ closed = true;
+ stopRead();
+ dirty();
+ entries.clear();
+ handles.clear();
+ listeners.clear();
+ },
+ diagnostics: () => ({
+ ...counters,
+ authors: entries.size,
+ handles: handles.size,
+ pending: !!pending,
+ error,
+ }),
+ };
+}
+/** Bare status plus legacy JSON; never treat unknown content as Online. */
+export function parsePresenceStatus(
+ event: Pick,
+): Exclude | undefined {
+ try {
+ if (event.content.length > 2048) return;
+ const status: unknown = ["online", "away", "offline"].includes(
+ event.content,
+ )
+ ? event.content
+ : JSON.parse(event.content)?.status;
+ if (status === "online" || status === "away" || status === "offline")
+ return status;
+ } catch {
+ /* Invalid presence is not an Offline signal. */
+ }
+}
diff --git a/src/features/presence/publisher.test.ts b/src/features/presence/publisher.test.ts
new file mode 100644
index 00000000..7c0f83ef
--- /dev/null
+++ b/src/features/presence/publisher.test.ts
@@ -0,0 +1,187 @@
+import { afterEach, beforeEach, expect, it, vi } from "vitest";
+import {
+ createPresencePublisher,
+ type PresencePublisherLock,
+} from "./publisher";
+import type { ActivitySnapshot, PresenceActivity } from "./activity";
+
+beforeEach(() => {
+ vi.useFakeTimers();
+ vi.setSystemTime(1000000);
+});
+afterEach(() => vi.useRealTimers());
+function activity() {
+ let value: ActivitySnapshot = { status: "online", visible: true };
+ const listeners = new Set<() => void>();
+ return {
+ source: {
+ snapshot: () => value,
+ subscribe(listener: () => void) {
+ listeners.add(listener);
+ return () => {
+ listeners.delete(listener);
+ };
+ },
+ } satisfies PresenceActivity,
+ update(next: ActivitySnapshot) {
+ value = next;
+ for (const listener of listeners) listener();
+ },
+ };
+}
+it("publishes initial/changed availability and slow renewals, never on visibility alone", async () => {
+ const a = activity();
+ const publish = vi.fn(
+ async (_status: "online" | "away", _signal: AbortSignal) => {},
+ );
+ const owner = createPresencePublisher({
+ activity: a.source,
+ publish,
+ random: () => 0,
+ });
+ owner.connection(true);
+ await vi.advanceTimersByTimeAsync(250);
+ expect(publish).toHaveBeenCalledTimes(1);
+ for (let i = 0; i < 100; i++)
+ a.update({ status: "online", visible: i % 2 === 0 });
+ await vi.advanceTimersByTimeAsync(59000);
+ expect(publish).toHaveBeenCalledTimes(1);
+ await vi.advanceTimersByTimeAsync(1000);
+ expect(publish).toHaveBeenCalledTimes(2);
+ a.update({ status: "away", visible: false });
+ await vi.advanceTimersByTimeAsync(1000);
+ expect(publish).toHaveBeenCalledTimes(3);
+ expect(publish.mock.calls[2]?.[0]).toBe("away");
+ owner.dispose();
+ await vi.advanceTimersByTimeAsync(180000);
+ expect(publish).toHaveBeenCalledTimes(3);
+});
+it("coalesces changes behind one in-flight publish and never replays missed ticks", async () => {
+ const a = activity();
+ const sends: { resolve(): void; signal: AbortSignal }[] = [];
+ const publish = vi.fn(
+ (_status: "online" | "away", signal: AbortSignal) =>
+ new Promise((resolve) => sends.push({ resolve, signal })),
+ );
+ const owner = createPresencePublisher({
+ activity: a.source,
+ publish,
+ random: () => 0,
+ });
+ owner.connection(true);
+ await vi.advanceTimersByTimeAsync(250);
+ a.update({ status: "away", visible: true });
+ await vi.advanceTimersByTimeAsync(600000);
+ expect(publish).toHaveBeenCalledTimes(1);
+ sends[0]?.resolve();
+ await vi.advanceTimersByTimeAsync(1000);
+ expect(publish).toHaveBeenCalledTimes(2);
+ expect(publish.mock.calls[1]?.[0]).toBe("away");
+ sends[1]?.resolve();
+ await vi.advanceTimersByTimeAsync(0);
+ vi.setSystemTime(Date.now() + 3600000);
+ await vi.advanceTimersByTimeAsync(60000);
+ expect(publish).toHaveBeenCalledTimes(3);
+ owner.dispose();
+});
+it("disconnect aborts old work and a late result cannot schedule work for its replacement", async () => {
+ const a = activity();
+ const sends: { resolve(): void; signal: AbortSignal }[] = [];
+ const publish = vi.fn(
+ (_status: "online" | "away", signal: AbortSignal) =>
+ new Promise((resolve) => sends.push({ resolve, signal })),
+ );
+ const owner = createPresencePublisher({
+ activity: a.source,
+ publish,
+ random: () => 0,
+ });
+ owner.connection(true);
+ await vi.advanceTimersByTimeAsync(250);
+ owner.connection(false);
+ expect(sends[0]?.signal.aborted).toBe(true);
+ owner.connection(true);
+ await vi.advanceTimersByTimeAsync(1000);
+ expect(publish).toHaveBeenCalledTimes(2);
+ sends[0]?.resolve();
+ await vi.advanceTimersByTimeAsync(0);
+ expect(owner.diagnostics().running).toBe(true);
+ sends[1]?.resolve();
+ await vi.advanceTimersByTimeAsync(0);
+ expect(owner.diagnostics()).toMatchObject({ running: false, accepted: 1 });
+ owner.dispose();
+});
+it("failure waits for the next renewal, rather than replaying a backlog", async () => {
+ const a = activity();
+ const publish = vi.fn(async () => {
+ throw new Error("rate limited");
+ });
+ const owner = createPresencePublisher({
+ activity: a.source,
+ publish,
+ random: () => 0,
+ });
+ owner.connection(true);
+ await vi.advanceTimersByTimeAsync(250);
+ expect(owner.diagnostics()).toMatchObject({ failures: 1, accepted: 0 });
+ await vi.advanceTimersByTimeAsync(59000);
+ expect(publish).toHaveBeenCalledTimes(1);
+ owner.dispose();
+});
+it("one scoped lifetime lock prevents duplicate publishers and passes ownership after disposal", async () => {
+ const waiting: (() => void)[] = [];
+ let locked = false;
+ const lock: PresencePublisherLock = async (signal, work) => {
+ if (locked) await new Promise((resolve) => waiting.push(resolve));
+ if (signal.aborted) return;
+ locked = true;
+ try {
+ await work();
+ } finally {
+ locked = false;
+ waiting.shift()?.();
+ }
+ };
+ const a = activity();
+ const publish = vi.fn(
+ async (_status: "online" | "away", _signal: AbortSignal) => {},
+ );
+ const first = createPresencePublisher({
+ activity: a.source,
+ publish,
+ lock,
+ random: () => 0,
+ });
+ const second = createPresencePublisher({
+ activity: a.source,
+ publish,
+ lock,
+ random: () => 0,
+ });
+ first.connection(true);
+ second.connection(true);
+ await vi.advanceTimersByTimeAsync(250);
+ expect(publish).toHaveBeenCalledTimes(1);
+ first.dispose();
+ await vi.advanceTimersByTimeAsync(250);
+ expect(publish).toHaveBeenCalledTimes(2);
+ second.dispose();
+});
+it("failed lock acquisition does not silently switch to uncoordinated publishing", async () => {
+ const a = activity();
+ const publish = vi.fn(
+ async (_status: "online" | "away", _signal: AbortSignal) => {},
+ );
+ const owner = createPresencePublisher({
+ activity: a.source,
+ publish,
+ lock: async () => {
+ throw new Error("lock unavailable");
+ },
+ });
+ owner.connection(true);
+ await vi.advanceTimersByTimeAsync(60000);
+ expect(publish).not.toHaveBeenCalled();
+ expect(owner.diagnostics().error).toBe("lock unavailable");
+ owner.dispose();
+});
diff --git a/src/features/presence/publisher.ts b/src/features/presence/publisher.ts
new file mode 100644
index 00000000..8225aa82
--- /dev/null
+++ b/src/features/presence/publisher.ts
@@ -0,0 +1,151 @@
+import type { PresenceActivity } from "./activity";
+export type PresencePublisherLock = (
+ signal: AbortSignal,
+ work: () => Promise,
+) => Promise;
+export function browserPresencePublisherLock(
+ scope: string,
+): PresencePublisherLock | undefined {
+ if (typeof navigator === "undefined" || !navigator.locks) return;
+ return (signal, work) =>
+ navigator.locks.request(
+ `buzz-presence:${scope}`,
+ { mode: "exclusive", signal },
+ work,
+ );
+}
+/** A renewable, lossy signal; exactly one in-flight operation, never a durable queue. */
+export function createPresencePublisher({
+ activity,
+ publish,
+ lock,
+ random = Math.random,
+}: {
+ activity: PresenceActivity;
+ publish(status: "online" | "away", signal: AbortSignal): Promise;
+ lock?: PresencePublisherLock | undefined;
+ random?: () => number;
+}) {
+ let closed = false;
+ let connected = false;
+ let owner: AbortController | undefined;
+ let leader = false;
+ let running = false;
+ let timer: ReturnType | undefined;
+ let desired = activity.snapshot().status;
+ let lastStart = -Infinity;
+ let due = 0;
+ const counters = { attempts: 0, accepted: 0, failures: 0 };
+ let error: string | undefined;
+ function schedule() {
+ if (closed || !connected || !leader || running || timer) return;
+ timer = setTimeout(
+ () => {
+ timer = undefined;
+ void send();
+ },
+ Math.max(0, due - Date.now(), lastStart + 1000 - Date.now()),
+ );
+ }
+ async function send() {
+ const current = owner;
+ if (
+ !current ||
+ current.signal.aborted ||
+ !leader ||
+ !connected ||
+ closed ||
+ running
+ )
+ return;
+ const status = desired;
+ lastStart = Date.now();
+ running = true;
+ counters.attempts++;
+ try {
+ await publish(
+ status,
+ AbortSignal.any([current.signal, AbortSignal.timeout(10000)]),
+ );
+ if (owner !== current || current.signal.aborted) return;
+ counters.accepted++;
+ error = undefined;
+ } catch (failure) {
+ if (owner !== current || current.signal.aborted) return;
+ counters.failures++;
+ error =
+ failure instanceof Error
+ ? failure.message
+ : "Presence publication unavailable";
+ } finally {
+ // A retired generation must not control a newer publisher's in-flight state.
+ if (owner === current) {
+ running = false;
+ due =
+ Date.now() + (status === desired ? 60000 + random() * 5000 : 1000);
+ schedule();
+ }
+ }
+ }
+ function stop() {
+ owner?.abort();
+ owner = undefined;
+ leader = running = false;
+ clearTimeout(timer);
+ timer = undefined;
+ }
+ function start() {
+ if (closed || !connected || owner) return;
+ const owned = new AbortController();
+ owner = owned;
+ const work = async () => {
+ if (closed || owned.signal.aborted || owner !== owned) return;
+ leader = true;
+ desired = activity.snapshot().status;
+ due = Date.now() + 250 + random() * 750; // Startup/reconnect yields and spreads scopes.
+ schedule();
+ await new Promise((resolve) =>
+ owned.signal.addEventListener("abort", () => resolve(), { once: true }),
+ );
+ };
+ void (lock ? lock(owned.signal, work) : work()).catch(
+ (failure: unknown) => {
+ if (owned.signal.aborted || owner !== owned) return;
+ error =
+ failure instanceof Error
+ ? failure.message
+ : "Presence ownership unavailable";
+ stop(); // No unlocked fallback that would silently create duplicate publishers.
+ },
+ );
+ }
+ const unsubscribe = activity.subscribe(() => {
+ const next = activity.snapshot().status;
+ if (desired === next) return;
+ desired = next;
+ due = Date.now();
+ clearTimeout(timer);
+ timer = undefined;
+ schedule();
+ });
+ return {
+ connection(active: boolean) {
+ if (closed || connected === active) return;
+ connected = active;
+ if (active) start();
+ else stop();
+ },
+ dispose() {
+ closed = true;
+ stop();
+ unsubscribe();
+ },
+ diagnostics: () => ({
+ ...counters,
+ leader,
+ coordinated: !!lock,
+ running,
+ error,
+ }),
+ };
+}
diff --git a/src/features/presence/react.tsx b/src/features/presence/react.tsx
new file mode 100644
index 00000000..7ef71603
--- /dev/null
+++ b/src/features/presence/react.tsx
@@ -0,0 +1,147 @@
+import {
+ useCallback,
+ useEffect,
+ useSyncExternalStore,
+ type RefObject,
+} from "react";
+import type { PresenceQueries } from "./directory";
+import styles from "./Presence.module.css";
+
+/** A selector, not a demand owner. Surfaces batch authors once for all their avatars. */
+export function PresenceIndicator({
+ presence,
+ author,
+ label = false,
+}: {
+ presence: PresenceQueries;
+ author: string;
+ label?: boolean;
+}) {
+ const subscribe = useCallback(
+ (listener: () => void) => presence.subscribe(author, listener),
+ [presence, author],
+ );
+ const snapshot = useCallback(() => presence.get(author), [presence, author]);
+ const status = useSyncExternalStore(subscribe, snapshot, snapshot);
+ const text =
+ status === "unknown"
+ ? "Presence unknown"
+ : status === "away"
+ ? "Away"
+ : status === "online"
+ ? "Online"
+ : "Offline";
+ return (
+
+
+ {label && {text} }
+
+ );
+}
+
+/** One demand handle per rendered surface. Intersection work never runs per heartbeat. */
+export function usePresenceSurface(
+ presence: PresenceQueries,
+ root: RefObject,
+) {
+ useEffect(() => {
+ const element = root.current;
+ if (!element) return;
+ const demand = presence.demand();
+ const visible = new Set();
+ const watched = new Set();
+ let timer: ReturnType | undefined;
+ let closed = false;
+ const flush = () => {
+ timer = undefined;
+ if (closed) return;
+ const authors =
+ document.visibilityState === "hidden"
+ ? []
+ : [...visible].flatMap((row) => {
+ const author = row.getAttribute("data-presence-author");
+ return author && element.contains(row) ? [author] : [];
+ });
+ const accepted = demand.update(authors);
+ if (accepted) element.removeAttribute("data-presence-limited");
+ else element.setAttribute("data-presence-limited", "true");
+ };
+ const schedule = () => {
+ if (!closed && !timer) timer = setTimeout(flush, 100);
+ };
+ const observer =
+ typeof IntersectionObserver === "undefined"
+ ? undefined
+ : new IntersectionObserver(
+ (changes) => {
+ for (const change of changes) {
+ if (change.isIntersecting) visible.add(change.target);
+ else visible.delete(change.target);
+ }
+ schedule();
+ },
+ { root: element, rootMargin: "160px" },
+ );
+ const scan = () => {
+ const rows = new Set(element.querySelectorAll("[data-presence-author]"));
+ for (const row of watched)
+ if (!rows.has(row)) {
+ watched.delete(row);
+ visible.delete(row);
+ observer?.unobserve(row);
+ }
+ for (const row of rows)
+ if (!watched.has(row)) {
+ watched.add(row);
+ if (observer) observer.observe(row);
+ else visible.add(row); // Bounded rendered-window fallback, never complete history.
+ }
+ schedule();
+ };
+ const mutation = new MutationObserver((changes) => {
+ if (
+ changes.some((change) =>
+ [...change.addedNodes, ...change.removedNodes].some(
+ (node) =>
+ node instanceof Element &&
+ (node.matches("[data-presence-author]") ||
+ node.querySelector("[data-presence-author]")),
+ ),
+ )
+ )
+ scan();
+ });
+ mutation.observe(element, { childList: true, subtree: true });
+ document.addEventListener("visibilitychange", schedule);
+ scan();
+ return () => {
+ closed = true;
+ clearTimeout(timer);
+ observer?.disconnect();
+ mutation.disconnect();
+ document.removeEventListener("visibilitychange", schedule);
+ demand.dispose();
+ element.removeAttribute("data-presence-limited");
+ };
+ }, [presence, root]);
+}
+
+export function usePresenceDemand(presence: PresenceQueries, author: string) {
+ useEffect(() => {
+ const demand = presence.demand();
+ const update = () =>
+ demand.update(document.visibilityState === "hidden" ? [] : [author]);
+ update();
+ document.addEventListener("visibilitychange", update);
+ return () => {
+ document.removeEventListener("visibilitychange", update);
+ demand.dispose();
+ };
+ }, [presence, author]);
+}
diff --git a/src/features/relay/presence-session.test.ts b/src/features/relay/presence-session.test.ts
new file mode 100644
index 00000000..acd382a0
--- /dev/null
+++ b/src/features/relay/presence-session.test.ts
@@ -0,0 +1,234 @@
+import { afterEach, beforeEach, expect, it, vi } from "vitest";
+import { createRelaySession } from "./session";
+import type { LiveCallbacks } from "./live";
+import { keypair, scriptedTransport, signed } from "./testing";
+import type { ActivitySnapshot, PresenceActivity } from "../presence/activity";
+
+beforeEach(() => {
+ vi.useFakeTimers();
+ vi.setSystemTime(1000000);
+});
+afterEach(() => vi.useRealTimers());
+function fixture(
+ options: { presence?: boolean | undefined; live?: boolean } = {
+ presence: true,
+ },
+) {
+ const relay = keypair(),
+ viewer = keypair(),
+ author = keypair();
+ const wire = scriptedTransport(viewer.pubkey, relay.pubkey);
+ let callbacks!: LiveCallbacks;
+ const update = vi.fn();
+ const observe = vi.fn();
+ const publish = vi.fn(
+ async (_status: "online" | "away", _signal: AbortSignal) => {},
+ );
+ let activity: ActivitySnapshot = { status: "online", visible: true };
+ const listeners = new Set<() => void>();
+ const source: PresenceActivity = {
+ snapshot: () => activity,
+ subscribe(listener) {
+ listeners.add(listener);
+ return () => {
+ listeners.delete(listener);
+ };
+ },
+ };
+ const owner = createRelaySession(
+ {
+ ...wire.transport,
+ agentActivity: true,
+ ...(options.presence === undefined ? {} : { presence: options.presence }),
+ ...(options.live === false
+ ? {}
+ : {
+ subscribe(value: LiveCallbacks) {
+ callbacks = value;
+ return {
+ update() {},
+ retry() {},
+ dispose() {},
+ observe,
+ presence: { update, publish },
+ };
+ },
+ }),
+ },
+ { presenceActivity: source },
+ );
+ const mount = () => {
+ const demand = owner.session.presence.demand();
+ callbacks.state({ status: "connected", routes: [] });
+ demand.update([author.pubkey]);
+ callbacks.presenceState?.({ status: "ready", authors: [author.pubkey] });
+ return demand;
+ };
+ return {
+ owner,
+ wire,
+ author,
+ relay,
+ callbacks,
+ update,
+ observe,
+ publish,
+ mount,
+ activity(next: ActivitySnapshot) {
+ activity = next;
+ for (const listener of listeners) listener();
+ },
+ snapshot: () =>
+ signed(relay, {
+ kind: 20001,
+ content: "online",
+ tags: [["p", author.pubkey]],
+ }),
+ };
+}
+it("production session routes presence outside retained projections and starts no message repair", async () => {
+ const f = fixture();
+ const demand = f.mount();
+ await vi.advanceTimersByTimeAsync(100);
+ const query = f.wire.next();
+ expect(query.filters).toEqual([
+ { kinds: [20001], authors: [f.author.pubkey], limit: 1 },
+ ]);
+ query.respond([f.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(f.owner.session.presence.get(f.author.pubkey)).toBe("online");
+ const view = f.owner.session.observe([{ kinds: [20001], limit: 10 }]);
+ expect(view.snapshot().events).toEqual([]);
+ f.callbacks.presence?.([
+ signed(f.author, { kind: 20001, content: "online", tags: [] }),
+ ]);
+ expect(view.snapshot().events).toEqual([]);
+ // Wrongly routed ephemeral frames must also be excluded, not stored by generic acceptance.
+ f.callbacks.receive([
+ signed(f.author, {
+ kind: 20001,
+ content: "online",
+ tags: [],
+ created_at: 2,
+ }),
+ ]);
+ expect(view.snapshot().events).toEqual([]);
+ expect(f.wire.pending).toEqual([]);
+ expect(f.owner.session.outbox).toBeUndefined();
+ demand.dispose();
+ view.dispose();
+ f.owner.dispose();
+});
+it("host activity suspends observers but keeps one availability publisher; disposal aborts both", async () => {
+ const f = fixture();
+ f.mount();
+ await vi.advanceTimersByTimeAsync(1100);
+ const request = f.wire.next();
+ expect(f.publish).toHaveBeenCalledTimes(1);
+ f.activity({ status: "online", visible: false });
+ expect(f.update).toHaveBeenLastCalledWith([]);
+ expect(request.signal?.aborted).toBe(true);
+ expect(f.owner.session.presence.get(f.author.pubkey)).toBe("unknown");
+ await vi.advanceTimersByTimeAsync(66000);
+ expect(f.publish).toHaveBeenCalledTimes(2);
+ expect(f.wire.pending).toHaveLength(0);
+ f.owner.dispose();
+ await vi.advanceTimersByTimeAsync(180000);
+ expect(f.publish).toHaveBeenCalledTimes(2);
+});
+it("session access revocation clears prior presence and rejects in-flight snapshot authority", async () => {
+ const f = fixture();
+ f.mount();
+ await vi.advanceTimersByTimeAsync(100);
+ const first = f.wire.next();
+ first.respond([f.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(f.owner.session.presence.get(f.author.pubkey)).toBe("online");
+ // Actual signed roster removes viewer; the session owns access invalidation.
+ const { roster } = await import("./testing");
+ const member = f.owner.session as typeof f.owner.session;
+ f.callbacks.receive([roster(f.relay, "a", [f.wire.transport.viewer])]);
+ f.callbacks.receive([roster(f.relay, "a", [], 1700000001)]);
+ expect(member.presence.get(f.author.pubkey)).toBe("unknown");
+ f.owner.dispose();
+});
+
+it("merged session keeps observer telemetry and presence independently owned across clear and disposal", async () => {
+ const f = fixture();
+ const demand = f.mount();
+ const release = f.owner.session.agentActivity.activate();
+ const generation = f.observe.mock.lastCall?.[0];
+ expect(generation).toBeTypeOf("number");
+ f.callbacks.state({
+ status: "connected",
+ routes: [{ id: "observer", status: "live", replay: "unknown" }],
+ });
+ expect(f.owner.session.agentActivity.snapshot().status).toBe("listening");
+ await vi.advanceTimersByTimeAsync(1100);
+ f.wire.next().respond([f.snapshot()]);
+ await vi.advanceTimersByTimeAsync(0);
+ expect(f.owner.session.presence.get(f.author.pubkey)).toBe("online");
+ const frame = {
+ id: "f".repeat(64),
+ agent: f.author.pubkey,
+ createdAt: Math.floor(Date.now() / 1000),
+ plaintext: JSON.stringify({
+ kind: "turn_started",
+ turnId: "one",
+ channelId: null,
+ timestamp: new Date().toISOString(),
+ }),
+ };
+ f.callbacks.observer?.(frame, generation);
+ expect(f.owner.session.agentActivity.snapshot().turns[0]?.state).toBe(
+ "working",
+ );
+ const view = f.owner.session.observe([{ kinds: [20001, 24200], limit: 10 }]);
+ expect(view.snapshot().events).toEqual([]);
+ demand.dispose();
+ expect(f.update).toHaveBeenLastCalledWith([]);
+ expect(f.owner.session.agentActivity.snapshot().records).toHaveLength(1);
+ f.callbacks.state({ status: "retrying", routes: [] });
+ expect(f.owner.session.agentActivity.snapshot().turns[0]?.state).toBe(
+ "unknown",
+ );
+ await f.owner.clearCache();
+ f.callbacks.observer?.(frame, generation);
+ expect(f.owner.session.agentActivity.snapshot().records).toHaveLength(0);
+ release();
+ expect(f.observe).toHaveBeenLastCalledWith(null);
+ const publications = f.publish.mock.calls.length;
+ f.owner.dispose();
+ f.callbacks.observer?.({ ...frame, id: "e".repeat(64) }, generation);
+ await vi.advanceTimersByTimeAsync(180000);
+ expect(f.publish).toHaveBeenCalledTimes(publications);
+ expect(f.owner.session.agentActivity.snapshot().status).toBe("unavailable");
+ expect(f.owner.session.agentActivity.snapshot().records).toHaveLength(0);
+ expect(view.snapshot().events).toEqual([]);
+ view.dispose();
+});
+
+it.each([{ presence: false }, {}, { presence: true, live: false }])(
+ "unsupported transport starts no presence reads, controls, publisher or renewals (%j)",
+ async (options) => {
+ const f = fixture(options);
+ try {
+ if (options.live !== false) f.mount();
+ else f.owner.session.presence.demand().update([f.author.pubkey]);
+ f.activity({ status: "away", visible: false });
+ f.activity({ status: "online", visible: true });
+ f.callbacks?.presenceState?.({
+ status: "ready",
+ authors: [f.author.pubkey],
+ });
+ await vi.advanceTimersByTimeAsync(180000);
+ expect(f.wire.pending).toEqual([]);
+ expect(f.update).not.toHaveBeenCalled();
+ expect(f.publish).not.toHaveBeenCalled();
+ expect(f.owner.diagnostics().presence.publisher).toBeUndefined();
+ expect(f.owner.session.presence.get(f.author.pubkey)).toBe("unknown");
+ } finally {
+ f.owner.dispose();
+ }
+ },
+);
diff --git a/src/features/relay/service.ts b/src/features/relay/service.ts
index 02bde1dd..f3888f63 100644
--- a/src/features/relay/service.ts
+++ b/src/features/relay/service.ts
@@ -1,3 +1,4 @@
+import type { PresenceActivity } from "../presence/activity";
import type { Context } from "@deepseek-ai/cordis";
import { createRelaySession, type RelaySession } from "./session";
import { createHeadPersistence } from "./persistence";
@@ -29,6 +30,7 @@ declare module "@deepseek-ai/cordis" {
export function provideRelay(
ctx: Context,
connect?: (signal: AbortSignal) => Promise,
+ presenceActivity?: PresenceActivity,
) {
let disposed = false;
let generation = 0;
@@ -86,6 +88,7 @@ export function provideRelay(
store.dispose();
store = createRelaySession(transport, {
prepared: true,
+ ...(presenceActivity ? { presenceActivity } : {}),
persistence: createHeadPersistence(
transport.viewer,
transport.scope ?? transport.relayAuthor,
diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts
index 3b5831f4..fed71c09 100644
--- a/src/features/relay/session.ts
+++ b/src/features/relay/session.ts
@@ -4,6 +4,13 @@ import {
type ReadOptions,
type RelayReader,
} from "./reader";
+import { createPresenceDirectory } from "../presence/directory";
+import {
+ createPresencePublisher,
+ browserPresencePublisherLock,
+ type PresencePublisherLock,
+} from "../presence/publisher";
+import type { PresenceActivity } from "../presence/activity";
import { createAgentActivity } from "../agents/activity";
import { OBSERVER_KIND } from "../agents/observer";
import { createAgentLibrary } from "../agents/library";
@@ -65,6 +72,8 @@ export function createRelaySession(
readStateStorage?: ReadStateStorage;
readPublisherLock?: ReadPublisherLock;
deliveryTimeoutMs?: number;
+ presenceActivity?: PresenceActivity;
+ presencePublisherLock?: PresencePublisherLock;
} = {},
) {
let closed = false;
@@ -72,6 +81,39 @@ export function createRelaySession(
const profiling =
options.profiling ?? transport?.profiling ?? createRelayProfiler();
const requests = createRelayReader(transport, { profiling });
+ let traffic: LiveSubscription | undefined;
+ const presenceSupported =
+ transport?.presence === true && !!transport.subscribe;
+ const presence = createPresenceDirectory({
+ reader: requests.reader,
+ relayAuthor: transport?.relayAuthor ?? "",
+ supported: presenceSupported,
+ updateInterests: (authors) => traffic?.presence?.update(authors),
+ notify: (listener) => notify(listener),
+ });
+ const presencePublisher =
+ options.presenceActivity && transport && presenceSupported
+ ? createPresencePublisher({
+ activity: options.presenceActivity,
+ publish: (status, signal) => {
+ if (!traffic?.presence)
+ return Promise.reject(
+ new Error("Presence publication unsupported"),
+ );
+ return traffic.presence.publish(status, signal);
+ },
+ lock:
+ options.presencePublisherLock ??
+ browserPresencePublisherLock(
+ `${transport.scope ?? transport.relayAuthor}:${transport.viewer}`,
+ ),
+ })
+ : undefined;
+ const presenceActivity = () =>
+ presence.visibility(options.presenceActivity?.snapshot().visible ?? true);
+ const stopPresenceActivity =
+ options.presenceActivity?.subscribe(presenceActivity);
+ presenceActivity();
let revision = 0;
let accessEpoch = 0;
let cacheClearEpoch = 0;
@@ -172,6 +214,7 @@ export function createRelaySession(
for (const id of revoked) recent.delete(id);
channels.purgeAccess((events) => events.filter(visibility(events)));
writes?.purgeConfirmed((event) => event.kind !== 0 && visible(event));
+ presence.clear();
profiles.clear();
emoji.clear();
agentLibrary.clear();
@@ -352,7 +395,6 @@ export function createRelaySession(
viewer: transport?.viewer ?? "",
notify,
});
- let traffic: LiveSubscription | undefined;
const liveListeners = new Set<() => void>();
let liveSnapshot: LiveSnapshot = Object.freeze({
status: transport?.subscribe ? "connecting" : "unavailable",
@@ -590,6 +632,7 @@ export function createRelaySession(
};
},
unread: unread.capability,
+ presence: presence.queries,
sidebarPreferences: sidebarPreferences.queries,
live,
profiling,
@@ -928,9 +971,14 @@ export function createRelaySession(
}
}
traffic = transport?.subscribe?.({
+ presence: (events) => presence.receive(events),
+ presenceState: (state) => presence.route(state),
observer: (frame, generation) => activity.receive(frame, generation),
receive(events, provenance) {
if (closed) return;
+ // Ephemeral traffic must never enter generic retention, even from a bad route.
+ presence.receive(events.filter((event) => event.kind === 20001));
+ events = events.filter((event) => event.kind !== 20001);
const candidates = new Set(
provenance?.phase === "live" && provenance.channelId
? events
@@ -995,6 +1043,8 @@ export function createRelaySession(
},
state(snapshot) {
if (closed) return;
+ presence.connection(snapshot.status === "connected");
+ presencePublisher?.connection(snapshot.status === "connected");
activity.state(snapshot);
if (
snapshot.status !== "connected" &&
@@ -1060,6 +1110,7 @@ export function createRelaySession(
cacheClearEpoch++;
activity.clear();
sidebarPreferences.clear();
+ presence.clear();
// New windows must not yield to or receive errors from retired owners.
catchups.clear();
catchupQueue.clear();
@@ -1077,6 +1128,9 @@ export function createRelaySession(
dispose() {
closed = true;
lifetime.abort();
+ stopPresenceActivity?.();
+ presencePublisher?.dispose();
+ presence.dispose();
activity.dispose();
sidebarPreferences.dispose();
stopInterests();
@@ -1099,6 +1153,10 @@ export function createRelaySession(
diagnostics: () => ({
...channels.diagnostics(),
profiles: profiles.stats(),
+ presence: {
+ ...presence.diagnostics(),
+ publisher: presencePublisher?.diagnostics(),
+ },
}),
};
}
diff --git a/tests/browser/broker-evidence.mjs b/tests/browser/broker-evidence.mjs
index 33c8f341..6af2048b 100644
--- a/tests/browser/broker-evidence.mjs
+++ b/tests/browser/broker-evidence.mjs
@@ -4,6 +4,7 @@ import assert from "node:assert/strict";
// is delayed, substituted or interpreted as delivery by this observer.
export function brokerEvidence(report, retiredStreams) {
const publications = [];
+ const retirements = new Map();
report.presencePublicationResponses = publications;
return {
middleware(req, res, next) {
@@ -16,6 +17,7 @@ export function brokerEvidence(report, retiredStreams) {
report.brokerRequests.push(record);
const url = new URL(req.url, `http://${req.headers.host}`);
const publishing = url.pathname.endsWith("/stream-presence-publish");
+ const retiredAtArrival = publishing ? new Map(retirements) : undefined;
const streaming = url.pathname.endsWith("/stream");
let streamId;
if (streaming) {
@@ -35,12 +37,12 @@ export function brokerEvidence(report, retiredStreams) {
try {
record.streamId = JSON.parse(body).streamId;
} catch {
- // Invalid/unreadable request evidence cannot classify a 503.
+ // Invalid/unreadable request evidence cannot classify an error.
}
});
const end = res.end;
res.end = function (chunk, ...args) {
- if (this.statusCode === 503) {
+ if (this.statusCode === 404 || this.statusCode === 503) {
let value;
try {
value = JSON.parse(String(chunk));
@@ -51,13 +53,26 @@ export function brokerEvidence(report, retiredStreams) {
url: url.href,
streamId: record.streamId,
at: performance.now(),
- status: 503,
+ status: this.statusCode,
body: value ?? null,
disposed:
+ this.statusCode === 503 &&
/^[0-9a-f]{32}$/.test(record.streamId ?? "") &&
value?.error === "Presence publication unconfirmed" &&
value?.code === "presence_owner_disposed" &&
Object.keys(value).length === 2,
+ // Retirement must precede arrival, not merely the eventual response
+ // or fixture teardown. A different relay's stream cannot classify.
+ retired:
+ this.statusCode === 404 &&
+ /^[0-9a-f]{32}$/.test(record.streamId ?? "") &&
+ retiredAtArrival.get(record.streamId) ===
+ url.pathname.replace(
+ /\/stream-presence-publish$/,
+ "/stream",
+ ) &&
+ value?.error === "Live stream no longer available" &&
+ Object.keys(value).length === 1,
});
}
return end.call(this, chunk, ...args);
@@ -74,37 +89,48 @@ export function brokerEvidence(report, retiredStreams) {
});
res.once("close", () => {
record.close = snapshot();
- if (streaming && streamId) retiredStreams.add(streamId);
- // A destroyed/truncated 503 without end() is still an unclassified failure.
- if (publishing && res.statusCode === 503 && !res.writableEnded)
+ if (streaming && /^[0-9a-f]{32}$/.test(streamId ?? "")) {
+ retiredStreams.add(streamId);
+ retirements.set(streamId, url.pathname);
+ }
+ // A destroyed/truncated error without end() remains unclassified.
+ if (
+ publishing &&
+ [404, 503].includes(res.statusCode) &&
+ !res.writableEnded
+ )
publications.push({
url: url.href,
streamId: record.streamId,
- status: 503,
+ status: res.statusCode,
disposed: false,
+ retired: false,
});
});
next();
},
assertPublications() {
assert.deepEqual(
- publications.filter((item) => !item.disposed),
+ publications.filter((item) => !item.disposed && !item.retired),
[],
- "Unclassified presence publication 503",
+ "Unclassified presence publication 404/503",
);
},
consoleFilter() {
// Each classified response permits at most one endpoint-qualified console
// diagnostic. Independent publication validation prevents same-URL masking.
- const remaining = publications.filter((item) => item.disposed);
+ const remaining = publications.filter(
+ (item) => item.disposed || item.retired,
+ );
return (message, location) => {
- if (
- !/^Failed to load resource: the server responded with a status of 503/.test(
+ const match =
+ /^Failed to load resource: the server responded with a status of (404|503)\b/.exec(
message,
- )
- )
- return false;
- const index = remaining.findIndex((item) => item.url === location);
+ );
+ if (!match) return false;
+ const index = remaining.findIndex(
+ (item) => item.url === location && item.status === Number(match[1]),
+ );
if (index < 0) return false;
remaining.splice(index, 1);
return true;
diff --git a/tests/browser/build.mjs b/tests/browser/build.mjs
index 6a16e28e..5ebbf5d0 100644
--- a/tests/browser/build.mjs
+++ b/tests/browser/build.mjs
@@ -10,7 +10,10 @@ const root = fileURLToPath(new URL("../../", import.meta.url));
// Playwright owns this worker-scoped build. Only compiled assets are shared;
// each test still owns its server, identities, relay state and browser storage.
-export async function buildApp({ developmentReact, pluginFixtures }, use) {
+export async function buildApp(
+ { developmentReact, pluginFixtures, withoutPresence },
+ use,
+) {
const directory = await mkdtemp(join(tmpdir(), "buzz-browser-build-"));
try {
const config = {
@@ -20,6 +23,26 @@ export async function buildApp({ developmentReact, pluginFixtures }, use) {
logLevel: "error",
plugins: [
react(),
+ ...(withoutPresence
+ ? [
+ {
+ name: "fixture-no-presence-control",
+ transform(code, id) {
+ if (id !== join(root, "src/features/relay/session.ts"))
+ return;
+ // Equivalent app, connection and read paths. Disable both optional
+ // owners through their shared support decision, not admission substitutes.
+ const support =
+ "transport?.presence === true && !!transport.subscribe";
+ if (code.split(support).length !== 2)
+ throw new Error(
+ `Missing presence control seam: ${support}`,
+ );
+ return code.replace(support, "false");
+ },
+ },
+ ]
+ : []),
...(pluginFixtures
? [
{
diff --git a/tests/browser/live.spec.mjs b/tests/browser/live.spec.mjs
index 2108b5b4..dc507032 100644
--- a/tests/browser/live.spec.mjs
+++ b/tests/browser/live.spec.mjs
@@ -16,11 +16,27 @@ const heads = (app, channel) =>
filter.until === undefined,
);
async function ready(page, app) {
- await open(page, app);
- await expect.poll(() => app.relay.hasRoute("primary", "alpha")).toBe(true);
- // The first head may already start after stream establishment; a duplicate
- // initial read is not required. Recovery below has its own positive gap control.
- await expect.poll(() => heads(app, "alpha").length).toBeGreaterThanOrEqual(1);
+ // Force the pre-establishment read and its subsequent catch-up apart. Initial
+ // rows/EOSE alone do not prove that a queued startup head has reached the UI.
+ app.relay.holdEose("alpha");
+ try {
+ await open(page, app);
+ await expect.poll(() => app.relay.hasRoute("primary", "alpha")).toBe(true);
+ expect(heads(app, "alpha")).toHaveLength(1);
+ const missed = app.append(
+ "primary",
+ "alpha",
+ "Startup catch-up complete",
+ false,
+ );
+ app.relay.releaseEose("alpha");
+ await expect(
+ history(page).locator(`[data-message-id="${missed.id}"]`),
+ ).toBeVisible();
+ expect(heads(app, "alpha")).toHaveLength(2);
+ } finally {
+ app.relay.releaseEose("alpha");
+ }
await expect.poll(() => app.relay.hasRoute("primary", "beta")).toBe(true);
await expect(retry(page)).toHaveCount(0);
await settle(page);
diff --git a/tests/browser/playwright.config.mjs b/tests/browser/playwright.config.mjs
index e2b1268a..577913b0 100644
--- a/tests/browser/playwright.config.mjs
+++ b/tests/browser/playwright.config.mjs
@@ -1,6 +1,11 @@
import { defineConfig } from "@playwright/test";
-const measurementFiles = ["channel-opening.spec.mjs", "scroll.spec.mjs"];
+const measurementFiles = [
+ "channel-opening.spec.mjs",
+ "scroll.spec.mjs",
+ "presence-contention.spec.mjs",
+ "presence-control.spec.mjs",
+];
export default defineConfig({
testDir: ".",
diff --git a/tests/browser/presence-contention.mjs b/tests/browser/presence-contention.mjs
new file mode 100644
index 00000000..ad9f3b0f
--- /dev/null
+++ b/tests/browser/presence-contention.mjs
@@ -0,0 +1,300 @@
+import { readFile } from "node:fs/promises";
+import { test, expect } from "./fixture.mjs";
+import { open } from "./timeline.mjs";
+
+const snapshots = (app) =>
+ app.report.queries.filter(({ filter }) => filter.kinds?.includes(20001));
+const ordinaryStarts = (app) =>
+ app.report.queries.filter(({ filter }) => !filter.kinds?.includes(20001));
+const headRequest = (request) =>
+ new URL(request.url()).pathname.endsWith("/query") &&
+ request
+ .postDataJSON()
+ ?.some(
+ (f) =>
+ f.kinds?.includes(9) &&
+ f["#h"]?.[0] === "beta" &&
+ f.top_level &&
+ f.until === undefined,
+ );
+
+// Optimistic text alone is not delivery: only verified observation removes the
+// exact event from the outbox. Keep this same criterion in the race controls.
+export async function confirmedSend(page, event) {
+ await expect(page.locator(`[data-message-id="${event.id}"]`)).toContainText(
+ event.content,
+ );
+ await expect(
+ page.locator("summary").filter({ hasText: /^Outbox ·/ }),
+ ).toHaveText("Outbox · 0 items");
+}
+
+export function contentionTests(withoutPresence) {
+ test.describe(
+ withoutPresence ? "no-presence control" : "held presence",
+ () => {
+ for (const action of ["send", "open"]) {
+ test(`presence contention: ${action} does not inherit optional HTTP pacing`, async ({
+ page,
+ app,
+ }) => {
+ await open(page, app);
+ const composer = page.getByRole("textbox", {
+ name: "Message #Alpha",
+ exact: true,
+ });
+ if (action === "send")
+ await composer.fill("Foreground contention probe");
+ if (!withoutPresence) {
+ await expect(
+ page
+ .locator(
+ '[aria-label="Channel message history"] [data-presence-status="online"]',
+ )
+ .first(),
+ ).toBeVisible();
+ expect(snapshots(app).length).toBeGreaterThan(0);
+ } else {
+ expect(snapshots(app)).toHaveLength(0);
+ expect(app.report.presencePublications).toHaveLength(0);
+ }
+ // Compare the same idle ordinary lane, never a preloaded foreground lane
+ // against an idle control. Presence must be the only new pacing input.
+ await expect
+ .poll(() => performance.now() - ordinaryStarts(app).at(-1).at)
+ .toBeGreaterThan(550);
+ const before = snapshots(app).length;
+ let held;
+ if (!withoutPresence) {
+ const started = new Promise((resolve) =>
+ app.relay.holdPresence(resolve),
+ );
+ app.presence("away"); // real signed conflict -> directory -> reader -> broker
+ held = await started;
+ expect(snapshots(app)).toHaveLength(before + 1);
+ expect(held.pending).toBe(true);
+ }
+ const phaseAt = held?.at ?? performance.now();
+ const requestMatch =
+ action === "open"
+ ? headRequest
+ : (request) =>
+ new URL(request.url()).pathname.endsWith("/publish");
+ // A verified live echo legitimately aborts an outstanding send ACK.
+ // Reads still require the response; sends prove delivery independently.
+ const response =
+ action === "open"
+ ? page.waitForResponse((response) =>
+ requestMatch(response.request()),
+ )
+ : undefined;
+ let browserRequestObservedAt, browserRequestUrl;
+ const browserOutcomes = [];
+ const observe = (request) => {
+ if (requestMatch(request)) {
+ browserRequestObservedAt = performance.now();
+ browserRequestUrl = request.url();
+ }
+ };
+ const received = (reply) => {
+ if (requestMatch(reply.request()))
+ browserOutcomes.push({
+ status: reply.status(),
+ at: performance.now(),
+ });
+ };
+ const failed = (request) => {
+ if (requestMatch(request))
+ browserOutcomes.push({
+ error: request.failure()?.errorText,
+ at: performance.now(),
+ });
+ };
+ page.on("request", observe);
+ page.on("response", received);
+ page.on("requestfailed", failed);
+ const priorBroker = app.report.brokerRequests.length;
+ const priorPublications = app.report.publications.length;
+ let sentEvent;
+ try {
+ if (action === "send") {
+ await composer.evaluate((input) => {
+ performance.mark("presence-foreground-intent");
+ input.form.requestSubmit();
+ });
+ } else {
+ expect(
+ ordinaryStarts(app).filter(
+ ({ filter }) =>
+ filter.top_level && filter["#h"]?.[0] === "beta",
+ ),
+ ).toHaveLength(0);
+ await page
+ .getByRole("button", { name: "Beta", exact: true })
+ .evaluate((button) => {
+ performance.mark("presence-foreground-intent");
+ button.click();
+ });
+ }
+ const reply = await response;
+ if (reply) expect(reply.status()).toBe(200);
+ if (action === "send") {
+ await expect
+ .poll(() => app.report.publications.length)
+ .toBe(priorPublications + 1);
+ sentEvent = app.report.publications[priorPublications].event;
+ await confirmedSend(page, sentEvent);
+ }
+ const candidates = app.report.brokerRequests
+ .slice(priorBroker)
+ .filter(({ url }) =>
+ url.endsWith(action === "send" ? "/publish" : "/query"),
+ );
+ expect(
+ candidates,
+ "foreground host timing must be unambiguous",
+ ).toHaveLength(1);
+ const broker = candidates[0];
+ const upstream =
+ action === "send"
+ ? app.report.publications.at(-1)
+ : ordinaryStarts(app).find(
+ ({ filter, at }) =>
+ at >= phaseAt &&
+ filter.top_level &&
+ filter["#h"]?.[0] === "beta",
+ );
+ // Host completion survives browser cancellation, without claiming
+ // that a closed socket or an optimistic row means acceptance.
+ await expect.poll(() => broker.finish?.finished).toBe(true);
+ expect(broker.finish.status).toBe(200);
+ const serverTiming = broker.finish.serverTiming;
+ if (reply)
+ expect(await reply.headerValue("server-timing")).toBe(
+ serverTiming,
+ );
+ const admissionMs = Number(
+ serverTiming?.match(/(?:^|,\s*)admission;dur=([\d.]+)/)?.[1],
+ );
+ const browserTiming = await page.evaluate(
+ ({ url, received }) => {
+ const intent = performance
+ .getEntriesByName("presence-foreground-intent")
+ .at(-1).startTime;
+ const resource = performance.getEntriesByName(url).at(-1);
+ return {
+ intent,
+ requestStart: resource?.startTime,
+ responseEnd: resource?.responseEnd,
+ intentToRequestMs: resource
+ ? resource.startTime - intent
+ : null,
+ intentToResponseMs:
+ received && resource ? resource.responseEnd - intent : null,
+ };
+ },
+ {
+ url: browserRequestUrl,
+ received: browserOutcomes.some(
+ (outcome) => outcome.status === 200,
+ ),
+ },
+ );
+ app.report.measurements.push({
+ action,
+ withoutPresence,
+ phaseAt,
+ browserRequestObservedAt,
+ note: "Browser request event observed on runner clock; not a hardware input or renderer timestamp",
+ brokerAt: broker.at,
+ upstreamAt: upstream.at,
+ arrivalAfterSnapshotMs: broker.at - phaseAt,
+ brokerToUpstreamMs: upstream.at - broker.at,
+ admissionMs,
+ serverTiming,
+ browserTiming,
+ browserOutcomes,
+ hostResponse: { ...broker.finish },
+ ...(sentEvent ? { confirmedEventId: sentEvent.id } : {}),
+ held: held && { ...held },
+ });
+ // This phase check is crucial: arriving after 500 ms would make old code pass.
+ expect(broker.at - phaseAt).toBeGreaterThanOrEqual(0);
+ expect(broker.at - phaseAt).toBeLessThan(200);
+ expect(Number.isFinite(admissionMs)).toBe(true);
+ expect(admissionMs).toBeLessThan(100);
+ expect(upstream.at - broker.at).toBeLessThan(150);
+ if (held) {
+ expect(held.pending || held.aborted).toBe(true);
+ if (held.aborted)
+ expect(held.completedAt).toBeGreaterThanOrEqual(phaseAt);
+ // Sending must not depend on a snapshot completion/cancellation.
+ if (action === "send") expect(held.pending).toBe(true);
+ }
+ if (action === "send") {
+ await expect(
+ page.getByText("Foreground contention probe", { exact: true }),
+ ).toBeVisible();
+ await expect(composer).toHaveValue("");
+ } else {
+ await expect(
+ page.getByRole("textbox", {
+ name: "Message #Beta",
+ exact: true,
+ }),
+ ).toBeVisible();
+ await expect(
+ page.locator("[data-message-id]").first(),
+ ).toBeVisible();
+ }
+ expect(app.report.quotaCharges.length).toBeGreaterThan(0);
+ expect(
+ app.report.quotaCharges.every((charge) => charge.accepted),
+ ).toBe(true);
+ expect(app.report.quotaRefusals).toEqual([]);
+ // Existing production profiler records read.queue/fetch/verify and
+ // write stages separately. Export after the timed action, not during it.
+ await page
+ .locator("summary")
+ .filter({ hasText: /^Relay timings$/ })
+ .evaluate((element) => {
+ for (
+ let parent = element.parentElement;
+ parent;
+ parent = parent.parentElement
+ )
+ if (parent.tagName === "DETAILS") parent.open = true;
+ });
+ const download = page.waitForEvent("download");
+ await page.getByRole("button", { name: "Export timings" }).click();
+ app.report.relayTimings = JSON.parse(
+ await readFile(await (await download).path(), "utf8"),
+ );
+ if (sentEvent)
+ expect(app.report.relayTimings).toContainEqual(
+ expect.objectContaining({
+ id: sentEvent.id,
+ stage: "send.delivery",
+ outcome: "ok",
+ }),
+ );
+ if (withoutPresence) {
+ expect(snapshots(app)).toHaveLength(0);
+ expect(app.report.presencePublications).toHaveLength(0);
+ expect(
+ app.report.liveRequests.filter(({ filter }) =>
+ filter.kinds.includes(20001),
+ ),
+ ).toHaveLength(0);
+ }
+ } finally {
+ page.off("request", observe);
+ page.off("response", received);
+ page.off("requestfailed", failed);
+ app.relay.releasePresence();
+ }
+ });
+ }
+ },
+ );
+}
diff --git a/tests/browser/presence-contention.spec.mjs b/tests/browser/presence-contention.spec.mjs
new file mode 100644
index 00000000..aa22b900
--- /dev/null
+++ b/tests/browser/presence-contention.spec.mjs
@@ -0,0 +1,8 @@
+import { test } from "./fixture.mjs";
+import { contentionTests } from "./presence-contention.mjs";
+test.use({
+ productionBroker: true,
+ composerPublication: true,
+ enforceQuotas: true,
+});
+contentionTests(false);
diff --git a/tests/browser/presence-control.spec.mjs b/tests/browser/presence-control.spec.mjs
new file mode 100644
index 00000000..b328b8de
--- /dev/null
+++ b/tests/browser/presence-control.spec.mjs
@@ -0,0 +1,141 @@
+import { stripVTControlCharacters } from "node:util";
+import { test, expect } from "./fixture.mjs";
+import { contentionTests, confirmedSend } from "./presence-contention.mjs";
+import { open } from "./timeline.mjs";
+test.use({
+ productionBroker: true,
+ composerPublication: true,
+ enforceQuotas: true,
+ withoutPresence: true,
+});
+contentionTests(true);
+
+// These are settlement controls, not performance samples: deliberately hold the
+// browser ACK while retaining the real composer, broker and verified live path.
+test("verified echo completes a send before its browser HTTP acknowledgement", async ({
+ page,
+ app,
+}) => {
+ await page.addInitScript(() => {
+ const original = window.fetch;
+ window.__sendAbortEvidence = [];
+ window.fetch = function (input, init) {
+ if (String(input).endsWith("/publish")) {
+ const event = JSON.parse(init.body);
+ init.signal.addEventListener(
+ "abort",
+ () => {
+ window.__sendAbortEvidence.push({
+ id: event.id,
+ reason: init.signal.reason?.name,
+ });
+ },
+ { once: true },
+ );
+ }
+ return original.call(this, input, init);
+ };
+ });
+ await open(page, app);
+ let release;
+ const held = new Promise((resolve) => {
+ release = resolve;
+ });
+ let hostReceipt,
+ hostStatus,
+ browserResponses = 0;
+ const received = (response) => {
+ if (response.url().endsWith("/publish")) browserResponses++;
+ };
+ page.on("response", received);
+ await page.route("**/api/relay/primary/publish", async (route) => {
+ const response = await route.fetch();
+ hostStatus = response.status();
+ hostReceipt = await response.json();
+ await held;
+ // The browser may already have cancelled; fulfilling an obsolete route is
+ // cleanup only and cannot be used as evidence of delivery.
+ await route.fulfill({ response }).catch(() => {});
+ });
+ try {
+ const composer = page.getByRole("textbox", {
+ name: "Message #Alpha",
+ exact: true,
+ });
+ await composer.fill("Echo before ACK control");
+ await composer.evaluate((input) => input.form.requestSubmit());
+ await expect.poll(() => app.report.publications.length).toBe(1);
+ const event = app.report.publications[0].event;
+ await confirmedSend(page, event);
+ await expect
+ .poll(() => hostReceipt)
+ .toEqual({ accepted: true, event_id: event.id });
+ expect(hostStatus).toBe(200);
+ await expect
+ .poll(() => page.evaluate(() => window.__sendAbortEvidence))
+ .toEqual([{ id: event.id, reason: "AbortError" }]);
+ expect(browserResponses).toBe(0);
+ app.report.echoBeforeAck = {
+ eventId: event.id,
+ hostReceipt,
+ hostStatus,
+ browserResponses,
+ abort: await page.evaluate(() => window.__sendAbortEvidence),
+ };
+ } finally {
+ release();
+ await page.unrouteAll({ behavior: "wait" });
+ page.off("response", received);
+ }
+});
+
+test("send confirmation rejects an optimistic row without verified relay observation", async ({
+ page,
+ app,
+}) => {
+ await open(page, app);
+ let release, submitted;
+ const held = new Promise((resolve) => {
+ release = resolve;
+ });
+ await page.route("**/api/relay/primary/publish", async (route) => {
+ submitted = route.request().postDataJSON();
+ await held;
+ await route.abort("aborted").catch(() => {});
+ });
+ try {
+ const composer = page.getByRole("textbox", {
+ name: "Message #Alpha",
+ exact: true,
+ });
+ await composer.fill("Unconfirmed optimistic control");
+ await composer.evaluate((input) => input.form.requestSubmit());
+ await expect.poll(() => submitted?.id).toBeTruthy();
+ await expect(
+ page.locator(`[data-message-id="${submitted.id}"]`),
+ ).toContainText(submitted.content);
+ await expect(composer).toHaveValue("");
+ // The same assertion used by the measurement must reject this tempting false
+ // positive. The default bounded expect timeout is unchanged.
+ const failure = await confirmedSend(page, submitted).then(
+ () => undefined,
+ (error) => stripVTControlCharacters(error.message),
+ );
+ expect(failure).toContain('Expected: "Outbox · 0 items"');
+ expect(failure).toContain('Received: "Outbox · 1 items"');
+ const reconciliation = app.report.queries.filter(
+ ({ filter }) => filter.ids?.[0] === submitted.id,
+ );
+ expect(reconciliation.length).toBeGreaterThan(0);
+ expect(app.report.publications).toHaveLength(0);
+ app.report.unconfirmedControl = {
+ eventId: submitted.id,
+ optimisticRowVisible: true,
+ confirmationRejected: true,
+ reconciliationReads: reconciliation.length,
+ };
+ } finally {
+ release();
+ await page.unrouteAll({ behavior: "wait" });
+ }
+});
diff --git a/tests/browser/presence-integration.spec.mjs b/tests/browser/presence-integration.spec.mjs
new file mode 100644
index 00000000..a8e62531
--- /dev/null
+++ b/tests/browser/presence-integration.spec.mjs
@@ -0,0 +1,91 @@
+import { test, expect } from "./fixture.mjs";
+import { open } from "./timeline.mjs";
+const history = (page) =>
+ page.getByRole("region", { name: "Channel message history" });
+
+// A real reading journey can publish durable read intent while presence runs.
+test.use({ productionBroker: true, readState: true });
+test("presence: production conversation routes, seed, conflicting live repair and navigation", async ({
+ page,
+ app,
+}, testInfo) => {
+ await open(page, app);
+ const presenceRoute = () =>
+ app.relay.sockets
+ .filter((s) => s.readyState === 1 && s.community === "primary")
+ .flatMap((s) => [...s.routes.values()])
+ .filter((f) => f.kinds.includes(20001));
+ // The fixture's peer owns the history; inject only after live delivery exists.
+ const author = app.histories.get("primary/alpha")[0].pubkey;
+ await expect
+ .poll(() => presenceRoute().some((f) => f.authors.includes(author)))
+ .toBe(true);
+ app.presence("online");
+ const online = history(page).locator('[data-presence-status="online"]');
+ await expect(online.first()).toBeVisible();
+ const reads = () =>
+ app.report.queries.filter(({ filter }) => filter.kinds?.includes(20001));
+ expect(reads().length).toBeGreaterThan(0);
+ expect(reads().every(({ filter }) => filter.authors.length <= 256)).toBe(
+ true,
+ );
+ const count = reads().length;
+ // The real profile panel shares the timeline's author, not another read owner.
+ const sockets = app.relay.sockets.length;
+ await history(page)
+ .getByRole("button", { name: /^View .* profile$/ })
+ .first()
+ .click();
+ const profile = page.getByRole("complementary", {
+ name: "Profile",
+ exact: true,
+ });
+ await expect(
+ profile.getByRole("img", { name: "Online", exact: true }),
+ ).toBeVisible();
+ await page.screenshot({ path: testInfo.outputPath("presence-profile.png") });
+ await profile.getByRole("button", { name: "Close channel panel" }).click();
+ await expect(profile).toHaveCount(0);
+ expect(app.relay.sockets).toHaveLength(sockets);
+ expect(reads()).toHaveLength(count);
+ app.presence("online");
+ await expect(online.first()).toBeVisible();
+ expect(reads()).toHaveLength(count);
+ app.presence("away");
+ await expect(
+ history(page).locator('[data-presence-status="away"]').first(),
+ ).toBeVisible();
+ // A delayed conflicting heartbeat cannot turn a snapshot-confirmed Away green.
+ app.presence("online", false);
+ await expect(online).toHaveCount(0);
+ await expect(
+ history(page).locator('[data-presence-status="away"]').first(),
+ ).toBeVisible();
+ // Exercise the real dwell -> journal -> encrypted HTTP publication alongside
+ // presence, rather than making success depend on finishing before its debounce.
+ await history(page).focus();
+ await expect
+ .poll(() => app.report.readPublications.length, { timeout: 12000 })
+ .toBeGreaterThan(0);
+ const { community, event, blob } = app.report.readPublications[0];
+ expect(community).toBe("primary");
+ expect(event.kind).toBe(30078);
+ const messageContexts = Object.keys(blob.contexts).filter((key) =>
+ key.startsWith("msg:"),
+ );
+ expect(messageContexts.length).toBeGreaterThan(0);
+ const ids = new Set(app.histories.get("primary/alpha").map(({ id }) => id));
+ for (const key of messageContexts) {
+ expect(ids.has(key.slice(4))).toBe(true);
+ expect(event.content).not.toContain(key.slice(4));
+ }
+ await page.getByRole("button", { name: "Home", exact: true }).first().click();
+ await expect.poll(() => presenceRoute().length).toBe(0);
+ // Agent Activity's separately owned route survives presence demand teardown.
+ expect(app.relay.hasRoute("primary", "observer")).toBe(true);
+ expect(app.relay.sockets).toHaveLength(sockets);
+ expect(app.report.presencePublications.length).toBeGreaterThan(0);
+ const starts = reads().map((r) => r.at);
+ for (let i = 1; i < starts.length; i++)
+ expect(starts[i] - starts[i - 1]).toBeGreaterThanOrEqual(4900);
+});
diff --git a/tests/browser/presence.spec.mjs b/tests/browser/presence.spec.mjs
new file mode 100644
index 00000000..0d193aa2
--- /dev/null
+++ b/tests/browser/presence.spec.mjs
@@ -0,0 +1,219 @@
+import { test, expect } from "@playwright/test";
+import { createServer } from "./vite-server.mjs";
+import react from "@vitejs/plugin-react";
+import { fileURLToPath } from "node:url";
+
+test("presence: viewport demand, equal-status silence, conflict repair and teardown", async ({
+ page,
+}, testInfo) => {
+ const server = await createServer({
+ root: fileURLToPath(new URL("../../", import.meta.url)),
+ configFile: false,
+ envFile: false,
+ plugins: [react()],
+ logLevel: "error",
+ server: { host: "127.0.0.1", port: 0 },
+ });
+ const errors = [];
+ page.on("pageerror", (error) => errors.push(String(error)));
+ try {
+ await server.listen();
+ const started = Date.now();
+ await page.goto(
+ `http://127.0.0.1:${server.httpServer.address().port}/tests/fixtures/presence.html`,
+ );
+ const timeline = page.getByRole("region", { name: "Presence timeline" });
+ await expect(timeline).toBeVisible();
+ await expect(
+ timeline.locator('[data-presence-status="online"]').first(),
+ ).toBeVisible();
+ const initial = await page.evaluate(() => ({
+ ...window.presenceFixture.report,
+ ...window.presenceFixture.diagnostics(),
+ }));
+ expect(initial.handles).toBe(1);
+ expect(initial.authors).toBeGreaterThan(0);
+ expect(initial.authors).toBeLessThanOrEqual(20);
+ const cohorts = await page.evaluate(() => window.presenceFixture.cohorts);
+ expect(initial.currentAuthors.length).toBeGreaterThan(0);
+ expect(
+ initial.currentAuthors.every((author) => cohorts[0].includes(author)),
+ ).toBe(true);
+ const onlineDot = timeline
+ .locator('[data-presence-status="online"] [data-status]')
+ .first();
+ await expect(onlineDot).toHaveCSS("border-radius", "50%");
+ await expect(onlineDot).toHaveCSS("width", "8px");
+ await expect(onlineDot).toHaveCSS("height", "8px");
+ const before = initial.notifications;
+ await page.evaluate(() => window.presenceFixture.heartbeat(10));
+ const renewed = await page.evaluate(() => ({
+ ...window.presenceFixture.report,
+ ...window.presenceFixture.diagnostics(),
+ }));
+ expect(renewed.notifications).toBe(before);
+ expect(renewed.reads).toBe(initial.reads);
+ await page.evaluate(() => {
+ window.presenceFixture.status("away");
+ window.presenceFixture.heartbeat(1, "away");
+ });
+ const awayDot = timeline
+ .locator('[data-presence-status="away"] [data-status]')
+ .first();
+ await expect(awayDot).toBeVisible();
+ // Test rendered geometry, not just status/ARIA text: Away must not be another circle.
+ for (const mode of ["light", "dark"]) {
+ await page.evaluate((value) => {
+ document.documentElement.dataset.colorMode = value;
+ }, mode);
+ await expect(awayDot).toHaveCSS("border-radius", "0px");
+ await expect(awayDot).toHaveCSS("width", "8px");
+ await expect(awayDot).toHaveCSS("height", "8px");
+ await testInfo.attach(`away-${mode}.png`, {
+ body: await timeline.locator("[data-message-id]").first().screenshot(),
+ contentType: "image/png",
+ });
+ }
+ await page.evaluate(() => {
+ window.presenceFixture.status("offline");
+ window.presenceFixture.heartbeat(1, "offline");
+ });
+ await expect(
+ timeline.locator('[data-presence-status="offline"]').first(),
+ ).toBeVisible();
+ await page.evaluate(() => window.presenceFixture.heartbeat(1, "online"));
+ await expect(
+ timeline.locator('[data-presence-status="online"]'),
+ ).toHaveCount(0);
+ await expect(
+ timeline.locator('[data-presence-status="offline"]').first(),
+ ).toBeVisible();
+ await timeline.evaluate((element) => {
+ element.scrollTop = element.scrollHeight;
+ });
+ await expect
+ .poll(() =>
+ page.evaluate(() => {
+ const { report, cohorts } = window.presenceFixture;
+ return (
+ report.currentAuthors.length > 0 &&
+ report.currentAuthors.every((author) => cohorts[1].includes(author))
+ );
+ }),
+ )
+ .toBe(true);
+ const scrolled = await page.evaluate(
+ () => window.presenceFixture.report.currentAuthors,
+ );
+ expect(
+ scrolled.some((author) => initial.currentAuthors.includes(author)),
+ ).toBe(false);
+ await page.getByRole("button", { name: "Toggle timeline" }).click();
+ await expect(timeline).toHaveCount(0);
+ const disposed = await page.evaluate(() =>
+ window.presenceFixture.diagnostics(),
+ );
+ expect(disposed.handles).toBe(0);
+ expect(disposed.authors).toBe(0);
+ await page.getByRole("button", { name: "Toggle timeline" }).click();
+ await expect(timeline).toBeVisible();
+ await expect
+ .poll(() =>
+ page.evaluate(() => window.presenceFixture.diagnostics().handles),
+ )
+ .toBe(1);
+ await testInfo.attach("presence-counters.json", {
+ body: JSON.stringify(
+ {
+ elapsedMs: Date.now() - started,
+ initial,
+ renewed,
+ disposed,
+ final: await page.evaluate(() => window.presenceFixture.report),
+ },
+ null,
+ 2,
+ ),
+ contentType: "application/json",
+ });
+ expect(errors).toEqual([]);
+ await page.evaluate(() => window.presenceFixture.dispose());
+ } finally {
+ await server.close();
+ }
+});
+
+test("presence: real same-origin Web Lock and cross-window activity handoff", async ({
+ context,
+}) => {
+ const server = await createServer({
+ root: fileURLToPath(new URL("../../", import.meta.url)),
+ configFile: false,
+ envFile: false,
+ logLevel: "error",
+ server: { host: "127.0.0.1", port: 0 },
+ });
+ let first, second;
+ try {
+ await server.listen();
+ first = await context.newPage();
+ second = await context.newPage();
+ const url = `http://127.0.0.1:${server.httpServer.address().port}/tests/fixtures/presence-publisher.html`;
+ const time = new Date("2026-09-12T00:00:00Z");
+ await first.clock.install({ time });
+ await second.clock.install({ time });
+ await first.goto(url);
+ await expect
+ .poll(() =>
+ first.evaluate(() => window.publisherFixture.diagnostics().leader),
+ )
+ .toBe(true);
+ await first.clock.runFor(1000);
+ await second.goto(url);
+ await second.clock.runFor(1000);
+ expect(
+ await first.evaluate(() => window.publisherFixture.publications),
+ ).toEqual(["online"]);
+ expect(
+ await second.evaluate(() => window.publisherFixture.publications),
+ ).toEqual([]);
+ expect(
+ await second.evaluate(
+ () => window.publisherFixture.diagnostics().coordinated,
+ ),
+ ).toBe(true);
+ // Timer suspension is lossy, not 600 seconds of catch-up renewals.
+ await first.clock.fastForward(600000);
+ await second.clock.fastForward(600000);
+ await expect
+ .poll(() => first.evaluate(() => window.publisherFixture.state().status))
+ .toBe("away");
+ await second.getByRole("button", { name: "Activity" }).click();
+ await expect
+ .poll(() => first.evaluate(() => window.publisherFixture.state().status))
+ .toBe("online");
+ await first.clock.runFor(1000);
+ expect(
+ await first.evaluate(() => window.publisherFixture.publications.length),
+ ).toBeLessThanOrEqual(4);
+ expect(
+ await second.evaluate(() => window.publisherFixture.publications),
+ ).toEqual([]);
+ await first.evaluate(() => window.publisherFixture.dispose());
+ await expect
+ .poll(() =>
+ second.evaluate(() => window.publisherFixture.diagnostics().leader),
+ )
+ .toBe(true);
+ await second.clock.runFor(1000);
+ expect(
+ await second.evaluate(() => window.publisherFixture.publications),
+ ).toEqual(["online"]);
+ } finally {
+ try {
+ await Promise.all([first?.close(), second?.close()]);
+ } finally {
+ await server.close();
+ }
+ }
+});
diff --git a/tests/fixtures/presence-publisher.html b/tests/fixtures/presence-publisher.html
new file mode 100644
index 00000000..d03ec086
--- /dev/null
+++ b/tests/fixtures/presence-publisher.html
@@ -0,0 +1 @@
+Presence owner fixture Activity
diff --git a/tests/fixtures/presence-publisher.ts b/tests/fixtures/presence-publisher.ts
new file mode 100644
index 00000000..8dad1510
--- /dev/null
+++ b/tests/fixtures/presence-publisher.ts
@@ -0,0 +1,28 @@
+import { createPresenceActivity } from "../../src/features/presence/activity";
+import {
+ browserPresencePublisherLock,
+ createPresencePublisher,
+} from "../../src/features/presence/publisher";
+const activity = createPresenceActivity();
+const publications: string[] = [];
+const owner = createPresencePublisher({
+ activity: activity.activity,
+ lock: browserPresencePublisherLock("isolated-browser-presence-fixture"),
+ async publish(status, signal) {
+ signal.throwIfAborted();
+ publications.push(status);
+ },
+ random: () => 0,
+});
+owner.connection(true);
+Object.assign(window, {
+ publisherFixture: {
+ publications,
+ state: () => activity.activity.snapshot(),
+ diagnostics: () => owner.diagnostics(),
+ dispose() {
+ owner.dispose();
+ activity.dispose();
+ },
+ },
+});
diff --git a/tests/fixtures/presence.html b/tests/fixtures/presence.html
new file mode 100644
index 00000000..3fcd12f3
--- /dev/null
+++ b/tests/fixtures/presence.html
@@ -0,0 +1 @@
+Presence fixture
diff --git a/tests/fixtures/presence.tsx b/tests/fixtures/presence.tsx
new file mode 100644
index 00000000..314a79fb
--- /dev/null
+++ b/tests/fixtures/presence.tsx
@@ -0,0 +1,156 @@
+// Actual session and production MessageRow/surface hooks, synthetic signed transport only.
+import { StrictMode, useRef, useState } from "react";
+import { createRoot } from "react-dom/client";
+import { createRelaySession } from "../../src/features/relay/session";
+import type { LiveCallbacks } from "../../src/features/relay/live";
+import { keypair, signed } from "../../src/features/relay/testing";
+import { usePresenceSurface } from "../../src/features/presence/react";
+import { MessageRow } from "../../src/features/messages/MessageRow";
+import type { ChannelMessage } from "../../src/features/relay/contracts";
+import "../../src/shared/styles/globals.css";
+const relay = keypair(),
+ viewer = keypair();
+const people = Array.from({ length: 40 }, () => keypair());
+let callbacks!: LiveCallbacks;
+let status = "online";
+let serial = 0;
+const report = {
+ reads: 0,
+ updates: 0,
+ publishes: 0,
+ maximumAuthors: 0,
+ currentAuthors: [] as string[],
+};
+const owner = createRelaySession({
+ presence: true,
+ viewer: viewer.pubkey,
+ relayAuthor: relay.pubkey,
+ media: () => undefined,
+ async query(filters) {
+ if (!filters[0]?.kinds?.includes(20001)) return [];
+ report.reads++;
+ const authors = filters[0].authors ?? [];
+ return status === "offline"
+ ? []
+ : authors.map((author) =>
+ signed(relay, {
+ kind: 20001,
+ content: status,
+ tags: [["p", author]],
+ }),
+ );
+ },
+ subscribe(value) {
+ callbacks = value;
+ queueMicrotask(() => callbacks.state({ status: "connected", routes: [] }));
+ return {
+ update() {},
+ retry() {},
+ dispose() {},
+ presence: {
+ update(authors) {
+ report.updates++;
+ report.currentAuthors = [...authors];
+ report.maximumAuthors = Math.max(
+ report.maximumAuthors,
+ authors.length,
+ );
+ queueMicrotask(() =>
+ callbacks.presenceState?.({
+ status: authors.length ? "ready" : "idle",
+ authors,
+ }),
+ );
+ },
+ async publish() {
+ report.publishes++;
+ },
+ },
+ };
+ },
+});
+// Disjoint top/bottom cohorts make observing every mounted row a test failure.
+const personIndex = (index: number) => (index < 500 ? 0 : 20) + (index % 20);
+function row(index: number): ChannelMessage {
+ return {
+ id: index.toString(16).padStart(64, "0"),
+ authorId: people[personIndex(index)]?.pubkey ?? viewer.pubkey,
+ channelId: "fixture",
+ createdAt: 1700000000 + index,
+ content: `Message ${index}`,
+ attachments: [],
+ reactions: [],
+ replyCount: 0,
+ participants: [],
+ mentions: [],
+ emoji: [],
+ };
+}
+function Surface() {
+ const scroller = useRef(null);
+ usePresenceSurface(owner.session.presence, scroller);
+ return (
+
+ {Array.from({ length: 1000 }, (_, index) => (
+ undefined}
+ onOpenLink={() => false}
+ day={false}
+ retry={undefined}
+ />
+ ))}
+
+ );
+}
+function App() {
+ const [shown, show] = useState(true);
+ return (
+
+ show(!shown)}>
+ Toggle timeline
+
+ {shown && }
+
+ );
+}
+Object.assign(window, {
+ presenceFixture: {
+ report,
+ cohorts: [people.slice(0, 20), people.slice(20)].map((group) =>
+ group.map((person) => person.pubkey),
+ ),
+ diagnostics: () => owner.diagnostics().presence,
+ heartbeat(count = 1, value = status) {
+ for (let n = 0; n < count; n++)
+ callbacks.presence?.(
+ people.map((person) =>
+ signed(person, {
+ kind: 20001,
+ content: value,
+ tags: [],
+ created_at: ++serial,
+ }),
+ ),
+ );
+ },
+ status(value: string) {
+ status = value;
+ },
+ dispose: () => owner.dispose(),
+ },
+});
+const root = document.getElementById("root");
+if (!root) throw new Error("Missing fixture root");
+createRoot(root).render(
+
+
+ ,
+);