diff --git a/src/lib/jetstream.ts b/src/lib/jetstream.ts index edcac5f..33e7d4d 100644 --- a/src/lib/jetstream.ts +++ b/src/lib/jetstream.ts @@ -81,6 +81,7 @@ export class JetstreamConnection { this.reconnectAttempts = 0; this.stopPolling(); this.options.onConnectionChange?.(true); + void this.catchUp(); }; this.ws.onmessage = (event) => { @@ -132,6 +133,29 @@ export class JetstreamConnection { this.reconnectTimer = setTimeout(() => this.connect(), delay); } + /** + * Read the opponent's record once, on every connect. + * + * A socket starts at the live tail, so a move committed between the page's + * read and this moment arrives nowhere. Replaying from a past cursor would + * cover that, but jetstream drains a backlog at close to real time, leaving + * every later frame as far behind as the cursor was old. One read costs a + * round trip and keeps the stream at the tail. + */ + private async catchUp(): Promise { + const { opponentDid } = this.options; + const opponentRkey = this.opponentRkey; + if (!opponentDid || !opponentRkey) return; + + const record = await getGamePublic(opponentDid, opponentRkey); + if (!record || this.destroyed) return; + this.options.onGameUpdate( + record as unknown as Record, + opponentDid, + opponentRkey + ); + } + private startPolling(): void { const { opponentDid } = this.options; const opponentRkey = this.opponentRkey; diff --git a/src/routes/game/[did]/[rkey]/+page.svelte b/src/routes/game/[did]/[rkey]/+page.svelte index 6c21877..f6a4330 100644 --- a/src/routes/game/[did]/[rkey]/+page.svelte +++ b/src/routes/game/[did]/[rkey]/+page.svelte @@ -48,7 +48,6 @@ let dismissShare = $state(false); let rematchOffer: { did: string; rkey: string } | null = $state(null); let rematchDismissed = $state(false); - let loadStartedAt = Date.now(); onMount(() => { return () => jsConnections.forEach(c => c.destroy()); @@ -76,10 +75,6 @@ }); async function loadGame() { - // Jetstream opens after these reads finish; rewind past them so a move - // made in between is replayed rather than lost. - loadStartedAt = Date.now(); - // Read the canonical record -- use authenticated agent if available, public otherwise const record = auth.agent ? await getGame(auth.agent, ownerDid, rkey) @@ -400,7 +395,6 @@ const js = new JetstreamConnection({ opponentDid, opponentRkey, - initialCursor: (loadStartedAt - 5_000) * 1000, onGameUpdate: async (record, did, eventRkey) => { // The opponent's other games land on this socket too. if (!isEventForGame(did, eventRkey, record, scope)) return; @@ -486,7 +480,6 @@ const js = new JetstreamConnection({ opponentDid: playerDid, opponentRkey: playerDid === ownerDid ? rkey : childRkey, - initialCursor: (loadStartedAt - 5_000) * 1000, onGameUpdate: (record, did, eventRkey) => { if (!isEventForGame(did, eventRkey, record, scope)) return; const pgn = record.pgn as string; diff --git a/tests/lib/jetstream.test.ts b/tests/lib/jetstream.test.ts index e350de7..4310bae 100644 --- a/tests/lib/jetstream.test.ts +++ b/tests/lib/jetstream.test.ts @@ -189,6 +189,66 @@ describe('message handling', () => { }); }); +describe('catch-up on connect', () => { + it('reads the opponent record once the socket opens', async () => { + vi.mocked(getGamePublic).mockResolvedValue({ pgn: '1. e4 e5' } as never); + const onGameUpdate = vi.fn(); + make({ opponentDid: OPPONENT, opponentRkey: 'opp-rkey', onGameUpdate }).connect(); + + latest().onopen?.(); + await vi.advanceTimersByTimeAsync(0); + + expect(getGamePublic).toHaveBeenCalledWith(OPPONENT, 'opp-rkey'); + expect(onGameUpdate).toHaveBeenCalledWith({ pgn: '1. e4 e5' }, OPPONENT, 'opp-rkey'); + }); + + it('subscribes at the live tail, not from a past cursor', () => { + make({ opponentDid: OPPONENT, opponentRkey: 'opp-rkey', onGameUpdate: vi.fn() }).connect(); + + expect(new URL(latest().url).searchParams.has('cursor')).toBe(false); + }); + + it('catches up again after a reconnect, since the gap repeats', async () => { + vi.mocked(getGamePublic).mockResolvedValue({ pgn: '1. e4' } as never); + const js = make({ opponentDid: OPPONENT, opponentRkey: 'opp-rkey', onGameUpdate: vi.fn() }); + js.connect(); + latest().onopen?.(); + await vi.advanceTimersByTimeAsync(0); + expect(getGamePublic).toHaveBeenCalledTimes(1); + + latest().close(); + await vi.advanceTimersByTimeAsync(1000); + latest().onopen?.(); + await vi.advanceTimersByTimeAsync(0); + + expect(getGamePublic).toHaveBeenCalledTimes(2); + }); + + it('does not read anything when the opponent record is unknown', async () => { + make({ opponentDid: '', onGameUpdate: vi.fn() }).connect(); + + latest().onopen?.(); + await vi.advanceTimersByTimeAsync(0); + + expect(getGamePublic).not.toHaveBeenCalled(); + }); + + it('drops a catch-up that lands after destroy', async () => { + let release: (v: unknown) => void = () => {}; + vi.mocked(getGamePublic).mockReturnValue(new Promise((r) => (release = r)) as never); + const onGameUpdate = vi.fn(); + const js = make({ opponentDid: OPPONENT, opponentRkey: 'opp-rkey', onGameUpdate }); + js.connect(); + latest().onopen?.(); + + js.destroy(); + release({ pgn: '1. e4' }); + await vi.advanceTimersByTimeAsync(0); + + expect(onGameUpdate).not.toHaveBeenCalled(); + }); +}); + describe('reconnect', () => { it('backs off exponentially up to the cap', () => { const js = make({ opponentDid: OPPONENT, onGameUpdate: vi.fn() }); @@ -305,9 +365,13 @@ describe('polling fallback', () => { await vi.advanceTimersByTimeAsync(3000); expect(getGamePublic).toHaveBeenCalledTimes(1); + // Reopening reads once to catch up; after that the interval is gone. latest().onopen?.(); + await vi.advanceTimersByTimeAsync(0); + const afterReopen = vi.mocked(getGamePublic).mock.calls.length; + await vi.advanceTimersByTimeAsync(9000); - expect(getGamePublic).toHaveBeenCalledTimes(1); + expect(getGamePublic).toHaveBeenCalledTimes(afterReopen); }); });