diff --git a/app/src/main.ts b/app/src/main.ts index 4572159..8ecf6a8 100644 --- a/app/src/main.ts +++ b/app/src/main.ts @@ -58,6 +58,7 @@ import { activityAt, formatActivityTime, presenceText, previewLine, sortByActivi import { roomProject, setRoomProject } from './room-projects.js' import { SharedProjectsPanel } from './shared-projects.js' import { RoomWatch } from './room-watch.js' +import { BrowserRoomArchiveStorage, RoomArchive, deleteRoomArchive, reseedRelays, type ReseedTarget } from './room-archive.js' import { PresenceAnnouncements } from './presence-announcements.js' import { readAgentRequestStatuses, type RequestAgent } from './agent-request-status.js' import { RoomBookmarks, accountRoomStore } from './room-bookmarks.js' @@ -95,6 +96,9 @@ import { localIdentity, sanitiseDisplayName, MAX_CHAT_TEXT_LENGTH, + CHAT_RETENTION_SECONDS, + MAX_CHAT_MESSAGES, + KINDS, type ParticipantIdentity, type DeviceCredential, type RoomPolicy, @@ -231,6 +235,10 @@ function slotWords(): string { return QUIET_SLOT >= 60 ? `${Math.ceil(QUIET_SLOT / 60)} minutes` : `${QUIET_SLOT} seconds` } const chatScroll = new ChatScroll(document.getElementById('chatLog')!, document.getElementById('newMessages') as HTMLButtonElement) +/** This browser's copy of every room it has been in: the original signed + * events, sealed under a device key. Undefined where the browser keeps no + * storage, and the rooms then read from relays alone. See room-archive.ts. */ +const roomArchive = (() => { try { return new RoomArchive(new BrowserRoomArchiveStorage()) } catch { return undefined } })() const conversationSearch = new ConversationSearch(document, selectChannel) const messageActions = new MessageActions() installReactionHold($('chatLog')) @@ -957,7 +965,7 @@ async function forgetThisBrowser(): Promise { title: 'Forget this browser?', message: 'This removes, from this browser only:\n' + '- your visitor identity\n' - + '- every room kept here, and its keys\n' + + '- every room kept here, its keys and the history this browser kept of it\n' + '- contacts and verified people\n' + '- text size and volume choices\n' + '- the saved connection to a Nostr signer\n\n' @@ -998,6 +1006,8 @@ async function forgetThisBrowser(): Promise { if (key.startsWith('kithmoot.')) storage.removeItem(key) } } + // The rooms' kept history goes with the rooms. + await deleteRoomArchive() history.replaceState(null, '', joinLinkBase()) approvedReload() @@ -6470,6 +6480,50 @@ function renderChat(messages: ChatMessage[]): void { updateConversationSearch() if (currentChannel === undefined) noteChatRead(messages) markConversationRead() + requestAnimationFrame(() => pageBackFromArchive(true)) +} + +/** + * Step back through this device's archive: when the reader reaches the top + * of the conversation, or when what is shown does not fill the log, as in a + * room whose last word is older than the retention window. The log keeps + * the reader's place while older messages arrive above it; see chat-scroll.ts. + */ +function pageBackFromArchive(onlyToFill = false): void { + const log = currentChannel === undefined ? session?.chat : channelLogs.get(currentChannel) + if (!log?.hasOlder) return + const el = $('chatLog') + if ($('roomArea').hidden || el.clientHeight === 0) return + if (onlyToFill ? el.scrollHeight > el.clientHeight : el.scrollTop > 64) return + void log.loadOlder() +} +$('chatLog').addEventListener('scroll', () => pageBackFromArchive(), { passive: true }) + +/** + * Put back on the room's relays what they forgot and this device kept. + * + * Late, so it never competes with joining. Never a quiet room's chat, which + * must not appear on a relay as bare room events; its rekeys ride in the + * open like any room's, so they still count. Only the room's own pool is + * written to. See `reseedRelays`. + */ +const RESEED_DELAY_MS = 5_000 +function scheduleReseed(s: RoomSession, pool: NostrRelayPool, quiet: boolean, authority: string | undefined): void { + const archive = roomArchive + if (!archive) return + const alive = (): boolean => !pool.closed && (session === s || dockedCall?.session === s) + setTimeout(() => { + if (!alive()) return + const since = nowSeconds() - CHAT_RETENTION_SECONDS + const logs = [s.chat, ...[AGENT_CHANNEL, TRANSCRIPT_CHANNEL, MINUTES_CHANNEL, CONTROL_CHANNEL].map(name => s.channel(name))] + const targets: ReseedTarget[] = [ + ...(quiet ? [] : logs.map(log => ({ kind: KINDS.CHAT, d: log.stream, since, limit: MAX_CHAT_MESSAGES }))), + ...(authority ? [{ kind: KINDS.ROOM_REKEY, d: s.roomId, limit: 1_000, authors: [authority] }] : []), + ] + reseedRelays(pool, archive, targets, { alive }) + .then(report => { for (const [url, count] of report.reseeded) console.info(`room archive: returned ${count} event${count === 1 ? '' : 's'} to ${url}`) }) + .catch(error => console.warn('room archive reseed', error)) + }, RESEED_DELAY_MS) } function updateConversationSearch(): void { @@ -9159,6 +9213,7 @@ async function startSession(asVisitor = false, retry?: { deadline: number }): Pr trackRole: activeTrackRole, ...(chatOnly ? {} : { assist: currentAssistOffer, relay: peerRelay }), ...forwarderMediaOptions, + ...(roomArchive ? { archive: roomArchive } : {}), // Epochs: follow a rekey signed by the room's authority, and ask it // first if the responder said the room is ahead of the secret we // were handed. See src/epoch.ts and docs/decisions.md. @@ -9200,6 +9255,7 @@ async function startSession(asVisitor = false, retry?: { deadline: number }): Pr trackRole: activeTrackRole, ...(chatOnly ? {} : { assist: currentAssistOffer, relay: peerRelay }), ...forwarderMediaOptions, + ...(roomArchive ? { archive: roomArchive } : {}), // Epochs: follow a rekey signed by the room's authority, and ask it // first if the responder said the room is ahead of the secret we // were handed. See src/epoch.ts and docs/decisions.md. @@ -9357,6 +9413,7 @@ async function startSession(asVisitor = false, retry?: { deadline: number }): Pr control.onChange((messages) => { if (session === s) ingestControl(messages) }) ingestControl(control.messages()) control.send(encodeControl({ op: 'catalogue?' })).catch(() => {}) + scheduleReseed(s, pool, !!quietTransport, sessionAuthority) renderRoomLockState() renderHost() // Empty until the keeper answers the `catalogue?` above with its signed diff --git a/app/src/room-archive.test.ts b/app/src/room-archive.test.ts new file mode 100644 index 0000000..0ba6315 --- /dev/null +++ b/app/src/room-archive.test.ts @@ -0,0 +1,248 @@ +import { webcrypto } from 'node:crypto' +import { beforeEach, describe, expect, it } from 'vitest' +import { finalizeEvent, generateSecretKey, type Event } from 'nostr-tools/pure' +import type { Filter } from 'nostr-tools/filter' +import { RoomArchive, reseedRelays, resetReseedHistory, type ArchiveKeys, type ArchivedRecord, type ReseedPool, type RoomArchiveStorage } from './room-archive.js' +import type { RelayConfig } from '../../src/relay-pool.js' + +const crypt = webcrypto as unknown as Crypto +const ROOM = 'ab'.repeat(32) +const OTHER = 'cd'.repeat(32) +const sk = generateSecretKey() + +class MemoryStorage implements RoomArchiveStorage { + stored?: ArchiveKeys + records = new Map() + async keys(): Promise { return this.stored } + async adoptKeys(candidate: ArchiveKeys): Promise { return this.stored ??= candidate } + async stream(stream: string): Promise { return [...this.records.values()].filter(r => r.stream === stream) } + async put(records: readonly ArchivedRecord[], remove: readonly string[] = []): Promise { for (const k of remove) this.records.delete(k); for (const r of records) this.records.set(r.key, r) } +} + +const chat = (text: string, at: number, d = ROOM, kind = 1460): Event => + finalizeEvent({ kind, created_at: at, tags: [['d', d]], content: `ciphertext:${text}` }, sk) + +describe('room archive at rest', () => { + it('round-trips original events, newest first, by kind and conversation, with a strict cursor', async () => { + const storage = new MemoryStorage() + const archive = new RoomArchive(storage, crypt) + const events = [chat('a', 100), chat('b', 200), chat('c', 300), chat('elsewhere', 250, OTHER), chat('rekey', 150, ROOM, 1462)] + for (const event of events) archive.keep(event) + archive.keep(events[0]!) + await archive.flushed() + expect(storage.records.size).toBe(5) + expect((await archive.read({ kind: 1460, d: ROOM, limit: 10 })).map(e => e.content)).toEqual(['ciphertext:c', 'ciphertext:b', 'ciphertext:a']) + expect(await archive.read({ kind: 1460, d: ROOM, limit: 10 })).toEqual([events[2], events[1], events[0]].map(e => JSON.parse(JSON.stringify(e)))) + expect((await archive.read({ kind: 1460, d: ROOM, since: 150, limit: 10 })).map(e => e.created_at)).toEqual([300, 200]) + expect((await archive.read({ kind: 1460, d: ROOM, before: { at: 200, id: events[1]!.id }, limit: 10 })).map(e => e.created_at)).toEqual([100]) + expect((await archive.read({ kind: 1462, d: ROOM, limit: 10 })).map(e => e.content)).toEqual(['ciphertext:rekey']) + + // A second tab, or the next visit, opens the same records. + const reopened = new RoomArchive(storage, crypt) + expect(await reopened.read({ kind: 1460, d: OTHER, limit: 10 })).toHaveLength(1) + }) + + it('shows neither room ids, event ids, contents nor times outside the ciphertext, under non-extractable keys', async () => { + const storage = new MemoryStorage() + const archive = new RoomArchive(storage, crypt) + const event = chat('secret', 1_799_999_999) + archive.keep(event) + await archive.flushed() + expect(storage.stored!.seal.extractable).toBe(false) + expect(storage.stored!.name.extractable).toBe(false) + const record = [...storage.records.values()][0]! + const visible = JSON.stringify({ key: record.key, stream: record.stream, version: record.version }) + Buffer.from(record.ciphertext).toString('latin1') + for (const leak of [ROOM, event.id, event.pubkey, event.sig, 'ciphertext:secret', String(event.created_at)]) expect(visible).not.toContain(leak) + expect(Object.keys(record).sort()).toEqual(['ciphertext', 'key', 'nonce', 'stream', 'version']) + }) + + it('checks events on the way in and records on the way out', async () => { + const storage = new MemoryStorage() + const archive = new RoomArchive(storage, crypt) + const good = chat('good', 100) + archive.keep({ ...good, content: 'changed after signing' }) + archive.keep({ ...chat('no tag', 100), tags: [] }) + archive.keep(good) + await archive.flushed() + expect(storage.records.size).toBe(1) + + // Tampered ciphertext, and a record moved to another conversation, are + // not handed to anybody. + archive.keep(chat('other', 100, OTHER)) + await archive.flushed() + const [first, second] = [...storage.records.values()] + storage.records.set(first!.key, { ...first!, ciphertext: first!.ciphertext.slice(1) }) + storage.records.set(second!.key, { ...second!, stream: first!.stream }) + const fresh = new RoomArchive(storage, crypt) + expect(await fresh.read({ kind: 1460, d: ROOM, limit: 10 })).toEqual([]) + expect(await fresh.read({ kind: 1460, d: OTHER, limit: 10 })).toEqual([]) + }) + + it('agrees on one pair of device keys when two tabs start at once', async () => { + const storage = new MemoryStorage() + const [a, b] = [new RoomArchive(storage, crypt), new RoomArchive(storage, crypt)] + a.keep(chat('from a', 1)); b.keep(chat('from b', 2)) + await Promise.all([a.flushed(), b.flushed()]) + expect(await new RoomArchive(storage, crypt).read({ kind: 1460, d: ROOM, limit: 10 })).toHaveLength(2) + }) + + it('past the cap drops the oldest, on disk too, and keeps the newest', async () => { + const storage = new MemoryStorage() + const archive = new RoomArchive(storage, crypt, { perStream: 3 }) + for (let i = 1; i <= 5; i++) archive.keep(chat(`m${i}`, i)) + await archive.flushed() + archive.keep(chat('m6', 6)) + await archive.flushed() + expect((await archive.read({ kind: 1460, d: ROOM, limit: 10 })).map(e => e.created_at)).toEqual([6, 5, 4]) + expect(storage.records.size).toBe(3) + expect((await new RoomArchive(storage, crypt).read({ kind: 1460, d: ROOM, limit: 10 })).map(e => e.created_at)).toEqual([6, 5, 4]) + }) + + it('lets a conversation go from memory when released, and reads it back from disk', async () => { + const storage = new MemoryStorage() + const archive = new RoomArchive(storage, crypt) + archive.keep(chat('a', 1)) + expect(await archive.read({ kind: 1460, d: ROOM, limit: 10 })).toHaveLength(1) + let loads = 0 + const stream = storage.stream.bind(storage) + storage.stream = async (s: string) => { loads++; return stream(s) } + await archive.read({ kind: 1460, d: ROOM, limit: 10 }) + expect(loads).toBe(0) + archive.release({ kind: 1460, d: ROOM }) + expect(await archive.read({ kind: 1460, d: ROOM, limit: 10 })).toHaveLength(1) + expect(loads).toBe(1) + }) + + it('a reader waits for its own conversation to be written, not for every room', async () => { + const storage = new MemoryStorage() + const archive = new RoomArchive(storage, crypt) + archive.keep(chat('here', 1)) + // Another room's write that never finishes must not hold this read up. + const put = storage.put.bind(storage) + let stall: () => void = () => {} + const stalled = new Promise(resolve => { stall = resolve }) + let reached = false + await archive.read({ kind: 1460, d: ROOM, limit: 10 }) + storage.put = async (records, remove) => { reached = true; await stalled; return put(records, remove) } + archive.keep(chat('elsewhere', 2, OTHER)) + await expect.poll(() => reached).toBe(true) + expect(await archive.read({ kind: 1460, d: ROOM, limit: 10 })).toHaveLength(1) + stall() + await archive.flushed() + }) +}) + +/** A pool of named relays, each holding what it holds. A relay in `forgets` + * accepts and stores nothing; one in `caps` returns at most that many + * events per request, newest first; one in `silent` never finishes. */ +class FakePool implements ReseedPool { + published: { url: string; id: string }[] = [] + asked: Filter[] = [] + constructor(public relays: Map, public config: RelayConfig[], public forgets = new Set(), + public caps = new Map(), public silent = new Set()) {} + describe(): RelayConfig[] { return this.config } + async query(url: string, filters: Filter[]): Promise<{ events: Event[]; complete: boolean }> { + const f = filters[0]! + this.asked.push(f) + if (this.silent.has(url)) return { events: [], complete: false } + const matching = (this.relays.get(url) ?? []).filter(e => (!f.ids || f.ids.includes(e.id)) && + (!f.kinds || f.kinds.includes(e.kind)) && (!f['#d'] || e.tags.some(t => t[0] === 'd' && f['#d']!.includes(t[1]!))) && + (f.since === undefined || e.created_at >= f.since) && (f.until === undefined || e.created_at <= f.until) && + (!f.authors || f.authors.includes(e.pubkey))) + const limit = Math.min(f.limit ?? Infinity, this.caps.get(url) ?? Infinity) + return { events: matching.sort((a, b) => b.created_at - a.created_at).slice(0, limit), complete: true } + } + async publishQuietly(url: string, event: Event): Promise { + if (!this.config.some(r => r.url === url && r.write)) throw new Error('not a room relay') + this.published.push({ url, id: event.id }) + if (!this.forgets.has(url)) this.relays.get(url)!.push(event) + } +} + +describe('reseeding forgetful relays', () => { + beforeEach(() => resetReseedHistory()) + const both = (url: string): RelayConfig => ({ url, read: true, write: true }) + + async function kept(...events: (Event | [Event, { quiet: true }])[]): Promise { + const archive = new RoomArchive(new MemoryStorage(), crypt) + for (const event of events) Array.isArray(event) ? archive.keep(...event) : archive.keep(event) + await archive.flushed() + return archive + } + + it('hands a relay that returned fewer the originals it lacks, unchanged, and only to that relay', async () => { + const events = [chat('1', 100), chat('2', 200), chat('3', 300)] + const archive = await kept(...events) + const pool = new FakePool(new Map([['wss://full', [...events]], ['wss://forgot', [events[0]!]]]), [both('wss://full'), both('wss://forgot'), { url: 'wss://read-only', read: true, write: false }]) + const report = await reseedRelays(pool, archive, [{ kind: 1460, d: ROOM, since: 0, limit: 500 }], { alive: () => true, gapMs: 0 }) + expect(pool.published).toEqual([{ url: 'wss://forgot', id: events[2]!.id }, { url: 'wss://forgot', id: events[1]!.id }]) + expect(pool.relays.get('wss://forgot')!.map(e => e.id).sort()).toEqual(events.map(e => e.id).sort()) + expect(pool.relays.get('wss://forgot')![1]).toEqual(JSON.parse(JSON.stringify(events[2]))) + expect(report.reseeded).toEqual(new Map([['wss://forgot', 2]])) + + // Asked again at once, nobody is compared twice. + await reseedRelays(pool, archive, [{ kind: 1460, d: ROOM, since: 0, limit: 500 }], { alive: () => true, gapMs: 0 }) + expect(pool.published).toHaveLength(2) + }) + + it('reads a relay that caps its answers to the end, and republishes nothing it holds', async () => { + const events = Array.from({ length: 12 }, (_, i) => chat(`m${i}`, 100 + i)) + const archive = await kept(...events) + const pool = new FakePool(new Map([['wss://capped', [...events]]]), [both('wss://capped')], new Set(), new Map([['wss://capped', 5]])) + await reseedRelays(pool, archive, [{ kind: 1460, d: ROOM, since: 0, limit: 500 }], { alive: () => true, gapMs: 0 }) + expect(pool.published).toEqual([]) + expect(pool.asked.some(f => f.until !== undefined)).toBe(true) + + // Missing one in the middle, it is found through the pages and only it goes back. + resetReseedHistory() + const gappy = new FakePool(new Map([['wss://capped', events.filter((_, i) => i !== 3)]]), [both('wss://capped')], new Set(), new Map([['wss://capped', 5]])) + await reseedRelays(gappy, archive, [{ kind: 1460, d: ROOM, since: 0, limit: 500 }], { alive: () => true, gapMs: 0 }) + expect(gappy.published.map(p => p.id)).toEqual([events[3]!.id]) + }) + + it('treats a relay that does not finish answering as unknown: nothing republished', async () => { + const archive = await kept(chat('a', 100), chat('b', 200)) + const pool = new FakePool(new Map([['wss://slow', []]]), [both('wss://slow')], new Set(), new Map(), new Set(['wss://slow'])) + await reseedRelays(pool, archive, [{ kind: 1460, d: ROOM, limit: 500 }], { alive: () => true, gapMs: 0 }) + expect(pool.published).toEqual([]) + }) + + it('stays inside the window and the budget, and skips anything whose signature fails', async () => { + const events = [chat('old', 10), chat('a', 100), chat('b', 200), chat('c', 300)] + const archive = await kept(...events) + const pool = new FakePool(new Map([['wss://empty', []]]), [both('wss://empty')]) + await reseedRelays(pool, archive, [{ kind: 1460, d: ROOM, since: 50, limit: 500 }], { alive: () => true, gapMs: 0, budget: 2 }) + expect(pool.published.map(p => p.id)).toEqual([events[3]!.id, events[2]!.id]) + + resetReseedHistory() + const forged = await kept({ ...chat('x', 100), sig: 'f'.repeat(128) }) + const clean = new FakePool(new Map([['wss://empty', []]]), [both('wss://empty')]) + await reseedRelays(clean, forged, [{ kind: 1460, d: ROOM, limit: 500 }], { alive: () => true, gapMs: 0 }) + expect(clean.published).toEqual([]) + }) + + it('never hands back what came through a quiet room', async () => { + const open = chat('open', 100) + const archive = await kept(open, [chat('quiet', 200), { quiet: true }]) + const pool = new FakePool(new Map([['wss://r', []]]), [both('wss://r')]) + await reseedRelays(pool, archive, [{ kind: 1460, d: ROOM, limit: 500 }], { alive: () => true, gapMs: 0 }) + expect(pool.published.map(p => p.id)).toEqual([open.id]) + // Still read back for the room itself. + expect(await archive.read({ kind: 1460, d: ROOM, limit: 10 })).toHaveLength(2) + }) + + it('republishes a rekey only when the room authority signed it, and stops when the room closes', async () => { + const authority = finalizeEvent({ kind: 1462, created_at: 1, tags: [['d', ROOM]], content: '' }, sk).pubkey + const stranger = finalizeEvent({ kind: 1462, created_at: 2, tags: [['d', ROOM]], content: '' }, generateSecretKey()) + const genuine = chat('rekey', 3, ROOM, 1462) + const archive = await kept(genuine, stranger) + const pool = new FakePool(new Map([['wss://r', []]]), [both('wss://r')]) + await reseedRelays(pool, archive, [{ kind: 1462, d: ROOM, limit: 1000, authors: [authority] }], { alive: () => true, gapMs: 0 }) + expect(pool.published.map(p => p.id)).toEqual([genuine.id]) + + resetReseedHistory() + const closed = new FakePool(new Map([['wss://r', []]]), [both('wss://r')]) + await reseedRelays(closed, archive, [{ kind: 1462, d: ROOM, limit: 1000, authors: [authority] }], { alive: () => false, gapMs: 0 }) + expect(closed.published).toEqual([]) + }) +}) diff --git a/app/src/room-archive.ts b/app/src/room-archive.ts new file mode 100644 index 0000000..34d9da4 --- /dev/null +++ b/app/src/room-archive.ts @@ -0,0 +1,470 @@ +/** + * This device's archive of the rooms it has been in: the original signed + * events, as they arrived, so a room's history does not depend on a relay + * remembering it. See `src/archive.ts` for what goes in and how it is read. + * + * Nothing decrypted is stored. Each event, still room-encrypted and signed, + * is sealed again with AES-GCM under a device key, the pattern of + * `history-index.ts` with its own schema. Outside the ciphertext a record + * carries two keyed hashes and nothing else: which conversation it belongs to + * and which event it is, both HMACs under a second device key. + * + * What that protects, honestly: somebody reading the stored records without + * the keys learns no room id, event id, sender or time. The keys themselves + * sit in the same IndexedDB on the same disk; "non-extractable" only stops + * the page's own scripts exporting them, so whoever can copy the whole + * browser profile can open the archive. Records are not padded, so their + * sizes show roughly how long each event is, and how many records each + * conversation holds shows too. + * + * `reseedRelays` is the other half: a room relay that returns fewer of a + * conversation's events than this device holds is handed the originals back, + * unchanged, a few at a time. Whether a relay keeps chat at all is the + * relay health's finding, not this module's. + */ +import { getEventHash, type Event } from 'nostr-tools/pure' +import type { Filter } from 'nostr-tools/filter' +import { archiveTag, compareArchived, olderThan, reseedCandidates, MAX_RESEED_EVENTS, type ArchiveMeta, type ArchiveQuery, type EventArchive } from '../../src/archive.js' +import { verifyEventUncached } from '../../src/verify.js' +import type { RelayConfig } from '../../src/relay-pool.js' + +const DATABASE_VERSION = 1 +const RECORD_VERSION = 1 +const EVENT_AAD = 'kithmoot.room-archive.event.v1:' +const NONCE_BYTES = 12 +/** An event larger than this was never a room event a relay would take. */ +const MAX_EVENT_BYTES = 128 * 1024 +/** Per conversation. Well past any room today; past it the oldest events + * go. What "keep" means on a full device is an open question in the plan. */ +export const MAX_ARCHIVED_PER_STREAM = 50_000 + +export interface ArchivedRecord { + /** HMAC of the event id: stable, so a second copy replaces the first. */ + key: string + /** HMAC of the kind and `d` tag: which conversation, opaquely. */ + stream: string + version: 1 + nonce: ArrayBuffer + ciphertext: ArrayBuffer +} + +export interface ArchiveKeys { + /** AES-GCM, seals each event. */ + seal: CryptoKey + /** HMAC-SHA-256, names streams and records. */ + name: CryptoKey +} + +export interface RoomArchiveStorage { + keys(): Promise + /** Store `candidate` unless keys already exist, and return whichever + * pair is stored: two tabs opening a fresh archive agree on one. */ + adoptKeys(candidate: ArchiveKeys): Promise + stream(stream: string): Promise + /** Write `records` and delete the records named by `remove`, together. */ + put(records: readonly ArchivedRecord[], remove?: readonly string[]): Promise +} + +/** What a sealed record holds. */ +interface Sealed { event: Event; quiet?: true } + +export class RoomArchive implements EventArchive { + #keys?: Promise + /** Decrypted per conversation on first read, then kept current by `keep` + * until `release`. */ + readonly #streams = new Map>>() + #queue: Sealed[] = [] + #flushing?: Promise + /** Events queued and not yet written, per conversation, so a reader waits + * for its own conversation and never for every room's. */ + readonly #pending = new Map() + readonly #drained = new Map void)[]>() + readonly #released = new Set() + #warned = false + + readonly #perStream: number + + constructor(private readonly storage: RoomArchiveStorage, private readonly crypt: Crypto = globalThis.crypto, limits: { perStream?: number } = {}) { + if (!crypt?.subtle || !crypt.getRandomValues) throw new Error('This device cannot keep an encrypted room archive.') + this.#perStream = limits.perStream ?? MAX_ARCHIVED_PER_STREAM + } + + keep(event: Event, meta?: ArchiveMeta): void { + try { + if (!wellFormed(event)) return + const name = streamName(event.kind, archiveTag(event)!) + this.#pending.set(name, (this.#pending.get(name) ?? 0) + 1) + this.#queue.push({ event: copyEvent(event), ...(meta?.quiet ? { quiet: true as const } : {}) }) + this.#flushing ??= Promise.resolve().then(() => this.#flush()).finally(() => { this.#flushing = undefined }) + } catch { + // Keeping is best effort; the room goes on without it. + } + } + + /** Resolves once everything kept so far is on disk. */ + async flushed(): Promise { + while (this.#flushing) await this.#flushing + } + + async read(query: ArchiveQuery): Promise { + if (!Number.isSafeInteger(query.limit) || query.limit < 1) throw new Error('Use an archive read limit of at least one.') + const d = query.d.toLowerCase() + const name = streamName(query.kind, d) + this.#released.delete(name) + if (this.#pending.get(name)) await new Promise(resolve => this.#drained.set(name, [...(this.#drained.get(name) ?? []), resolve])) + const events = await this.#load(query.kind, d) + return [...events.values()] + .filter(kept => !(query.reseedable && kept.quiet)) + .map(kept => kept.event) + .filter(e => (query.since === undefined || e.created_at >= query.since) && (!query.before || olderThan(e, query.before))) + .sort(compareArchived) + .slice(0, query.limit) + .map(copyEvent) + } + + release(query: Pick): void { + const name = streamName(query.kind, query.d.toLowerCase()) + // Still being written: let it go once it is on disk. + if (this.#pending.get(name)) this.#released.add(name) + else this.#streams.delete(name) + } + + async #flush(): Promise { + while (this.#queue.length) { + const batch = this.#queue.splice(0, 200) + const added: [Map, string][] = [] + const touched = new Set>() + try { + const records: ArchivedRecord[] = [] + for (const kept of batch) { + const { event } = kept + const d = archiveTag(event)! + const events = await this.#load(event.kind, d) + if (events.has(event.id)) continue + events.set(event.id, kept) + added.push([events, event.id]) + touched.add(events) + records.push(await this.#seal(kept, d)) + } + // Past the cap, the oldest go: a room that keeps talking keeps its + // newest history rather than freezing on the day it filled up. + const remove: string[] = [] + for (const events of touched) { + if (events.size <= this.#perStream) continue + const oldest = [...events.values()].map(k => k.event).sort(compareArchived).slice(this.#perStream) + for (const event of oldest) { + events.delete(event.id) + const key = await this.#hash(`event\n${event.id}`) + const written = records.findIndex(r => r.key === key) + if (written >= 0) records.splice(written, 1) + else remove.push(key) + } + } + if (records.length || remove.length) await this.storage.put(records, remove) + } catch (error) { + // Not on disk, so not remembered as kept: a later copy tries again. + for (const [events, id] of added) events.delete(id) + if (!this.#warned) { this.#warned = true; console.warn('room archive could not keep events', error) } + } finally { + for (const { event } of batch) this.#settled(streamName(event.kind, archiveTag(event)!)) + } + } + } + + #settled(name: string): void { + const left = (this.#pending.get(name) ?? 1) - 1 + if (left > 0) { this.#pending.set(name, left); return } + this.#pending.delete(name) + for (const resolve of this.#drained.get(name) ?? []) resolve() + this.#drained.delete(name) + if (this.#released.delete(name)) this.#streams.delete(name) + } + + #load(kind: number, d: string): Promise> { + const name = streamName(kind, d) + let loading = this.#streams.get(name) + if (!loading) { + loading = this.#open(kind, d) + // A failed read is retried next time rather than remembered as empty. + loading.catch(() => { if (this.#streams.get(name) === loading) this.#streams.delete(name) }) + this.#streams.set(name, loading) + } + return loading + } + + async #open(kind: number, d: string): Promise> { + const stream = await this.#streamHandle(kind, d) + const { seal } = await this.#deviceKeys() + const events = new Map() + for (const record of await this.storage.stream(stream)) { + const kept = await this.#unseal(record, stream, seal) + // Checked again on the way out: the right conversation, and an id + // that is the hash of what it names. The signature is the reader's + // to check, as for any event off a relay. + if (kept && kept.event.kind === kind && archiveTag(kept.event) === d) events.set(kept.event.id, kept) + } + return events + } + + async #unseal(record: ArchivedRecord, stream: string, seal: CryptoKey): Promise { + if (record.version !== RECORD_VERSION || record.stream !== stream || !(record.nonce instanceof ArrayBuffer) || + !(record.ciphertext instanceof ArrayBuffer) || record.nonce.byteLength !== NONCE_BYTES) return undefined + try { + const plaintext = await this.crypt.subtle.decrypt({ name: 'AES-GCM', iv: record.nonce, additionalData: aad(stream) }, seal, record.ciphertext) + const sealed = JSON.parse(new TextDecoder().decode(plaintext)) as Sealed + if (!sealed || typeof sealed !== 'object' || !wellFormed(sealed.event)) return undefined + return { event: sealed.event, ...(sealed.quiet === true ? { quiet: true as const } : {}) } + } catch { + return undefined + } + } + + async #seal(kept: Sealed, d: string): Promise { + const stream = await this.#streamHandle(kept.event.kind, d) + const { seal } = await this.#deviceKeys() + const nonce = new Uint8Array(NONCE_BYTES); this.crypt.getRandomValues(nonce) + const plaintext = new TextEncoder().encode(JSON.stringify(kept)) + const ciphertext = await this.crypt.subtle.encrypt({ name: 'AES-GCM', iv: nonce, additionalData: aad(stream) }, seal, plaintext) + return { key: await this.#hash(`event\n${kept.event.id}`), stream, version: RECORD_VERSION, nonce: nonce.slice().buffer, ciphertext } + } + + #streamHandle(kind: number, d: string): Promise { return this.#hash(`stream\n${kind}\n${d}`) } + + async #hash(value: string): Promise { + const { name } = await this.#deviceKeys() + const mac = new Uint8Array(await this.crypt.subtle.sign('HMAC', name, new TextEncoder().encode(value))) + return [...mac].map(byte => byte.toString(16).padStart(2, '0')).join('') + } + + #deviceKeys(): Promise { + this.#keys ??= (async () => { + const existing = await this.storage.keys() + if (existing) return existing + const seal = await this.crypt.subtle.generateKey({ name: 'AES-GCM', length: 256 }, false, ['encrypt', 'decrypt']) + const name = await this.crypt.subtle.generateKey({ name: 'HMAC', hash: 'SHA-256' }, false, ['sign']) + return this.storage.adoptKeys({ seal, name }) + })() + this.#keys.catch(() => { this.#keys = undefined }) + return this.#keys + } +} + +function streamName(kind: number, d: string): string { return `${kind}:${d}` } + +function aad(stream: string): ArrayBuffer { return new TextEncoder().encode(EVENT_AAD + stream).slice().buffer } + +/** A complete event whose id is the hash of what it says. */ +function wellFormed(event: Event): boolean { + if (!event || typeof event !== 'object' || !Number.isSafeInteger(event.kind) || !Number.isSafeInteger(event.created_at) || + typeof event.content !== 'string' || !Array.isArray(event.tags) || !/^[0-9a-f]{64}$/.test(event.id) || + !/^[0-9a-f]{64}$/.test(event.pubkey) || !/^[0-9a-f]{128}$/.test(event.sig) || !archiveTag(event)) return false + if (event.content.length > MAX_EVENT_BYTES) return false + try { return getEventHash(event) === event.id } catch { return false } +} + +/** The signed fields only: no cached verdict or other property rides along. */ +function copyEvent(event: Event): Event { + return { id: event.id, pubkey: event.pubkey, created_at: event.created_at, kind: event.kind, tags: event.tags.map(tag => [...tag]), content: event.content, sig: event.sig } +} + +/** Browser persistence. The two device keys are structured-cloned + * CryptoKeys and stay non-extractable; records carry only keyed hashes + * outside their ciphertext. */ +export class BrowserRoomArchiveStorage implements RoomArchiveStorage { + #db: Promise | undefined + constructor(private readonly dbName = ROOM_ARCHIVE_DATABASE, private readonly factory: IDBFactory = globalThis.indexedDB) { + if (!factory) throw new Error('This browser does not provide storage for a room archive.') + } + async keys(): Promise { + return (await request((await this.#transaction('keys', 'readonly')).objectStore('keys').get('device'))) as ArchiveKeys | undefined + } + async adoptKeys(candidate: ArchiveKeys): Promise { + const transaction = await this.#transaction('keys', 'readwrite') + const store = transaction.objectStore('keys') + const existing = await request(store.get('device')) as ArchiveKeys | undefined + if (!existing) store.put(candidate, 'device') + await complete(transaction) + return existing ?? candidate + } + async stream(stream: string): Promise { + return await request((await this.#transaction('events', 'readonly')).objectStore('events').index('stream').getAll(stream)) + } + async put(records: readonly ArchivedRecord[], remove: readonly string[] = []): Promise { + const transaction = await this.#transaction('events', 'readwrite') + const store = transaction.objectStore('events') + for (const key of remove) store.delete(key) + for (const record of records) store.put(record) + await complete(transaction) + } + async #transaction(store: 'keys' | 'events', mode: IDBTransactionMode): Promise { + return (await this.#database()).transaction(store, mode) + } + #database(): Promise { + return this.#db ??= new Promise((resolve, reject) => { + const open = this.factory.open(this.dbName, DATABASE_VERSION) + open.onupgradeneeded = () => { + const db = open.result + if (!db.objectStoreNames.contains('keys')) db.createObjectStore('keys') + if (!db.objectStoreNames.contains('events')) db.createObjectStore('events', { keyPath: 'key' }).createIndex('stream', 'stream') + } + open.onsuccess = () => { + // Another tab deleting the archive (forgetting this browser) must + // not wait on this one. + open.result.onversionchange = () => open.result.close() + resolve(open.result) + } + open.onerror = () => reject(open.error) + }) + } +} + +export const ROOM_ARCHIVE_DATABASE = 'kithmoot-room-archive-v1' + +/** Remove this browser's whole archive: part of forgetting the browser. */ +export function deleteRoomArchive(factory: IDBFactory | undefined = globalThis.indexedDB): Promise { + if (!factory) return Promise.resolve() + return new Promise(resolve => { + const deleting = factory.deleteDatabase(ROOM_ARCHIVE_DATABASE) + deleting.onsuccess = deleting.onerror = deleting.onblocked = () => resolve() + }) +} + +function request(value: IDBRequest): Promise { return new Promise((resolve, reject) => { value.onsuccess = () => resolve(value.result); value.onerror = () => reject(value.error) }) } +function complete(transaction: IDBTransaction): Promise { + return new Promise((resolve, reject) => { + transaction.oncomplete = () => resolve() + transaction.onerror = () => reject(transaction.error) + transaction.onabort = () => reject(transaction.error ?? new Error('Room archive transaction was aborted.')) + }) +} + +// --------------------------------------------------------------------------- +// Reseeding forgetful relays +// --------------------------------------------------------------------------- + +/** One conversation to compare: the filter a reader subscribes with. */ +export interface ReseedTarget { + kind: number + d: string + since?: number + limit: number + /** Only these signers: the authority, for rekeys. */ + authors?: string[] +} + +/** What reseeding needs of a relay pool. `NostrRelayPool` has all of it. */ +export interface ReseedPool { + describe(): RelayConfig[] + query(url: string, filters: Filter[], timeoutMs?: number): Promise<{ events: Event[]; complete: boolean }> + publishQuietly(url: string, event: Event): Promise +} + +export interface ReseedOptions { + /** False once the room this is for has closed or been left. */ + alive: () => boolean + /** Pause between two republished events, so a relay sees a trickle. */ + gapMs?: number + /** The most events one call republishes, over every relay. */ + budget?: number + /** How long one relay has to answer one request. */ + timeoutMs?: number + now?: () => number +} + +export interface ReseedReport { + /** Per relay URL: how many originals it accepted back. */ + reseeded: Map +} + +/** When a conversation was last compared with a relay, by `url d`, so + * reopening a room over and over does not ask the same relay again. */ +const lastCompared = new Map() +export const RESEED_INTERVAL_MS = 10 * 60_000 +const DEFAULT_RESEED_BUDGET = 1_000 +/** How many times one comparison pages back with `until` before it settles + * for what it has. */ +const MAX_COMPARE_PAGES = 10 + +/** + * What one relay holds of a conversation, paged back with `until` so a relay + * that caps how many events it returns is still read to the end. Undefined + * when the relay did not answer in full: a slow or closed request says + * nothing about what it holds. `floor` is set when paging stopped before the + * relay ran out, and nothing older than it can be judged. + */ +async function relayHolds(pool: ReseedPool, url: string, filter: Filter, timeoutMs: number | undefined): Promise<{ ids: Set; floor?: number } | undefined> { + const ids = new Set() + let until: number | undefined + for (let page = 0; page < MAX_COMPARE_PAGES; page++) { + const answer = await pool.query(url, [{ ...filter, ...(until !== undefined ? { until } : {}) }], timeoutMs) + if (!answer.complete) return undefined + let fresh = 0 + for (const event of answer.events) if (!ids.has(event.id)) { ids.add(event.id); fresh++ } + if (fresh === 0) return { ids } + // Inclusive, because several events can share the oldest second. + until = Math.min(...answer.events.map(e => e.created_at)) + } + return { ids, floor: until } +} + +/** + * Compare each target with each of the room's read-and-write relays, alone, + * and hand a relay that returned fewer events than this device holds the + * ones it lacks. Only the pool's own relays are ever written to: the room's + * relays as this device uses them, never a person's own or anybody else's. + * A quiet room's chat is never handed back, whatever the targets say. + * Bounded by window, by `budget` and by `RESEED_INTERVAL_MS` per relay and + * conversation, paced by `gapMs`, and written through the pool's quiet path, + * which never marks the room's relay health or reopens its sockets. + */ +export async function reseedRelays(pool: ReseedPool, archive: EventArchive, targets: readonly ReseedTarget[], opts: ReseedOptions): Promise { + const report: ReseedReport = { reseeded: new Map() } + const gap = opts.gapMs ?? 100 + let budget = opts.budget ?? DEFAULT_RESEED_BUDGET + const now = opts.now ?? Date.now + const relays = pool.describe().filter(relay => relay.read && relay.write).map(relay => relay.url) + for (const target of targets) { + for (const url of relays) { + if (!opts.alive() || budget <= 0) return report + const key = `${url} ${target.kind}:${target.d}` + if (now() - (lastCompared.get(key) ?? -Infinity) < RESEED_INTERVAL_MS) continue + lastCompared.set(key, now()) + const filter: Filter = { kinds: [target.kind], '#d': [target.d], limit: target.limit, ...(target.since !== undefined ? { since: target.since } : {}), ...(target.authors ? { authors: target.authors } : {}) } + let held: Awaited> + let archived: Event[] + try { + ;[held, archived] = await Promise.all([ + relayHolds(pool, url, filter, opts.timeoutMs), + archive.read({ kind: target.kind, d: target.d, since: target.since, limit: target.limit, reseedable: true }), + ]) + } catch { + continue + } + // Unknown is not empty: no reseed and no verdict until it answers. + if (!held || !opts.alive()) continue + const candidates = reseedCandidates( + archived.filter(e => !target.authors || target.authors.includes(e.pubkey)), + held.ids, + Math.min(budget, MAX_RESEED_EVENTS), + held.floor, + ).filter(verifyEventUncached) + const sent: string[] = [] + for (const event of candidates) { + if (!opts.alive()) return report + budget-- + try { + await pool.publishQuietly(url, event) + sent.push(event.id) + } catch { + // Refused or unreachable: the next comparison will try again. + } + if (gap > 0) await new Promise(resolve => setTimeout(resolve, gap)) + } + if (sent.length) report.reseeded.set(url, (report.reseeded.get(url) ?? 0) + sent.length) + } + } + return report +} + +/** For tests: forget when conversations were last compared. */ +export function resetReseedHistory(): void { lastCompared.clear() } diff --git a/playwright.config.ts b/playwright.config.ts index a8a6d35..d790983 100644 --- a/playwright.config.ts +++ b/playwright.config.ts @@ -47,7 +47,7 @@ export default defineConfig({ // importantly, goes dark again on mute - an analyser that is never pulled // reports silence for ever with nothing in the console, so this feature // can fail by simply never happening. - testMatch: ['notification-settings.spec.ts', 'knock.spec.ts', 'private-room.spec.ts', 'confirmations.spec.ts', 'den-journey.spec.ts', 'assignments.spec.ts', 'relay-settings.spec.ts', 'agent-receipts.spec.ts', 'context.spec.ts', 'persistent-groups.spec.ts', 'workspace.spec.ts', 'share-viewer.spec.ts', 'screen-share-audio.spec.ts', 'chat-comfort.spec.ts', 'phone-chat.spec.ts', 'messages.spec.ts', 'phone-message-links.spec.ts', 'quiet.spec.ts', 'contact-card.spec.ts', 'model-shortcuts.spec.ts', 'e2e.spec.ts', 'media.spec.ts', 'no-google-ice.spec.ts', 'safari-ice.spec.ts', 'volume.spec.ts', 'soak.spec.ts', 'agent.spec.ts', 'effects.spec.ts', 'camera-effects-no-phone-home.spec.ts', 'relay-capability.spec.ts', 'peer-assist.spec.ts', 'rooms.spec.ts', 'speaking.spec.ts', 'speaking-quiet-device.spec.ts', 'verification.spec.ts', 'channels.spec.ts', 'chat-reliability.spec.ts', 'updates.spec.ts', 'nostr-rooms.spec.ts', 'room-switching.spec.ts', 'conversation-search.spec.ts', 'drafts.spec.ts', 'home.spec.ts', 'site.spec.ts', 'forget-this-browser.spec.ts', 'sign-out-clears-bunker-key.spec.ts', 'wake-lock.spec.ts', 'call-stability.spec.ts', 'call-dock.spec.ts', 'mute-badge.spec.ts', 'desktop-room-layout.spec.ts', 'portrait-share.spec.ts', 'call-layout.spec.ts', 'call-focus.spec.ts', 'call-chat-divider.spec.ts', 'call-bell.spec.ts', 'shortcuts.spec.ts', 'call-prefs.spec.ts', 'mobile-landscape.spec.ts'], + testMatch: ['notification-settings.spec.ts', 'knock.spec.ts', 'private-room.spec.ts', 'confirmations.spec.ts', 'den-journey.spec.ts', 'assignments.spec.ts', 'relay-settings.spec.ts', 'agent-receipts.spec.ts', 'context.spec.ts', 'persistent-groups.spec.ts', 'workspace.spec.ts', 'share-viewer.spec.ts', 'screen-share-audio.spec.ts', 'chat-comfort.spec.ts', 'phone-chat.spec.ts', 'messages.spec.ts', 'phone-message-links.spec.ts', 'quiet.spec.ts', 'contact-card.spec.ts', 'model-shortcuts.spec.ts', 'e2e.spec.ts', 'media.spec.ts', 'no-google-ice.spec.ts', 'safari-ice.spec.ts', 'volume.spec.ts', 'soak.spec.ts', 'agent.spec.ts', 'effects.spec.ts', 'camera-effects-no-phone-home.spec.ts', 'relay-capability.spec.ts', 'peer-assist.spec.ts', 'rooms.spec.ts', 'speaking.spec.ts', 'speaking-quiet-device.spec.ts', 'verification.spec.ts', 'channels.spec.ts', 'chat-reliability.spec.ts', 'room-archive.spec.ts', 'updates.spec.ts', 'nostr-rooms.spec.ts', 'room-switching.spec.ts', 'conversation-search.spec.ts', 'drafts.spec.ts', 'home.spec.ts', 'site.spec.ts', 'forget-this-browser.spec.ts', 'sign-out-clears-bunker-key.spec.ts', 'wake-lock.spec.ts', 'call-stability.spec.ts', 'call-dock.spec.ts', 'mute-badge.spec.ts', 'desktop-room-layout.spec.ts', 'portrait-share.spec.ts', 'call-layout.spec.ts', 'call-focus.spec.ts', 'call-chat-divider.spec.ts', 'call-bell.spec.ts', 'shortcuts.spec.ts', 'call-prefs.spec.ts', 'mobile-landscape.spec.ts'], // Public relays take a few seconds to round-trip a roster event, and the // join-last case waits on three of those in sequence: A's entry, B's, and // then A and B answering C's arrival. The stage-1 live test used similar @@ -85,12 +85,12 @@ export default defineConfig({ projects: [ { name: 'firefox', - testMatch: ['knock.spec.ts', 'private-room.spec.ts', 'confirmations.spec.ts', 'den-journey.spec.ts', 'assignments.spec.ts', 'relay-settings.spec.ts', 'agent-receipts.spec.ts', 'context.spec.ts', 'persistent-groups.spec.ts', 'workspace.spec.ts', 'chat-comfort.spec.ts', 'phone-chat.spec.ts', 'messages.spec.ts', 'quiet.spec.ts', 'contact-card.spec.ts', 'model-shortcuts.spec.ts', 'chat-reliability.spec.ts', 'nostr-rooms.spec.ts', 'room-switching.spec.ts', 'conversation-search.spec.ts', 'drafts.spec.ts', 'home.spec.ts', 'site.spec.ts'], + testMatch: ['knock.spec.ts', 'private-room.spec.ts', 'confirmations.spec.ts', 'den-journey.spec.ts', 'assignments.spec.ts', 'relay-settings.spec.ts', 'agent-receipts.spec.ts', 'context.spec.ts', 'persistent-groups.spec.ts', 'workspace.spec.ts', 'chat-comfort.spec.ts', 'phone-chat.spec.ts', 'messages.spec.ts', 'quiet.spec.ts', 'contact-card.spec.ts', 'model-shortcuts.spec.ts', 'chat-reliability.spec.ts', 'room-archive.spec.ts', 'nostr-rooms.spec.ts', 'room-switching.spec.ts', 'conversation-search.spec.ts', 'drafts.spec.ts', 'home.spec.ts', 'site.spec.ts'], use: { ...devices['Desktop Firefox'] }, }, { name: 'webkit', - testMatch: ['knock.spec.ts', 'private-room.spec.ts', 'confirmations.spec.ts', 'den-journey.spec.ts', 'assignments.spec.ts', 'relay-settings.spec.ts', 'agent-receipts.spec.ts', 'context.spec.ts', 'persistent-groups.spec.ts', 'workspace.spec.ts', 'chat-comfort.spec.ts', 'phone-chat.spec.ts', 'messages.spec.ts', 'quiet.spec.ts', 'contact-card.spec.ts', 'model-shortcuts.spec.ts', 'chat-reliability.spec.ts', 'nostr-rooms.spec.ts', 'room-switching.spec.ts', 'conversation-search.spec.ts', 'drafts.spec.ts', 'home.spec.ts', 'site.spec.ts', 'safari-ice.spec.ts', 'updates.spec.ts'], + testMatch: ['knock.spec.ts', 'private-room.spec.ts', 'confirmations.spec.ts', 'den-journey.spec.ts', 'assignments.spec.ts', 'relay-settings.spec.ts', 'agent-receipts.spec.ts', 'context.spec.ts', 'persistent-groups.spec.ts', 'workspace.spec.ts', 'chat-comfort.spec.ts', 'phone-chat.spec.ts', 'messages.spec.ts', 'quiet.spec.ts', 'contact-card.spec.ts', 'model-shortcuts.spec.ts', 'chat-reliability.spec.ts', 'room-archive.spec.ts', 'nostr-rooms.spec.ts', 'room-switching.spec.ts', 'conversation-search.spec.ts', 'drafts.spec.ts', 'home.spec.ts', 'site.spec.ts', 'safari-ice.spec.ts', 'updates.spec.ts'], use: { ...devices['Desktop Safari'] }, }, { diff --git a/src/archive.test.ts b/src/archive.test.ts new file mode 100644 index 0000000..c216b80 --- /dev/null +++ b/src/archive.test.ts @@ -0,0 +1,211 @@ +import { describe, expect, it } from 'vitest' +import { finalizeEvent, generateSecretKey, getPublicKey, type Event } from 'nostr-tools/pure' +import { SimRelay, SimTransport } from '../test/sim-relay.js' +import { archiveTag, compareArchived, olderThan, reseedCandidates, type ArchiveMeta, type ArchiveQuery, type EventArchive } from './archive.js' +import { CHAT_RETENTION_SECONDS, ChatLog, MAX_CHAT_MESSAGES, MAX_CHAT_MESSAGES_PER_MINUTE, encodeChatEvent, type ChatMessage } from './chat.js' +import type { Filter } from 'nostr-tools/filter' +import type { RoomPolicy } from './types.js' +import { createDeviceCredential } from './credential.js' +import { localIdentity } from './identity.js' +import { deriveRoom } from './room.js' +import { RoomSession } from './session.js' +import { KINDS } from './kinds.js' + +const NOW = 1_800_000_000 + +/** The contract, in memory: what `app/src/room-archive.ts` does on disk. */ +class MemoryArchive implements EventArchive { + readonly events = new Map() + readonly meta = new Map() + released: string[] = [] + keep(event: Event, meta?: ArchiveMeta): void { this.events.set(event.id, event); this.meta.set(event.id, meta) } + release(q: { kind: number; d: string }): void { this.released.push(q.d) } + async read(q: ArchiveQuery): Promise { + return [...this.events.values()] + .filter(e => e.kind === q.kind && archiveTag(e) === q.d && (q.since === undefined || e.created_at >= q.since) && (!q.before || olderThan(e, q.before))) + .sort(compareArchived).slice(0, q.limit) + } +} + +async function room(seed = 7) { + const { roomId, roomKey } = deriveRoom(new Uint8Array(32).fill(seed)) + const deviceSk = generateSecretKey() + const credential = await createDeviceCredential({ + identity: localIdentity(generateSecretKey()), devicePubkey: getPublicKey(deviceSk), roomId, expiresAt: NOW + 3600, + }) + const message = (text: string, sentAt = NOW, id = `m-${text}`): Event => { + const msg: ChatMessage = { id, participant: credential.pubkey, device: getPublicKey(deviceSk), credential, text, sentAt } + return encodeChatEvent(msg, { roomId, roomKey, deviceSk }) + } + return { roomId, roomKey, deviceSk, credential, message } +} + +describe('reseed decisions', () => { + const at = (n: number): Event => ({ id: n.toString(16).padStart(64, '0'), created_at: 1000 + n, kind: 1460, pubkey: '', tags: [['d', 'x']], content: '', sig: '' }) + + it('hands back only what a relay that returned fewer is missing, newest first and bounded', () => { + const archived = [at(1), at(2), at(3), at(4)] + expect(reseedCandidates(archived, new Set([at(2).id]))).toEqual([at(4), at(3), at(1)]) + expect(reseedCandidates(archived, new Set([at(2).id]), 2)).toEqual([at(4), at(3)]) + expect(reseedCandidates(archived, new Set())).toHaveLength(4) + }) + + it('leaves alone a relay that returned as many as the archive holds, whatever it holds', () => { + const archived = [at(1), at(2)] + expect(reseedCandidates(archived, new Set([at(1).id, at(2).id]))).toEqual([]) + expect(reseedCandidates(archived, new Set([at(1).id, at(9).id]))).toEqual([]) + expect(reseedCandidates([], new Set())).toEqual([]) + }) + + it('pages newest first with a strict cursor, a tie on the second broken on id', () => { + const a = { created_at: 5, id: 'b' }, b = { created_at: 5, id: 'a' }, c = { created_at: 4, id: 'z' } + expect([c, b, a].sort(compareArchived)).toEqual([a, b, c]) + expect(olderThan(b, { at: a.created_at, id: a.id })).toBe(true) + expect(olderThan(a, { at: a.created_at, id: a.id })).toBe(false) + expect(olderThan(c, { at: 5, id: '' })).toBe(true) + }) +}) + +describe('a chat log over an archive', () => { + it('opens on what the device kept when every relay has forgotten', async () => { + const r = await room() + const archive = new MemoryArchive() + for (const text of ['one', 'two', 'three']) archive.keep(r.message(text)) + const log = new ChatLog({ ...r, transport: new SimTransport(new SimRelay()), archive, now: () => NOW }) + await expect.poll(() => log.messages().length).toBe(3) + expect(log.messages().map(m => m.text).sort()).toEqual(['one', 'three', 'two']) + expect(log.messages()[0]!.lane).toBeUndefined() + log.close() + }) + + it('decodes an archived event by the same rules as a relay one: forged, foreign and gated are refused', async () => { + const r = await room() + const other = await room(8) + const archive = new MemoryArchive() + const good = r.message('good') + archive.keep(good) + // Same d tag, but under another room's key. + archive.keep({ ...other.message('foreign'), tags: [['d', r.roomId]] }) + // A valid event whose signature no longer matches its body. + const forged = r.message('forged') + archive.keep({ ...forged, content: good.content }) + // Signed by a device that the credential inside does not name. + const stranger = generateSecretKey() + archive.keep(finalizeEvent({ kind: KINDS.CHAT, created_at: NOW, tags: [['d', r.roomId]], content: forged.content }, stranger)) + const log = new ChatLog({ ...r, transport: new SimTransport(new SimRelay()), archive, now: () => NOW }) + await expect.poll(() => log.hasOlder).toBe(true) + expect(log.messages().map(m => m.text)).toEqual(['good']) + + // A gated room refuses the same archived message it would refuse from a relay. + const gated = new ChatLog({ ...r, transport: new SimTransport(new SimRelay()), archive, now: () => NOW, + policy: { tier: 'kith', admitted: [getPublicKey(generateSecretKey())] } }) + await expect.poll(() => gated.hasOlder).toBe(true) + expect(gated.messages()).toEqual([]) + log.close(); gated.close() + }) + + it('shows the lane once a relay copy of an archived message arrives', async () => { + const r = await room() + const archive = new MemoryArchive() + const event = r.message('seen before') + archive.keep(event) + const relay = new SimRelay() + const inner = new SimTransport(relay) + const transport = { + publish: (e: Event) => inner.publish(e), + subscribe: (fs: Filter[], on: (e: Event, via?: string) => void, eose?: () => void) => inner.subscribe(fs, e => on(e, 'wss://relay.example'), eose), + close: () => inner.close(), + describe: () => [{ url: 'wss://relay.example', read: true, write: true }], + } + const log = new ChatLog({ ...r, transport, archive, now: () => NOW }) + await expect.poll(() => log.messages().length).toBe(1) + expect(log.messages()[0]!.lane).toBeUndefined() + relay.publish(event) + expect(log.messages()[0]!.lane).toBe('public') + log.close() + expect(archive.released).toEqual([r.roomId]) + }) + + it('keeps nothing the rate limit refused, and marks what a quiet room said', async () => { + const r = await room() + const relay = new SimRelay() + const archive = new MemoryArchive() + const log = new ChatLog({ ...r, transport: new SimTransport(relay), archive, now: () => NOW }) + for (let i = 0; i < MAX_CHAT_MESSAGES_PER_MINUTE + 5; i++) relay.publish(r.message(`flood ${i}`)) + expect(log.messages()).toHaveLength(MAX_CHAT_MESSAGES_PER_MINUTE) + expect(archive.events.size).toBe(MAX_CHAT_MESSAGES_PER_MINUTE) + log.close() + + const quietArchive = new MemoryArchive() + const policy: RoomPolicy = { tier: 'open', quiet: true, members: [r.credential.pubkey] } + const quiet = new ChatLog({ ...r, transport: new SimTransport(new SimRelay()), archive: quietArchive, policy, now: () => NOW }) + await quiet.send('hush') + expect([...quietArchive.meta.values()]).toEqual([{ quiet: true }]) + quiet.close() + }) + + it('keeps every event it accepts, its own sends included, and nothing it refused', async () => { + const r = await room() + const relay = new SimRelay() + const archive = new MemoryArchive() + const log = new ChatLog({ ...r, transport: new SimTransport(relay), archive, now: () => NOW }) + await log.send('mine') + relay.publish(r.message('theirs')) + relay.publish(finalizeEvent({ kind: KINDS.CHAT, created_at: NOW, tags: [['d', r.roomId]], content: 'rubbish' }, generateSecretKey())) + expect([...archive.events.values()]).toHaveLength(2) + expect(log.messages().map(m => m.text).sort()).toEqual(['mine', 'theirs']) + log.close() + }) + + it('pages back past the retention window and the message cap, a page at a time', async () => { + const r = await room() + const archive = new MemoryArchive() + const old = NOW - CHAT_RETENTION_SECONDS - 90 * 24 * 60 * 60 + const total = MAX_CHAT_MESSAGES + 150 + // Spread over the window and far past it: 100 inside, the rest older. + // Three seconds apart, inside the per-sender rate every message obeys. + for (let i = 0; i < total; i++) archive.keep(r.message(`n${i}`, i < total - 100 ? old + 3 * i : NOW - 3 * (total - i))) + const log = new ChatLog({ ...r, transport: new SimTransport(new SimRelay()), archive, now: () => NOW }) + await expect.poll(() => log.messages().length, { timeout: 30_000 }).toBe(100) + expect(log.hasOlder).toBe(true) + let read = 0 + for (let step = 0; step < 20 && log.hasOlder; step++) read += await log.loadOlder() + expect(read).toBe(total - 100) + expect(log.hasOlder).toBe(false) + expect(log.messages()).toHaveLength(total) + expect(log.messages()[0]!.text).toBe('n0') + expect(await log.loadOlder()).toBe(0) + log.close() + // Six hundred signatures and credentials, twice: slow on a busy machine. + }, 60_000) + + it('a session that a relay forgot still reaches its epoch and its history from the archive', async () => { + const secret = new Uint8Array(32).fill(21) + const authoritySk = generateSecretKey() + const authority = getPublicKey(authoritySk) + const aliceKeys = { identity: localIdentity(generateSecretKey()), deviceSk: generateSecretKey() } + const archive = new MemoryArchive() + const base = { secret, now: () => NOW, announceJitterMs: 0, authority, epochSettleMs: 0 } + const relay = new SimRelay({ replay: true }) + const keeper = new RoomSession({ ...base, transport: new SimTransport(relay), identity: localIdentity(generateSecretKey()), deviceSk: generateSecretKey(), name: 'Keeper' }) + const alice = new RoomSession({ ...base, ...aliceKeys, transport: new SimTransport(relay), name: 'Alice', archive }) + await keeper.join([], {}) + await alice.join([], {}) + for (let i = 0; i < 10; i++) await new Promise(resolve => setTimeout(resolve, 0)) + await keeper.rekey({ authoritySk }) + for (let i = 0; i < 10; i++) await new Promise(resolve => setTimeout(resolve, 0)) + expect(alice.epoch).toBe(1) + await keeper.chat.send('said in epoch 1') + expect(alice.chat.messages().map(m => m.text)).toEqual(['said in epoch 1']) + expect([...archive.events.values()].some(e => e.kind === KINDS.ROOM_REKEY)).toBe(true) + alice.leave(); keeper.leave() + + // Every relay forgot: the rekey and the chat are only on Alice's device. + const again = new RoomSession({ ...base, ...aliceKeys, transport: new SimTransport(new SimRelay({ replay: true })), name: 'Alice', archive }) + await again.join([], {}) + await expect.poll(() => again.chat.messages().length).toBe(1) + expect(again.epoch).toBe(1) + expect(again.chat.messages().map(m => m.text)).toEqual(['said in epoch 1']) + again.leave() + }, 30_000) +}) diff --git a/src/archive.ts b/src/archive.ts new file mode 100644 index 0000000..01b1fa8 --- /dev/null +++ b/src/archive.ts @@ -0,0 +1,96 @@ +import type { Event } from 'nostr-tools/pure' + +/** + * A device's own copy of a room's original signed events. + * + * Room history used to be exactly as durable as the one relay that still + * stored it: every client read chat from relays every time, and a relay that + * forgot took the room's past with it. An archive keeps each event as it + * arrived, ciphertext and signature, never anything decrypted, so the copy is + * as safe to hold as the relay's was, can be checked again when it is read + * back, and can be handed to a relay unchanged. + * + * What goes in is only what the caller has already accepted: a chat event + * that decoded under the room's rules, a rekey whose signature checked. What + * comes out is untrusted again and goes through the same decoder as an event + * off a relay, so an archive can widen what a device remembers and never what + * it believes. The web implementation is `app/src/room-archive.ts`. + */ +export interface EventArchive { + /** Keep one accepted event. Idempotent, and never throws: it runs inside a + * relay subscription handler. Writing may finish later. */ + keep(event: Event, meta?: ArchiveMeta): void + /** Kept events of one kind under one `d` tag, newest first. */ + read(query: ArchiveQuery): Promise + /** Nobody is reading this conversation now: an archive that holds it in + * memory may let it go. */ + release?(query: Pick): void +} + +/** What an archive records beside an event, sealed with it. */ +export interface ArchiveMeta { + /** It came through a quiet room, where chat never appears on a relay as + * a bare room event. Such an event is never handed back to a relay, + * whatever the room is opened as later. See `quiet.ts`. */ + quiet?: boolean +} + +export interface ArchiveQuery { + kind: number + d: string + /** Oldest `created_at` to include. */ + since?: number + /** Only events strictly older than this position, in newest-first order: + * earlier `created_at`, or the same second and a smaller id. The cursor a + * reader pages back with. */ + before?: ArchiveCursor + limit: number + /** Only events that may be handed back to a relay: not a quiet room's. */ + reseedable?: boolean +} + +export interface ArchiveCursor { at: number; id: string } + +/** The `d` tag an event is filed under, or undefined. */ +export function archiveTag(event: Pick): string | undefined { + const d = event.tags.find(tag => tag[0] === 'd')?.[1] + return typeof d === 'string' && d.length > 0 && d.length <= 128 ? d.toLowerCase() : undefined +} + +/** Newest first; a tie on the second breaks on id, the same order a cursor + * pages through. */ +export function compareArchived(a: Pick, b: Pick): number { + return b.created_at - a.created_at || (a.id < b.id ? 1 : a.id > b.id ? -1 : 0) +} + +/** Whether `event` sits strictly after `cursor` in newest-first order. */ +export function olderThan(event: Pick, cursor: ArchiveCursor): boolean { + return event.created_at < cursor.at || (event.created_at === cursor.at && event.id < cursor.id) +} + +/** The most events one device republishes to one relay for one conversation + * in one pass: the window a reader renders, and no more. */ +export const MAX_RESEED_EVENTS = 500 + +/** + * Which of a device's archived events a relay should be handed back. + * + * Only when the relay returned fewer of the conversation's events than the + * archive holds for the same window: a relay that returned as many or more + * is not forgetful, whatever it is missing, and is left alone. Then the + * archived events it did not return, newest first, at most `limit`. + * + * `floor` is for an answer that may be cut short: a relay that caps how many + * events it returns hands back its newest and stops, so an archived event + * older than the oldest one it returned says nothing about whether the relay + * holds it. Such events are left out of the count and out of the result. + * + * The events are the archive's originals, unchanged; the caller verifies + * each signature again before it publishes anything. + */ +export function reseedCandidates(archived: readonly Event[], returned: ReadonlySet, limit = MAX_RESEED_EVENTS, floor?: number): Event[] { + if (!Number.isSafeInteger(limit) || limit < 0) throw new Error('reseed limit must be a whole number') + const judged = floor === undefined ? archived : archived.filter(event => event.created_at >= floor) + if (returned.size >= judged.length) return [] + return judged.filter(event => !returned.has(event.id)).sort(compareArchived).slice(0, limit) +} diff --git a/src/chat.ts b/src/chat.ts index 102944c..5d0124f 100644 --- a/src/chat.ts +++ b/src/chat.ts @@ -22,6 +22,8 @@ import { sanitiseDisplayName } from './display-name.js' import { evaluateAccess } from './access.js' import { inspectAgentOwnershipSignature, normaliseAgentOwnership, verifyAgentOwnership } from './ownership.js' import type { RelayTransport } from './relay-pool.js' +import { olderThan, type ArchiveCursor, type ArchiveMeta, type EventArchive } from './archive.js' +import { isQuietPolicy } from './quiet.js' import { laneOfRelayUrl, laneOfRelays, type Lane } from './lane.js' import type { AgentOwnership, DeviceCredential, KindredProof, RoomPolicy } from './types.js' @@ -29,6 +31,11 @@ export const MAX_CHAT_TEXT_LENGTH = 2_000 export const CHAT_RETENTION_SECONDS = 30 * 24 * 60 * 60 export const MAX_CHAT_MESSAGES = 500 export const MAX_CHAT_MESSAGES_PER_MINUTE = 30 +/** How many archived messages one step back through history reads. */ +export const CHAT_ARCHIVE_PAGE = 100 +/** Archived events decoded between two yields to the page: a decode is a + * signature, a credential and a decryption, a few milliseconds each. */ +const ARCHIVE_DECODE_CHUNK = 40 const CHANNEL_ID_INFO = 'kithmoot/v1/channel-id/' const CHANNEL_KEY_INFO = 'kithmoot/v1/channel-key/' @@ -643,6 +650,14 @@ export interface ChatLogOptions { owner?: AgentOwnership /** Historical signed claim; never current-owner authority. */ ownerClaim?: AgentOwnership + /** + * This device's own copy of the room's events. Given one, the log shows + * what it holds before any relay answers, keeps every event it accepts, + * and can page back past the retention window with `loadOlder`. What it + * reads from the archive is decoded by exactly the rules a relay's events + * are. See `archive.ts`. + */ + archive?: EventArchive } /** What `send` may say beyond the text. */ @@ -699,6 +714,20 @@ export class ChatLog { #unsub: () => void /** The epoch this log reads and writes. Undefined is epoch 0. */ #epoch?: EpochRoot + /** How many messages the log holds: `MAX_CHAT_MESSAGES`, plus a page for + * every step back through the archive a reader asked for. */ + #window = MAX_CHAT_MESSAGES + /** The oldest send time a reader has paged back to, which moves the + * retention cut for this log. Undefined until it pages. */ + #pagedTo?: number + /** Where the next step back through the archive starts; undefined until + * the archive has answered. */ + #cursor?: ArchiveCursor + #archiveDone = false + #paging?: Promise + /** Messages read from the archive, by event id, until a relay's copy says + * which lane they travelled. Bounded like `#decoded`. */ + readonly #laneless = new Map() constructor(opts: ChatLogOptions) { this.#opts = opts @@ -708,9 +737,18 @@ export class ChatLog { this.#unsub = this.#subscribe() } - #subscribe(): () => void { + /** The `d` tag this log reads and writes under now: public on the wire, + * and what a caller asks a relay or the archive for. */ + get stream(): string { const root = rootOf({ roomId: this.#opts.roomId, roomKey: this.#opts.roomKey, epoch: this.#epoch }) - const { id } = deriveChannel(root.id, root.key, this.#opts.channel) + return deriveChannel(root.id, root.key, this.#opts.channel).id + } + + #subscribe(): () => void { + const id = this.stream + this.#cursor = undefined + this.#archiveDone = !this.#opts.archive + if (this.#opts.archive) void this.#readArchive(id, this.#epoch) return this.#opts.transport.subscribe( // The newest the log can hold, not the whole retention window: the // rest would be decoded only to fall off the end. @@ -748,10 +786,99 @@ export class ChatLog { */ rekey(next: EpochRoot): void { this.#unsub() + this.#release() this.#epoch = next this.#unsub = this.#subscribe() } + /** What this device holds, first: the same window a relay is asked for, + * so a room opens on its history before any relay has answered. */ + async #readArchive(d: string, epoch: EpochRoot | undefined): Promise { + const since = this.#now() - CHAT_RETENTION_SECONDS + let events: Event[] = [] + try { + events = await this.#opts.archive!.read({ kind: KINDS.CHAT, d, since, limit: MAX_CHAT_MESSAGES }) + } catch { + // An archive that cannot be read is a device without one; relays still answer. + } + if (this.#closed || this.#epoch !== epoch) return + const oldest = events[events.length - 1] + this.#cursor = oldest ? { at: oldest.created_at, id: oldest.id } : { at: since, id: '' } + // Told even when nothing in the window was kept, so a reader can ask + // for what lies before it. + if (!await this.#ingestAll(events, epoch)) this.#notify() + } + + /** Whether the archive may hold messages older than the log shows. */ + get hasOlder(): boolean { + return !this.#archiveDone && this.#cursor !== undefined + } + + /** + * Step back through this device's archive: read the next page of older + * events and show them, past the retention window if need be. Resolves + * to how many events the archive handed over; zero means there are no + * more. The render window grows by what a reader asks for, never on its + * own. + */ + loadOlder(count = CHAT_ARCHIVE_PAGE): Promise { + if (this.#paging) return this.#paging + const archive = this.#opts.archive + const cursor = this.#cursor + if (!archive || !cursor || this.#archiveDone || this.#closed) return Promise.resolve(0) + const epoch = this.#epoch + const d = this.stream + const run = async (): Promise => { + let events: Event[] + try { + events = (await archive.read({ kind: KINDS.CHAT, d, before: cursor, limit: count })).filter(e => olderThan(e, cursor)) + } catch { + return 0 + } + if (this.#closed || this.#epoch !== epoch) return 0 + const oldest = events[events.length - 1] + if (!oldest) { + this.#archiveDone = true + return 0 + } + this.#cursor = { at: oldest.created_at, id: oldest.id } + this.#window += events.length + this.#pagedTo = Math.min(this.#pagedTo ?? Infinity, oldest.created_at) + await this.#ingestAll(events, epoch) + return events.length + } + this.#paging = run().finally(() => { this.#paging = undefined }) + return this.#paging + } + + /** Archived events in, a chunk at a time with the page given a turn in + * between, newest first, and one notification per chunk that showed + * something. Stops if the log closes or changes key meanwhile. Returns + * whether anything showed. */ + async #ingestAll(events: readonly Event[], epoch: EpochRoot | undefined): Promise { + let any = false + for (let i = 0; i < events.length; i += ARCHIVE_DECODE_CHUNK) { + if (i > 0) await new Promise(resolve => setTimeout(resolve, 0)) + if (this.#closed || this.#epoch !== epoch) return any + let changed = false + for (const event of events.slice(i, i + ARCHIVE_DECODE_CHUNK)) changed = this.#ingest(event, undefined, true) || changed + if (changed) this.#notify() + any ||= changed + } + return any + } + + /** Let the archive drop this conversation from memory. */ + #release(): void { + this.#opts.archive?.release?.({ kind: KINDS.CHAT, d: this.stream }) + } + + /** The oldest send time this log shows: the retention window, or further + * back when a reader paged there. */ + #floor(): number { + return Math.min(this.#now() - CHAT_RETENTION_SECONDS, this.#pagedTo ?? Infinity) + } + /** The channel this log is, or undefined for the main chat. */ get channel(): string | undefined { return this.#opts.channel @@ -862,6 +989,8 @@ export class ChatLog { throw new Error('This conversation has closed or changed its key. Copy your message into the current conversation to send it.') } await this.#opts.transport.publish(event) + // Kept once a relay has it, whether or not a relay echoes it back. + this.#opts.archive?.keep(event, this.#archiveMeta()) } } @@ -894,19 +1023,32 @@ export class ChatLog { close(): void { this.#closed = true this.#unsub() + this.#release() this.#listeners.clear() } - #ingest(event: Event, via?: string): void { - if (this.#decoded.has(event.id)) return + /** Returns whether the log changed. `fromArchive` events are neither kept + * again nor announced one by one; the caller notifies once. */ + #ingest(event: Event, via?: string, fromArchive = false): boolean { + if (this.#decoded.has(event.id)) { + // A message first read from the archive, now arriving from a relay: + // that relay is the lane it travelled. + const shown = fromArchive ? undefined : this.#laneless.get(event.id) + if (!shown) return false + this.#laneless.delete(event.id) + shown.lane = this.#laneOf(via) + if (shown.lane === undefined) return false + this.#notify() + return true + } // Older than everything a full log keeps: it would be decoded and // dropped straight away. `encodeChatEvent` writes `sentAt` as // `created_at`; a sender who puts an earlier one on the outside only // loses their own message. - const oldest = this.#messages.length >= MAX_CHAT_MESSAGES ? this.#messages[0] : undefined - if (oldest && event.created_at < oldest.sentAt) return + const oldest = this.#messages.length >= this.#window ? this.#messages[0] : undefined + if (oldest && event.created_at < oldest.sentAt) return false this.#decoded.add(event.id) - if (this.#decoded.size > MAX_CHAT_MESSAGES * 4) { + if (this.#decoded.size > this.#window * 4) { const first = this.#decoded.values().next().value if (first !== undefined) this.#decoded.delete(first) } @@ -918,32 +1060,51 @@ export class ChatLog { channel: this.#opts.channel, ...(this.#epoch ? { epoch: this.#epoch } : {}), }) - if (!msg) return - msg.lane = this.#laneOf(via) - if (msg.sentAt < this.#now() - CHAT_RETENTION_SECONDS) return - if (this.#seen.has(msg.id)) return + if (!msg) return false + // An archived event has no relay to name, and the lane it once took is + // not recorded, so it claims none until a relay's copy arrives. + msg.lane = fromArchive ? undefined : this.#laneOf(via) + const floor = this.#floor() + if (msg.sentAt < floor) return false + if (this.#seen.has(msg.id)) return false const senderTimes = (this.#senderTimes.get(msg.device) ?? []) - .filter((sentAt) => sentAt >= this.#now() - CHAT_RETENTION_SECONDS) + .filter((sentAt) => sentAt >= floor) if (senderTimes.filter((sentAt) => Math.abs(sentAt - msg.sentAt) < 60).length >= MAX_CHAT_MESSAGES_PER_MINUTE) { - return + return false } senderTimes.push(msg.sentAt) this.#senderTimes.set(msg.device, senderTimes) - while (this.#senderTimes.size > MAX_CHAT_MESSAGES) { + while (this.#senderTimes.size > this.#window) { const oldest = this.#senderTimes.keys().next().value if (oldest === undefined) break this.#senderTimes.delete(oldest) } this.#seen.add(msg.id) + // Accepted by every rule the room has, the rate limit included, so worth + // keeping whether or not this log has room to show it for long. + if (!fromArchive) this.#opts.archive?.keep(event, this.#archiveMeta()) + else { + this.#laneless.set(event.id, msg) + if (this.#laneless.size > this.#window) this.#laneless.delete(this.#laneless.keys().next().value!) + } this.#messages.push(msg) this.#messages.sort(compareMessages) - while (this.#messages.length > MAX_CHAT_MESSAGES) { + while (this.#messages.length > this.#window) { const removed = this.#messages.shift() if (removed) this.#seen.delete(removed.id) } + if (!fromArchive) this.#notify() + return true + } + + #archiveMeta(): ArchiveMeta | undefined { + return isQuietPolicy(this.#opts.policy) ? { quiet: true } : undefined + } + + #notify(): void { const snapshot = this.messages() // Guarded: decodeChatEvent is written never to throw precisely because // this runs inside a relay subscription handler, and a throwing caller diff --git a/src/index.ts b/src/index.ts index 19377f6..4ccd623 100644 --- a/src/index.ts +++ b/src/index.ts @@ -250,11 +250,14 @@ export { CHAT_RETENTION_SECONDS, MAX_CHAT_MESSAGES, MAX_CHAT_MESSAGES_PER_MINUTE, + CHAT_ARCHIVE_PAGE, MAX_CHAT_ATTACHMENTS, MAX_ATTACHMENT_URL_LENGTH, MAX_ATTACHMENT_NAME_LENGTH, normaliseAttachment, } from './chat.js' +export { archiveTag, compareArchived, olderThan, reseedCandidates, MAX_RESEED_EVENTS } from './archive.js' +export type { EventArchive, ArchiveQuery, ArchiveCursor, ArchiveMeta } from './archive.js' export type { ChatMessage, ChatMessageKind, diff --git a/src/relay-pool.test.ts b/src/relay-pool.test.ts index 71f6b38..53e5433 100644 --- a/src/relay-pool.test.ts +++ b/src/relay-pool.test.ts @@ -139,6 +139,37 @@ describe('NostrRelayPool', () => { expect(b.stored.map((e) => e.id)).toEqual([event.id]) }) + it('asks one relay alone and writes to one relay alone, quietly', async () => { + // A room's shared subscription hides which relay held what; a device + // putting history back has to ask, and write to, each one separately. + const onlyA = evt(1460, [['d', 'room']]) + const onBoth = evt(1460, [['d', 'room']]) + a.seed(onlyA); a.seed(onBoth); b.seed(onBoth) + const fromB = await pool.query(URL_B, [{ kinds: [1460], '#d': ['room'] }], 2_000) + expect(fromB.complete).toBe(true) + expect(fromB.events.map(e => e.id)).toEqual([onBoth.id]) + const fromA = await pool.query(URL_A, [{ kinds: [1460], '#d': ['room'] }], 2_000) + expect(fromA.complete).toBe(true) + expect(fromA.events.map(e => e.id).sort()).toEqual([onlyA.id, onBoth.id].sort()) + await pool.publishQuietly(URL_B, onlyA) + expect(b.stored.map(e => e.id)).toContain(onlyA.id) + expect(a.stored.filter(e => e.id === onlyA.id)).toHaveLength(1) + // Quiet: neither a refusal nor an acceptance is shown as the relay's health. + b.rejectPublishes = true + await expect(pool.publishQuietly(URL_B, evt(1460, [['d', 'room']]))).rejects.toBeTruthy() + expect(pool.health()[1]!.lastError).toBeUndefined() + expect(pool.health()[1]!.lastPublishedAt).toBeUndefined() + expect(pool.publishing).toBe(false) + await expect(pool.publishQuietly('wss://elsewhere.test', onlyA)).rejects.toThrow(/not a writable relay/) + await expect(pool.query('wss://elsewhere.test', [{ kinds: [1460] }])).rejects.toThrow(/not a readable relay/) + }) + + it('says a query that timed out is unknown, not empty', async () => { + a.silent = true + a.seed(evt(1460, [['d', 'room']])) + expect(await pool.query(URL_A, [{ kinds: [1460], '#d': ['room'] }], 200)).toEqual({ events: [], complete: false }) + }) + it('says when every relay has answered, so a caller can close without cutting a slow one off', async () => { await pool.publish(evt()) await pool.settled() diff --git a/src/relay-pool.ts b/src/relay-pool.ts index b27cbda..736fda6 100644 --- a/src/relay-pool.ts +++ b/src/relay-pool.ts @@ -391,6 +391,58 @@ export class NostrRelayPool implements RelayTransport { for (const sub of this.#subscriptions) if (sub.bindings.has(url)) this.#startRelay(sub, url) } + /** + * Ask one configured, readable relay alone and collect what it returns. + * A room's shared subscription cannot say which relay holds what, because + * the first relay to deliver an event hides every later copy; this can. + * The events are signature-checked like every other. + * + * `complete` is true only when the relay itself said it had sent + * everything. nostr-tools fakes an end-of-stored-events when its own wait + * runs out, so that wait is set well past `timeoutMs` and this function's + * own timer decides first: a relay that was slow, closed the request or + * dropped the socket gives an answer that is unknown, not empty. Nothing + * is resent and no health is marked. + */ + query(url: string, filters: Filter[], timeoutMs = 8_000): Promise<{ events: Event[]; complete: boolean }> { + if (this.#closed) return Promise.reject(new Error('pool is closed')) + const relay = this.#relays.find(r => r.url === normalizeURL(url)) + if (!relay?.read) return Promise.reject(new Error('not a readable relay of this pool')) + return new Promise(resolve => { + const events = new Map() + let settled = false + const finish = (complete: boolean): void => { + if (settled) return + settled = true + clearTimeout(timer) + handle.close() + resolve({ events: [...events.values()], complete }) + } + const handle = this.#pool.subscribeMap(filters.map(filter => ({ url: relay.url, filter: { ...filter } })), { + abort: this.#abort.signal, + maxWait: timeoutMs * 2 + 5_000, + onevent: event => { events.set(event.id, event) }, + oneose: () => finish(true), + onclose: () => finish(false), + }) + const timer = setTimeout(() => finish(false), timeoutMs) + }) + } + + /** + * Hand one configured, writable relay an event, quietly: one attempt, no + * reconnect on a timeout, no health marks and not counted as a publish in + * flight. For background work such as putting back history a relay + * forgot, whose refusals are not the relay failing the room and whose + * retries must never disturb the room's own sockets and subscriptions. + */ + async publishQuietly(url: string, event: Event): Promise { + if (this.#closed) throw new Error('pool is closed') + const relay = this.#relays.find(r => r.url === normalizeURL(url)) + if (!relay?.write) throw new Error('not a writable relay of this pool') + await this.#pool.publish([relay.url], event, { abort: this.#abort.signal })[0] + } + describe(): RelayConfig[] { return this.#relays.map(relay => { if (!this.circleAtUse) return { ...relay } diff --git a/src/session.ts b/src/session.ts index 3206845..582b164 100644 --- a/src/session.ts +++ b/src/session.ts @@ -24,6 +24,7 @@ import { encodeDescriptorEvent, decodeDescriptorEvent } from './descriptor.js' import { encodeCallBellEvent, type CallBellState } from './call-bell.js' import { ChatLog } from './chat.js' import type { EpochRoot } from './chat.js' +import type { EventArchive } from './archive.js' import { EpochRefusedError, decodeRekeyEvent, @@ -131,6 +132,12 @@ export interface CallView { export interface RoomSessionBaseOptions { transport: RelayTransport + /** This device's own copy of the room's events: every chat log keeps and + * reads through it, and the authority's rekeys are kept and replayed at + * join, so a room a relay forgot still opens where this device left it. + * Everything read back is checked by the same rules as a relay's events. + * See `archive.ts`. */ + archive?: EventArchive secret: Uint8Array /** This endpoint's own key. Never the participant's. */ deviceSk: Uint8Array @@ -614,6 +621,10 @@ export class RoomSession { (event) => this.#ingestRekey(event), ) try { + // The rekeys this device kept, through the same door as a relay's: + // a relay that forgot them cannot put this device back in epoch 0. + const kept = await this.#opts.archive?.read({ kind: KINDS.ROOM_REKEY, d: this.roomId, limit: 1_000 }).catch(() => []) + for (const event of kept ?? []) this.#ingestRekey(event) await this.#settleEpoch() } catch (err) { this.#unsubRekey?.() @@ -749,6 +760,7 @@ export class RoomSession { ...(this.#epochRoot() ? { epoch: this.#epochRoot() } : {}), ...(this.#ownerToCarry() ? { owner: this.#ownerToCarry() } : {}), ...(this.#ownerClaimToCarry() ? { ownerClaim: this.#ownerClaimToCarry() } : {}), + ...(this.#opts.archive ? { archive: this.#opts.archive } : {}), }) // Presence is live state, so it has to be restated and it has to lapse - @@ -896,7 +908,9 @@ export class RoomSession { #ingestRekey(event: Event): void { if (!this.#opts.authority) return const epoch = peekRekeyEvent(event, { roomId: this.roomId, authority: this.#opts.authority }) - if (epoch === null || epoch <= this.#epoch.epoch) return + if (epoch === null) return + this.#opts.archive?.keep(event) + if (epoch <= this.#epoch.epoch) return this.#pendingRekeys.set(epoch, event) // During join the settle step drains; after it, every rekey is acted on // the moment it arrives. @@ -1735,6 +1749,7 @@ export class RoomSession { ...(this.#epochRoot() ? { epoch: this.#epochRoot() } : {}), ...(this.#ownerToCarry() ? { owner: this.#ownerToCarry() } : {}), ...(this.#ownerClaimToCarry() ? { ownerClaim: this.#ownerClaimToCarry() } : {}), + ...(this.#opts.archive ? { archive: this.#opts.archive } : {}), }) this.#channels.set(name, log) return log diff --git a/test/room-archive.spec.ts b/test/room-archive.spec.ts new file mode 100644 index 0000000..1ddd968 --- /dev/null +++ b/test/room-archive.spec.ts @@ -0,0 +1,79 @@ +import { test, expect, type Browser, type BrowserContext, type Page } from '@playwright/test' +import { deriveRoom, encodeJoinUrl, generateRoomSecret } from '../src/room.js' +import { TEST_RELAY_HTTP } from './relays.js' + +// A room's history is as durable as the devices that saw it, not the relay +// that happened to keep it. On 27 September 2026 a two-person room showed +// 425 messages on one device and none on another, because only one relay +// still held the chat and no client kept a copy. Here the relay forgets the +// room outright: the device that was there still opens on its history, and +// puts it back, so a newcomer reads it from the relay again. + +function testRelay(baseURL: string): string { + const url = new URL('/__test-relay', baseURL) + url.protocol = 'wss:' + return url.href +} + +async function contextFor(browser: Browser, baseURL: string): Promise { + const context = await browser.newContext({ ignoreHTTPSErrors: true, serviceWorkers: 'block' }) + await context.route('**/turn', r => r.fulfill({ status: 503, body: '' })) + await context.routeWebSocket(url => url.protocol === 'wss:' && url.href !== testRelay(baseURL), ws => ws.close()) + return context +} + +async function join(page: Page, url: string, name: string): Promise { + await page.goto(url) + await page.locator('#displayName').fill(name) + await page.locator('#join').click() + await expect(page.locator('#roomArea')).toBeVisible() +} + +/** How many of the room's main chat events the relay holds; `forget` + * drops them all first. */ +async function stored(roomId: string, forget = false): Promise { + const response = await fetch(`${TEST_RELAY_HTTP}/__stored?kind=1460&d=${roomId}`, { method: forget ? 'DELETE' : 'GET' }) + return (await response.json() as { count: number }).count +} + +test('a room its relay forgot still opens on its history, and the device that kept it puts it back', async ({ browser, baseURL }) => { + const secret = generateRoomSecret() + const { roomId } = deriveRoom(secret) + const url = encodeJoinUrl(baseURL!, secret, [testRelay(baseURL!)]) + const kept = await contextFor(browser, baseURL!) + const newcomer = await contextFor(browser, baseURL!) + try { + const first = await kept.newPage() + await join(first, url, 'Keeper of history') + const said = ['The first thing said', 'The second thing said', 'The third thing said'] + for (const text of said) { + await first.locator('#chatInput').fill(text) + await first.locator('#chatInput').press('Enter') + } + await expect(first.locator('#chatLog .msg')).toHaveCount(said.length) + await expect.poll(() => stored(roomId)).toBeGreaterThanOrEqual(said.length) + const held = await stored(roomId) + await first.close() + + // The relay forgets the room's chat, as a public relay does. + expect(await stored(roomId, true)).toBe(0) + + // The same browser comes back: its history is on screen with no relay + // holding any of it. + const again = await kept.newPage() + await join(again, url, 'Keeper of history') + for (const text of said) await expect(again.locator('#chatLog')).toContainText(text) + await expect(again.locator('#chatLog .msg')).toHaveCount(said.length) + + // And it hands the originals back to the relay. + await expect.poll(() => stored(roomId), { timeout: 60_000 }).toBe(held) + + // So somebody who was never here reads them off the relay. + const reader = await newcomer.newPage() + await join(reader, url, 'Newcomer') + for (const text of said) await expect(reader.locator('#chatLog')).toContainText(text) + } finally { + await kept.close() + await newcomer.close() + } +}) diff --git a/test/workspace.spec.ts b/test/workspace.spec.ts index f16375d..2c322a2 100644 --- a/test/workspace.spec.ts +++ b/test/workspace.spec.ts @@ -378,7 +378,8 @@ test('switching rooms restores independent reading places after delayed history hold = true await switchTo('Reading room') await expect.poll(() => held.length).toBeGreaterThan(0) - await expect(log.locator('.msg')).toHaveCount(0) + // The room's saved history is on screen before the relays answer. + await expect(log.locator('.msg')).toHaveCount(12) await expect(page.locator('#newMessages')).toBeVisible() release() await expect(log.locator('.msg')).toHaveCount(13) diff --git a/test/ws-relay.mjs b/test/ws-relay.mjs index 5b416a1..1743dbf 100644 --- a/test/ws-relay.mjs +++ b/test/ws-relay.mjs @@ -19,7 +19,8 @@ // RELAY_PORT=7778 node test/ws-relay.mjs // // A plain HTTP GET answers 200, which is what lets Playwright's `webServer` -// wait on it. The same test-only server accepts encrypted Blossom uploads in +// wait on it. `/__stored?kind=&d=` counts stored events, and DELETE on it +// forgets them. The same test-only server accepts encrypted Blossom uploads in // memory, capped at the production endpoint's 270 MiB request limit. Nothing // here is a product or a persistent store. @@ -104,6 +105,21 @@ const http = createServer((req, res) => { res.end(bytes) return } + // Test-only: make this relay forget, or count, the stored events of one + // kind under one `d` tag, as a public relay that drops a room's history + // does. `room-archive.spec.ts` uses it to prove a device puts back what a + // relay forgot. + if (url.pathname === '/__stored') { + const kind = Number(url.searchParams.get('kind')) + const d = url.searchParams.get('d') + const matching = (e) => e.kind === kind && (!d || e.tags.some((t) => t[0] === 'd' && t[1] === d)) + if (req.method === 'DELETE') { + for (let i = stored.length - 1; i >= 0; i--) if (matching(stored[i])) stored.splice(i, 1) + } + res.writeHead(200, { 'content-type': 'application/json' }) + res.end(JSON.stringify({ count: stored.filter(matching).length })) + return + } res.writeHead(200, { 'content-type': 'text/plain' }) res.end('kithmoot test relay\n') })