From f19fd77d97253d40e43b2e99f1c41a8f820a3a64 Mon Sep 17 00:00:00 2001 From: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> Date: Thu, 1 Oct 2026 09:49:36 -0600 Subject: [PATCH 1/5] fix(relay): use writer reads for channel confirmations Signed-off-by: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> --- dev/relay-broker-api.test.mjs | 21 ++++++ docs/relay-queries.md | 20 ++++++ .../channel-members/administration.test.ts | 23 ++++++- .../channel-members/administration.ts | 1 + src/features/relay/channel-details.test.ts | 17 ++++- src/features/relay/channel-details.ts | 1 + src/features/relay/channel-lifecycle.test.ts | 38 ++++++++++- src/features/relay/channel-lifecycle.ts | 22 +++++-- src/features/relay/contracts.ts | 15 ++++- src/features/relay/direct-messages.test.ts | 25 +++++++- src/features/relay/direct-messages.ts | 1 + src/features/relay/live-demand.test.ts | 10 ++- src/features/relay/mentions.test.ts | 23 ++++++- src/features/relay/session.ts | 26 ++++++-- src/features/relay/store.ts | 36 +++++++++-- src/features/relay/work-sessions.test.ts | 64 +++++++++++++++---- src/features/relay/work-sessions.ts | 32 ++++++++-- 17 files changed, 330 insertions(+), 45 deletions(-) diff --git a/dev/relay-broker-api.test.mjs b/dev/relay-broker-api.test.mjs index 1b1b6b0ab..ab9e3499e 100644 --- a/dev/relay-broker-api.test.mjs +++ b/dev/relay-broker-api.test.mjs @@ -955,6 +955,27 @@ test("local capacity is explicitly unsent, not relay quota; unknown upstream pub } }, 10000); +test("strong channel confirmation filters reach the upstream query unchanged", async () => { + const h = await harness(() => Response.json([])); + try { + const transport = await connectBrokerTransport(h.base); + const filters = [ + { + kinds: [39002], + "#d": ["11111111-1111-4111-8111-111111111111"], + limit: 1, + consistency: "strong", + }, + ]; + await transport.query(filters); + expect(h.calls).toHaveLength(1); + expect(h.calls[0].url).toBe(`${fixtureRelayUrl}/query`); + expect(h.calls[0].body).toEqual(filters); + } finally { + await h.close(); + } +}); + // Reader-to-host priority propagation control contributed by Brain. test("reader and transport start foreground work without waiting for background completion", async () => { let release; diff --git a/docs/relay-queries.md b/docs/relay-queries.md index aa7ced52d..6370ea8a4 100644 --- a/docs/relay-queries.md +++ b/docs/relay-queries.md @@ -66,6 +66,26 @@ 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. +## Read-your-writes consistency + +`fresh: true` prevents sharing an older in-flight request; it does not select +the writer. Channel creation/admission and recovery, DM opening, channel edits, +member administration, lifecycle confirmations, and mention preflights request +`consistency: "strong"` on their authoritative filters. Signed evidence remains +required; an accepted command is not membership, and writer routing does not +wait for asynchronous relay side effects. Existing cancellation and bounded +confirmation retries are unchanged. + +Channel discovery accepts an explicit consistency option for post-write exact +reads and the full-roster fallback. A queued writer-backed refresh survives an +older in-flight pass or quota pause, then reverts to ordinary routing. Signed +membership hints request that same full writer-backed pass, including metadata; +ordinary startup, browsing, reconnect, and DM visibility refresh stay replica- +eligible. Details/member-admin dialogs use writer-backed state for their shared +load/preflight/confirmation reads. Work-session membership preflights (including +session sends and canvas saves) also use the writer, so a just-added member does +not fail the next operation. + ## Community emoji `session.emoji` owns the current community's kind-30030 `d=buzz:custom-emoji` diff --git a/src/features/channel-members/administration.test.ts b/src/features/channel-members/administration.test.ts index 4f3266631..e73ed86e2 100644 --- a/src/features/channel-members/administration.test.ts +++ b/src/features/channel-members/administration.test.ts @@ -1,6 +1,6 @@ import { afterEach, expect, it, vi } from "vitest"; import { finalizeEvent, getPublicKey, type EventTemplate } from "nostr-tools"; -import type { RelayEvent } from "../relay/events"; +import type { ReadFilter, RelayEvent } from "../relay/events"; import { PublishRejected } from "../relay/outbox"; import { createRelaySession } from "../relay/session"; import { matchesEvent } from "../relay/projection"; @@ -65,7 +65,9 @@ function harness(actor = "owner", role: string | undefined = "member") { ...roster(actor, role), ]; let access = true; - const read = vi.fn(async (_filters?: unknown, _options?: unknown) => events); + const read = vi.fn( + async (_filters: readonly ReadFilter[], _options?: unknown) => events, + ); const sign = vi.fn(async (template: EventTemplate) => finalizeEvent(structuredClone(template), key), ); @@ -110,6 +112,22 @@ function harness(actor = "owner", role: string | undefined = "member") { }; } +it.each(["admin", "remove"] as const)( + "confirms %s while the replica retains the previous role", + async (role) => { + const h = harness(); + const replica = h.events(); + h.read.mockImplementation(async (filters) => + filters.every((filter) => filter.consistency === "strong") + ? h.events() + : replica, + ); + await h.owner.capability.run(id, { ...change, role }); + expect(h.owner.capability.snapshot(id).operation?.status).toBe("confirmed"); + expect(h.publish).toHaveBeenCalledOnce(); + }, +); + it.each(["admin", "member", "guest", "remove"] as const)( "confirms %s from exact fresh signed state, never acknowledgement alone", async (role) => { @@ -123,6 +141,7 @@ it.each(["admin", "member", "guest", "remove"] as const)( expect(h.read.mock.calls[0]).toEqual([ [39000, 39001, 39002].map((kind) => ({ kinds: [kind], + consistency: "strong", authors: [author], "#d": [id], limit: 1, diff --git a/src/features/channel-members/administration.ts b/src/features/channel-members/administration.ts index 35657b1fd..309687b84 100644 --- a/src/features/channel-members/administration.ts +++ b/src/features/channel-members/administration.ts @@ -70,6 +70,7 @@ export function createMemberAdministration({ const events = await reader.read( [39000, 39001, 39002].map((kind) => ({ kinds: [kind], + consistency: "strong" as const, authors: [relayAuthor], "#d": [id], limit: 1, diff --git a/src/features/relay/channel-details.test.ts b/src/features/relay/channel-details.test.ts index f5af603b8..dc2e364f4 100644 --- a/src/features/relay/channel-details.test.ts +++ b/src/features/relay/channel-details.test.ts @@ -59,7 +59,7 @@ function harness(role = "owner") { record(39001, role === "member" ? [] : [["p", viewer, role]]), record(39002, [["p", viewer, "", role]]), ]; - const read = vi.fn(async () => events); + const read = vi.fn(async (_filters: readonly ReadFilter[]) => events); const sign = vi.fn(async (event: EventTemplate) => finalizeEvent(structuredClone(event), key), ); @@ -114,6 +114,7 @@ it("saves only after two fresh authority checks and projects confirmed readback" expect(call).toEqual([ [39000, 39001, 39002].map((kind) => ({ kinds: [kind], + consistency: "strong", authors: [relayAuthor], "#d": [id], limit: 1, @@ -125,6 +126,20 @@ it("saves only after two fresh authority checks and projects confirmed readback" expect(h.acceptDiscovery).toHaveBeenCalledWith([h.events()[0]]); expect(h.owner.capability.snapshot(id)).toBeUndefined(); }); +it("confirms channel edits while the replica still has the old metadata", async () => { + const h = harness(); + const replica = h.events(); + h.read.mockImplementation(async (filters) => + filters.every((filter) => filter.consistency === "strong") + ? h.events() + : replica, + ); + const base = await h.owner.capability.load(id); + await h.owner.capability.save(base, draft); + expect(h.publish).toHaveBeenCalledOnce(); + expect(h.acceptDiscovery).toHaveBeenLastCalledWith([h.events()[0]]); + expect(h.owner.capability.snapshot(id)).toBeUndefined(); +}); it.each(["owner", "admin", "member"])( "derives current direct %s authority", async (role) => { diff --git a/src/features/relay/channel-details.ts b/src/features/relay/channel-details.ts index de9e04b68..ff98dd271 100644 --- a/src/features/relay/channel-details.ts +++ b/src/features/relay/channel-details.ts @@ -86,6 +86,7 @@ export function createChannelDetails({ const events = await reader.read( [39000, 39001, 39002].map((kind) => ({ kinds: [kind], + consistency: "strong" as const, authors: [relayAuthor], "#d": [id], limit: 1, diff --git a/src/features/relay/channel-lifecycle.test.ts b/src/features/relay/channel-lifecycle.test.ts index 04e012048..cd53e42b8 100644 --- a/src/features/relay/channel-lifecycle.test.ts +++ b/src/features/relay/channel-lifecycle.test.ts @@ -1,7 +1,7 @@ import { createHash } from "node:crypto"; import { schnorr } from "@noble/curves/secp256k1.js"; import { bytesToHex } from "nostr-tools/utils"; -import { describe, expect, it, vi } from "vitest"; +import { assert, describe, expect, it, vi } from "vitest"; import { finalizeEvent, getPublicKey, type EventTemplate } from "nostr-tools"; import { PublishRejected } from "./outbox"; import { @@ -140,6 +140,41 @@ function harness(role = "owner", type = "stream", owners = 1) { }; } +it.each(["archive", "delete", "leave", "hide"] as const)( + "confirms %s while ordinary reads still see pre-write state", + async (action) => { + const h = harness("owner", action === "hide" ? "dm" : "stream", 2); + const replica = h.getEvents(); + const currentRead = h.read.getMockImplementation(); + assert(currentRead); + h.read.mockImplementation(async (filters, options) => { + if (filters.every((filter) => filter.consistency === "strong")) + return currentRead(filters, options); + if (filters[0]?.kinds?.[0] === 30622) return []; + return replica.filter((event) => + filters.some((filter) => filter.kinds?.includes(event.kind)), + ); + }); + try { + await h.owner.capability.refreshVisibility(); + expect(h.read.mock.calls[0]?.[0][0]).not.toHaveProperty("consistency"); + await h.owner.capability.run(action, id); + for (const [filters] of h.read.mock.calls.slice(1, 3)) + for (const filter of filters) + expect(filter).not.toHaveProperty("consistency"); + expect(h.publish).toHaveBeenCalledOnce(); + expect(h.read.mock.calls.at(-1)?.[0][0]?.consistency).toBe("strong"); + if (action === "hide") + expect(h.owner.capability.snapshot().hidden).toEqual([id]); + else if (action === "archive") + expect(h.acceptDiscovery).toHaveBeenCalledOnce(); + else expect(h.removed).toHaveBeenCalledExactlyOnceWith(id); + } finally { + h.owner.dispose(); + } + }, +); + describe("type and role boundaries", () => { it.each([ ["owner", "stream", 1, true, true, false, false], @@ -295,6 +330,7 @@ it.each(["archive", "delete", "leave", "hide"] as const)( expect(h.read.mock.calls.at(-1)?.[0]).toEqual([ { kinds: [39002], + consistency: "strong", authors: [relayAuthor], "#d": [id], "#p": [viewer], diff --git a/src/features/relay/channel-lifecycle.ts b/src/features/relay/channel-lifecycle.ts index 73e4ca619..0896f88da 100644 --- a/src/features/relay/channel-lifecycle.ts +++ b/src/features/relay/channel-lifecycle.ts @@ -108,12 +108,14 @@ export function createChannelLifecycle({ id: string, signal: AbortSignal, member = false, + strong = false, ) { if (!reader) throw new Error("Channel actions are unavailable on this connection"); const events = await reader.read( kinds.map((kind) => ({ kinds: [kind], + ...(strong ? { consistency: "strong" as const } : {}), authors: [relayAuthor], "#d": [id], limit: 1, @@ -187,8 +189,14 @@ export function createChannelLifecycle({ return Object.freeze({ ...settings, deleteUnavailable: true }); } } - async function readVisibility(signal: AbortSignal) { - const events = await read([DM_VISIBILITY_KIND], viewer, signal, true); + async function readVisibility(signal: AbortSignal, strong = false) { + const events = await read( + [DM_VISIBILITY_KIND], + viewer, + signal, + true, + strong, + ); const record = lifecycleRecord( events, DM_VISIBILITY_KIND, @@ -314,10 +322,16 @@ export function createChannelLifecycle({ signal.addEventListener("abort", abort, { once: true }); }); if (action === "hide") { - if ((await readVisibility(signal)).includes(id)) return; + if ((await readVisibility(signal, true)).includes(id)) return; } else { const kind = action === "leave" ? 39002 : 39000; - const events = await read([kind], id, signal, action === "leave"); + const events = await read( + [kind], + id, + signal, + action === "leave", + true, + ); const record = lifecycleRecord(events, kind, id, relayAuthor); if (action === "archive") { if ( diff --git a/src/features/relay/contracts.ts b/src/features/relay/contracts.ts index 098012a03..d9be69e4c 100644 --- a/src/features/relay/contracts.ts +++ b/src/features/relay/contracts.ts @@ -1,3 +1,4 @@ +import type { ReadFilter } from "./events"; import type { CustomEmoji } from "./emoji"; import type { ReadOptions } from "./reader"; import type { Delivery } from "./outbox"; @@ -147,17 +148,25 @@ export type ChannelWindow = Readonly<{ freshness?: "cached" | "verified"; historyLimited?: boolean; }>; +/** Post-write discovery opts into writer reads; browsing keeps the default. */ +export type ChannelReadOptions = ReadOptions & Pick; /** Reads are side-effect-free; snapshots retain identity until their value changes. * Commands are idempotent requests; the store decides whether network work is needed. */ export interface ChannelQueries { list(): ChannelList; /** Bounded discovery lookup; never inserts public previews into list(). */ get?(channelId: string): ChannelSummary | undefined; - resolve?(channelIds: readonly string[], options?: ReadOptions): Promise; + resolve?( + channelIds: readonly string[], + options?: ChannelReadOptions, + ): Promise; /** Exact re-read of one already-listed channel's roster, merged into the * ready list. `resolve` admits channels the list lacks; this confirms a * membership change on one it already carries, without a full rediscovery. */ - refreshRoster?(channelId: string, options?: ReadOptions): Promise; + refreshRoster?( + channelId: string, + options?: ChannelReadOptions, + ): Promise; subscribeList(listener: () => void): () => void; window(channelId: string): ChannelWindow; subscribeWindow(channelId: string, listener: () => void): () => void; @@ -171,5 +180,5 @@ export interface ChannelQueries { /** Roster warming is optional for fixture-only query implementations. The * caller supplies preferred (e.g. starred) ids; the store orders the rest. */ warm?(preferred: readonly string[]): void; - refreshList?(): void; + refreshList?(options?: Pick): void; } diff --git a/src/features/relay/direct-messages.test.ts b/src/features/relay/direct-messages.test.ts index 2cdbee18d..1f5bf8629 100644 --- a/src/features/relay/direct-messages.test.ts +++ b/src/features/relay/direct-messages.test.ts @@ -6,7 +6,9 @@ import type { ReadFilter, RelayEvent } from "./events"; import type { RelayWriter } from "./transport"; const id = "11111111-1111-4111-8111-111111111111"; -function setup(options: { badRoster?: boolean; untrusted?: boolean } = {}) { +function setup( + options: { badRoster?: boolean; untrusted?: boolean; lagging?: boolean } = {}, +) { const viewer = keypair(), other = keypair(), relay = keypair(), @@ -37,7 +39,10 @@ function setup(options: { badRoster?: boolean; untrusted?: boolean } = {}) { about: "Unused biography", }), ]; - return discovery; + return options.lagging && + !filters.every((filter) => filter.consistency === "strong") + ? [] + : discovery; }); const publish = vi.fn(async () => {}); const openDirectMessage = vi.fn(async () => id); @@ -99,6 +104,22 @@ it("opens using signed exact membership, then confirms the regular outbox messag t.owner.dispose(); } }); +it("opens a DM before its signed discovery reaches the replica", async () => { + const t = setup({ lagging: true }); + try { + await expect( + t.dm.open([t.other.pubkey], new AbortController().signal), + ).resolves.toBe(id); + expect(t.owner.session.channels.get?.(id)?.members).toEqual( + [t.viewer.pubkey, t.other.pubkey].sort(), + ); + expect(t.query.mock.calls.at(-1)?.[0]).toEqual([ + expect.objectContaining({ consistency: "strong", "#d": [id] }), + ]); + } finally { + t.owner.dispose(); + } +}); it.each([{ badRoster: true }, { untrusted: true }])( "refuses a receipt without trusted exact participants: %j", async (options) => { diff --git a/src/features/relay/direct-messages.ts b/src/features/relay/direct-messages.ts index b3336b3d3..6d47ca1c3 100644 --- a/src/features/relay/direct-messages.ts +++ b/src/features/relay/direct-messages.ts @@ -83,6 +83,7 @@ export function createDirectMessages( [ { kinds: [39000, 39002], + consistency: "strong", authors: [transport.relayAuthor], "#d": [id], limit: 2, diff --git a/src/features/relay/live-demand.test.ts b/src/features/relay/live-demand.test.ts index a275b910f..93cbf922c 100644 --- a/src/features/relay/live-demand.test.ts +++ b/src/features/relay/live-demand.test.ts @@ -219,6 +219,7 @@ it("coalesces hints during a busy roster without letting retry clicks create wor h.live.established(); await flush(); const first = h.wire.next(); + expect(first.filters[0]).not.toHaveProperty("consistency"); h.owner.session.live.retry(); h.owner.session.live.retry(); first.respond([]); @@ -237,7 +238,9 @@ it("coalesces hints during a busy roster without letting retry clicks create wor second.respond([]); await flush(); expect(h.wire.pending).toHaveLength(1); - h.wire.next().respond([]); + const confirmation = h.wire.next(); + expect(confirmation.filters[0]?.consistency).toBe("strong"); + confirmation.respond([]); await flush(); expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); expect(h.wire.pending).toHaveLength(0); @@ -253,6 +256,7 @@ it("queued membership hints and reconnect cannot bypass a learned roster pause", h.live.established(); await flush(); const first = h.wire.next(); + expect(first.filters[0]).not.toHaveProperty("consistency"); h.live.receive([ signed(h.relay, { kind: 44100, @@ -274,7 +278,9 @@ it("queued membership hints and reconnect cannot bypass a learned roster pause", h.owner.session.live.retry(); await flush(); expect(h.wire.pending).toHaveLength(1); - h.wire.next().respond([]); + const confirmation = h.wire.next(); + expect(confirmation.filters[0]?.consistency).toBe("strong"); + confirmation.respond([]); await flush(); expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); expect(h.wire.pending).toHaveLength(0); diff --git a/src/features/relay/mentions.test.ts b/src/features/relay/mentions.test.ts index 406fcf739..0d2bdd575 100644 --- a/src/features/relay/mentions.test.ts +++ b/src/features/relay/mentions.test.ts @@ -30,6 +30,7 @@ function setup(sessionMode = false) { return signed(viewer, template); }); let currentRoster: RelayEvent | undefined; + let replicaRoster: RelayEvent | undefined; const profileEvents: RelayEvent[] = []; const owner = createRelaySession( { @@ -40,7 +41,13 @@ function setup(sessionMode = false) { filters[0]?.kinds?.[0] === 39002 && filters[0]?.limit === 1 ) - return Promise.resolve(currentRoster ? [currentRoster] : []); + return Promise.resolve( + replicaRoster && filters[0]?.consistency !== "strong" + ? [replicaRoster] + : currentRoster + ? [currentRoster] + : [], + ); if ( sessionMode && filters.every((filter) => filter.kinds?.every((kind) => kind === 0)) @@ -98,6 +105,9 @@ function setup(sessionMode = false) { ...wire, ...owner, members, + lagReplica: () => { + replicaRoster = currentRoster; + }, async agentProfile(key: typeof honey) { profileEvents.push( signed(key, { @@ -145,6 +155,17 @@ it.each([false, true])( expect(h.session.outbox?.snapshot()[0]?.delivery).toBe("accepted"); }, ); +it("mentions a just-added member despite an older replica roster", async () => { + const h = setup(); + await h.members([viewer.pubkey]); + h.lagReplica(); + await h.members([viewer.pubkey, honey.pubkey], 1700000001); + const id = h.session.messages.send("c", "@Honey help", [honey.pubkey]); + await vi.waitFor(() => expect(h.publish).toHaveBeenCalledOnce()); + expect(h.publish.mock.calls[0]?.[0].id).toBe(id); + expect(h.publish.mock.calls[0]?.[0].tags).toContainEqual(["p", honey.pubkey]); +}); + it("typed names create no recipient tags; unconfirmed or forged membership cannot grant mention permission", async () => { const h = setup(); expect(() => h.session.messages.send("c", "@Honey", [honey.pubkey])).toThrow( diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts index 92e20a96c..b231a99d7 100644 --- a/src/features/relay/session.ts +++ b/src/features/relay/session.ts @@ -943,6 +943,7 @@ export function createRelaySession( [ { kinds: [39002], + consistency: "strong", authors: [transport.relayAuthor], "#d": [channelId], limit: 1, @@ -999,7 +1000,15 @@ export function createRelaySession( // Confirm only this viewer's exact creation receipt. Discovery may be // incomplete; this never admits the channel or grants content access. const events = await requests.reader.read( - [{ kinds: [9007], ids: [id], authors: [transport.viewer], limit: 1 }], + [ + { + kinds: [9007], + ids: [id], + authors: [transport.viewer], + limit: 1, + consistency: "strong", + }, + ], { signal: lifetime.signal, fresh: true }, ); return events.some( @@ -1846,12 +1855,19 @@ export function createRelaySession( } } let rosterTimer: ReturnType | undefined; - function refreshRoster() { - if (closed || rosterTimer) return; + let strongRosterRefresh = false; + function refreshRoster(strong = false) { + if (closed) return; + strongRosterRefresh ||= strong; + if (rosterTimer) return; const timer = setTimeout(() => { timers.delete(timer); rosterTimer = undefined; - if (!closed) channels.queries.refreshList?.(); + const consistency = strongRosterRefresh + ? { consistency: "strong" as const } + : {}; + strongRosterRefresh = false; + if (!closed) channels.queries.refreshList?.(consistency); }, 0); rosterTimer = timer; timers.add(timer); @@ -2029,7 +2045,7 @@ export function createRelaySession( ), ) ) - refreshRoster(); + refreshRoster(true); const epoch = accessEpoch; const generation = liveGeneration; const visible = accept(events); diff --git a/src/features/relay/store.ts b/src/features/relay/store.ts index cd63b769d..2acef1819 100644 --- a/src/features/relay/store.ts +++ b/src/features/relay/store.ts @@ -5,6 +5,7 @@ import { createRelayProfiler, type RelayProfiler } from "./profiling"; import { ReadError, readErrorKind } from "./errors"; import type { ChannelList, + ChannelReadOptions, ChannelMessage, ChannelQueries, ChannelWindow, @@ -137,6 +138,7 @@ export function createChannelStore( epoch = 0, listBusy = false; let listAgain = false; + let strongListAgain = false; let listRetryAt = 0; type RosterRefresh = Readonly<{ state: "idle" | "pending" | "verified" | "deferred" | "error"; @@ -974,11 +976,12 @@ export function createChannelStore( } /** Apply roster authority as soon as it succeeds; names are a separate, * optional read and cannot delay revocation or overwrite newer live grants. */ - async function discover(force = false) { + async function discover(force = false, strong = false) { if (disposed || !transport || !discovery || options.cachedOnly) return; // Hints/establishment during a read require a later read. During a quota // pause they retain an obligation, not another request with a deadline. if (force) listAgain = true; + if (strong) strongListAgain = true; if ( listBusy || performance.now() < listRetryAt || @@ -988,6 +991,11 @@ export function createChannelStore( ) return; listAgain = false; + // A post-write request queued behind an older pass must keep writer routing. + const consistency = strongListAgain + ? { consistency: "strong" as const } + : {}; + strongListAgain = false; listRetryAt = 0; listBusy = true; rosterRefresh = Object.freeze({ state: "pending" }); @@ -1025,6 +1033,7 @@ export function createChannelStore( [ { kinds: [39002], + ...consistency, "#p": [transport.viewer], limit: DISCOVERY_LIMIT, ...(cursor @@ -1112,6 +1121,7 @@ export function createChannelStore( [ { kinds: [39002], + ...consistency, authors: [transport.relayAuthor], "#d": batch, "#p": [transport.viewer], @@ -1165,6 +1175,7 @@ export function createChannelStore( [ { kinds: [39000], + ...consistency, "#d": wanted.slice(offset, offset + DISCOVERY_LIMIT), limit: DISCOVERY_LIMIT, }, @@ -1245,7 +1256,7 @@ export function createChannelStore( * the canonical regression test. */ async function resolve( channelIds: readonly string[], - settings?: ReadOptions, + settings?: ChannelReadOptions, ) { if (disposed || !transport || !discovery || options.cachedOnly) throw new Error("Relay is unavailable"); @@ -1265,12 +1276,18 @@ export function createChannelStore( [ { kinds: [39000], + ...(settings?.consistency + ? { consistency: settings.consistency } + : {}), authors: [transport.relayAuthor], "#d": ids, limit: ids.length + 1, }, { kinds: [39002], + ...(settings?.consistency + ? { consistency: settings.consistency } + : {}), authors: [transport.relayAuthor], "#d": ids, "#p": [transport.viewer], @@ -1377,7 +1394,10 @@ export function createChannelStore( * - The read is viewer-scoped (`#p`), so the relay never answers with a roster * this viewer is absent from. An omitted roster changes nothing: revocation * by omission stays with the complete viewer-roster pass and live traffic. */ - async function refreshRoster(channelId: string, settings?: ReadOptions) { + async function refreshRoster( + channelId: string, + settings?: ChannelReadOptions, + ) { if (disposed || !transport || !discovery || options.cachedOnly) throw new Error("Relay is unavailable"); if (list.status !== "ready") @@ -1388,6 +1408,9 @@ export function createChannelStore( [ { kinds: [39002], + ...(settings?.consistency + ? { consistency: settings.consistency } + : {}), authors: [transport.relayAuthor], "#d": [channelId], "#p": [transport.viewer], @@ -1541,9 +1564,10 @@ export function createChannelStore( if (!persistence?.readStartup) void discover(); else void restore().then(() => discover()); }, - refreshList() { - if (!persistence?.readStartup) void discover(true); - else void restore().then(() => discover(true)); + refreshList(settings?: Pick) { + const strong = settings?.consistency === "strong"; + if (!persistence?.readStartup) void discover(true, strong); + else void restore().then(() => discover(true, strong)); }, /** Background roster warm. The caller supplies preferred ids (e.g. starred); * the rest follow by recency of their retained head, never-fetched last. */ diff --git a/src/features/relay/work-sessions.test.ts b/src/features/relay/work-sessions.test.ts index 209ec7020..c7a0525c8 100644 --- a/src/features/relay/work-sessions.test.ts +++ b/src/features/relay/work-sessions.test.ts @@ -580,18 +580,26 @@ it.each<["open" | "private", string[][]]>([ request.filters.some((filter) => filter["#d"]), ); expect(exact?.filters).toEqual([ - { kinds: [39000], authors: [relay.pubkey], "#d": [id], limit: 2 }, + { + kinds: [39000], + authors: [relay.pubkey], + "#d": [id], + limit: 2, + consistency: "strong", + }, { kinds: [39002], authors: [relay.pubkey], "#d": [id], "#p": [viewer.pubkey], limit: 2, + consistency: "strong", }, ]); for (const request of wire.pending.splice(0)) request.respond( - request === exact + request === exact && + request.filters.every((filter) => filter.consistency === "strong") ? [ metadata(relay, id, "Release notes", undefined, tags), roster(relay, id, [viewer.pubkey]), @@ -657,14 +665,16 @@ it("creates a channel during initial discovery without committing a list of only records .find(({ event }) => event.kind === 9007) ?.event.tags.find(([name]) => name === "h")?.[1] ?? ""; - // Serve every pending read: the viewer roster page lists `joined`, and - // metadata answers whatever ids discovery asks for. - const serve = (joined: readonly string[]) => { + // Ordinary reads remain frozen before creation; only writer reads see it. + // Metadata answers the IDs requested by each discovery pass. + const serve = () => { for (const request of wire.pending.splice(0)) { const [filter] = request.filters; request.respond( filter?.kinds?.includes(39002) && filter["#p"] - ? joined.map((channel) => roster(relay, channel, [viewer.pubkey])) + ? (filter.consistency === "strong" ? [other, id] : [other]).map( + (channel) => roster(relay, channel, [viewer.pubkey]), + ) : filter?.kinds?.includes(39000) ? (filter["#d"] ?? []).map((channel) => metadata( @@ -689,7 +699,7 @@ it("creates a channel during initial discovery without committing a list of only ), ).toBe(false); // The in-flight page was served before the relay accepted the create. - serve([other]); + serve(); await vi.waitFor(() => expect( snapshots.find((snapshot) => snapshot.status === "ready")?.channels, @@ -698,7 +708,7 @@ it("creates a channel during initial discovery without committing a list of only // The first pass finishes with its metadata read before the forced second // pass starts; that pass carries the new channel's roster. await vi.waitFor(() => expect(wire.pending).toHaveLength(1)); - serve([other]); + serve(); await vi.waitFor(() => expect( wire.pending.some((request) => @@ -706,7 +716,8 @@ it("creates a channel during initial discovery without committing a list of only ), ).toBe(true), ); - serve([other, id]); + expect(wire.pending[0]?.filters[0]?.consistency).toBe("strong"); + serve(); await expect(creating).resolves.toBe(id); expect(owner.session.channels.list()).toMatchObject({ status: "ready", @@ -715,6 +726,16 @@ it("creates a channel during initial discovery without committing a list of only expect.objectContaining({ id, members: [viewer.pubkey] }), ]), }); + // Finish metadata and prove the one-shot writer selection does not stick. + await vi.waitFor(() => expect(wire.pending).toHaveLength(1)); + expect(wire.pending[0]?.filters[0]?.consistency).toBe("strong"); + serve(); + await vi.waitFor(() => + expect(owner.session.live.snapshot().roster.state).toBe("verified"), + ); + owner.session.channels.refreshList?.(); + await vi.waitFor(() => expect(wire.pending).toHaveLength(1)); + expect(wire.pending[0]?.filters[0]).not.toHaveProperty("consistency"); // Canonical check, cited from the store's `resolve` docstring: no ready // snapshot ever held only the new channel. expect( @@ -780,9 +801,21 @@ it("confirms an agent added to a session and its parent with two concurrent exac roster(relay, parent, members.get(parent) ?? [], clock), roster(relay, child, members.get(child) ?? [], clock), ]; - return events.filter((event) => - filters.some((filter) => matchesEvent(event, filter)), + const replica = events.map((event) => + event.kind === 39002 + ? roster( + relay, + event.tags.find(([key]) => key === "d")?.[1] ?? "", + [viewer.pubkey], + 1_700_000_000, + ) + : event, ); + return ( + filters.every((filter) => filter.consistency === "strong") + ? events + : replica + ).filter((event) => filters.some((filter) => matchesEvent(event, filter))); }); const owner = createRelaySession( { @@ -833,6 +866,7 @@ it("confirms an agent added to a session and its parent with two concurrent exac "#d": [id], "#p": [viewer.pubkey], limit: 2, + consistency: "strong", }, ]); for (const { release } of held.splice(0)) release(); @@ -922,6 +956,7 @@ it.each([true, false])( [ { kinds: [9007], + consistency: "strong", ids: [creation.id], authors: [viewer.pubkey], limit: 1, @@ -1084,6 +1119,7 @@ it("admits a new channel from its exact lookup without a full discovery", async ).resolves.toBeUndefined(); expect(test.resolve).toHaveBeenCalledExactlyOnceWith([id], { signal: expect.any(AbortSignal), + consistency: "strong", }); expect(lookupSignal(test).aborted).toBe(true); expect(test.refreshList).not.toHaveBeenCalled(); @@ -1162,6 +1198,7 @@ it.each([ await vi.waitFor(() => expect(test.refreshList).toHaveBeenCalledOnce()); expect(test.resolve).toHaveBeenCalledExactlyOnceWith([id], { signal: expect.any(AbortSignal), + consistency: "strong", }); expect(test.resolve.mock.invocationCallOrder[0]).toBeLessThan( test.refreshList.mock.invocationCallOrder[0] ?? 0, @@ -1207,6 +1244,7 @@ it("confirms an agent added to a listed channel from one exact roster read witho expect(test.resolve).not.toHaveBeenCalled(); expect(test.refreshRoster).toHaveBeenCalledExactlyOnceWith(id, { signal: expect.any(AbortSignal), + consistency: "strong", }); expect(lookupSignal(test, test.refreshRoster).aborted).toBe(true); expect(test.refreshList).not.toHaveBeenCalled(); @@ -1278,6 +1316,7 @@ it.each([ await vi.waitFor(() => expect(test.refreshList).toHaveBeenCalledOnce()); expect(test.refreshRoster).toHaveBeenCalledExactlyOnceWith(id, { signal: expect.any(AbortSignal), + consistency: "strong", }); expect(test.refreshRoster.mock.invocationCallOrder[0]).toBeLessThan( test.refreshList.mock.invocationCallOrder[0] ?? 0, @@ -1511,10 +1550,11 @@ it("recovers a lost normal-channel creation acknowledgment only with its exact v [ { kinds: [39000, 39002], + consistency: "strong", "#d": ["11111111-1111-4111-8111-111111111111"], limit: 2, }, - { ids: [creation.id], limit: 1 }, + { ids: [creation.id], limit: 1, consistency: "strong" }, ], expect.anything(), ); diff --git a/src/features/relay/work-sessions.ts b/src/features/relay/work-sessions.ts index 7839f6f5a..a05ce032c 100644 --- a/src/features/relay/work-sessions.ts +++ b/src/features/relay/work-sessions.ts @@ -57,7 +57,10 @@ export function createWorkSessions( check(); return; } - const events = await reader.read([{ ids: [id], limit: 1 }], { signal }); + const events = await reader.read( + [{ ids: [id], limit: 1, consistency: "strong" }], + { signal }, + ); if (events.some((event) => event.id === id)) { check(); return; @@ -83,8 +86,13 @@ export function createWorkSessions( // by a stale roster. The verified reader applies signed discovery first. const events = await reader.read( [ - { kinds: [39000, 39002], "#d": [channelId], limit: 2 }, - { ids: [id], limit: 1 }, + { + kinds: [39000, 39002], + "#d": [channelId], + limit: 2, + consistency: "strong", + }, + { ids: [id], limit: 1, consistency: "strong" }, ], { signal }, ); @@ -100,7 +108,10 @@ export function createWorkSessions( } check(); if (!retry) { - const events = await reader.read([{ ids: [id], limit: 1 }], { signal }); + const events = await reader.read( + [{ ids: [id], limit: 1, consistency: "strong" }], + { signal }, + ); if (events.some((event) => event.id === id)) return; throw new Error( existing.error ?? "The operation could not be confirmed.", @@ -251,6 +262,7 @@ export function createWorkSessions( try { const options = { signal: AbortSignal.any([signal, exact.signal]), + consistency: "strong" as const, }; await (listed ? channels.refreshRoster?.(id, options) @@ -261,7 +273,7 @@ export function createWorkSessions( })(), ]); } - if (!settled) channels.refreshList?.(); + if (!settled) channels.refreshList?.({ consistency: "strong" }); await wait; } async function refreshMembership(id: string) { @@ -269,7 +281,15 @@ export function createWorkSessions( throw new Error("Channel membership is unavailable."); if (!relayAuthor) throw new Error("Channel membership is unavailable."); const events = await reader.read( - [{ kinds: [39002], authors: [relayAuthor], "#d": [id], limit: 1 }], + [ + { + kinds: [39002], + authors: [relayAuthor], + "#d": [id], + limit: 1, + consistency: "strong", + }, + ], { signal, fresh: true, priority: "foreground" }, ); const channel = From 813651348ac404bb82e4ac19f7bc4307954e6e90 Mon Sep 17 00:00:00 2001 From: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> Date: Thu, 1 Oct 2026 10:50:04 -0600 Subject: [PATCH 2/5] fix(relay): confirm template setup and agent removal on writer Signed-off-by: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> --- docs/relay-queries.md | 5 ++ .../agent-selection.test.tsx | 69 ++++++++++++++++++- .../profiles/ProfileAgentArchive.test.tsx | 44 ++++++++++++ src/features/channel-templates/capability.ts | 6 +- src/features/relay/session.ts | 11 +-- src/features/relay/work-sessions.ts | 20 +++++- 6 files changed, 144 insertions(+), 11 deletions(-) diff --git a/docs/relay-queries.md b/docs/relay-queries.md index 6370ea8a4..cb5404f2e 100644 --- a/docs/relay-queries.md +++ b/docs/relay-queries.md @@ -85,6 +85,11 @@ eligible. Details/member-admin dialogs use writer-backed state for their shared load/preflight/confirmation reads. Work-session membership preflights (including session sends and canvas saves) also use the writer, so a just-added member does not fail the next operation. +Template setup also confirms exact Canvas/member events and selected Canvas heads +against the writer without replaying accepted commands. Agent deletion discovers +member channels and confirms each removal with writer-backed rosters; unreadable +rosters still fail closed. Ordinary Canvas browsing and standalone Canvas/recipe +save confirmation routing are unchanged by this policy. ## Community emoji diff --git a/src/bundled/channel-templates/agent-selection.test.tsx b/src/bundled/channel-templates/agent-selection.test.tsx index e993b565a..6e4dc6973 100644 --- a/src/bundled/channel-templates/agent-selection.test.tsx +++ b/src/bundled/channel-templates/agent-selection.test.tsx @@ -26,7 +26,7 @@ import { createAgentControl } from "../../features/agents/control"; import { controlFixture } from "../../features/agents/control-testing"; import { keypair, signed, roster } from "../../features/relay/testing"; import { matchesEvent } from "../../features/relay/projection"; -import type { RelayEvent } from "../../features/relay/events"; +import type { ReadFilter, RelayEvent } from "../../features/relay/events"; import { createOutbox, type OutgoingEvent } from "../../features/relay/outbox"; import type { EventTemplate } from "nostr-tools"; import type { KitRecord } from "../../features/channel-templates/model"; @@ -105,6 +105,8 @@ function harness( }); const stored = new Map(); const published: RelayEvent[] = []; + const replica = { lag: false }; + const reads: ReadFilter[] = []; const sign = vi.fn(async (value: EventTemplate) => signed(viewer, value)); let journal: readonly OutgoingEvent[] = []; const channels = new Map([ @@ -177,6 +179,7 @@ function harness( }, }, query: async (filters) => { + reads.push(...filters); if ( filters.some((filter) => filter.kinds?.some((kind) => [39000, 39002].includes(kind)), @@ -211,7 +214,13 @@ function harness( ), ]; return events.filter((event) => - filters.some((filter) => matchesEvent(event, filter)), + filters.some( + (filter) => + matchesEvent(event, filter) && + (!replica.lag || + filter.consistency === "strong" || + ![9007, 9000, 40100, 39000, 39002].includes(event.kind)), + ), ); }, }, @@ -231,6 +240,8 @@ function harness( native, fixture, published, + replica, + reads, sign, viewer, journal: () => journal, @@ -1256,6 +1267,60 @@ it.each([ }, ); +it.each(["", "# Seed plan"])( + "finishes template setup against the writer with stale replicas and no live echoes (Canvas: %s)", + async (canvas) => { + const test = harness(); + test.replica.lag = true; + try { + const id = await test.owner.session.channelCreation.create({ + name: "Writer-confirmed setup", + visibility: "private", + setup: { + agents: [test.fixture.agent.pubkey], + canvas, + groupId: "", + templateId: "saved", + }, + }); + const receipt = `buzz-channel-setup.v2:https://relay.example.test:${test.viewer.pubkey}:${id}`; + // Receipt retirement is the completion barrier, not admission or publish ACK. + await waitFor(() => expect(localStorage.getItem(receipt)).toBeNull()); + expect(test.owner.session.channelCreation.notices()).toEqual([]); + expect(test.published.map((event) => event.kind)).toEqual( + canvas ? [9007, 40100, 9000] : [9007, 9000], + ); + expect(test.sign).toHaveBeenCalledTimes(test.published.length); + expect(test.journal()).toEqual([]); + expect(test.owner.session.channels.get?.(id)?.members).toContain( + test.fixture.agent.pubkey, + ); + for (const event of test.published.filter((event) => event.kind !== 9007)) + expect(test.reads).toContainEqual({ + ids: [event.id], + limit: 1, + consistency: "strong", + }); + if (canvas) { + const heads = test.reads.filter((filter) => + filter.kinds?.includes(40100), + ); + expect(heads).toHaveLength(3); // Before seeding, after save, before members. + expect(heads.every((filter) => filter.consistency === "strong")).toBe( + true, + ); + // Ordinary browsing still uses the lagging replica, not the writer. + await expect( + test.owner.session.canvas.read(id), + ).resolves.toBeUndefined(); + expect(test.reads.at(-1)?.consistency).toBeUndefined(); + } + } finally { + test.dispose(); + } + }, +); + it("retains a failed group placement independently and permits another Create", async () => { const test = harness(); try { diff --git a/src/bundled/profiles/ProfileAgentArchive.test.tsx b/src/bundled/profiles/ProfileAgentArchive.test.tsx index b1ce35de7..12c135df8 100644 --- a/src/bundled/profiles/ProfileAgentArchive.test.tsx +++ b/src/bundled/profiles/ProfileAgentArchive.test.tsx @@ -19,6 +19,7 @@ import { afterEach, expect, it, vi } from "vitest"; import { createRelaySession } from "../../features/relay/session"; import type { LiveCallbacks } from "../../features/relay/live"; import type { RelayData } from "../../features/relay/service"; +import { matchesEvent } from "../../features/relay/projection"; import { archiveRelay, keypair, @@ -418,6 +419,49 @@ it("verified owner deletes through the base confirmation: confirmed channel remo expect(fixture.instance.session.outbox?.snapshot()).toEqual([]), ); }); +it("deletes using writer rosters when replicas miss a recent channel and retain removed members", async () => { + const fixture = mountManaged(viewer, true, { + [room]: [viewer.pubkey, agent.pubkey], + }); + const query = fixture.transport.query; + const replica = await query( + [ + { + kinds: [39002], + authors: [relay.pubkey], + "#p": [agent.pubkey], + limit: 500, + }, + ], + new AbortController().signal, + ); + // A recent addition exists only on the writer and is not in the loaded list. + fixture.channels[hidden] = [agent.pubkey]; + expect(fixture.instance.session.channels.list().channels).not.toContainEqual( + expect.objectContaining({ id: hidden }), + ); + fixture.transport.query = async (filters, ...rest) => + ( + await Promise.all( + filters.map((filter) => + filter.kinds?.includes(39002) && filter.consistency !== "strong" + ? replica.filter((event) => matchesEvent(event, filter)) + : query([filter], ...rest), + ), + ) + ).flat(); + await confirmDelete(userEvent.setup()); + await vi.waitFor(() => expect(fixture.close).toHaveBeenCalledOnce()); + expect(fixture.published.map((event) => event.kind)).toEqual([ + 9001, 9001, 9035, + ]); + expect(fixture.channels).toEqual({ [room]: [viewer.pubkey], [hidden]: [] }); + expect(deletes(fixture)).toHaveLength(1); + await vi.waitFor(() => + expect(fixture.instance.session.outbox?.snapshot()).toEqual([]), + ); +}); + it("cancel leaves the agent untouched", async () => { const fixture = mountManaged(); const user = userEvent.setup(); diff --git a/src/features/channel-templates/capability.ts b/src/features/channel-templates/capability.ts index ce6024a62..753618e4f 100644 --- a/src/features/channel-templates/capability.ts +++ b/src/features/channel-templates/capability.ts @@ -227,12 +227,14 @@ export function createChannelKit({ }); const canvas = Object.freeze({ available: !!outbox?.supports(40100), - async read(channel: string) { + async read(channel: string, options: Pick = {}) { signal.throwIfAborted(); if (!canWrite(channel)) throw new Error("Canvas is unavailable after channel access changed"); return selectedHead( - await fresh([{ kinds: [40100], "#h": [channel], limit: 1 }]), + await fresh([ + { kinds: [40100], "#h": [channel], limit: 1, ...options }, + ]), ); }, async save(channel: string, content: string, expected: string | undefined) { diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts index b231a99d7..4c5a37edf 100644 --- a/src/features/relay/session.ts +++ b/src/features/relay/session.ts @@ -1190,16 +1190,17 @@ export function createRelaySession( await workSessions.delivered(id, undefined, false, false); if ( !( - await verified.read([{ ids: [id], limit: 1 }], { - signal: lifetime.signal, - fresh: true, - }) + await verified.read( + [{ ids: [id], limit: 1, consistency: "strong" }], + { signal: lifetime.signal, fresh: true }, + ) ).some((event) => event.id === id) ) throw new Error("Setup is awaiting exact relay confirmation"); }, refresh: (id, member) => workSessions.refresh(id, { member }, false), - canvasHead: async (id) => (await channelKit.canvas.read(id))?.id, + canvasHead: async (id) => + (await channelKit.canvas.read(id, { consistency: "strong" }))?.id, async preflight(input) { if ( !input.name.trim() || diff --git a/src/features/relay/work-sessions.ts b/src/features/relay/work-sessions.ts index a05ce032c..3d9e1e876 100644 --- a/src/features/relay/work-sessions.ts +++ b/src/features/relay/work-sessions.ts @@ -316,7 +316,15 @@ export function createWorkSessions( if (!relayAuthor) throw new Error("Channel membership is unavailable."); const limit = 500; const events = await discovery.read( - [{ kinds: [39002], authors: [relayAuthor], "#p": [pubkey], limit }], + [ + { + kinds: [39002], + authors: [relayAuthor], + "#p": [pubkey], + limit, + consistency: "strong", + }, + ], { signal: AbortSignal.any([signal, caller]), fresh: true, @@ -343,7 +351,15 @@ export function createWorkSessions( async function listsMember(id: string, pubkey: string, caller: AbortSignal) { if (!relayAuthor) throw new Error("Channel membership is unavailable."); const events = await discovery.read( - [{ kinds: [39002], authors: [relayAuthor], "#d": [id], limit: 1 }], + [ + { + kinds: [39002], + authors: [relayAuthor], + "#d": [id], + limit: 1, + consistency: "strong", + }, + ], { signal: AbortSignal.any([signal, caller]), fresh: true, From eb33025862d97ac40e038dee003952bdb7348248 Mon Sep 17 00:00:00 2001 From: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> Date: Thu, 1 Oct 2026 11:31:22 -0600 Subject: [PATCH 3/5] fix(relay): retain writer routing after refused roster reads Signed-off-by: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> --- docs/relay-queries.md | 13 +++-- src/features/relay/live-demand.test.ts | 72 ++++++++++++++++++++++++++ src/features/relay/store.ts | 4 ++ 3 files changed, 84 insertions(+), 5 deletions(-) diff --git a/docs/relay-queries.md b/docs/relay-queries.md index cb5404f2e..f19dbcfca 100644 --- a/docs/relay-queries.md +++ b/docs/relay-queries.md @@ -78,13 +78,16 @@ confirmation retries are unchanged. Channel discovery accepts an explicit consistency option for post-write exact reads and the full-roster fallback. A queued writer-backed refresh survives an -older in-flight pass or quota pause, then reverts to ordinary routing. Signed +older in-flight pass or quota pause. If the writer-backed roster pass itself +fails, its explicit retry retains writer routing and the existing cooldown; +a successful pass returns later refreshes to ordinary routing. Signed membership hints request that same full writer-backed pass, including metadata; ordinary startup, browsing, reconnect, and DM visibility refresh stay replica- -eligible. Details/member-admin dialogs use writer-backed state for their shared -load/preflight/confirmation reads. Work-session membership preflights (including -session sends and canvas saves) also use the writer, so a just-added member does -not fail the next operation. +eligible. The details editor and member-administration capability use writer-backed +state for their shared load/preflight/confirmation reads; the member dialog's +separate display-roster load remains replica-eligible. Work-session membership +preflights (including session sends and canvas saves) also use the writer, so a +just-added member does not fail the next operation. Template setup also confirms exact Canvas/member events and selected Canvas heads against the writer without replaying accepted commands. Agent deletion discovers member channels and confirms each removal with writer-backed rosters; unreadable diff --git a/src/features/relay/live-demand.test.ts b/src/features/relay/live-demand.test.ts index 93cbf922c..91bd66b1d 100644 --- a/src/features/relay/live-demand.test.ts +++ b/src/features/relay/live-demand.test.ts @@ -289,6 +289,78 @@ it("queued membership hints and reconnect cannot bypass a learned roster pause", } }); +it.each([ + { name: "rate-limited", retryAfterMs: 60_000 }, + { name: "refused without a cooldown", retryAfterMs: undefined }, +])( + "a $name strong roster pass keeps writer routing for explicit retry", + async ({ retryAfterMs }) => { + const h = setup(); + try { + h.connected(); + h.live.established(); + await flush(); + h.wire.next().respond([]); + await flush(); + h.live.receive([ + signed(h.relay, { + kind: 44100, + content: "", + tags: [["p", h.viewer.pubkey]], + }), + ]); + await flush(); + const refused = h.wire.next(); + expect(refused.filters[0]?.consistency).toBe("strong"); + refused.fail( + new ReadError( + "unavailable", + "Read refused", + retryAfterMs === undefined ? undefined : 429, + retryAfterMs, + ), + ); + await flush(); + expect(h.owner.session.live.snapshot().roster.state).toBe("error"); + expect(h.wire.pending).toHaveLength(0); + if (retryAfterMs !== undefined) { + h.owner.session.live.retry(); + await flush(); + expect(h.wire.pending).toHaveLength(0); + vi.spyOn(performance, "now").mockReturnValue( + performance.now() + retryAfterMs + 1, + ); + } + h.owner.session.live.retry(); + await flush(); + const retry = h.wire.next(); + expect(retry.filters[0]?.consistency).toBe("strong"); + const granted = [ + roster(h.relay, "new", [h.viewer.pubkey]), + metadata(h.relay, "new", "New channel"), + ]; + // Only the writer has the newly granted membership; no live roster echo. + retry.respond(retry.filters[0]?.consistency === "strong" ? granted : []); + await flush(); + expect(h.owner.session.channels.list().channels.map((c) => c.id)).toEqual( + ["new"], + ); + expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); + expect(h.wire.pending).toHaveLength(0); + h.owner.session.channels.refreshList?.(); + await flush(); + const ordinary = h.wire.next(); + expect(ordinary.filters[0]).not.toHaveProperty("consistency"); + ordinary.respond(granted); + await flush(); + expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); + expect(h.wire.pending).toHaveLength(0); + } finally { + h.owner.dispose(); + } + }, +); + it.each(["unavailable", "denied"] as const)( "a %s metadata read retains roster authority and remains recoverable", async (kind) => { diff --git a/src/features/relay/store.ts b/src/features/relay/store.ts index 2acef1819..02d2e1f54 100644 --- a/src/features/relay/store.ts +++ b/src/features/relay/store.ts @@ -1189,6 +1189,10 @@ export function createChannelStore( if (!disposed && generation === epoch) outcome = { state: "verified" }; } catch (error) { if (disposed || generation !== epoch) return; + // A refused roster has not fulfilled this pass's writer requirement. + // Retain it for deliberate retry without scheduling more work. + if (readingRoster && consistency.consistency === "strong") + strongListAgain = true; outcome = isAbort(error) ? { state: "deferred" } : { state: "error", error: describe(error) }; From 77104349abea5ed96ca3c3133eba81a17068b683 Mon Sep 17 00:00:00 2001 From: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> Date: Thu, 1 Oct 2026 11:36:38 -0600 Subject: [PATCH 4/5] fix(relay): retain writer routing through metadata failures Signed-off-by: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> --- docs/relay-queries.md | 8 +++---- src/features/relay/live-demand.test.ts | 30 ++++++++++++++++++-------- src/features/relay/store.ts | 7 +++--- 3 files changed, 28 insertions(+), 17 deletions(-) diff --git a/docs/relay-queries.md b/docs/relay-queries.md index f19dbcfca..52a464b06 100644 --- a/docs/relay-queries.md +++ b/docs/relay-queries.md @@ -78,10 +78,10 @@ confirmation retries are unchanged. Channel discovery accepts an explicit consistency option for post-write exact reads and the full-roster fallback. A queued writer-backed refresh survives an -older in-flight pass or quota pause. If the writer-backed roster pass itself -fails, its explicit retry retains writer routing and the existing cooldown; -a successful pass returns later refreshes to ordinary routing. Signed -membership hints request that same full writer-backed pass, including metadata; +older in-flight pass or quota pause. If the writer-backed pass itself fails, +including its metadata phase, its explicit retry retains writer routing and the +existing cooldown; a successful pass returns later refreshes to ordinary routing. +Signed membership hints request that same full writer-backed pass, including metadata; ordinary startup, browsing, reconnect, and DM visibility refresh stay replica- eligible. The details editor and member-administration capability use writer-backed state for their shared load/preflight/confirmation reads; the member dialog's diff --git a/src/features/relay/live-demand.test.ts b/src/features/relay/live-demand.test.ts index 91bd66b1d..32837062e 100644 --- a/src/features/relay/live-demand.test.ts +++ b/src/features/relay/live-demand.test.ts @@ -290,11 +290,13 @@ it("queued membership hints and reconnect cannot bypass a learned roster pause", }); it.each([ - { name: "rate-limited", retryAfterMs: 60_000 }, - { name: "refused without a cooldown", retryAfterMs: undefined }, + { name: "rate-limited roster", retryAfterMs: 60_000, metadata: false }, + { name: "refused roster", retryAfterMs: undefined, metadata: false }, + { name: "rate-limited metadata", retryAfterMs: 60_000, metadata: true }, + { name: "refused metadata", retryAfterMs: undefined, metadata: true }, ])( - "a $name strong roster pass keeps writer routing for explicit retry", - async ({ retryAfterMs }) => { + "a $name read keeps writer routing for explicit retry", + async ({ retryAfterMs, metadata: failMetadata }) => { const h = setup(); try { h.connected(); @@ -310,8 +312,21 @@ it.each([ }), ]); await flush(); - const refused = h.wire.next(); + const membership = roster(h.relay, "new", [h.viewer.pubkey]); + let refused = h.wire.next(); expect(refused.filters[0]?.consistency).toBe("strong"); + if (failMetadata) { + refused.respond([membership]); + await flush(); + expect( + h.owner.session.channels.list().channels.map((c) => c.id), + ).toEqual(["new"]); + refused = h.wire.next(); + expect(refused.filters[0]).toMatchObject({ + kinds: [39000], + consistency: "strong", + }); + } refused.fail( new ReadError( "unavailable", @@ -335,10 +350,7 @@ it.each([ await flush(); const retry = h.wire.next(); expect(retry.filters[0]?.consistency).toBe("strong"); - const granted = [ - roster(h.relay, "new", [h.viewer.pubkey]), - metadata(h.relay, "new", "New channel"), - ]; + const granted = [membership, metadata(h.relay, "new", "New channel")]; // Only the writer has the newly granted membership; no live roster echo. retry.respond(retry.filters[0]?.consistency === "strong" ? granted : []); await flush(); diff --git a/src/features/relay/store.ts b/src/features/relay/store.ts index 02d2e1f54..fd9039edc 100644 --- a/src/features/relay/store.ts +++ b/src/features/relay/store.ts @@ -1189,10 +1189,9 @@ export function createChannelStore( if (!disposed && generation === epoch) outcome = { state: "verified" }; } catch (error) { if (disposed || generation !== epoch) return; - // A refused roster has not fulfilled this pass's writer requirement. - // Retain it for deliberate retry without scheduling more work. - if (readingRoster && consistency.consistency === "strong") - strongListAgain = true; + // A failed pass has not fulfilled its writer requirement. Even a names + // failure retries the roster, which must not fall back to a stale replica. + if (consistency.consistency === "strong") strongListAgain = true; outcome = isAbort(error) ? { state: "deferred" } : { state: "error", error: describe(error) }; From c7d773ce3a5b2098693deb0e3c9de4b584c90f3e Mon Sep 17 00:00:00 2001 From: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> Date: Fri, 2 Oct 2026 10:29:19 -0600 Subject: [PATCH 5/5] fix(relay): retain writer routing when hint confirmations retire Signed-off-by: Carl <32a2e2c9d428ee08902cab75d956da2c1d235a22d4766b0dd4138bf6e2e5db1d@buzz.block.builderlab.xyz> --- docs/relay-queries.md | 5 +- src/features/relay/live-demand.test.ts | 202 ++++++++++++------------- src/features/relay/session.ts | 8 +- src/features/relay/store.ts | 5 + 4 files changed, 112 insertions(+), 108 deletions(-) diff --git a/docs/relay-queries.md b/docs/relay-queries.md index 0bb7db706..575f4d67e 100644 --- a/docs/relay-queries.md +++ b/docs/relay-queries.md @@ -83,7 +83,10 @@ is interrupted, including its metadata phase, its next pass retains writer routi existing cooldown; a successful pass returns later refreshes to ordinary routing. Signed membership hints use writer-backed exact reads or the existing full-roster fallback, including metadata. A replica pass superseding pending exact hint -confirmations queues a writer-backed pass to settle those grants. +confirmations queues a writer-backed pass to settle those grants. If list failure, +disconnect or cache clear retires queued or in-flight hint confirmations, the store +retains their writer requirement for the next deliberate refresh, Retry or +establishment, without starting an automatic recovery pass or bypassing cooldown. Ordinary startup, browsing, reconnect, and DM visibility refresh stay replica- eligible. The details editor and member-administration capability use writer-backed state for their shared load/preflight/confirmation reads; the member dialog's diff --git a/src/features/relay/live-demand.test.ts b/src/features/relay/live-demand.test.ts index f7d84078a..5a9707a8e 100644 --- a/src/features/relay/live-demand.test.ts +++ b/src/features/relay/live-demand.test.ts @@ -291,6 +291,37 @@ function exactFilters(h: ReturnType, ids: string[]) { ]; } +// A retired hint must survive to deliberate recovery, without making ordinary +// refreshes writer-backed once that obligation has been fulfilled. +async function recoverRetiredHint(h: ReturnType) { + expect(h.wire.pending).toHaveLength(1); + const recovery = h.wire.next(); + expect(recovery.filters[0]?.kinds).toEqual([39002]); + expect(recovery.filters[0]?.["#d"]).toBeUndefined(); + expect(recovery.filters[0]?.consistency).toBe("strong"); + recovery.respond( + recovery.filters[0]?.consistency === "strong" + ? [ + metadata(h.relay, "a", "Alpha"), + roster(h.relay, "a", [h.viewer.pubkey]), + ] + : [], + ); + await flush(); + expect(h.owner.session.channels.list()).toMatchObject({ + status: "ready", + channels: [{ id: "a", name: "Alpha" }], + }); + expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); + expect(h.wire.pending).toHaveLength(0); + h.owner.session.channels.refreshList?.(); + const ordinary = h.wire.next(); + expect(ordinary.filters[0]).not.toHaveProperty("consistency"); + ordinary.respond([]); + await flush(); + expect(h.wire.pending).toHaveLength(0); +} + it("a named member-added hint confirms only that channel and preserves verified roster state", async () => { const h = setup(); try { @@ -590,97 +621,82 @@ it("membership hints from another signer or for another viewer do not read disco } }); -it("cache-clear retires an exact membership confirmation without automatic recovery", async () => { - const h = setup(); - try { - await readyRoster(h); - h.live.receive([membershipHint(h, "a")]); - await flush(); - const exact = h.wire.next(); - expect(exact.filters).toEqual(exactFilters(h, ["a"])); - await h.owner.clearCache(); - exact.respond([ - metadata(h.relay, "a", "Alpha"), - roster(h.relay, "a", [h.viewer.pubkey]), - ]); - await flush(); - await flush(); - expect(exact.signal?.aborted).toBe(true); - expect(h.owner.session.channels.get?.("a")).toBeUndefined(); - expect(h.wire.pending).toHaveLength(0); - expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); - } finally { - h.owner.dispose(); - } -}); - -it("disconnect retires an exact membership confirmation until new establishment", async () => { - const h = setup(); - try { - await readyRoster(h); - h.live.receive([membershipHint(h, "a")]); - await flush(); - const exact = h.wire.next(); - h.live.state({ status: "retrying", routes: [] }); - exact.respond([ - metadata(h.relay, "a", "Alpha"), - roster(h.relay, "a", [h.viewer.pubkey]), - ]); - await flush(); - await flush(); - expect(exact.signal?.aborted).toBe(true); - expect(h.owner.session.channels.get?.("a")).toBeUndefined(); - expect(h.wire.pending).toHaveLength(0); - h.connected(); - h.live.established(); - await flush(); - const full = h.wire.next(); - expect(full.filters[0]?.["#d"]).toBeUndefined(); - full.respond([]); - await flush(); - expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); - expect(h.wire.pending).toHaveLength(0); - } finally { - h.owner.dispose(); - } -}); - -it.each(["cache-clear", "disconnect"])( - "a %s before a queued hint dispatches retires it without an exact or fallback read", - async (reset) => { +it.each([ + { reset: "cache-clear", inFlight: false }, + { reset: "cache-clear", inFlight: true }, + { reset: "disconnect", inFlight: false }, + { reset: "disconnect", inFlight: true }, +])( + "$reset retires a hint (in flight: $inFlight) but preserves writer routing for deliberate recovery", + async ({ reset, inFlight }) => { const h = setup(); try { await readyRoster(h); - // The broker drains traffic before connection state in one loop, so the - // hint is still queued behind its timer when the reset arrives. The reset - // leaves the list ready; only the queue itself can stop the dispatch. h.live.receive([membershipHint(h, "a")]); + if (inFlight) await flush(); + const exact = inFlight ? h.wire.next() : undefined; + if (exact) expect(exact.filters).toEqual(exactFilters(h, ["a"])); if (reset === "cache-clear") await h.owner.clearCache(); else h.live.state({ status: "retrying", routes: [] }); + if (exact) { + expect(exact.signal?.aborted).toBe(true); + exact.respond([ + metadata(h.relay, "a", "Alpha"), + roster(h.relay, "a", [h.viewer.pubkey]), + ]); + } await flush(); await flush(); + // Retirement must neither admit a delayed grant nor start a recovery pass. expect(h.wire.pending).toHaveLength(0); expect(h.owner.session.channels.get?.("a")).toBeUndefined(); expect(h.owner.session.channels.list().status).toBe("ready"); + expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); if (reset === "disconnect") { h.connected(); h.live.established(); - await flush(); - const full = h.wire.next(); - expect(full.filters[0]?.kinds).toEqual([39002]); - expect(full.filters[0]?.["#d"]).toBeUndefined(); - full.respond([ - metadata(h.relay, "a", "Alpha"), - roster(h.relay, "a", [h.viewer.pubkey]), - ]); - await flush(); - expect(h.owner.session.channels.list()).toMatchObject({ - status: "ready", - channels: [{ id: "a", name: "Alpha" }], - }); - } - expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); + } else h.owner.session.channels.refreshList?.(); + await flush(); + await recoverRetiredHint(h); + } finally { + h.owner.dispose(); + } + }, +); + +it.each(["cache-clear", "disconnect"])( + "%s during a replica pass retains a later hint without automatically rerunning the interrupted pass", + async (reset) => { + const h = setup(); + try { + await readyRoster(h); + h.owner.session.channels.refreshList?.(); + const full = h.wire.next(); + expect(full.filters[0]).not.toHaveProperty("consistency"); + h.live.receive([membershipHint(h, "a")]); + await flush(); + const exact = h.wire.next(); + expect(exact.filters).toEqual(exactFilters(h, ["a"])); + if (reset === "cache-clear") await h.owner.clearCache(); + else h.live.state({ status: "retrying", routes: [] }); + expect(exact.signal?.aborted).toBe(true); + expect(full.signal?.aborted).toBe(true); + exact.respond([ + metadata(h.relay, "a", "Alpha"), + roster(h.relay, "a", [h.viewer.pubkey]), + ]); + full.respond([]); + await flush(); + await flush(); + expect(h.owner.session.live.snapshot().roster.state).toBe("deferred"); + expect(h.owner.session.channels.get?.("a")).toBeUndefined(); expect(h.wire.pending).toHaveLength(0); + if (reset === "disconnect") { + h.connected(); + h.live.established(); + } else h.owner.session.live.retry(); + await flush(); + await recoverRetiredHint(h); } finally { h.owner.dispose(); } @@ -714,22 +730,14 @@ it("an exact membership confirmation cannot overwrite a concurrent roster error" }); expect(h.owner.session.live.snapshot().roster.state).toBe("error"); expect(h.wire.pending).toHaveLength(0); - vi.spyOn(performance, "now").mockReturnValue(performance.now() + 61000); h.owner.session.live.retry(); await flush(); - h.wire - .next() - .respond([ - metadata(h.relay, "a", "Alpha"), - roster(h.relay, "a", [h.viewer.pubkey]), - ]); - await flush(); - expect(h.owner.session.channels.list()).toMatchObject({ - status: "ready", - channels: [{ id: "a", name: "Alpha" }], - }); - expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); expect(h.wire.pending).toHaveLength(0); + expect(h.owner.session.channels.list().status).toBe("error"); + vi.spyOn(performance, "now").mockReturnValue(performance.now() + 61000); + h.owner.session.live.retry(); + await flush(); + await recoverRetiredHint(h); } finally { h.owner.dispose(); } @@ -769,21 +777,7 @@ it("a concurrent roster failure without a cooldown retires the exact confirmatio expect(h.wire.pending).toHaveLength(0); h.owner.session.live.retry(); await flush(); - expect(h.wire.pending).toHaveLength(1); - const recovery = h.wire.next(); - expect(recovery.filters[0]?.kinds).toEqual([39002]); - expect(recovery.filters[0]?.["#d"]).toBeUndefined(); - recovery.respond([ - metadata(h.relay, "a", "Alpha"), - roster(h.relay, "a", [h.viewer.pubkey]), - ]); - await flush(); - expect(h.owner.session.channels.list()).toMatchObject({ - status: "ready", - channels: [{ id: "a", name: "Alpha" }], - }); - expect(h.owner.session.live.snapshot().roster.state).toBe("verified"); - expect(h.wire.pending).toHaveLength(0); + await recoverRetiredHint(h); } finally { h.owner.dispose(); } diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts index a92f8b4e7..1f0c92102 100644 --- a/src/features/relay/session.ts +++ b/src/features/relay/session.ts @@ -1905,10 +1905,12 @@ export function createRelaySession( } /** A cache clear, a disconnect or a list that leaves ready drops hints with * the rest of the session's reader work, before a queued batch's timer can - * dispatch it into the new epoch, and releases any inherited grants: a ready - * commit must not hide a failed discovery, and the next establishment's full - * pass or Retry owns recovery. */ + * dispatch it into the new epoch. Transfer their writer requirement to the + * store without scheduling work: a ready commit must not hide a failed + * discovery, and the next establishment's full pass or Retry owns recovery. */ function dropHintConfirmations() { + if (hintQueue.size > 0 || hintReads.size > 0) + channels.requireStrongListRead(); rosterOwesGrants = false; retireHintConfirmations(); } diff --git a/src/features/relay/store.ts b/src/features/relay/store.ts index 155f6a5f1..160f005f8 100644 --- a/src/features/relay/store.ts +++ b/src/features/relay/store.ts @@ -2017,6 +2017,11 @@ export function createChannelStore( return { queries, roster: () => rosterRefresh, + /** Preserve a retired exact confirmation's routing without scheduling work. + * The next refresh/Retry owns dispatch and the existing cooldown. */ + requireStrongListRead() { + if (!disposed) strongListAgain = true; + }, retryList() { if (rosterRefresh.state === "error" || rosterRefresh.state === "deferred") void discover(true);