diff --git a/docs/inbox.md b/docs/inbox.md new file mode 100644 index 000000000..7c4612e7d --- /dev/null +++ b/docs/inbox.md @@ -0,0 +1,107 @@ +# Inbox: in-progress port + +Inbox is currently a bundled placeholder page (`buzz.inbox/inbox`). This +stacked PR's session evidence is not displayed until the dependent Inbox UI PR. +No new navigation behavior is part of this evidence slice. + +## Evidence slice (stacked with Inbox UI) + +PR3 adds no Inbox UI. The subsequent Inbox page uses `session.unread.inbox()` as +its single conversation projection and `session.inboxFeed` as bounded addressed +history demand. These changes are one launch batch, not a separately shippable +backend feature. Ordinary channel and thread readers, read-state storage and +outbox retain their existing ownership. No projects, approvals or reminders. + +## Ownership and limits + +`session.unread.inbox()` / `subscribeInbox()` own retained verified unread +evidence and read actions. `session.inboxFeed` owns finite, verified addressed +history demand and exact incomplete-target metadata, not a row cache. Finite +results and live arrivals, edits and deletions use the existing session admission +and shared unread fold, so own/deleted messages stay absent and unresolved roots +never become duplicate conversations. Current membership gates the feed's +completeness targets and unread's rows; joining alone does not materialize +pre-membership history without fresh shared admission. Inbox renders those shared +conversation rows directly; there is no second project/approval row merge or +feed-owned reconciliation buffer and deletion-count abort. PR4 owns only +presentation, filtering and selection. No parallel signing or persistence is +added. Optional profile enrichment belongs to PR4; access, cache clear and +session retirement fence these projections. Opening Inbox does not mark rows +read; selecting an unread row does. Canonical Messages keeps its own reading +behavior. + +DM read clears the channel through its newest retained evidence. Thread read +advances the thread prefix, including earlier unshown replies, plus individually +represented top-level mentions and local message marks, but not unrelated +messages. Participation follows the shared unread owner's direct-parent policy, +including its existing bounded lookups for replies whose membership is undecided. +Fetched roots and participation witnesses remain structural, never extra Inbox rows. +A lookup-only root uses the existing exact-message manual-unread target until counted +root evidence arrives; its known root still supplies grouping and read-through. +Relevant replies remain thread activity even when the root cannot be fetched. +Multiple steps are not atomic: failures leave remaining evidence retryable. Manual unread is local to this device. Hosts without frontier-sync +disable read mutations. Saved frontiers are not proof of remote reconciliation. + +This is **bounded recent evidence**, not a complete historical inbox. Unread +retains at most 4,096 events / 8 MiB, observing up to 500 recent events per +128-channel roster batch. A lazy addressed query returns up to 50 kind-9/40002 +messages. For those exact IDs and unresolved failed targets, general `#e` +queries page signed kind-40003 edits and kind-5/9005 deletions, then deletions +of the edits; `include_aux` on an ordinary `#p` query is not a supported relay +contract. Both auxiliary stages must finish before the feed is ready. Retained +unread evidence can be provisionally admitted between queries, but exact target +IDs are marked incomplete **before** unread subscribers are notified. The +consumer must not show their body as current while incomplete; failures retain +that metadata and offer retry. This is a completeness signal, not atomic content +admission or a second message fold. A target outside the latest 50 addressed +rows is not discovered; if a failed target falls outside a later page its +bounded auxiliary check still runs before its incomplete flag is cleared. +Incomplete obligations survive disconnect and unrelated access revocation while +the corresponding readable unread evidence survives; full cache/session retirement +clears both. Tombstone checks include retained author edits even if a later relay +query omits their soft-deleted rows. If reference visibility withholds an auxiliary +event, the finite attempt stays failed/incomplete rather than treating the filtered +page as exhausted history. Explicit Retry can settle it once existing shared +readers have admitted the missing reference; Inbox adds no reference-resolution loop. +Auxiliary reads cap retained results at 2,000 events / 4 MiB per stage; the shared reader +keeps its existing per-request deadline and cancellation. Missing roots, +participation or older activity can omit rows; an empty Inbox does not prove +complete history. No polling, independent row source or channel window is +opened for the feed. Access/cache/disconnect/disposal fence pending reads; local +read intent remains with unread. Packaged/native acceptance and human visual +feedback remain separate. + +Reminders and their NIP-ER lifecycle are **not included** in this change. +The unfinished reminder prototype is preserved separately for later work, +not shipped in the Inbox source or broker routes. Follow/mute Inbox policy, +full backlog discovery and nonchat activity are outside this slice. Project and +approval queries, grouping, routing and detail presentation are deliberately +excluded rather than presented as partial parity. Native/ACP packaged acceptance and human visual feedback +remain open. Do not treat this as the full OG Inbox port. + +## Verification status + +`inbox-feed.test.ts` exercises the real session reader/visibility/unread owners +with signed ephemeral events, including the first-admission subscriber ordering, +held edit and tombstone reads, failure/retry, reentrant access removal and cache +reset, root regrouping, and an edited addressed target older than 500 ordinary +messages. `unread.test.ts` retains the current main read/catch-up behavior and +adds Inbox projection/read-state cases. Neither file establishes browser paint, +real relay persistence or packaged/native acceptance. PR4 must cover the +per-row pending/failed presentation contract and exact inline previews. + + +### Review repairs + +Incomplete obligations now survive lifecycle changes that retain readable evidence. +Closure includes retained author edits omitted by later relay queries and fails +visibly on withheld auxiliary evidence. Shared unread remains the sole row owner; +the unused feed row copy and its 70-deletion abort were removed with explicit +approval. A prepared channel-read intent lets the dependent UI retry the original +cutoff and manual-clear keys. Tests retain all prior scenarios, distinguish both +finite retention guards, and prove 71 admitted author deletions do not resurrect +rows or require a second reconciliation buffer. + +Prior local tests and review cover the repairs on the combined tree; this PR's +updated head requires its own checks. No human/live/native acceptance or shipping +readiness is claimed. diff --git a/src/features/relay/inbox-feed.test.ts b/src/features/relay/inbox-feed.test.ts new file mode 100644 index 000000000..3bf19aa34 --- /dev/null +++ b/src/features/relay/inbox-feed.test.ts @@ -0,0 +1,1480 @@ +import { afterEach, expect, it, vi } from "vitest"; +import { createRelaySession } from "./session"; +import { keypair, message, metadata, roster, signed } from "./testing"; +import type { RelayEvent, ReadFilter } from "./events"; +import type { LiveCallbacks } from "./live"; +import { byteSize } from "./budget"; +import { createInboxFeed } from "./inbox-feed"; + +// General #e reads return an empty terminal page. The ordinary #p response is +// still explicitly gated by each test; no incidental timer orders admission. +const owners: ReturnType[] = []; +afterEach(() => { + for (const owner of owners.splice(0)) owner.dispose(); + vi.useRealTimers(); +}); +function deferred() { + let resolve!: (value: T) => void; + let reject!: (error: Error) => void; + const promise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + return { promise, resolve, reject }; +} +function setup() { + const viewer = keypair(), + relay = keypair(), + alice = keypair(); + let live!: LiveCallbacks; + const calls: { + filters: readonly ReadFilter[]; + pending: ReturnType>; + }[] = []; + const query = vi.fn((filters: readonly ReadFilter[]) => { + if (filters[0]?.["#e"]) return Promise.resolve([] as RelayEvent[]); + const pending = deferred(); + calls.push({ filters, pending }); + return pending.promise; + }); + const owner = createRelaySession({ + viewer: viewer.pubkey, + relayAuthor: relay.pubkey, + query, + media: () => undefined, + subscribe(callbacks) { + live = callbacks; + return { update() {}, retry() {}, dispose() {} }; + }, + }); + owners.push(owner); + const admit = (members: string[], at = 10) => + live.receive([ + roster(relay, "room", members, at), + metadata(relay, "room", "Room", at), + ]); + return { ...owner, viewer, alice, relay, admit, calls, live, query }; +} +async function take(h: ReturnType, kind: number) { + await vi.waitFor(() => + expect(h.calls.some((call) => call.filters[0]?.kinds?.includes(kind))).toBe( + true, + ), + ); + const index = h.calls.findIndex((call) => + call.filters[0]?.kinds?.includes(kind), + ); + const call = h.calls.splice(index, 1)[0]; + if (!call) throw new Error("Missing expected feed query"); + return call.pending; +} +it("reads addressed history independently of the retained unread window, never a fake complete archive", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const addressed = message(h.alice, "room", "older mention", 20, [ + ["p", h.viewer.pubkey], + ]); + const other = message(h.alice, "room", "not addressed", 21); + const work = h.session.inboxFeed.ensure(); + const mentions = await take(h, 9); + expect(mentions.promise).toBeDefined(); + mentions.resolve([addressed, other]); + await work; + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + // The session reconciles authorized history into shared unread evidence, not a shadow store. + expect(h.session.unread.inbox().items.map((item) => item.messageId)).toEqual([ + addressed.id, + ]); + expect(h.session.channels.window("room").rows).toEqual([]); + expect(h.query).toHaveBeenCalledTimes(2); // addressed #p and exact #e terminal page + expect(h.query.mock.calls[1]?.[0]).toEqual([ + { + kinds: [40003, 5, 9005], + "#e": [addressed.id], + limit: 500, + }, + ]); + expect(h.query.mock.calls[0]?.[0]).toEqual([ + { kinds: [40002, 9], "#p": [h.viewer.pubkey], limit: 50 }, + ]); +}); +it("reconciles live addressed updates and drops membership-revoked history before listeners", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const work = h.session.inboxFeed.ensure(); + const mentions = await take(h, 9); + const first = addressed(h, "first", 20); + mentions.resolve([first]); + await work; + expect(rows(h).map((row) => row.messageId)).toEqual([first.id]); + const updates: string[][] = []; + h.session.unread.subscribeInbox(() => + updates.push(rows(h).map((row) => row.messageId)), + ); + const later = addressed(h, "later", 22); + h.live.receive([later]); + expect(rows(h).map((row) => row.messageId)).toEqual([later.id, first.id]); + expect(updates).toEqual([[later.id, first.id]]); + h.admit([h.alice.pubkey], 30); + expect(rows(h)).toEqual([]); + // The revocation notification itself must already expose empty rows. + expect(updates).toEqual([[later.id, first.id], []]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [], + }); +}); +it("cache clear fences late completions; fresh demand recovers", async () => { + const h = setup(); + h.admit([h.viewer.pubkey]); + const old = addressed(h, "old", 20); + h.live.receive([old]); + expect(rows(h).map((row) => row.messageId)).toEqual([old.id]); + const work = h.session.inboxFeed.ensure(); + const mentions = await take(h, 9); + await h.clearCache(); + expect(rows(h)).toEqual([]); + mentions.resolve([old]); + await work; + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [], + }); + expect(rows(h)).toEqual([]); + expect(h.query).toHaveBeenCalledTimes(1); + const again = h.session.inboxFeed.ensure(); + (await take(h, 9)).resolve([]); + await again; + expect(h.session.inboxFeed.snapshot().status).toBe("ready"); +}); +it("surfaces read denial without treating it as empty, then retries", async () => { + const h = setup(); + h.admit([h.viewer.pubkey]); + const work = h.session.inboxFeed.ensure(); + (await take(h, 9)).reject(new Error("restricted: feed unavailable")); + await work; + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "error", + error: "restricted: feed unavailable", + }); + const again = h.session.inboxFeed.refresh(); + (await take(h, 9)).resolve([]); + await again; + expect(h.session.inboxFeed.snapshot().status).toBe("ready"); +}); + +async function discover(h: ReturnType) { + h.session.channels.ensureList(); + (await take(h, 39002)).resolve([ + roster(h.relay, "room", [h.viewer.pubkey, h.alice.pubkey], 10), + metadata(h.relay, "room", "Room", 10), + ]); + await vi.waitFor(() => + expect(h.session.channels.list().status).toBe("ready"), + ); +} +const rows = (h: ReturnType) => h.session.unread.inbox().items; +const addressed = (h: ReturnType, text: string, at: number) => + message(h.alice, "room", text, at, [["p", h.viewer.pubkey]]); +it("complete initial roster remains lazy; first demand, invalidation and fresh demand really read", async () => { + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); + const h = setup(); + await discover(h); + expect( + h.query.mock.calls.some(([filters]) => + filters.some((filter) => !!filter["#p"] && filter.kinds?.includes(9)), + ), + ).toBe(false); + expect(h.session.inboxFeed.snapshot().status).toBe("idle"); + const first = h.session.inboxFeed.ensure(); + const gate = await take(h, 9); + expect(h.session.inboxFeed.snapshot().status).toBe("loading"); + const issue = addressed(h, "historical mention", 15); + gate.resolve([issue]); + await first; + expect(rows(h).map((row) => row.preview)).toContain("historical mention"); + const addressedReads = () => + h.query.mock.calls.filter(([filters]) => + filters.some((filter) => !!filter["#p"] && filter.kinds?.includes(9)), + ).length; + const count = addressedReads(); + expect(count).toBe(1); + h.admit([h.alice.pubkey], 30); + expect(rows(h)).toEqual([]); + // Admission/listeners run synchronously; drain microtasks and the existing + // zero-delay scheduling boundary before granting any fresh Inbox demand. + await vi.advanceTimersByTimeAsync(0); + expect(addressedReads()).toBe(count); + expect(h.session.inboxFeed.snapshot().status).toBe("idle"); + const next = h.session.inboxFeed.ensure(); + const fresh = await take(h, 9); + expect(addressedReads()).toBe(count + 1); + fresh.resolve([issue]); + await next; + expect(h.session.inboxFeed.snapshot().status).toBe("ready"); + expect(rows(h)).toEqual([]); +}); +it("reconnect recovers demanded feed but never starts an unopened feed", async () => { + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); + const h = setup(); + await discover(h); + h.live.state({ status: "connected", routes: [] }); + h.live.state({ status: "retrying", routes: [] }); + h.live.state({ status: "connected", routes: [] }); + h.live.established(); + await vi.advanceTimersByTimeAsync(0); + // Complete the roster and the existing zero-delay reconnect callback explicitly. + (await take(h, 39002)).resolve([ + roster(h.relay, "room", [h.viewer.pubkey, h.alice.pubkey], 10), + metadata(h.relay, "room", "Room", 10), + ]); + await vi.waitFor(() => + expect(h.session.channels.list().status).toBe("ready"), + ); + expect(h.session.inboxFeed.snapshot().status).toBe("idle"); + const first = h.session.inboxFeed.ensure(); + (await take(h, 9)).resolve([]); + await first; + h.live.state({ status: "retrying", routes: [] }); + expect(h.session.inboxFeed.snapshot().status).toBe("idle"); + h.live.state({ status: "connected", routes: [] }); + h.live.established(); + await vi.advanceTimersByTimeAsync(0); + (await take(h, 39002)).resolve([ + roster(h.relay, "room", [h.viewer.pubkey, h.alice.pubkey], 10), + metadata(h.relay, "room", "Room", 10), + ]); + (await take(h, 9)).resolve([addressed(h, "reconnected mention", 30)]); + await vi.waitFor(() => + expect(h.session.inboxFeed.snapshot().status).toBe("ready"), + ); + expect(rows(h).map((row) => row.preview)).toContain("reconnected mention"); +}); +it("finite completion merges concurrent live chat arrivals and author deletions", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const old = addressed(h, "old", 20), + live = addressed(h, "live", 30); + // Establish deletion admission before holding the finite mention read. + h.live.receive([old]); + const first = h.session.inboxFeed.ensure(); + const mention = await take(h, 9); + try { + h.live.receive([live]); + expect(rows(h).map((row) => row.preview)).toContain("live"); + h.live.receive([ + signed(h.alice, { + kind: 5, + tags: [["e", old.id]], + content: "", + created_at: 31, + }), + ]); + } finally { + mention.resolve([old]); + } + await first; + expect(rows(h).map((row) => row.preview)).toEqual(["live"]); + const again = h.session.inboxFeed.refresh(); + const stale = await take(h, 9); + h.admit([h.alice.pubkey], 40); + stale.resolve([old, live]); + await again; + expect(rows(h)).toEqual([]); +}); +it("71 admitted author deletions keep shared rows deleted without aborting a held finite attempt", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const chats = Array.from({ length: 71 }, (_, i) => + addressed(h, `delete ${i}`, 20 + i), + ); + // Establish every target's actual row and shared reference admission first. + h.live.receive(chats); + expect( + rows(h) + .map((row) => row.messageId) + .sort(), + ).toEqual(chats.map((chat) => chat.id).sort()); + const feed = h.session.inboxFeed; + const work = feed.ensure(); + const held = await take(h, 9); + const deletions = chats.map((chat, i) => + signed(h.alice, { + kind: i % 2 ? 9005 : 5, + tags: [["e", chat.id]], + content: "", + created_at: 100 + i, + }), + ); + expect(new Set(deletions.map((event) => event.id)).size).toBe(71); + try { + h.live.receive(deletions); + expect(rows(h)).toEqual([]); + // The retired guard failed here on deletion 71 despite remaining well + // inside shared unread and finite auxiliary budgets. + expect(feed.snapshot().status).toBe("loading"); + } finally { + held.resolve(chats.slice(0, 50)); + await work; + } + expect(feed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + error: undefined, + }); + expect(rows(h)).toEqual([]); // Late addressed history cannot resurrect them. + expect(h.query).toHaveBeenCalledTimes(2); + expect(h.query.mock.calls[1]?.[0]).toEqual([ + { + kinds: [40003, 5, 9005], + "#e": chats + .slice(0, 50) + .map((chat) => chat.id) + .sort(), + limit: 500, + }, + ]); +}); +it("channel feed history uses the shared fold: own/deleted rows stay absent and unresolved roots never duplicate", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const root = message(h.viewer, "room", "root", 15); + const reply = message(h.alice, "room", "reply", 20, [ + ["p", h.viewer.pubkey], + ["e", root.id, "", "reply"], + ]); + const own = message(h.viewer, "room", "own addressed", 21, [ + ["p", h.viewer.pubkey], + ]); + const work = h.session.inboxFeed.ensure(); + (await take(h, 9)).resolve([reply, own]); + await work; + expect(rows(h).map((row) => row.id)).toEqual([`room:${reply.id}`]); + h.live.receive([root]); + expect(rows(h).map((row) => row.id)).toEqual([`room:${root.id}`]); + h.live.receive([ + signed(h.alice, { + kind: 5, + tags: [ + ["h", "room"], + ["e", reply.id], + ], + content: "", + created_at: 30, + }), + ]); + expect(rows(h)).toEqual([]); + const refresh = h.session.inboxFeed.refresh(); + (await take(h, 9)).resolve([reply, own]); + await refresh; + expect(rows(h)).toEqual([]); +}); + +it("a live deletion arriving before a held finite result suppresses its late row even while loading", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const chat = addressed(h, "late deleted", 20); + // Establish shared target visibility first; orphan reference-only deletions + // are deliberately not admitted by the session access owner. + h.live.receive([chat]); + const work = h.session.inboxFeed.ensure(); + const mention = await take(h, 9); + h.live.receive([ + signed(h.alice, { + kind: 5, + tags: [["e", chat.id]], + content: "", + created_at: 30, + }), + ]); + mention.resolve([chat]); + await work; + expect(rows(h)).toEqual([]); +}); + +it.each([false, true])( + "disposal erases populated unread rows and feed metadata and rejects late completion (pending=%s)", + async (pending) => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const feed = h.session.inboxFeed; + const unread = h.session.unread; + const mention = addressed(h, "retired mention", 20); + const initial = feed.ensure(); + (await take(h, 9)).resolve([mention]); + await initial; + expect(feed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + expect(unread.inbox().items.map((row) => row.messageId)).toEqual([ + mention.id, + ]); + const work = pending ? feed.refresh() : undefined; + const held = pending ? await take(h, 9) : undefined; + const notify = vi.fn(); + const notifyRows = vi.fn(); + feed.subscribe(notify); + unread.subscribeInbox(notifyRows); + try { + h.dispose(); + expect(feed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [], + }); + expect(unread.inbox().items).toEqual([]); + expect(notify).not.toHaveBeenCalled(); + expect(notifyRows).not.toHaveBeenCalled(); + } finally { + held?.resolve([mention]); + await work; + } + expect(feed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [], + }); + expect(unread.inbox().items).toEqual([]); + expect(notify).not.toHaveBeenCalled(); + expect(notifyRows).not.toHaveBeenCalled(); + }, +); + +it("joining after public pre-membership history needs fresh shared admission, not another feed read", async () => { + const h = setup(); + h.live.receive([ + roster(h.relay, "room", [h.alice.pubkey], 10), + metadata(h.relay, "room", "Room", 10, [["public"]]), + ]); + const chat = addressed(h, "Addressed chat", 20); + const work = h.session.inboxFeed.ensure(); + (await take(h, 9)).resolve([chat]); + await work; + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + expect(rows(h)).toEqual([]); + const reads = h.query.mock.calls.length; + expect(reads).toBe(1); + expect(h.query.mock.calls.some(([filters]) => filters[0]?.["#e"])).toBe( + false, + ); + h.admit([h.viewer.pubkey, h.alice.pubkey], 30); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + // Unread never admitted the pre-membership chat. Joining alone cannot create + // a row or completeness obligation from that unretained evidence. + expect(rows(h)).toEqual([]); + await h.session.inboxFeed.ensure(); + expect(h.query).toHaveBeenCalledTimes(reads); + h.live.receive([chat]); + expect(rows(h).map((row) => row.messageId)).toEqual([chat.id]); +}); +it("excludes nonchat and approvals from unread rows and completeness targets and coalesces one bounded chat request", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const chat = addressed(h, "Chat mention", 20); + const nonchat = [ + 1, 45001, 45003, 1618, 1619, 1621, 1630, 1631, 1632, 1633, 46010, 46011, + 46012, + ].map((kind) => + signed(h.alice, { + kind, + content: `Excluded ${kind}`, + created_at: 21, + tags: [ + ["h", "room"], + ["p", h.viewer.pubkey], + ["a", `30617:${h.alice.pubkey}:repo`], + ], + }), + ); + const first = h.session.inboxFeed.ensure(), + second = h.session.inboxFeed.refresh(); + (await take(h, 9)).resolve([chat, ...nonchat]); + await Promise.all([first, second]); + expect(h.query).toHaveBeenCalledTimes(2); // one coalesced #p, one #e closure + expect(h.query.mock.calls[1]?.[0]).toEqual([ + { kinds: [40003, 5, 9005], "#e": [chat.id], limit: 500 }, + ]); + expect(rows(h).map((row) => row.messageId)).toEqual([chat.id]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + const before = h.session.unread.inbox(); + h.live.receive(nonchat); + expect(h.session.unread.inbox()).toBe(before); + expect(rows(h).map((row) => row.messageId)).toEqual([chat.id]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); +}); + +it("first verified admission marks exact target incomplete before unread subscribers see original", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const old = addressed(h, "OLD BODY", 20); + const sequence: { + preview: string | undefined; + status: string; + incomplete: readonly string[]; + }[] = []; + const observe = () => + sequence.push({ + preview: rows(h)[0]?.preview, + status: h.session.inboxFeed.snapshot().status, + incomplete: h.session.inboxFeed.snapshot().incomplete, + }); + const stop = h.session.unread.subscribeInbox(observe); + const stopFeed = h.session.inboxFeed.subscribe(observe); + try { + const work = h.session.inboxFeed.ensure(); + (await take(h, 9)).resolve([old]); + await work; + expect(sequence.filter(({ preview }) => preview === "OLD BODY")).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + status: "loading", + incomplete: [old.id], + }), + ]), + ); + expect( + sequence.filter( + ({ preview, incomplete, status }) => + preview === "OLD BODY" && + !incomplete.includes(old.id) && + status !== "ready", + ), + ).toEqual([]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + } finally { + stop(); + stopFeed(); + } +}); + +it("settles a signed old addressed edit after 500 newer ordinary rows without opening detail or live replay", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const old = addressed(h, "OLD BODY", 20); + const edit = signed(h.alice, { + kind: 40003, + created_at: 21, + content: "CURRENT BODY", + tags: [ + ["h", "room"], + ["e", old.id], + ], + }); + const newer = Array.from({ length: 500 }, (_, i) => + message(h.alice, "room", `ordinary ${i}`, 100 + i), + ); + const addressedQuery = deferred(); + const filters: ReadFilter[] = []; + h.query.mockImplementation(async ([filter]) => { + if (!filter) return []; + filters.push(filter); + if (filter["#p"]) return addressedQuery.promise; + if (filter["#h"]) return newer; + if (filter["#e"]?.includes(old.id)) { + // General #e, not include_aux on #p, returns author edits of exact IDs. + return filter.until === undefined ? [edit] : []; + } + if (filter["#e"]?.includes(edit.id)) return []; + return []; + }); + const unread = h.session.unread.ensure(); + await unread; + expect(rows(h)).toEqual([]); + const work = h.session.inboxFeed.ensure(); + await vi.waitFor(() => + expect(filters.some((filter) => !!filter["#p"])).toBe(true), + ); + addressedQuery.resolve([old]); + await work; + expect(rows(h).map((item) => item.preview)).toEqual(["CURRENT BODY"]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + expect(filters.filter((filter) => filter["#e"])).toEqual([ + { kinds: [40003, 5, 9005], "#e": [old.id], limit: 500 }, + { + kinds: [40003, 5, 9005], + "#e": [old.id], + limit: 500, + until: 21, + before_id: edit.id, + }, + { kinds: [5, 9005], "#e": [edit.id], limit: 500 }, + ]); + expect(filters.find((filter) => filter["#p"])).toEqual({ + kinds: [40002, 9], + "#p": [h.viewer.pubkey], + limit: 50, + }); + expect(h.session.channels.window("room").rows).toEqual([]); +}); + +it("keeps exact incomplete evidence across held edit and deletion reads, failed retry and recovery", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const old = addressed(h, "OLD BODY", 20); + const edit = signed(h.alice, { + kind: 40003, + created_at: 21, + content: "CURRENT BODY", + tags: [ + ["h", "room"], + ["e", old.id], + ], + }); + const hold = deferred(); + let editReads = 0; + const tombstone = deferred(); + const retry = deferred(); + const retryTombstone = deferred(); + const deletion = signed(h.alice, { + kind: 5, + created_at: 22, + content: "", + tags: [["e", edit.id]], + }); + let tombstoneReads = 0; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [old]; + if (filter?.["#e"]?.includes(old.id)) { + if (filter.until !== undefined) return []; + editReads++; + return editReads === 1 ? hold.promise : retry.promise; + } + if (filter?.["#e"]?.includes(edit.id)) { + if (filter.until !== undefined) return []; + tombstoneReads++; + return tombstoneReads === 1 ? tombstone.promise : retryTombstone.promise; + } + return []; + }); + const feed = h.session.inboxFeed; + const work = feed.ensure(); + try { + await vi.waitFor(() => expect(editReads).toBe(1)); + expect(feed.snapshot()).toMatchObject({ + status: "loading", + incomplete: [old.id], + }); + expect(rows(h)[0]?.preview).toBe("OLD BODY"); // PR4 must show placeholder instead + hold.resolve([edit]); + await vi.waitFor(() => expect(rows(h)[0]?.preview).toBe("CURRENT BODY")); + expect(feed.snapshot()).toMatchObject({ + status: "loading", + incomplete: [old.id], + }); + tombstone.reject(new Error("edit tombstones unavailable")); + await work; + expect(feed.snapshot()).toMatchObject({ + status: "error", + incomplete: [old.id], + error: "edit tombstones unavailable", + }); + const next = feed.refresh(); + await vi.waitFor(() => expect(editReads).toBe(2)); + expect(feed.snapshot()).toMatchObject({ + status: "loading", + incomplete: [old.id], + }); + retry.resolve([]); // Soft-deleted edits are absent from ordinary relay queries. + await vi.waitFor(() => expect(tombstoneReads).toBe(2)); + expect(feed.snapshot()).toMatchObject({ + status: "loading", + incomplete: [old.id], + }); + expect(rows(h)[0]?.preview).toBe("CURRENT BODY"); + retryTombstone.resolve([deletion]); + await next; + expect(rows(h)[0]?.preview).toBe("OLD BODY"); + expect(feed.snapshot()).toMatchObject({ status: "ready", incomplete: [] }); + } finally { + hold.resolve([]); + retry.resolve([]); + tombstone.resolve([]); + retryTombstone.resolve([]); + } +}); + +it("tombstones a stored edit, then restores its original target body only after closure settles", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "original", 20); + const edit = signed(h.alice, { + kind: 40003, + created_at: 21, + content: "superseded", + tags: [ + ["h", "room"], + ["e", target.id], + ], + }); + const deletion = signed(h.alice, { + kind: 5, + created_at: 22, + content: "", + tags: [["e", edit.id]], + }); + let release!: (events: RelayEvent[]) => void; + const held = new Promise((resolve) => { + release = resolve; + }); + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [target]; + if (filter?.["#e"]?.includes(target.id)) + return filter.until === undefined ? [edit] : []; + if (filter?.["#e"]?.includes(edit.id)) + return filter.until === undefined ? held : []; + return []; + }); + const work = h.session.inboxFeed.ensure(); + try { + await vi.waitFor(() => expect(rows(h)[0]?.preview).toBe("superseded")); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "loading", + incomplete: [target.id], + }); + release([deletion]); + await work; + expect(rows(h)[0]?.preview).toBe("original"); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + } finally { + release([]); + } +}); + +it("access removal while pre-admission subscriber runs cannot admit the pending target", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "must not appear", 20); + const stop = h.session.inboxFeed.subscribe(() => { + if (h.session.inboxFeed.snapshot().incomplete.includes(target.id)) + h.admit([h.alice.pubkey], 30); + }); + try { + const work = h.session.inboxFeed.ensure(); + (await take(h, 9)).resolve([target]); + await work; + expect(rows(h)).toEqual([]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [], + }); + expect(h.query.mock.calls.some(([filters]) => filters[0]?.["#e"])).toBe( + false, + ); + } finally { + stop(); + } +}); + +it("a cache reset fences an old pending edit read and its late result", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "old", 20); + let release!: (events: RelayEvent[]) => void; + const held = new Promise((resolve) => { + release = resolve; + }); + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [target]; + if (filter?.["#e"]) return held; + return []; + }); + const work = h.session.inboxFeed.ensure(); + try { + await vi.waitFor(() => + expect(h.session.inboxFeed.snapshot().incomplete).toEqual([target.id]), + ); + expect(h.session.inboxFeed.snapshot().status).toBe("loading"); + expect(rows(h)[0]?.preview).toBe("old"); + await h.clearCache(); + release([ + signed(h.alice, { + kind: 40003, + content: "late", + created_at: 21, + tags: [ + ["h", "room"], + ["e", target.id], + ], + }), + ]); + await work; + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [], + }); + expect(rows(h)).toEqual([]); + } finally { + release([]); + } +}); + +it("pending exact group member survives verified root regrouping and representative changes", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const root = message(h.viewer, "room", "my root", 15); + const reply = message(h.alice, "room", "old reply", 20, [ + ["p", h.viewer.pubkey], + ["e", root.id, "", "reply"], + ]); + let release!: (events: RelayEvent[]) => void; + const held = new Promise((resolve) => { + release = resolve; + }); + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [reply]; + if (filter?.["#e"]) return held; + return []; + }); + const work = h.session.inboxFeed.ensure(); + try { + await vi.waitFor(() => + expect(h.session.inboxFeed.snapshot().incomplete).toEqual([reply.id]), + ); + expect(rows(h)[0]).toMatchObject({ + id: `room:${reply.id}`, + messageIds: [reply.id], + }); + h.live.receive([root]); + expect(rows(h)[0]).toMatchObject({ + id: `room:${root.id}`, + messageId: reply.id, + messageIds: [reply.id], + }); + expect( + rows(h)[0]?.messageIds.some((id) => + h.session.inboxFeed.snapshot().incomplete.includes(id), + ), + ).toBe(true); + } finally { + release([]); + await work; + } +}); + +it("a synchronous loading subscriber clears demand without stranding a stale work promise", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + let once = true; + const feed = h.session.inboxFeed; + const stop = feed.subscribe(() => { + if (once && feed.snapshot().status === "loading") { + once = false; + feed.clear(); + } + }); + try { + const cancelled = feed.ensure(); + // Explicitly settle the fixture's aborted request, without a test timeout. + (await take(h, 9)).resolve([]); + await cancelled; + expect(feed.snapshot()).toMatchObject({ status: "idle", incomplete: [] }); + const fresh = feed.ensure(); + (await take(h, 9)).resolve([]); + await fresh; + expect(feed.snapshot().status).toBe("ready"); + } finally { + stop(); + } +}); + +it("retry verifies an old incomplete target even if a new addressed page excludes it", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const old = addressed(h, "old", 20); + const newer = Array.from({ length: 50 }, (_, i) => + addressed(h, `new ${i}`, 100 + i), + ); + let attempt = 0; + const exactReads: readonly string[][] = []; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return ++attempt === 1 ? [old] : newer; + if (filter?.["#e"] && filter.until === undefined) { + (exactReads as string[][]).push([...filter["#e"]]); + if (attempt === 1) throw new Error("overlays unavailable"); + } + return []; + }); + await h.session.inboxFeed.ensure(); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "error", + incomplete: [old.id], + }); + await h.session.inboxFeed.refresh(); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + expect(exactReads[1]).toContain(old.id); + expect(exactReads[1]).toHaveLength(51); +}); + +it("a non-advancing edit page remains retryable instead of clearing incomplete evidence", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "original", 20); + const edits = Array.from({ length: 2 }, (_, i) => + signed(h.alice, { + kind: 40003, + created_at: 30 + i, + content: `revision ${i}`, + tags: [ + ["h", "room"], + ["e", target.id], + ], + }), + ); + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [target]; + if (filter?.["#e"]?.includes(target.id)) return edits; + return []; + }); + await h.session.inboxFeed.ensure(); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "error", + incomplete: [target.id], + }); + expect(h.session.inboxFeed.snapshot().error).toBe( + "Inbox message updates did not advance. Retry inbox.", + ); + expect(h.session.inboxFeed.snapshot().status).not.toBe("ready"); +}); + +it("ordinary live receive after cache clearing does not throw from the finite callback fence", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + await h.clearCache(); + expect(() => h.live.receive([addressed(h, "ordinary", 20)])).not.toThrow(); +}); + +it("a reentrant feed clear during exact pre-admission never publishes the old body", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "never publish", 20); + const feed = h.session.inboxFeed; + const observed: string[] = []; + const stopUnread = h.session.unread.subscribeInbox(() => { + observed.push(...rows(h).map((row) => row.preview)); + }); + const stop = feed.subscribe(() => { + if (feed.snapshot().incomplete.includes(target.id)) feed.clear(); + }); + try { + const work = feed.ensure(); + (await take(h, 9)).resolve([target]); + await work; + expect(observed).not.toContain("never publish"); + expect(rows(h)).toEqual([]); + expect(feed.snapshot()).toMatchObject({ status: "idle", incomplete: [] }); + } finally { + stop(); + stopUnread(); + } +}); + +it.each(["disconnect", "unrelated revocation"] as const)( + "%s preserves unresolved retained targets and retries them outside the next addressed page", + async (change) => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + h.live.receive([ + roster(h.relay, "other", [h.viewer.pubkey], 10), + metadata(h.relay, "other", "Other", 10), + ]); + h.live.state({ status: "connected", routes: [] }); + const old = addressed(h, "unsettled", 20); + const newer = Array.from({ length: 50 }, (_, i) => + addressed(h, `new ${i}`, 100 + i), + ); + const initial = deferred(); + const retry = deferred(); + const exactReads: readonly string[][] = []; + let attempt = 0; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return ++attempt === 1 ? [old] : newer; + if (filter?.["#e"]) { + (exactReads as string[][]).push([...filter["#e"]]); + return attempt === 1 ? initial.promise : retry.promise; + } + return []; + }); + const feed = h.session.inboxFeed; + const snapshots: { retained: boolean; incomplete: readonly string[] }[] = + []; + const observe = () => + snapshots.push({ + retained: rows(h).some((row) => row.messageIds.includes(old.id)), + incomplete: feed.snapshot().incomplete, + }); + const stopUnread = h.session.unread.subscribeInbox(observe); + const stopFeed = feed.subscribe(observe); + const first = feed.ensure(); + try { + await vi.waitFor(() => expect(exactReads).toHaveLength(1)); + expect(rows(h).map((row) => row.preview)).toContain("unsettled"); + if (change === "disconnect") + h.live.state({ status: "retrying", routes: [] }); + else h.live.receive([roster(h.relay, "other", [], 30)]); + initial.resolve([]); + await first; + expect(feed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [old.id], + }); + expect(rows(h).map((row) => row.preview)).toContain("unsettled"); + expect( + snapshots.filter(({ retained }) => retained).length, + ).toBeGreaterThan(0); + expect( + snapshots.filter( + ({ retained, incomplete }) => + retained && !incomplete.includes(old.id), + ), + ).toEqual([]); + if (change === "disconnect") + h.live.state({ status: "connected", routes: [] }); + const next = feed.refresh(); + await vi.waitFor(() => expect(exactReads).toHaveLength(2)); + expect(exactReads[1]).toContain(old.id); + expect(exactReads[1]).toHaveLength(51); + expect(feed.snapshot().incomplete).toContain(old.id); + retry.resolve([]); + await next; + expect(feed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + } finally { + stopUnread(); + stopFeed(); + initial.resolve([]); + retry.resolve([]); + } + }, +); + +it("refresh checks a previously settled retained edit omitted after disconnect", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + h.live.state({ status: "connected", routes: [] }); + const target = addressed(h, "original", 20); + const edit = signed(h.alice, { + kind: 40003, + created_at: 21, + content: "deleted edit", + tags: [ + ["h", "room"], + ["e", target.id], + ], + }); + const deletion = signed(h.alice, { + kind: 9005, + created_at: 22, + content: "", + tags: [["e", edit.id]], + }); + const held = deferred(); + let attempt = 0; + let tombstoneReads = 0; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) { + attempt++; + return [target]; + } + if (filter?.until !== undefined) return []; + if (filter?.["#e"]?.includes(target.id)) return attempt === 1 ? [edit] : []; + if (filter?.["#e"]?.includes(edit.id)) { + tombstoneReads++; + return attempt === 1 ? [] : held.promise; + } + return []; + }); + const feed = h.session.inboxFeed; + await feed.ensure(); + expect(feed.snapshot()).toMatchObject({ status: "ready", incomplete: [] }); + expect(rows(h)[0]?.preview).toBe("deleted edit"); + h.live.state({ status: "retrying", routes: [] }); + h.live.state({ status: "connected", routes: [] }); + const next = feed.refresh(); + try { + await vi.waitFor(() => expect(tombstoneReads).toBe(2)); + expect(feed.snapshot()).toMatchObject({ + status: "loading", + incomplete: [target.id], + }); + held.resolve([deletion]); + await next; + expect(rows(h)[0]?.preview).toBe("original"); + expect(feed.snapshot()).toMatchObject({ status: "ready", incomplete: [] }); + } finally { + held.resolve([]); + } +}); + +it("a fully visibility-filtered auxiliary page is incomplete, not a false terminal page", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "original", 20); + const unretained = message(h.alice, "room", "another removal target", 21); + const deletion = signed(h.alice, { + kind: 5, + created_at: 30, + content: "", + tags: [ + ["h", "room"], + ["e", target.id], + ["e", unretained.id], + ], + }); + const olderEdit = signed(h.alice, { + kind: 40003, + created_at: 25, + content: "older applicable update", + tags: [ + ["h", "room"], + ["e", target.id], + ], + }); + const mentions = deferred(); + const auxiliary: ReadFilter[] = []; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return mentions.promise; + if (filter?.["#e"]?.includes(target.id)) { + auxiliary.push(filter); + if (filter.until === undefined) return [deletion]; + return filter.until === 30 ? [olderEdit] : []; + } + return []; + }); + const feed = h.session.inboxFeed; + const work = feed.ensure(); + try { + await vi.waitFor(() => expect(h.query).toHaveBeenCalledTimes(1)); + expect(rows(h)).toEqual([]); + mentions.resolve([target]); + await work; + // Fail closed rather than claiming the unadmitted page proved exhaustion. + expect(feed.snapshot()).toMatchObject({ + status: "error", + incomplete: [target.id], + error: + "Inbox message updates could not be verified for current access. Retry inbox.", + }); + expect(rows(h)[0]?.preview).toBe("original"); + expect(auxiliary).toHaveLength(1); + // Existing shared admission can later supply the missing reference. Retry + // must then traverse the older applicable update and a real empty terminal. + h.live.receive([unretained]); + await feed.refresh(); + expect(auxiliary.map((filter) => filter.until)).toEqual([ + undefined, + undefined, + 30, + 25, + ]); + expect(feed.snapshot()).toMatchObject({ status: "ready", incomplete: [] }); + expect(rows(h)).toEqual([]); // Author's bulk deletion is now fully admitted. + } finally { + mentions.resolve([]); + } +}); + +it("a signed short-page walk settles edits and checks their deletion dependencies", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "original", 20); + const edits = Array.from({ length: 3 }, (_, i) => + signed(h.alice, { + kind: 40003, + created_at: 33 - i, + content: `revision ${i}`, + tags: [ + ["h", "room"], + ["e", target.id], + ], + }), + ); + const pages: ReadFilter[] = []; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [target]; + if (filter?.["#e"]?.includes(target.id)) { + pages.push(filter); + return edits + .filter( + (edit) => + filter.until === undefined || edit.created_at < filter.until, + ) + .slice(0, 1); + } + return []; + }); + await h.session.inboxFeed.ensure(); + expect(pages.map((filter) => filter.until)).toEqual([undefined, 33, 32, 31]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + expect(rows(h)[0]?.preview).toBe("revision 0"); + expect( + h.query.mock.calls.some( + ([filters]) => + filters[0]?.kinds?.join(",") === "5,9005" && + filters[0]?.["#e"]?.length === 3, + ), + ).toBe(true); +}); + +// Retention budgets belong to the feed, not signing/session folding. Synthetic +// DTOs model the already-verified reader boundary here; signed admission, edits, +// and short-page closure remain covered through the real session above. +it.each(["events", "bytes"] as const)( + "walks advancing auxiliary pages with a discriminating %s budget outcome", + async (budget) => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "original", 20); + const count = budget === "events" ? 2001 : 9; + const edits: RelayEvent[] = Array.from({ length: count }, (_, i) => ({ + ...target, + id: (i + 1).toString(16).padStart(64, "0"), + kind: 40003, + created_at: 30 + count - i, + content: budget === "bytes" ? "x".repeat(512 * 1024) : `revision ${i}`, + tags: [ + ["h", "room"], + ["e", target.id], + ], + })); + expect(edits.length > 2000).toBe(budget === "events"); + expect(byteSize(edits) > 4 * 1024 * 1024).toBe(budget === "bytes"); + const pageSize = budget === "events" ? 500 : 3; + const pages: ReadFilter[] = []; + const reader = { + read: vi.fn(async (filters: readonly ReadFilter[]) => { + const filter = filters[0]; + if (!filter?.["#e"]?.includes(target.id)) return []; + pages.push(filter); + return edits + .filter( + (edit) => + filter.until === undefined || edit.created_at < filter.until, + ) + .slice(0, pageSize); + }), + }; + const feed = createInboxFeed({ + viewer: h.viewer.pubkey, + channels: h.session.channels, + reader, + async addressedRead(_filter, _signal, prepare) { + prepare([target]); + return [target]; + }, + retainedEvent: (id) => (id === target.id ? target : undefined), + retainedEditIds: () => [], + }); + try { + await feed.ensure(); + expect(pages.length).toBe(budget === "events" ? 5 : 3); + expect(feed.snapshot()).toMatchObject({ + status: "error", + incomplete: [target.id], + error: "Inbox message updates exceed the read budget. Retry inbox.", + }); + // Failure must be retryable; no blanket pagination rejection may pass. + pages.length = 0; + const first = edits[0]; + if (!first) throw Error("Missing budget fixture edit"); + edits.splice(0, edits.length, { ...first, content: "bounded retry" }); + await feed.refresh(); + expect(pages).toHaveLength(2); + expect(feed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + expect(reader.read.mock.calls.at(-1)?.[0][0]?.["#e"]).toEqual([first.id]); + } finally { + feed.dispose(); + } + }, +); + +it("access purge drops denied obligations but retains the still-readable target", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + h.live.receive([ + roster(h.relay, "other", [h.viewer.pubkey], 10), + metadata(h.relay, "other", "Other", 10), + ]); + const kept = addressed(h, "kept", 20); + const denied = message(h.alice, "other", "denied", 21, [ + ["p", h.viewer.pubkey], + ]); + const held = deferred(); + let reads = 0; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [kept, denied]; + if (filter?.["#e"]) { + reads++; + return held.promise; + } + return []; + }); + const feed = h.session.inboxFeed; + const work = feed.ensure(); + try { + await vi.waitFor(() => expect(reads).toBe(1)); + expect(feed.snapshot().incomplete).toEqual( + expect.arrayContaining([kept.id, denied.id]), + ); + expect(rows(h)).toHaveLength(2); + h.live.receive([roster(h.relay, "other", [], 30)]); + held.resolve([]); + await work; + expect(rows(h).map((row) => row.messageId)).toEqual([kept.id]); + expect(feed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [kept.id], + }); + await h.clearCache(); + expect(rows(h)).toEqual([]); + expect(feed.snapshot().incomplete).toEqual([]); + } finally { + held.resolve([]); + } +}); + +it.each(["revoke-regrant", "dispose"] as const)( + "an auxiliary response cannot re-admit evidence after %s", + async (change) => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "original", 20); + const edit = signed(h.alice, { + kind: 40003, + created_at: 21, + content: "stale update", + tags: [ + ["h", "room"], + ["e", target.id], + ], + }); + const held = deferred(); + let reads = 0; + h.query.mockImplementation(async ([filter]) => { + if (filter?.["#p"]) return [target]; + if (filter?.["#e"]) { + reads++; + return held.promise; + } + return []; + }); + const feed = h.session.inboxFeed; + const work = feed.ensure(); + try { + await vi.waitFor(() => expect(reads).toBe(1)); + expect(rows(h)[0]?.preview).toBe("original"); + if (change === "revoke-regrant") { + h.admit([h.alice.pubkey], 30); + h.admit([h.viewer.pubkey, h.alice.pubkey], 31); + } else await h[change](); + held.resolve([edit]); + await work; + expect(rows(h)).toEqual([]); + expect(feed.snapshot()).toMatchObject({ + status: "idle", + incomplete: [], + }); + expect(reads).toBe(1); + } finally { + held.resolve([]); + } + }, +); + +it("retained dependency closure checks only live author edits of the exact addressed target", async () => { + const h = setup(); + h.admit([h.viewer.pubkey, h.alice.pubkey]); + const target = addressed(h, "original", 20); + const unrelated = message(h.alice, "room", "not addressed", 21); + const edit = ( + author: typeof h.alice, + id: string, + content: string, + at: number, + ) => + signed(author, { + kind: 40003, + created_at: at, + content, + tags: [ + ["h", "room"], + ["e", id], + ], + }); + const originalEdit = edit(h.alice, target.id, "surviving older edit", 22); + const newestEdit = edit(h.alice, target.id, "newest edit", 23); + const removedEdit = edit(h.alice, target.id, "already deleted", 24); + const wrongAuthor = edit(h.viewer, target.id, "not an author edit", 25); + const unrelatedEdit = edit(h.alice, unrelated.id, "another target", 26); + const deletion = signed(h.alice, { + kind: 5, + created_at: 27, + content: "", + tags: [["e", removedEdit.id]], + }); + // All are already retained; the relay omits them from this finite query. + h.live.receive([ + target, + unrelated, + originalEdit, + newestEdit, + removedEdit, + wrongAuthor, + unrelatedEdit, + deletion, + ]); + expect(rows(h)[0]?.preview).toBe("newest edit"); + const filters: ReadFilter[] = []; + h.query.mockImplementation(async ([filter]) => { + if (!filter) return []; + filters.push(filter); + return filter["#p"] ? [target] : []; + }); + await h.session.inboxFeed.ensure(); + expect(filters.filter((filter) => filter["#e"])).toEqual([ + { kinds: [40003, 5, 9005], "#e": [target.id], limit: 500 }, + { + kinds: [5, 9005], + "#e": [originalEdit.id, newestEdit.id].sort(), + limit: 500, + }, + ]); + expect(h.session.inboxFeed.snapshot()).toMatchObject({ + status: "ready", + incomplete: [], + }); + expect(rows(h)[0]?.preview).toBe("newest edit"); +}); diff --git a/src/features/relay/inbox-feed.ts b/src/features/relay/inbox-feed.ts new file mode 100644 index 000000000..601ce4a97 --- /dev/null +++ b/src/features/relay/inbox-feed.ts @@ -0,0 +1,253 @@ +import { byteSize } from "./budget"; +import type { ChannelQueries } from "./contracts"; +import type { RelayEvent, ReadFilter } from "./events"; +import type { RelayReader } from "./reader"; + +const mentionKinds = [9, 40002]; +const cursorOf = (event: RelayEvent) => ({ + until: event.created_at, + before_id: event.id, +}); +const older = (event: RelayEvent, cursor: ReturnType) => + event.created_at < cursor.until || + (event.created_at === cursor.until && event.id > cursor.before_id); +export type InboxFeedSnapshot = Readonly<{ + status: "idle" | "loading" | "ready" | "error"; + error?: string | undefined; + /** Exact addressed targets whose stored edits/deletions are not settled. */ + incomplete: readonly string[]; +}>; + +/** The session owns finite, verified, viewer-addressed history. No new live route, + * polling, or parallel channel store. Current membership gates completeness targets. */ +export function createInboxFeed({ + viewer, + channels, + reader, + addressedRead, + retainedEvent, + retainedEditIds, + notify = (listener) => listener(), +}: { + viewer: string; + channels: ChannelQueries; + reader: RelayReader; + /** Private lookups in the existing bounded unread evidence owner. */ + retainedEvent: (id: string) => RelayEvent | undefined; + retainedEditIds: (ids: readonly string[]) => readonly string[]; + /** Verified session admission calls prepare before publishing unread evidence. */ + addressedRead: ( + filter: ReadFilter, + signal: AbortSignal, + prepare: (events: readonly RelayEvent[]) => void, + ) => Promise; + notify?: (listener: () => void) => void; +}) { + let closed = false; + let epoch = 0; + let work: Promise | undefined; + let controller: AbortController | undefined; + let requested = false; + let snapshot: InboxFeedSnapshot = Object.freeze({ + status: "idle", + incomplete: [], + }); + const listeners = new Set<() => void>(); + const admitted = (event: RelayEvent) => { + if (!event.tags.some(([name, value]) => name === "p" && value === viewer)) + return false; + const hs = event.tags.filter(([name]) => name === "h"); + return ( + hs.length === 0 || + (hs.length === 1 && + channels + .list() + .channels.some( + (channel) => + channel.id === hs[0]?.[1] && + !channel.cached && + channel.members?.includes(viewer), + )) + ); + }; + function project(events: readonly RelayEvent[]) { + return Object.freeze( + events.filter( + (event) => mentionKinds.includes(event.kind) && admitted(event), + ), + ); + } + function publish(patch: Partial) { + if (closed) return; + snapshot = Object.freeze({ + ...snapshot, + ...patch, + }); + for (const listener of listeners) notify(listener); + } + async function overlays( + ids: readonly string[], + kinds: readonly number[], + signal: AbortSignal, + ) { + if (!ids.length) return []; + const collected: RelayEvent[] = []; + let cursor: ReturnType | undefined; + for (;;) { + const page = [ + ...(await reader.read( + [ + { + kinds, + "#e": ids, + limit: 500, + ...(cursor ?? {}), + }, + ], + { signal }, + )), + ].sort((a, b) => b.created_at - a.created_at || a.id.localeCompare(b.id)); + const previous = cursor; + if (previous && page.some((event) => !older(event, previous))) + throw new Error("Inbox message updates did not advance. Retry inbox."); + collected.push(...page); + if (collected.length > 2000 || byteSize(collected) > 4 * 1024 * 1024) + throw new Error( + "Inbox message updates exceed the read budget. Retry inbox.", + ); + const last = page.at(-1); + if (!last) return collected; + cursor = cursorOf(last); + // Authorized responses can have short pages. Only empty ends the walk. + } + } + async function refresh() { + if (closed) return; + if (work) return work; + requested = true; + const generation = epoch; + const owned = new AbortController(); + controller = owned; + // Do not clear incomplete targets on retry: the new first read can re-admit + // the old body before its stored overlays are checked again. + work = (async () => { + try { + // Bounded addressed chat history; nonchat and approvals belong elsewhere. + const mentionFilter: ReadFilter = { + kinds: mentionKinds, + "#p": [viewer], + limit: 50, + }; + const mentions = await addressedRead( + mentionFilter, + owned.signal, + (events) => { + if (closed || generation !== epoch || owned.signal.aborted) return; + const ids = project(events).map((event) => event.id); + if (ids.length) + publish({ + incomplete: Object.freeze([ + ...new Set([...snapshot.incomplete, ...ids]), + ]), + }); + }, + ); + if (closed || generation !== epoch) return; + // Retry older failed targets even when newer addressed rows pushed them + // beyond the latest page; never clear an unchecked exact ID. + const targets = [ + ...new Set([ + ...snapshot.incomplete, + ...project(mentions).map((event) => event.id), + ]), + ]; + const updates = await overlays(targets, [40003, 5, 9005], owned.signal); + const edits = updates.filter((event) => event.kind === 40003); + await overlays( + [ + ...new Set([ + ...retainedEditIds(targets), + ...edits.map((event) => event.id), + ]), + ], + [5, 9005], + owned.signal, + ); + if (closed || generation !== epoch) return; + publish({ + status: "ready", + incomplete: [], + error: undefined, + }); + } catch (error) { + if (closed || generation !== epoch || owned.signal.aborted) return; + publish({ + status: "error", + error: + error instanceof Error + ? error.message + : "Inbox history unavailable", + }); + } + })().finally(() => { + if (controller === owned) { + controller = undefined; + work = undefined; + } + }); + if (generation === epoch && controller === owned) + publish({ status: "loading", error: undefined }); + return work; + } + function clear(incomplete: readonly string[] = []) { + epoch++; + controller?.abort(); + controller = undefined; + work = undefined; + // Previously demanded data recovers explicitly or through reconnect. + publish({ status: "idle", incomplete, error: undefined }); + } + const retainedIncomplete = () => + Object.freeze(snapshot.incomplete.filter((id) => retainedEvent(id))); + return Object.freeze({ + snapshot: () => snapshot, + subscribe(listener: () => void) { + listeners.add(listener); + return () => { + listeners.delete(listener); + }; + }, + ensure: () => + work ?? (snapshot.status === "idle" ? refresh() : Promise.resolve()), + refresh, + stale() { + epoch++; + controller?.abort(); + controller = undefined; + work = undefined; + publish({ + status: "idle", + incomplete: retainedIncomplete(), + error: undefined, + }); + }, + reconnect() { + if (requested) void refresh(); + }, + clear: () => clear(), + purge: () => clear(retainedIncomplete()), + dispose() { + closed = true; + epoch++; + controller?.abort(); + controller = undefined; + work = undefined; + // Retained capabilities expose no retired content; do not notify dead consumers. + snapshot = Object.freeze({ + status: "idle", + incomplete: [], + }); + listeners.clear(); + }, + }); +} diff --git a/src/features/relay/inbox.ts b/src/features/relay/inbox.ts new file mode 100644 index 000000000..2652c2c9d --- /dev/null +++ b/src/features/relay/inbox.ts @@ -0,0 +1,29 @@ +import type { ReadTarget } from "./read-state-model"; + +/** A conversation projected from the unread owner's bounded verified evidence. */ +export type InboxItem = Readonly<{ + id: string; + channelId: string; + target: ReadTarget; + /** Oldest observed unread message, otherwise newest relevant message. */ + messageId: string; + latestMessageId: string; + /** Exact verified group members, for pending preview evidence across regrouping. */ + messageIds: readonly string[]; + rootId?: string; + authorId: string; + preview: string; + createdAt: number; + mentioned: boolean; + thread: boolean; + unreadCount: number; + manual: boolean; + /** Explicit prefixes; a thread prefix never acknowledges its top-level root. */ + readThrough: readonly Readonly<{ target: ReadTarget; messageId: string }>[]; +}>; +export type InboxSnapshot = Readonly<{ + status: "idle" | "loading" | "ready" | "error"; + items: readonly InboxItem[]; + freshness: "unknown" | "observed" | "stale"; + error?: string | undefined; +}>; diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts index 2d0f8fc81..5198aee24 100644 --- a/src/features/relay/session.ts +++ b/src/features/relay/session.ts @@ -51,6 +51,7 @@ import { } from "./read-state-storage"; import { createTyping } from "./typing"; import { createUnread } from "./unread"; +import { createInboxFeed } from "./inbox-feed"; import type { IncomingListener, IncomingMessage } from "./incoming"; import { objectBody } from "./body"; import type { ChannelList, ChannelSummary } from "./contracts"; @@ -389,6 +390,7 @@ export function createRelaySession( for (const purge of views.values()) purge(); commit(); unread.purge(); + inboxFeed.purge(); } finally { if (--revoking === 0) { const pending = [...notifications]; @@ -401,6 +403,10 @@ export function createRelaySession( filters: readonly ReadFilter[], settings?: ReadOptions, channelTraffic = true, + beforeInbox?: ( + events: readonly RelayEvent[], + raw: readonly RelayEvent[], + ) => void, ) { if ( options.cachedOnly || @@ -474,7 +480,12 @@ export function createRelaySession( if (closed || epoch !== accessEpoch) throw new DOMException("Stale relay search", "AbortError"); } - const visible = accept(events, channelTraffic); + const visible = accept( + events, + channelTraffic, + beforeInbox, + settings?.signal, + ); // Discovery must see signed grants/removals even when their channel is // currently denied; only the store interprets roster completeness. return channelTraffic @@ -488,6 +499,11 @@ export function createRelaySession( function accept( events: readonly RelayEvent[], channelTraffic = true, + beforeInbox?: ( + events: readonly RelayEvent[], + raw: readonly RelayEvent[], + ) => void, + signal?: AbortSignal, ): readonly RelayEvent[] { if (closed) return []; // Authority precedes every projection, even in a batch containing both a @@ -508,6 +524,12 @@ export function createRelaySession( ) .filter(visibility(events, true)); const epoch = accessEpoch; + beforeInbox?.(visible, events); + if ( + beforeInbox && + (closed || epoch !== accessEpoch || signal?.aborted || cacheClearing) + ) + throw new DOMException("Stale Inbox read", "AbortError"); typing.accept(visible); if (closed || epoch !== accessEpoch) return []; profiling.measure( @@ -792,6 +814,33 @@ export function createRelaySession( viewer: transport?.viewer ?? "", notify, }); + const inboxFeed = createInboxFeed({ + // A withheld auxiliary event is not proof of an exhausted history page. + // Fail closed at raw/admitted admission; do not relax reference visibility. + reader: { + read: (filters, settings) => + readVerified(filters, settings, true, (visible, raw) => { + const admitted = new Set(visible.map((event) => event.id)); + if ( + raw.some( + (event) => + [40003, 5, 9005].includes(event.kind) && + !admitted.has(event.id), + ) + ) + throw new Error( + "Inbox message updates could not be verified for current access. Retry inbox.", + ); + }), + }, + retainedEvent: unread.event, + retainedEditIds: unread.retainedEditIds, + addressedRead: (filter, signal, prepare) => + readVerified([filter], { signal }, true, prepare), + channels: channels.queries, + viewer: transport?.viewer ?? "", + notify, + }); let traffic: LiveSubscription | undefined; const liveListeners = new Set<() => void>(); let liveSnapshot: LiveSnapshot = Object.freeze({ @@ -1677,6 +1726,7 @@ export function createRelaySession( statuses: statuses.queries, agentLibrary: agentLibrary.queries, agentChoices, + inboxFeed, workflows: workflows.capability, projects, projectGit: transport?.projectGit @@ -2242,6 +2292,7 @@ export function createRelaySession( workflows.interrupt(); channels.staleHeads(); unread.stale(); + inboxFeed.stale(); } // Access-revoked CLOSED is a refresh hint, not signed archive/membership // authority. Aggregate snapshots repeat failures; only react to a new one. @@ -2289,6 +2340,7 @@ export function createRelaySession( emoji.reconnect(); statuses.reconnect(); unread.reconnect(); + inboxFeed.reconnect(); for (const refresh of refreshers) void refresh(); } }, 0); @@ -2389,6 +2441,7 @@ export function createRelaySession( for (const clear of views.values()) clear(true); recent.clear(); unread.clear(); + inboxFeed.clear(); requests.invalidate(); profiles.clear(); emoji.clear(); @@ -2427,6 +2480,7 @@ export function createRelaySession( for (const timer of timers) clearTimeout(timer); for (const dispose of [...views.keys()]) dispose(); unread.dispose(); + inboxFeed.dispose(); writes?.dispose(); requests.dispose(); channels.dispose(); diff --git a/src/features/relay/unread-lookups.test.ts b/src/features/relay/unread-lookups.test.ts index ba719c5f2..d4561518f 100644 --- a/src/features/relay/unread-lookups.test.ts +++ b/src/features/relay/unread-lookups.test.ts @@ -204,3 +204,103 @@ it("lookup results survive a full-window reset and refresh without asking again" expect(attention()).toMatchObject({ category: "thread", unread: true }); expect(lookups()).toBe(count); }); + +it("Inbox-only demand resolves direct conversation participation without counting lookup witnesses", async () => { + const { owner, store, reader, lookups } = await setup(); + const parent = event(peer, 9, 10, []); + const mine = reply(viewer, 11, parent); + const sibling = reply(peer, 20, parent); + const nestedParent = reply(peer, 21, parent); + const nested = reply(peer, 22, nestedParent); + store.push(parent, mine, nestedParent); + owner.accept([sibling, nested]); + const published: string[][] = []; + const stop = owner.capability.subscribeInbox(() => { + published.push( + owner.capability.inbox().items.flatMap((item) => item.messageIds), + ); + }); + try { + expect(owner.capability.inbox().items).toEqual([]); + await vi.waitFor(() => { + const items = owner.capability.inbox().items; + expect(items).toHaveLength(1); + expect(items[0]).toMatchObject({ + id: `c0:${parent.id}`, + messageId: sibling.id, + messageIds: [sibling.id], + rootId: parent.id, + thread: true, + unreadCount: 1, + target: { kind: "message", channelId: "c0", messageId: sibling.id }, + }); + }); + expect(published).toContainEqual([sibling.id]); + expect(lookups()).toBeGreaterThan(0); + expect( + reader.read.mock.calls.some(([[filter]]) => + filter?.authors?.includes(viewer), + ), + ).toBe(true); + // A fetched parent is structural, not a counted message eligible for a + // root-scoped manual action. The existing message fallback is actionable. + const item = owner.capability.inbox().items[0]; + assert(item); + expect(item.readThrough).toEqual([ + { + target: { kind: "thread", channelId: "c0", rootId: parent.id }, + messageId: sibling.id, + }, + ]); + await owner.capability.markUnreadLocal(item.target); + expect(owner.capability.inbox().items[0]?.manual).toBe(true); + await owner.capability.clearUnreadLocal(item.target); + expect(owner.capability.inbox().items[0]).toMatchObject({ + unreadCount: 1, + manual: false, + }); + // Fresh counted root evidence can safely promote the context-menu target. + owner.accept([parent]); + expect(owner.capability.inbox().items[0]?.target).toEqual({ + kind: "thread", + channelId: "c0", + rootId: parent.id, + }); + } finally { + stop(); + } +}); + +it("Inbox retains a relevant direct reply as thread activity when the root cannot be fetched", async () => { + const { owner, store } = await setup(); + const root = event(peer, 9, 10, []); + const parent = reply(viewer, 11, root); + const response = event(peer, 9, 20, [ + ["e", root.id, "", "root"], + ["e", parent.id, "", "reply"], + ]); + store.push(parent); // Root genuinely unavailable, not another counted row. + owner.accept([response]); + const stop = owner.capability.subscribeInbox(() => {}); + try { + await vi.waitFor(() => + expect(owner.capability.inbox().items).toHaveLength(1), + ); + const item = owner.capability.inbox().items[0]; + assert(item); + expect(item).toMatchObject({ + id: `c0:${response.id}`, + messageIds: [response.id], + thread: true, + unreadCount: 1, + }); + expect(item.rootId).toBeUndefined(); + expect(item.target).toEqual({ + kind: "message", + channelId: "c0", + messageId: response.id, + }); + } finally { + stop(); + } +}); diff --git a/src/features/relay/unread.test.ts b/src/features/relay/unread.test.ts index 4e38b9d96..d72c21303 100644 --- a/src/features/relay/unread.test.ts +++ b/src/features/relay/unread.test.ts @@ -2249,3 +2249,305 @@ it.each(["channel", "thread", "message"] as const)( expect(h.journal()?.localUnread).toEqual({}); }, ); + +it("inbox groups relevant conversations, preserves read rows and exact unread resume points", async () => { + const h = setup(); + h.grant("room"); + h.grant("dm"); + h.emit([metadata(h.relay, "dm", "DM", 11, [["t", "dm"]])]); + const root = message(h.viewer, "room", "My thread", 20); + const mention = message(h.alice, "room", "Mention", 21, [ + ["p", h.viewer.pubkey], + ]); + const reply = message(h.alice, "room", "First reply", 22, [ + ["e", root.id, "", "reply"], + ]); + const newer = message(h.alice, "room", "Mentioned reply", 23, [ + ["e", root.id, "", "reply"], + ["p", h.viewer.pubkey], + ]); + const dm = message(h.alice, "dm", "Direct hello", 24); + h.emit([ + root, + mention, + reply, + newer, + dm, + message(h.alice, "room", "Not relevant", 25), + ]); + const unread = h.session.unread; + const before = unread.inbox(); + expect(before.items).toHaveLength(3); + const thread = before.items.find((item) => item.thread); + assert(thread); + expect(thread).toMatchObject({ + messageId: reply.id, + latestMessageId: newer.id, + mentioned: true, + unreadCount: 2, + preview: "First reply", + }); + expect(before.items[0]).toMatchObject({ + channelId: "dm", + target: { kind: "channel", channelId: "dm" }, + }); + expect(unread.inbox()).toBe(before); + const listener = vi.fn(); + const stop = unread.subscribeInbox(listener); + h.emit([newer]); + expect(listener).not.toHaveBeenCalled(); + expect(unread.inbox()).toBe(before); + await unread.markThrough(thread.target, thread.latestMessageId); + const read = unread.inbox().items.find((item) => item.id === thread.id); + assert(read); + expect(read).toMatchObject({ unreadCount: 0, messageId: newer.id }); + expect(unread.activity("room").items).toHaveLength(0); + await unread.markUnreadLocal(thread.target); + expect( + unread.inbox().items.find((item) => item.id === thread.id), + ).toMatchObject({ unreadCount: 0, manual: true }); + await unread.markThrough(thread.target, thread.latestMessageId); + expect( + unread.inbox().items.find((item) => item.id === thread.id)?.manual, + ).toBe(false); + expect(unread.inbox().items).toHaveLength(3); + expect( + unread.snapshot({ kind: "channel", channelId: "room" }).observedCount, + ).toBe(2); + stop(); +}); + +it("inbox folds edits and deletions and revokes all evidence before a reentrant subscriber", async () => { + const h = setup(); + h.grant("room"); + const row = message(h.alice, "room", "original", 20, [ + ["p", h.viewer.pubkey], + ]); + h.emit([row]); + expect(h.session.unread.inbox().items[0]?.preview).toBe("original"); + h.emit([ + signed(h.alice, { + kind: 40003, + created_at: 21, + content: "edited", + tags: [ + ["h", "room"], + ["e", row.id], + ], + }), + ]); + expect(h.session.unread.inbox().items[0]?.preview).toBe("edited"); + const noticed: number[] = []; + h.session.unread.subscribe(h.target, () => + noticed.push(h.session.unread.inbox().items.length), + ); + h.emit([roster(h.relay, "room", [], 30)]); + expect(noticed.at(-1)).toBe(0); + expect(h.session.unread.inbox().items).toHaveLength(0); + h.grant("room", 31); + expect(h.session.unread.inbox().items).toHaveLength(0); + h.emit([row]); + h.emit([ + signed(h.alice, { + kind: 5, + created_at: 32, + content: "", + tags: [["e", row.id]], + }), + ]); + expect(h.session.unread.inbox().items).toHaveLength(0); + h.dispose(); + expect(h.session.unread.inbox().items).toHaveLength(0); +}); + +it("inbox keeps unresolved mentions exact, joins a verified root, and never clears unrelated roots", async () => { + const h = setup(); + h.grant("room"); + const root = message(h.alice, "room", "Mentioned root", 20, [ + ["p", h.viewer.pubkey], + ]); + const reply = message(h.alice, "room", "Mentioned reply", 21, [ + ["p", h.viewer.pubkey], + ["e", root.id, "", "reply"], + ]); + const other = message(h.alice, "room", "Other mention", 22, [ + ["p", h.viewer.pubkey], + ]); + h.emit([reply, other]); + expect( + h.session.unread.inbox().items.find((item) => item.messageId === reply.id) + ?.target.kind, + ).toBe("message"); + h.emit([root]); + const item = h.session.unread.inbox().items.find((item) => item.thread); + assert(item); + expect(h.session.unread.inbox().items).toHaveLength(2); + expect(item.readThrough.map((step) => step.target.kind)).toEqual([ + "message", + "thread", + ]); + await h.session.unread.markUnreadLocal({ + kind: "message", + channelId: "room", + messageId: root.id, + }); + for (const step of item.readThrough) + await h.session.unread.markThrough(step.target, step.messageId); + expect( + h.session.unread.inbox().items.find((row) => row.id === item.id), + ).toMatchObject({ unreadCount: 0, manual: false }); + expect( + h.session.unread.inbox().items.find((row) => row.messageId === other.id) + ?.unreadCount, + ).toBe(1); +}); + +it("inbox observation exposes empty, failure, recovery and cache clear without a new read owner", async () => { + const h = setup(); + h.grant("room"); + expect(h.session.unread.inbox().status).toBe("idle"); + let release!: () => void; + const hold = new Promise((resolve) => { + release = resolve; + }); + h.query.mockImplementation(async (filters) => { + if (filters[0]?.kinds?.includes(9)) { + await hold; + throw new Error("offline"); + } + return []; + }); + const work = h.session.unread.ensure(); + try { + await vi.waitFor(() => + expect( + h.query.mock.calls.some(([filters]) => filters[0]?.kinds?.includes(9)), + ).toBe(true), + ); + expect(h.session.unread.inbox().status).toBe("loading"); + } finally { + release(); + } + await work; + expect(h.session.unread.inbox()).toMatchObject({ + status: "error", + error: "offline", + }); + h.query.mockResolvedValue([]); + await h.session.unread.refresh(); + expect(h.session.unread.inbox()).toMatchObject({ + status: "ready", + items: [], + }); + h.emit([message(h.alice, "room", "fresh", 20, [["p", h.viewer.pubkey]])]); + expect(h.session.unread.inbox().items).toHaveLength(1); + await h.clearCache(); + expect(h.session.unread.inbox().items).toHaveLength(0); +}); + +it("a reply surviving root deletion can be marked unread and then cleared", async () => { + const h = setup(); + h.grant("room"); + const root = message(h.alice, "room", "root", 20); + const reply = message(h.alice, "room", "surviving mention", 21, [ + ["e", root.id, "", "reply"], + ["p", h.viewer.pubkey], + ]); + h.emit([ + root, + reply, + signed(h.alice, { + kind: 5, + content: "", + created_at: 22, + tags: [["e", root.id]], + }), + ]); + const item = h.session.unread.inbox().items[0]; + assert(item); + expect(item).toMatchObject({ + target: { kind: "message", messageId: reply.id }, + rootId: root.id, + }); + for (const step of item.readThrough) + await h.session.unread.markThrough(step.target, step.messageId); + await h.session.unread.markUnreadLocal(item.target); + const marked = h.session.unread.inbox().items[0]; + assert(marked); + expect(marked.manual).toBe(true); + for (const step of marked.readThrough) + await h.session.unread.markThrough(step.target, step.messageId); + expect(h.session.unread.inbox().items[0]).toMatchObject({ + manual: false, + unreadCount: 0, + }); +}); + +it("a prepared channel read retries its captured cut, not later arrivals or the retry clock", async () => { + const h = setup(); + h.grant("room"); + const before = message(h.alice, "room", "before click", 11); + h.emit([before]); + const unread = h.session.unread; + await unread.markUnreadLocal(h.target); + clock(20); + const retry = unread.prepareChannelRead("room"); + h.failSave(); + await expect(retry()).rejects.toThrow("disk full"); + expect(h.snapshot()).toMatchObject({ + observedCount: 1, + manual: "local-only", + }); + const later = message(h.alice, "room", "after click", 25); + h.emit([later]); + clock(30); + await retry(); + expect(h.journal()?.state.frontiers.room).toBe(20); + expect(unread.attention("room", before.id).unread).toBe(false); + expect(unread.attention("room", later.id).unread).toBe(true); + expect(h.snapshot()).toMatchObject({ observedCount: 1, manual: "none" }); + await unread.markChannelRead("room"); + expect(h.journal()?.state.frontiers.room).toBe(30); + expect(h.snapshot().observedCount).toBe(0); +}); + +it("prepared channel reads serialize each invocation with existing channel mutations", async () => { + const h = setup(); + h.grant("room"); + h.emit([message(h.alice, "room", "before click", 11)]); + const unread = h.session.unread; + await unread.ensure(); + clock(20); + const retry = unread.prepareChannelRead("room"); + const held = h.holdSaveStarted(); + const first = retry(); + try { + await held.started; + const mark = unread.markUnreadLocal(h.target); + const last = retry(); + held.release(); + await Promise.all([first, mark, last]); + expect(h.journal()?.state.frontiers.room).toBe(20); + expect(h.snapshot()).toMatchObject({ observedCount: 0, manual: "none" }); + } finally { + held.release(); + } +}); + +it.each(["clearCache", "dispose", "revoke-regrant"] as const)( + "a prepared channel read cannot retry past %s", + async (change) => { + const h = setup(); + h.grant("room"); + h.emit([message(h.alice, "room", "before click", 11)]); + await h.session.unread.markUnreadLocal(h.target); + const before = h.journal(); + const retry = h.session.unread.prepareChannelRead("room"); + if (change === "revoke-regrant") { + h.emit([roster(h.relay, "room", [], 20)]); + h.grant("room", 21); + } else await h[change](); + await expect(retry()).rejects.toThrow(); + expect(h.journal()).toEqual(before); + }, +); diff --git a/src/features/relay/unread.ts b/src/features/relay/unread.ts index d8a5ae098..d02bec4ce 100644 --- a/src/features/relay/unread.ts +++ b/src/features/relay/unread.ts @@ -1,3 +1,4 @@ +import type { InboxItem, InboxSnapshot } from "./inbox"; import type { ChannelQueries } from "./contracts"; import type { RelayEvent } from "./events"; import { @@ -69,6 +70,8 @@ export type ReadingHandle = Readonly<{ dispose(): void; }>; export interface UnreadCapability { + inbox(): InboxSnapshot; + subscribeInbox(listener: () => void): () => void; snapshot(target: ReadTarget): UnreadSnapshot; /** Same verified attention/frontier policy as badges, not a notification event source. */ attention(channelId: string, messageId: string): MessageAttention; @@ -80,6 +83,10 @@ export interface UnreadCapability { ensure(): Promise; refresh(): Promise; retrySync(): Promise; + /** Durable local intent revision, not sync-health changes. */ + revision(): number; + /** Access/cache/connection retirement fence. */ + generation(): number; reading(channelId: string): ReadingHandle; /** Explicit prefix intent, unlike individual-message visibility observations. */ markThrough( @@ -101,6 +108,8 @@ export interface UnreadCapability { enterChannel(channelId: string): Promise; /** Explicit channel prefix through retained verified evidence, including replies. */ markChannelRead(channelId: string): Promise; + /** Capture one prefix/time cut; each invocation queues that same intent for retry. */ + prepareChannelRead(channelId: string): () => Promise; /** `markChannelRead` serialised over every accessible listed channel that still * shows unread evidence or a local mark. Channels with nothing to clear are * skipped, so an already-read community costs no writes. One failing channel @@ -189,6 +198,10 @@ export function createUnread({ const activityListeners = new Map void>>(); const activitySnapshots = new Map(); const activityDirty = new Set(); + const inboxListeners = new Set<() => void>(); + let inboxSnapshot: InboxSnapshot | undefined; + let inboxDirty = true; + let inboxLoading = false; const handles = new Set<() => void>(); const views = new Map< () => void, @@ -724,6 +737,140 @@ export function createUnread({ snapshots.set(key, value); return value; } + function inbox(): InboxSnapshot { + if (inboxSnapshot && !inboxDirty) return inboxSnapshot; + inboxDirty = false; + indexEvidence(); + const state = reads.state(); + const items: InboxItem[] = []; + if (!closed) + for (const channel of channels.list().channels) { + if (channel.cached || !channel.members?.includes(viewer)) continue; + const dm = channel.channelType === "dm"; + const groups = new Map(); + for (const entry of byChannel.get(channel.id) ?? []) { + want(entry, dm); + if (entry.event.pubkey === viewer || !category(entry, dm)) continue; + const id = dm ? channel.id : (entry.rootId ?? entry.event.id); + const group = groups.get(id) ?? []; + group.push(entry); + groups.set(id, group); + } + if (!groups.size) continue; + const content = new Map( + foldMessages(channel.id, "", [...events.values()], { + includeReplies: true, + }).map((row) => [row.id, row.content]), + ); + for (const [id, entries] of groups) { + entries.sort( + (a, b) => + a.event.created_at - b.event.created_at || + a.event.id.localeCompare(b.event.id), + ); + const latest = entries[entries.length - 1]; + if (!latest) continue; + const unread = entries.filter((entry) => isUnread(entry, state, dm)); + const representative = unread[0] ?? latest; + const replies = entries.filter((entry) => entry.rootId !== undefined); + const lastReply = replies[replies.length - 1]; + const target: ReadTarget = dm + ? { kind: "channel", channelId: channel.id } + : lastReply?.rootId && + events.has(lastReply.rootId) && + !tombstones.has(lastReply.rootId) + ? { + kind: "thread", + channelId: channel.id, + rootId: lastReply.rootId, + } + : { + kind: "message", + channelId: channel.id, + messageId: latest.event.id, + }; + const readThrough: { target: ReadTarget; messageId: string }[] = dm + ? [] + : entries + .filter( + (entry) => + !entry.rootId || + !!reads.localUnread(`msg:${entry.event.id}`), + ) + .map((entry) => ({ + target: { + kind: "message" as const, + channelId: channel.id, + messageId: entry.event.id, + }, + messageId: entry.event.id, + })); + if (!dm && lastReply?.rootId) + readThrough.push({ + target: { + kind: "thread", + channelId: channel.id, + rootId: lastReply.rootId, + }, + messageId: lastReply.event.id, + }); + items.push( + Object.freeze({ + id: `${channel.id}:${id}`, + channelId: channel.id, + target: Object.freeze(target), + messageId: representative.event.id, + latestMessageId: latest.event.id, + messageIds: Object.freeze(entries.map((entry) => entry.event.id)), + ...(representative.rootId + ? { rootId: representative.rootId } + : {}), + authorId: representative.event.pubkey, + preview: + content.get(representative.event.id) ?? + representative.event.content, + createdAt: latest.event.created_at, + mentioned: entries.some((entry) => entry.mentioned), + thread: entries.some((entry) => entry.parentId !== undefined), + unreadCount: unread.length, + manual: + !!reads.localUnread(targetKey(target)) || + readThrough.some( + ({ target }) => !!reads.localUnread(targetKey(target)), + ), + readThrough: Object.freeze( + readThrough.map((step) => + Object.freeze({ + ...step, + target: Object.freeze(step.target), + }), + ), + ), + }), + ); + } + } + items.sort((a, b) => b.createdAt - a.createdAt || a.id.localeCompare(b.id)); + const next: InboxSnapshot = Object.freeze({ + status: error + ? "error" + : inboxLoading + ? "loading" + : freshness === "unknown" + ? "idle" + : "ready", + items: Object.freeze(items), + freshness, + ...(error ? { error } : {}), + }); + // Keep React's snapshot stable for irrelevant evidence and duplicate deliveries. + if ( + !inboxSnapshot || + JSON.stringify(inboxSnapshot) !== JSON.stringify(next) + ) + inboxSnapshot = next; + return inboxSnapshot; + } function addActivityListener(channelId: string, listener: () => void) { activity(channelId); const set = activityListeners.get(channelId) ?? new Set(); @@ -736,6 +883,9 @@ export function createUnread({ } function publish(channelIds?: ReadonlySet) { if (closed) return; + inboxDirty = true; + const previousInbox = inboxSnapshot; + const nextInbox = inboxListeners.size ? inbox() : undefined; const changed: string[] = []; const changedActivity: string[] = []; for (const [key, old] of snapshots) { @@ -764,6 +914,8 @@ export function createUnread({ } } // Replace/invalidate ALL affected projections before any reentrant callback. + if (nextInbox && nextInbox !== previousInbox) + for (const listener of inboxListeners) notify(listener); for (const key of changed) for (const listener of listeners.get(key) ?? []) notify(listener); for (const channelId of changedActivity) @@ -1088,6 +1240,7 @@ export function createUnread({ if (closed) return; if (refresh) return refresh; repairAgain = false; + inboxLoading = true; const generation = epoch; refresh = (async () => { await reads.ensure(); @@ -1098,7 +1251,11 @@ export function createUnread({ (channel) => !channel.cached && channel.members?.includes(viewer), ) .map((channel) => channel.id); - if (!ids.length) return; + if (!ids.length) { + freshness = "observed"; + error = undefined; + return; + } try { // The relay caps aggregate explicit #h values at 128 per request. // Keep roster scope: an unscoped read also includes unjoined open channels. @@ -1135,8 +1292,11 @@ export function createUnread({ } })().finally(() => { refresh = undefined; + inboxLoading = false; + publish(); if (!closed && repairAgain) void repair(); }); + publish(); return refresh; } /** The viewer's deletions end lookup memberships they were evidence for, @@ -1313,6 +1473,12 @@ export function createUnread({ }; } const capability: UnreadCapability = Object.freeze({ + inbox, + subscribeInbox(listener) { + inbox(); + inboxListeners.add(listener); + return () => inboxListeners.delete(listener); + }, snapshot, attention, subscribe(target, listener) { @@ -1342,6 +1508,8 @@ export function createUnread({ await reads.refresh("foreground"); await repair(); }, + revision: reads.revision, + generation: () => epoch, retrySync: async () => { await reads.refresh(); await reads.flush(); @@ -1561,8 +1729,11 @@ export function createUnread({ }); }, async markChannelRead(channelId) { + return capability.prepareChannelRead(channelId)(); + }, + prepareChannelRead(channelId) { const read = channelReadIntent(channelId, Math.floor(Date.now() / 1000)); - return serialize(channelId, read); + return () => serialize(channelId, read); }, async markAllChannelsRead() { if (closed) throw new Error("Read target unavailable"); @@ -1671,6 +1842,30 @@ export function createUnread({ ); return owners && [...owners].every(allowed) ? event : undefined; }, + // Edits omitted by later relay queries can still affect the shared fold. + // Read their deletion evidence too; no second edit cache or projection. + retainedEditIds(ids: readonly string[]) { + if (closed) return []; + const targets = new Set(ids); + const owners = channelOwnership((id) => events.get(id)); + return [...events.values()].flatMap((event) => { + if (event.kind !== 40003 || deleted(event)) return []; + const channels = owners(event); + if (!channels || ![...channels].every(allowed)) return []; + return event.tags.some(([name, id]) => { + const target = id && targets.has(id) ? events.get(id) : undefined; + return ( + name === "e" && + target && + contentKind(target) && + target.pubkey === event.pubkey && + !deleted(target) + ); + }) + ? [event.id] + : []; + }); + }, accept, purge, reconnect() { @@ -1715,6 +1910,9 @@ export function createUnread({ activitySnapshots.clear(); activityDirty.clear(); events.clear(); + inboxSnapshot = undefined; + inboxDirty = true; + inboxListeners.clear(); forcedMessages.clear(); entered.clear(); reads.dispose(); diff --git a/tests/browser/thread-unread.spec.mjs b/tests/browser/thread-unread.spec.mjs index abfa03a00..34e43d453 100644 --- a/tests/browser/thread-unread.spec.mjs +++ b/tests/browser/thread-unread.spec.mjs @@ -266,9 +266,14 @@ test("thread buttons show observed unread independently, clear only after readin await expect(dot(first)).toHaveCount(0); await expect(dot(other)).toBeVisible(); await expect(other).toHaveAccessibleName(/Observed unread replies/); // No channel-wide shortcut. + // The close transition owns channel width; unread receipt does not mean its + // resulting Virtua remeasurement and native scrolling have finished. await page .getByRole("button", { name: "Close Thread tab", exact: true }) .click(); + await expect( + page.locator('[data-panel-dock][data-closing="true"]'), + ).toHaveCount(0); const own = app.reply(roots[0].id, true); // Barrier: the session has indexed the reply, so its unread effect is final. await expect @@ -283,7 +288,28 @@ test("thread buttons show observed unread independently, clear only after readin await expect(first).toHaveAccessibleName("View thread: 23 replies"); app.reply(roots[0].id); await expect(first).toHaveAccessibleName(/Observed unread replies/); + await virtuaIdle(page); + // Remeasurement while geometry settles may start another native scroll. + await expect(page.locator("[data-channel-timeline] ol")).toHaveCSS( + "pointer-events", + "auto", + ); + const reopenAttempt = await page.evaluate( + () => window.fixtureNavigation.snapshot().attempt.id, + ); await first.click(); + await expect + .poll(() => + page.evaluate((before) => { + const { status, entry, attempt } = window.fixtureNavigation.snapshot(); + return { + status, + freshAttempt: attempt.id !== before, + messageId: entry.target.messageId, + }; + }, reopenAttempt), + ) + .toEqual({ status: "opened", freshAttempt: true, messageId: roots[0].id }); await expect( panel.getByText("New peer reply", { exact: true }), ).toBeVisible();