diff --git a/e2e/harness/index.ts b/e2e/harness/index.ts index 5537f10d..ccf86465 100644 --- a/e2e/harness/index.ts +++ b/e2e/harness/index.ts @@ -21,6 +21,7 @@ import { formatTranscript, installWebSocketRoutes, silenceRelay, + withholdEose, } from "./net/websocket" import {watchFaults} from "./faults" import {boot} from "./app/boot" @@ -47,6 +48,7 @@ export { getPublishedEvents, getTranscript, silenceRelay, + withholdEose, } from "./net/websocket" export {readCachedEvents} from "./app/cache" export { @@ -137,6 +139,8 @@ export type PageOptions = { hosting?: HostingFixtures // Relay urls that take the socket and answer nothing; silenceRelay does the same mid-test. silent?: string[] + // Relay urls that serve events and never eose, so every request runs out its own deadline. + eoseless?: string[] } export type Harness = { @@ -232,6 +236,10 @@ export const test = base.extend({ silenceRelay(context, url) } + for (const url of options.eoseless ?? []) { + withholdEose(context, url) + } + await mockRelayInfo(context, options.relayInfo ?? {}) await mockAnalytics(context) await mockDufflepud(context) diff --git a/e2e/harness/net/websocket.ts b/e2e/harness/net/websocket.ts index 1a226abd..563c7870 100644 --- a/e2e/harness/net/websocket.ts +++ b/e2e/harness/net/websocket.ts @@ -28,6 +28,7 @@ type Traffic = { leaks: Set forgotten: Set silenced: Set + eoseless: Set } const trafficStore = makeContextStore("installWebSocketRoutes") @@ -51,6 +52,19 @@ const openEmptyRelay = (): RelayConnection => { } } +// A relay that serves what it has and never says it is done, so a request runs out its own deadline. +const openEoselessRelay = (connection: RelayConnection): RelayConnection => ({ + onMessage(listener) { + connection.onMessage(message => { + if (message[0] !== RelayMessageType.Eose) { + listener(message) + } + }) + }, + send: message => connection.send(message), + close: () => connection.close(), +}) + // A relay that takes the socket and says nothing, so the only way out is the caller's own deadline. const openSilentRelay = (): RelayConnection => ({ onMessage() {}, @@ -72,7 +86,9 @@ const serve = (traffic: Traffic, zooid: Zooid, route: WebSocketRoute) => { const relay = zooid.relays.get(url) if (relay) { - return relay.connect() + const connection = relay.connect() + + return traffic.eoseless.has(url) ? openEoselessRelay(connection) : connection } traffic.leaks.add(url) @@ -104,6 +120,7 @@ export const installWebSocketRoutes = (context: BrowserContext, zooid: Zooid) => leaks: new Set(), forgotten: new Set(), silenced: new Set(), + eoseless: new Set(), }) return context.routeWebSocket( @@ -136,6 +153,10 @@ export const forgetRelay = (context: BrowserContext, url: string) => export const silenceRelay = (context: BrowserContext, url: string) => trafficStore.get(context).silenced.add(normalizeRelayUrl(url)) +// Resolved at open too, so a page that boots into it passes `eoseless` to `as`. +export const withholdEose = (context: BrowserContext, url: string) => + trafficStore.get(context).eoseless.add(normalizeRelayUrl(url)) + // Every frame in both directions, oldest first. Attach it to a failing test. export const formatTranscript = (context: BrowserContext) => getTranscript(context) diff --git a/e2e/specs/rooms.spec.ts b/e2e/specs/rooms.spec.ts index 9aedea77..94380ee2 100644 --- a/e2e/specs/rooms.spec.ts +++ b/e2e/specs/rooms.spec.ts @@ -1040,3 +1040,20 @@ test("US-115 connect a wallet without losing the zap you were composing", async await expect(zap.getByRole("button", {name: "Send Zap"})).toBeVisible() await expect(amount).toHaveValue("210") }) + +test("fills a room from a relay that never says it is done", async ({seed, as}) => { + const scenario = await seed(({relay, user, at}) => { + const space = relay("space") + + space.room("general", {name: "General"}) + space.join(user.alice, "general") + space.join(user.bob, "general") + space.message(user.bob, "general", "the buoy is back on station", at(3, HOUR)) + }) + + const {url} = scenario.space("space") + const page = await as(users.alice, roomPath(url, "general"), {eoseless: [url]}) + + // Nothing will tell the feed its page is finished, so a room that waits for that stays empty. + await expect(message(page, "the buoy is back on station")).toBeVisible() +}) diff --git a/src/app/feeds.ts b/src/app/feeds.ts index ae0f5431..1b4d6215 100644 --- a/src/app/feeds.ts +++ b/src/app/feeds.ts @@ -547,7 +547,7 @@ export const makeFeed = ({ const unsubscribers = syncFeed({relays, filters, addEvents, removeEvents}) // Each relay answers a limit for itself, so a page runs out where the relay that gave least ran out. - const loadSpan = async (extension: Filter) => { + const loadSpan = async (extension: Filter, onAnswer?: (event: TrustedEvent) => void) => { let complete = false const pages = new Map() @@ -561,11 +561,15 @@ export const makeFeed = ({ } else { pages.set(url, {count: 1, lowest: event.created_at}) } + + onAnswer?.(event) } const found = await network.get().request({ relays, autoClose: true, + // The socket sends five per 100ms first come first served, so the feed on screen goes ahead of the space's background pulls. + priority: 1, signal: spanSignal(controller.signal), threshold: SPAN_THRESHOLD, filters: filters.map(filter => ({...filter, ...extension})), @@ -587,7 +591,11 @@ export const makeFeed = ({ const until = oldest const since = until - olderInterval - const {found, complete, pages} = await loadSpan({since, until, limit: PAGE_SIZE}) + + // A page comes back newest first, so each event proves the stretch above it has been answered. + const {found, complete, pages} = await loadSpan({since, until, limit: PAGE_SIZE}, event => + reach(event.created_at), + ) // A relay that answered with less than its limit has covered its whole span. let edge: Maybe @@ -632,7 +640,17 @@ export const makeFeed = ({ return {found: found.length, complete, exhausted: false} } - addEvents(relays.flatMap(url => Array.from(getEventsForUrl(url, filters)))) + // What the repository already holds is an answer too, so a page of it renders before any relay replies. + const stored = relays + .flatMap(url => Array.from(getEventsForUrl(url, filters))) + .sort(compareEventsAsc) + .slice(-PAGE_SIZE) + + if (stored.length > 0) { + reach(stored[0].created_at) + } + + addEvents(stored) return { events, @@ -725,6 +743,7 @@ export const makeCalendarFeed = ({ const found = await network.get().request({ relays, autoClose: true, + priority: 1, signal: spanSignal(controller.signal), threshold: SPAN_THRESHOLD, filters: [{kinds: [EVENT_TIME], "#D": daysBetween(since, until).map(String)}],