From 46fc83d211ff26ade01acba731f41fa22a7831d6 Mon Sep 17 00:00:00 2001 From: weishu Date: Wed, 5 Aug 2026 14:55:02 +0800 Subject: [PATCH] feat: cut relay tunnel bandwidth with compression and SSE replay Relay-metered traffic drops on every channel that carried avoidable bytes; binary payloads (attachments, voice audio) are unchanged. Hub: - gzip /api/* JSON responses, gated by q-aware Accept-Encoding negotiation (explicit gzip;q=0 beats a wildcard in either order; hono's compress() alone matches by substring) - enable WebSocket permessage-deflate and default flagless ws.send() to compressed frames - Bun negotiates the extension but compresses nothing unless each send opts in, and @socket.io/bun-engine never passes the flag (measured 96 KB terminal payload -> 608 B on wire) - replay missed SSE events on reconnect: 256-event/2MB ring buffer, per-process epoch ids bound to the authenticated namespace so a token swap can never resume from a foreign cursor, standard Last-Event-ID header preferred over the ?lastEventId fallback, live broadcasts queued until the replay flushes to preserve order Web: - skip the full sessions/details/messages resync when the hub answers resume:ok - a phone unlock now costs a handshake plus the gap delta instead of refetching everything - own every EventSource retry path: take over browser-native CONNECTING retries, defer reconnects while the tab is hidden, and raise the backoff ceiling to 5 min after repeated failures - drop the 30s skills/slash-commands polling; refresh on demand without blocking the suggestion menu behind a stalled CLI RPC - remember the websocket upgrade across terminal socket reconnects --- hub/src/sse/sseManager.test.ts | 215 ++++++++++++++++++++++ hub/src/sse/sseManager.ts | 165 ++++++++++++++++- hub/src/web/routes/events.replay.test.ts | 183 ++++++++++++++++++ hub/src/web/routes/events.ts | 53 ++++-- hub/src/web/server.ts | 34 +++- hub/src/web/sseCompression.test.ts | 33 +++- hub/src/web/sseCompression.ts | 28 ++- hub/src/web/wsCompression.test.ts | 62 +++++++ hub/src/web/wsCompression.ts | 25 +++ shared/src/schemas.ts | 9 +- web/src/App.tsx | 17 +- web/src/hooks/queries/useSkills.ts | 19 +- web/src/hooks/queries/useSlashCommands.ts | 19 +- web/src/hooks/useAgentTerminalSocket.ts | 6 + web/src/hooks/useSSE.ts | 85 ++++++++- web/src/hooks/useTerminalSocket.ts | 3 + 16 files changed, 905 insertions(+), 51 deletions(-) create mode 100644 hub/src/web/routes/events.replay.test.ts create mode 100644 hub/src/web/wsCompression.test.ts create mode 100644 hub/src/web/wsCompression.ts diff --git a/hub/src/sse/sseManager.test.ts b/hub/src/sse/sseManager.test.ts index 209c278f..fb15d2e7 100644 --- a/hub/src/sse/sseManager.test.ts +++ b/hub/src/sse/sseManager.test.ts @@ -119,3 +119,218 @@ describe('SSEManager namespace filtering', () => { expect(received[0]?.id).toBe('visible') }) }) + +describe('SSEManager reconnect replay', () => { + type Sent = { event: SyncEvent; eventId: string | undefined } + + function subscribeCollecting( + manager: SSEManager, + id: string, + options: { namespace?: string; sessionId?: string | null; all?: boolean; resumeFrom?: string | null } = {} + ): { sent: Sent[]; resume: 'ok' | 'gap'; replay: Array<{ event: SyncEvent; eventId: string }> } { + const sent: Sent[] = [] + const result = manager.subscribe({ + id, + namespace: options.namespace ?? 'alpha', + all: options.all ?? !options.sessionId, + sessionId: options.sessionId ?? null, + resumeFrom: options.resumeFrom ?? null, + send: (event, eventId) => { + sent.push({ event, eventId }) + }, + sendHeartbeat: () => {} + }) + return { sent, resume: result.resume, replay: result.replay } + } + + it('assigns monotonic ids to broadcast events', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + + expect(sent).toHaveLength(2) + const [first, second] = sent + expect(first?.eventId).toMatch(/^[0-9a-f-]{8}:1:[0-9a-f]{8}$/) + expect(second?.eventId).toMatch(/^[0-9a-f-]{8}:2:[0-9a-f]{8}$/) + expect(first?.eventId?.split(':')[0]).toBe(second?.eventId?.split(':')[0]) + expect(first?.eventId?.split(':')[2]).toBe(second?.eventId?.split(':')[2]) + }) + + it('resumes with a filtered replay of missed events', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + manager.broadcast({ type: 'session-updated', sessionId: 's2', namespace: 'alpha' }) + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'beta' }) + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + const cursor = sent[0]?.eventId + manager.unsubscribe('a') + + const reconnect = subscribeCollecting(manager, 'b', { sessionId: 's1', resumeFrom: cursor }) + + expect(reconnect.resume).toBe('ok') + // s2 filtered out (session mismatch), beta filtered out (namespace) + expect(reconnect.replay).toHaveLength(1) + expect(reconnect.replay[0]?.event).toMatchObject({ sessionId: 's1', namespace: 'alpha' }) + expect(reconnect.replay[0]?.eventId?.split(':')[1]).toBe('4') + }) + + it('resumes ok with empty replay when nothing was missed', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + manager.unsubscribe('a') + + const reconnect = subscribeCollecting(manager, 'b', { resumeFrom: sent[0]?.eventId }) + + expect(reconnect.resume).toBe('ok') + expect(reconnect.replay).toHaveLength(0) + + // no pending queue: live events flow immediately + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + expect(reconnect.sent).toHaveLength(1) + }) + + it('reports a gap for foreign or malformed cursors', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + + for (const cursor of ['deadbeef:1', 'not-a-cursor', ':', 'deadbeef:', '']) { + const { resume, replay } = subscribeCollecting(manager, `c-${cursor}`, { resumeFrom: cursor || null }) + expect(resume).toBe('gap') + expect(replay).toHaveLength(0) + } + }) + + it('reports a gap for cursors from the future', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + const [epoch, , tag] = sent[0]?.eventId?.split(':') ?? [] + manager.unsubscribe('a') + + const reconnect = subscribeCollecting(manager, 'b', { resumeFrom: `${epoch}:999:${tag}` }) + expect(reconnect.resume).toBe('gap') + }) + + it('rejects a cursor issued under a different namespace', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a', { namespace: 'alpha' }) + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + const cursor = sent[0]?.eventId + manager.unsubscribe('a') + + // Same subscription shape, different authenticated namespace (token + // swap on the same hub): the alpha cursor must not vouch for beta. + const foreign = subscribeCollecting(manager, 'b', { namespace: 'beta', resumeFrom: cursor }) + expect(foreign.resume).toBe('gap') + expect(foreign.replay).toHaveLength(0) + + // The same cursor is still valid for its own namespace. + const home = subscribeCollecting(manager, 'c', { namespace: 'alpha', resumeFrom: cursor }) + expect(home.resume).toBe('ok') + expect(home.replay).toHaveLength(1) + }) + + it('reports a gap when the cursor has been evicted from the ring', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + // Overflow the 256-entry ring so the first event is evicted + for (let i = 0; i < 300; i++) { + manager.broadcast({ type: 'session-updated', sessionId: `s${i}`, namespace: 'alpha' }) + } + manager.unsubscribe('a') + + const stale = subscribeCollecting(manager, 'b', { resumeFrom: sent[0]?.eventId }) + expect(stale.resume).toBe('gap') + + const fresh = subscribeCollecting(manager, 'c', { resumeFrom: sent[298]?.eventId }) + expect(fresh.resume).toBe('ok') + expect(fresh.replay).toHaveLength(1) + expect(fresh.replay[0]?.event).toMatchObject({ sessionId: 's299' }) + }) + + it('queues live broadcasts during replay and drains them in order', async () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + manager.broadcast({ type: 'session-updated', sessionId: 's2', namespace: 'alpha' }) + manager.unsubscribe('a') + + const reconnect = subscribeCollecting(manager, 'b', { resumeFrom: sent[0]?.eventId }) + expect(reconnect.replay).toHaveLength(1) + + // Fires while the caller is still writing the replay: must be queued, + // not delivered. + manager.broadcast({ type: 'session-updated', sessionId: 's3', namespace: 'alpha' }) + expect(reconnect.sent).toHaveLength(0) + + await manager.drainPending('b') + expect(reconnect.sent).toHaveLength(1) + expect(reconnect.sent[0]?.event).toMatchObject({ sessionId: 's3' }) + expect(reconnect.sent[0]?.eventId?.split(':')[1]).toBe('3') + + // After the drain the connection is live + manager.broadcast({ type: 'session-updated', sessionId: 's4', namespace: 'alpha' }) + expect(reconnect.sent).toHaveLength(2) + }) + + it('drains events that arrive while the drain itself is awaiting', async () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'alpha' }) + manager.broadcast({ type: 'session-updated', sessionId: 's2', namespace: 'alpha' }) + manager.unsubscribe('a') + + const order: string[] = [] + let injected = false + manager.subscribe({ + id: 'b', + namespace: 'alpha', + all: true, + resumeFrom: sent[0]?.eventId, + send: async (event) => { + const sessionId = 'sessionId' in event ? event.sessionId : '?' + order.push(String(sessionId)) + if (!injected) { + injected = true + // Simulates a broadcast racing in mid-drain + manager.broadcast({ type: 'session-updated', sessionId: 's-mid', namespace: 'alpha' }) + await Promise.resolve() + } + }, + sendHeartbeat: () => {} + }) + manager.broadcast({ type: 'session-updated', sessionId: 's3', namespace: 'alpha' }) + + await manager.drainPending('b') + + expect(order).toEqual(['s3', 's-mid']) + }) + + it('evicts by byte budget while keeping at least one event', () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const { sent } = subscribeCollecting(manager, 'a') + const bigMetadata = 'x'.repeat(1_500_000) + manager.broadcast({ type: 'session-updated', sessionId: 'big-1', namespace: 'alpha', data: { summary: bigMetadata } as never }) + manager.broadcast({ type: 'session-updated', sessionId: 'big-2', namespace: 'alpha', data: { summary: bigMetadata } as never }) + manager.broadcast({ type: 'session-updated', sessionId: 'big-3', namespace: 'alpha', data: { summary: bigMetadata } as never }) + manager.unsubscribe('a') + + // Only big-3 fits the 2MB budget, so big-2 was evicted UNSEEN by a + // cursor pointing at big-1 - that cursor cannot resume... + const stale = subscribeCollecting(manager, 'b', { resumeFrom: sent[0]?.eventId }) + expect(stale.resume).toBe('gap') + + // ...while a cursor at big-2 only needs big-3, which survives. An + // evicted event the client has already seen never blocks resume. + const fresh = subscribeCollecting(manager, 'c', { resumeFrom: sent[1]?.eventId }) + expect(fresh.resume).toBe('ok') + expect(fresh.replay).toHaveLength(1) + expect(fresh.replay[0]?.event).toMatchObject({ sessionId: 'big-3' }) + }) +}) diff --git a/hub/src/sse/sseManager.ts b/hub/src/sse/sseManager.ts index e4fec1ad..cabd4f7a 100644 --- a/hub/src/sse/sseManager.ts +++ b/hub/src/sse/sseManager.ts @@ -1,3 +1,4 @@ +import { createHash, randomUUID } from 'node:crypto' import type { SyncEvent } from '../sync/syncEngine' import type { VisibilityState } from '../visibility/visibilityTracker' import type { VisibilityTracker } from '../visibility/visibilityTracker' @@ -11,15 +12,50 @@ export type SSESubscription = { } type SSEConnection = SSESubscription & { - send: (event: SyncEvent) => void | Promise + send: (event: SyncEvent, eventId?: string) => void | Promise sendHeartbeat: () => void | Promise + /** + * Non-null while a resumed connection is still writing its replay: live + * broadcasts land here instead of on the wire so replayed events keep + * their original order. `drainPending` flushes and clears it. + */ + pending: Array<{ event: SyncEvent; eventId: string }> | null } +export type SSEResumeResult = { + /** 'ok': `replay` holds every missed event; client may skip its resync. */ + resume: 'ok' | 'gap' + replay: Array<{ event: SyncEvent; eventId: string }> +} + +/** How many broadcast events are kept for reconnect replay. */ +const EVENT_BUFFER_CAPACITY = 256 +/** + * Byte budget for the replay buffer (sum of JSON-encoded events). Bounds + * memory when large message payloads flow; evicting early just means older + * cursors resync via REST like they always did. + */ +const EVENT_BUFFER_MAX_BYTES = 2 * 1024 * 1024 + export class SSEManager { private readonly connections: Map = new Map() private heartbeatTimer: NodeJS.Timeout | null = null private readonly heartbeatMs: number private readonly visibilityTracker: VisibilityTracker + /** + * Replay ring buffer. Event ids are `${epoch}:${seq}:${nsTag}`: the epoch + * is per-process so a cursor from before a hub restart can never match a + * fresh buffer, seq grows monotonically within the process, and nsTag + * binds the cursor to the namespace it was issued under - a client whose + * token swap changed its namespace must NOT resume from the old cursor + * (its verdict would skip the resync the new namespace needs), even + * though replay filtering alone would never leak foreign events. + */ + private readonly epoch = randomUUID().slice(0, 8) + private nextSeq = 1 + private readonly eventBuffer: Array<{ seq: number; event: SyncEvent; bytes: number }> = [] + private eventBufferBytes = 0 + private readonly namespaceTags = new Map() constructor(heartbeatMs = 30_000, visibilityTracker: VisibilityTracker) { this.heartbeatMs = heartbeatMs @@ -33,9 +69,11 @@ export class SSEManager { sessionId?: string | null machineId?: string | null visibility?: VisibilityState - send: (event: SyncEvent) => void | Promise + /** Last event id the client saw; enables replay instead of resync. */ + resumeFrom?: string | null + send: (event: SyncEvent, eventId?: string) => void | Promise sendHeartbeat: () => void | Promise - }): SSESubscription { + }): SSESubscription & SSEResumeResult { const subscription: SSEConnection = { id: options.id, namespace: options.namespace, @@ -43,7 +81,15 @@ export class SSEManager { sessionId: options.sessionId ?? null, machineId: options.machineId ?? null, send: options.send, - sendHeartbeat: options.sendHeartbeat + sendHeartbeat: options.sendHeartbeat, + pending: null + } + + const { resume, replay } = this.resolveResume(subscription, options.resumeFrom ?? null) + if (replay.length > 0) { + // Live broadcasts must not overtake the replay the caller is about + // to write; queue them until drainPending. + subscription.pending = [] } this.connections.set(subscription.id, subscription) @@ -58,10 +104,104 @@ export class SSEManager { namespace: subscription.namespace, all: subscription.all, sessionId: subscription.sessionId, - machineId: subscription.machineId + machineId: subscription.machineId, + resume, + replay } } + /** + * Flush events queued while the caller was writing a replay, then switch + * the connection to direct delivery. Must be awaited after the replay has + * been written; a no-op for connections that never queued. + */ + async drainPending(id: string): Promise { + const connection = this.connections.get(id) + if (!connection) { + return + } + while (connection.pending) { + const batch = connection.pending.splice(0) + if (batch.length === 0) { + // No await between this check and the assignment, so no + // broadcast can slip into the queue and be dropped. + connection.pending = null + break + } + for (const item of batch) { + try { + await connection.send(item.event, item.eventId) + } catch { + this.unsubscribe(connection.id) + return + } + } + } + } + + private namespaceTag(namespace: string): string { + const cached = this.namespaceTags.get(namespace) + if (cached) { + return cached + } + const tag = createHash('sha256').update(`${this.epoch}|${namespace}`).digest('hex').slice(0, 8) + this.namespaceTags.set(namespace, tag) + return tag + } + + private eventIdFor(seq: number, namespace: string): string { + return `${this.epoch}:${seq}:${this.namespaceTag(namespace)}` + } + + private resolveResume(connection: SSEConnection, resumeFrom: string | null): SSEResumeResult { + if (!resumeFrom) { + return { resume: 'gap', replay: [] } + } + const parts = resumeFrom.split(':') + if (parts.length !== 3 || parts[0] !== this.epoch) { + return { resume: 'gap', replay: [] } + } + if (parts[2] !== this.namespaceTag(connection.namespace)) { + // Cursor was issued under a different namespace (token swap on the + // same hub): its position says nothing about what THIS namespace + // has missed, so force the full resync. + return { resume: 'gap', replay: [] } + } + const seq = Number(parts[1]) + if (!Number.isSafeInteger(seq) || seq < 1 || seq >= this.nextSeq) { + return { resume: 'gap', replay: [] } + } + const oldestBuffered = this.eventBuffer[0]?.seq ?? this.nextSeq + if (seq < oldestBuffered - 1) { + // Events between the cursor and the buffer were evicted. + return { resume: 'gap', replay: [] } + } + const replay: Array<{ event: SyncEvent; eventId: string }> = [] + for (const entry of this.eventBuffer) { + if (entry.seq > seq && this.shouldSend(connection, entry.event)) { + replay.push({ event: entry.event, eventId: this.eventIdFor(entry.seq, connection.namespace) }) + } + } + return { resume: 'ok', replay } + } + + private recordEvent(event: SyncEvent): number { + const seq = this.nextSeq++ + const bytes = JSON.stringify(event).length + this.eventBuffer.push({ seq, event, bytes }) + this.eventBufferBytes += bytes + while ( + this.eventBuffer.length > EVENT_BUFFER_CAPACITY + || (this.eventBufferBytes > EVENT_BUFFER_MAX_BYTES && this.eventBuffer.length > 1) + ) { + const evicted = this.eventBuffer.shift() + if (evicted) { + this.eventBufferBytes -= evicted.bytes + } + } + return seq + } + unsubscribe(id: string): void { this.connections.delete(id) this.visibilityTracker.removeConnection(id) @@ -70,6 +210,10 @@ export class SSEManager { } } + hasSubscription(id: string): boolean { + return this.connections.has(id) + } + async sendToast(namespace: string, event: Extract): Promise { const deliveries: Array> = [] for (const connection of this.connections.values()) { @@ -105,12 +249,21 @@ export class SSEManager { } broadcast(event: SyncEvent): void { + const seq = this.recordEvent(event) for (const connection of this.connections.values()) { if (!this.shouldSend(connection, event)) { continue } - void Promise.resolve(connection.send(event)).catch(() => { + // The id is per-connection: same epoch and seq, but tagged with + // the receiving namespace so the cursor stays bound to it. + const eventId = this.eventIdFor(seq, connection.namespace) + if (connection.pending) { + connection.pending.push({ event, eventId }) + continue + } + + void Promise.resolve(connection.send(event, eventId)).catch(() => { this.unsubscribe(connection.id) }) } diff --git a/hub/src/web/routes/events.replay.test.ts b/hub/src/web/routes/events.replay.test.ts new file mode 100644 index 00000000..dc5c82c5 --- /dev/null +++ b/hub/src/web/routes/events.replay.test.ts @@ -0,0 +1,183 @@ +import { describe, expect, it } from 'bun:test' +import { Hono } from 'hono' +import { SSEManager } from '../../sse/sseManager' +import { VisibilityTracker } from '../../visibility/visibilityTracker' +import type { WebAppEnv } from '../middleware/auth' +import { createEventsRoutes } from './events' + +type Frame = { id: string | null; data: Record } + +/** + * Reads SSE frames from a live stream until `count` frames arrived, then + * disconnects. In-memory `app.request` streams need a specific teardown to + * reach hono's `stream.onAbort` (the route's release path): an + * AbortController signal does not propagate to `c.req.raw.signal`, and a + * bare `reader.cancel()` with no read in flight does not propagate either - + * cancellation only reaches the stream source while a read is parked. So: + * park a read, then cancel. + */ +async function collectFrames(response: Response, count: number): Promise { + const reader = response.body!.getReader() + const decoder = new TextDecoder() + const frames: Frame[] = [] + let buffer = '' + try { + while (frames.length < count) { + const { done, value } = await reader.read() + if (done) { + break + } + buffer += decoder.decode(value, { stream: true }) + let boundary = buffer.indexOf('\n\n') + while (boundary >= 0) { + const raw = buffer.slice(0, boundary) + buffer = buffer.slice(boundary + 2) + boundary = buffer.indexOf('\n\n') + let id: string | null = null + let data = '' + for (const line of raw.split('\n')) { + if (line.startsWith('id:')) { + id = line.slice(3).trim() + } else if (line.startsWith('data:')) { + data += line.slice(5).trim() + } + } + if (data) { + frames.push({ id, data: JSON.parse(data) as Record }) + } + } + } + } finally { + // The parked read engages the stream's pull on the next macrotask; + // cancelling before that happens is silently ignored. + const parked = reader.read().catch(() => null) + await new Promise((resolve) => setTimeout(resolve, 0)) + await reader.cancel().catch(() => {}) + await parked + } + return frames +} + +function buildApp(manager: SSEManager): Hono { + const app = new Hono() + app.use('*', async (c, next) => { + c.set('namespace', 'ns-test') + c.set('userId', 1) + await next() + }) + app.route('/api', createEventsRoutes(() => manager, () => null, () => null)) + return app +} + +async function openStream(app: Hono, query: string): Promise { + return await app.request(`/api/events?all=true${query}`) +} + +describe('GET /api/events replay', () => { + it('first connect gets a gap verdict and id-tagged live events', async () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const app = buildApp(manager) + + const res = await openStream(app, '') + + // Give the route a beat to subscribe before broadcasting + await new Promise((r) => setTimeout(r, 20)) + manager.broadcast({ type: 'session-updated', sessionId: 's1', namespace: 'ns-test' }) + manager.broadcast({ type: 'session-updated', sessionId: 's2', namespace: 'ns-test' }) + + const frames = await collectFrames(res, 3) + + expect(frames[0]?.data.type).toBe('connection-changed') + expect((frames[0]?.data.data as { resume?: string }).resume).toBe('gap') + expect(frames[0]?.id).toBeNull() + + expect(frames[1]?.data.sessionId).toBe('s1') + expect(frames[1]?.id).toMatch(/^[0-9a-f-]{8}:\d+:[0-9a-f]{8}$/) + expect(frames[2]?.data.sessionId).toBe('s2') + }) + + it('reconnect with lastEventId replays the missed events before live traffic', async () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const app = buildApp(manager) + + // First connection observes event 1, then drops. + const firstRes = await openStream(app, '') + await new Promise((r) => setTimeout(r, 20)) + manager.broadcast({ type: 'session-updated', sessionId: 'seen', namespace: 'ns-test' }) + const firstFrames = await collectFrames(firstRes, 2) + const cursor = firstFrames[1]?.id + expect(cursor).toBeTruthy() + + // Missed while disconnected. + manager.broadcast({ type: 'session-updated', sessionId: 'missed-1', namespace: 'ns-test' }) + manager.broadcast({ type: 'session-updated', sessionId: 'missed-2', namespace: 'ns-test' }) + + // Reconnect with the cursor. + const secondRes = await openStream(app, `&lastEventId=${encodeURIComponent(cursor!)}`) + await new Promise((r) => setTimeout(r, 20)) + manager.broadcast({ type: 'session-updated', sessionId: 'live-after', namespace: 'ns-test' }) + + const frames = await collectFrames(secondRes, 4) + + expect(frames[0]?.data.type).toBe('connection-changed') + expect((frames[0]?.data.data as { resume?: string }).resume).toBe('ok') + expect(frames.slice(1).map((f) => f.data.sessionId)).toEqual(['missed-1', 'missed-2', 'live-after']) + // Replayed frames keep their original ids so the cursor keeps advancing + expect(frames[1]?.id?.split(':')[1]).toBe('2') + expect(frames[2]?.id?.split(':')[1]).toBe('3') + }) + + it('reconnect with a stale cursor from another process gets a gap', async () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const app = buildApp(manager) + + const res = await openStream(app, '&lastEventId=00000000%3A42%3A00000000') + const frames = await collectFrames(res, 1) + + expect(frames[0]?.data.type).toBe('connection-changed') + expect((frames[0]?.data.data as { resume?: string }).resume).toBe('gap') + }) + + it('prefers the standard Last-Event-ID header over the query parameter', async () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const app = buildApp(manager) + + const firstRes = await openStream(app, '') + await new Promise((r) => setTimeout(r, 20)) + manager.broadcast({ type: 'session-updated', sessionId: 'e1', namespace: 'ns-test' }) + manager.broadcast({ type: 'session-updated', sessionId: 'e2', namespace: 'ns-test' }) + const firstFrames = await collectFrames(firstRes, 3) + const oldCursor = firstFrames[1]?.id + const freshCursor = firstFrames[2]?.id + expect(oldCursor && freshCursor).toBeTruthy() + + manager.broadcast({ type: 'session-updated', sessionId: 'e3', namespace: 'ns-test' }) + + // Native auto-reconnects keep the stale URL cursor but send the fresh + // one in the header - the header must win, replaying only e3. + const res = await app.request( + `/api/events?all=true&lastEventId=${encodeURIComponent(oldCursor!)}`, + { headers: { 'Last-Event-ID': freshCursor! } } + ) + const frames = await collectFrames(res, 2) + + expect((frames[0]?.data.data as { resume?: string }).resume).toBe('ok') + expect(frames[1]?.data.sessionId).toBe('e3') + }) + + it('releases the subscription when the client disconnects', async () => { + const manager = new SSEManager(0, new VisibilityTracker()) + const app = buildApp(manager) + + const res = await openStream(app, '') + const frames = await collectFrames(res, 1) + expect(frames[0]?.data.type).toBe('connection-changed') + + // After abort the connection must be unsubscribed: broadcasting to a + // dead stream would otherwise throw / leak. + await new Promise((r) => setTimeout(r, 50)) + const subscriptionId = (frames[0]?.data.data as { subscriptionId?: string }).subscriptionId + expect(subscriptionId).toBeTruthy() + expect(manager.hasSubscription(subscriptionId!)).toBe(false) + }) +}) diff --git a/hub/src/web/routes/events.ts b/hub/src/web/routes/events.ts index 202b9d56..0ccf28a6 100644 --- a/hub/src/web/routes/events.ts +++ b/hub/src/web/routes/events.ts @@ -52,6 +52,12 @@ export function createEventsRoutes( const machineId = parseOptionalId(query.machineId) const subscriptionId = randomUUID() const visibility = parseVisibility(query.visibility) + // Native EventSource auto-reconnects keep the original URL and carry + // the freshest cursor in the standard Last-Event-ID header; the query + // parameter covers manually created reconnects (new EventSource with + // a rebuilt URL). Header wins - it is never staler than the URL. + const resumeFrom = parseOptionalId(c.req.header('Last-Event-ID')) + ?? parseOptionalId(query.lastEventId) const namespace = c.get('namespace') let resolvedSessionId = sessionId @@ -79,14 +85,18 @@ export function createEventsRoutes( } const response = streamSSE(c, async (stream) => { - manager.subscribe({ + // Reconnects supply the last SSE id they saw; when the manager + // still has everything after it, the missed events are replayed + // here and the client skips its REST resync (`resume: 'ok'`). + const { resume, replay } = manager.subscribe({ id: subscriptionId, namespace, all, sessionId: resolvedSessionId, machineId, visibility, - send: (event) => stream.writeSSE({ data: JSON.stringify(event) }), + resumeFrom, + send: (event, eventId) => stream.writeSSE({ data: JSON.stringify(event), id: eventId }), sendHeartbeat: async () => { await stream.writeSSE({ data: JSON.stringify({ @@ -100,23 +110,34 @@ export function createEventsRoutes( } }) - await stream.writeSSE({ - data: JSON.stringify({ - type: 'connection-changed', - data: { - status: 'connected', - subscriptionId - } + try { + // Verdict first, replay second, live traffic third: the client + // decides whether to resync from the verdict, so it must never + // see replayed or live events before it. + await stream.writeSSE({ + data: JSON.stringify({ + type: 'connection-changed', + data: { + status: 'connected', + subscriptionId, + resume + } + }) }) - }) - await new Promise((resolve) => { - const done = () => resolve() - c.req.raw.signal.addEventListener('abort', done, { once: true }) - stream.onAbort(done) - }) + for (const item of replay) { + await stream.writeSSE({ data: JSON.stringify(item.event), id: item.eventId }) + } + await manager.drainPending(subscriptionId) - manager.unsubscribe(subscriptionId) + await new Promise((resolve) => { + const done = () => resolve() + c.req.raw.signal.addEventListener('abort', done, { once: true }) + stream.onAbort(done) + }) + } finally { + manager.unsubscribe(subscriptionId) + } }) return compressSseResponse(response, c.req.header('Accept-Encoding')) diff --git a/hub/src/web/server.ts b/hub/src/web/server.ts index 6ae596b5..79a629c4 100644 --- a/hub/src/web/server.ts +++ b/hub/src/web/server.ts @@ -1,4 +1,5 @@ import { Hono } from 'hono' +import { compress } from 'hono/compress' import { cors } from 'hono/cors' import { logger } from 'hono/logger' import { join } from 'node:path' @@ -29,6 +30,8 @@ import { createVoiceRoutes } from './routes/voice' import type { SSEManager } from '../sse/sseManager' import type { VisibilityTracker } from '../visibility/visibilityTracker' import type { Server as BunServer, ServerWebSocket } from 'bun' +import { applyDefaultWsCompression } from './wsCompression' +import { acceptsGzip } from './sseCompression' import type { Server as SocketEngine } from '@socket.io/bun-engine' import { jwtVerify } from 'jose' import type { WebSocketData } from '@socket.io/bun-engine' @@ -234,11 +237,34 @@ function createWebApp(options: { const corsMiddleware = cors({ origin: corsOriginOption, allowMethods: ['GET', 'POST', 'PATCH', 'DELETE', 'OPTIONS'], - allowHeaders: ['authorization', 'content-type'] + // last-event-id: browsers attach it to EventSource reconnects for + // SSE replay; allow it in case a browser preflights the request. + allowHeaders: ['authorization', 'content-type', 'last-event-id'] }) app.use('/api/*', corsMiddleware) app.use('/cli/*', corsMiddleware) + // Gzip JSON API responses. Over the relay tunnel every byte is metered + // twice (the SNI proxy copies in both directions), and API payloads are + // repetitive JSON that compresses to roughly a quarter of its size. + // + // This deliberately does not touch /api/events: streamSSE sets + // Transfer-Encoding, which hono's compress() skips, and that stream is + // already gzipped by compressSseResponse with an explicit sync flush. + // Binary uploads/downloads are skipped too - compress() only handles + // content types it knows are compressible. + // + // Gated on the q-aware parser because hono's compress() matches the + // Accept-Encoding value by substring: `gzip;q=0` - an explicit refusal - + // would otherwise still get a gzip body it cannot consume. + const gzipCompress = compress({ encoding: 'gzip' }) + app.use('/api/*', async (c, next) => { + if (acceptsGzip(c.req.header('Accept-Encoding'))) { + return gzipCompress(c, next) + } + return next() + }) + app.route('/cli', createCliRoutes(options.getSyncEngine)) app.route('/api', createAuthRoutes(options.jwtSecret, options.store)) @@ -408,7 +434,13 @@ export async function startWebServer(options: { maxRequestBodySize: Math.max(socketHandler.maxRequestBodySize, 68 * 1024 * 1024), websocket: { ...originalWsHandler, + // Advertise permessage-deflate. Negotiation alone compresses + // nothing in Bun — each send() opts in — so open() below also + // makes compression the default for flagless sends. See + // wsCompression.ts for the contract. + perMessageDeflate: true, open(ws: unknown) { + applyDefaultWsCompression(ws as ServerWebSocket) const wsAny = ws as ServerWebSocket<{ _qwenProxy?: boolean; _geminiProxy?: boolean }> if (wsAny.data?._geminiProxy) { geminiProxyHandler.open(wsAny) diff --git a/hub/src/web/sseCompression.test.ts b/hub/src/web/sseCompression.test.ts index c3308af6..6380da23 100644 --- a/hub/src/web/sseCompression.test.ts +++ b/hub/src/web/sseCompression.test.ts @@ -1,6 +1,6 @@ import { describe, expect, it } from 'bun:test' import zlib from 'node:zlib' -import { compressSseResponse } from './sseCompression' +import { acceptsGzip, compressSseResponse } from './sseCompression' function gunzip(data: Uint8Array): string { return zlib.gunzipSync(Buffer.from(data)).toString('utf8') @@ -209,3 +209,34 @@ describe('compressSseResponse backpressure', () => { await reader.cancel('done') }) }) + +describe('acceptsGzip', () => { + it('handles plain and missing headers', () => { + expect(acceptsGzip(undefined)).toBe(false) + expect(acceptsGzip('')).toBe(false) + expect(acceptsGzip('gzip')).toBe(true) + expect(acceptsGzip('gzip, deflate, br')).toBe(true) + expect(acceptsGzip('br, deflate')).toBe(false) + expect(acceptsGzip('*')).toBe(true) + expect(acceptsGzip('identity;q=1, *;q=0')).toBe(false) + }) + + it('honors q-values on explicit gzip entries', () => { + expect(acceptsGzip('gzip;q=0')).toBe(false) + expect(acceptsGzip('gzip;q=0.0')).toBe(false) + expect(acceptsGzip('gzip;q=0.5')).toBe(true) + expect(acceptsGzip('br, gzip;q=0')).toBe(false) + }) + + it('lets an explicit gzip entry override the wildcard in either order', () => { + expect(acceptsGzip('*;q=1, gzip;q=0')).toBe(false) + expect(acceptsGzip('gzip;q=0, *;q=1')).toBe(false) + expect(acceptsGzip('*;q=0, gzip;q=1')).toBe(true) + expect(acceptsGzip('gzip;q=1, *;q=0')).toBe(true) + }) + + it('treats an unparseable q as the default weight', () => { + expect(acceptsGzip('gzip;q=abc')).toBe(true) + expect(acceptsGzip('*;q=abc, gzip;q=0')).toBe(false) + }) +}) diff --git a/hub/src/web/sseCompression.ts b/hub/src/web/sseCompression.ts index bcfa9547..d0462ed0 100644 --- a/hub/src/web/sseCompression.ts +++ b/hub/src/web/sseCompression.ts @@ -4,27 +4,41 @@ import zlib from 'node:zlib' * True when the client is willing to receive gzip. * * `Accept-Encoding: gzip;q=0` means the opposite of what a substring match - * would suggest, so parse the q-value rather than looking for the word. + * would suggest, so parse the q-value rather than looking for the word. All + * entries are read before deciding because an explicit `gzip` entry takes + * precedence over `*` regardless of where it appears in the header (RFC 9110 + * §12.5.3): `*;q=1, gzip;q=0` refuses gzip, `*;q=0, gzip;q=1` accepts it. + * Exported because the API compression middleware needs the same q-aware + * negotiation (hono's compress() matches by substring and would gzip for + * clients that explicitly refuse it). */ -function acceptsGzip(acceptEncoding: string | undefined): boolean { +export function acceptsGzip(acceptEncoding: string | undefined): boolean { if (!acceptEncoding) { return false } + let gzipQ: number | null = null + let wildcardQ: number | null = null for (const part of acceptEncoding.split(',')) { const [rawName, ...params] = part.split(';') const name = rawName?.trim().toLowerCase() if (name !== 'gzip' && name !== '*') { continue } - const q = params + const qParam = params .map((param) => param.trim().toLowerCase()) .find((param) => param.startsWith('q=')) - if (q && Number(q.slice(2)) === 0) { - return false + // Absent or unparseable q counts as 1 (the header's default weight); + // repeated entries let the last one win. + const parsed = qParam ? Number(qParam.slice(2)) : 1 + const q = Number.isFinite(parsed) ? parsed : 1 + if (name === 'gzip') { + gzipQ = q + } else { + wildcardQ = q } - return true } - return false + const effective = gzipQ ?? wildcardQ + return effective !== null && effective > 0 } /** diff --git a/hub/src/web/wsCompression.test.ts b/hub/src/web/wsCompression.test.ts new file mode 100644 index 00000000..c778d806 --- /dev/null +++ b/hub/src/web/wsCompression.test.ts @@ -0,0 +1,62 @@ +import { describe, expect, test } from 'bun:test' +import { applyDefaultWsCompression } from './wsCompression' + +function makeFakeWs(): { ws: { send: (data: string | Bun.BufferSource, compress?: boolean) => number }; calls: Array<{ data: string | Bun.BufferSource; compress: boolean | undefined; self: unknown }> } { + const calls: Array<{ data: string | Bun.BufferSource; compress: boolean | undefined; self: unknown }> = [] + const ws = { + send(this: unknown, data: string | Bun.BufferSource, compress?: boolean): number { + calls.push({ data, compress, self: this }) + return typeof data === 'string' ? data.length : data.byteLength + } + } + return { ws, calls } +} + +describe('applyDefaultWsCompression', () => { + test('flagless send defaults to compress=true (bun-engine call shape)', () => { + const { ws, calls } = makeFakeWs() + applyDefaultWsCompression(ws) + + ws.send('payload') + + expect(calls).toHaveLength(1) + expect(calls[0]?.compress).toBe(true) + expect(calls[0]?.data).toBe('payload') + }) + + test('explicit flags are preserved in both directions', () => { + const { ws, calls } = makeFakeWs() + applyDefaultWsCompression(ws) + + ws.send('a', false) + ws.send('b', true) + + expect(calls[0]?.compress).toBe(false) + expect(calls[1]?.compress).toBe(true) + }) + + test('original send stays bound to the socket and returns its result', () => { + const { ws, calls } = makeFakeWs() + const original = ws.send + applyDefaultWsCompression(ws) + + const result = ws.send('12345') + + expect(result).toBe(5) + expect(ws.send).not.toBe(original) + // bind() target must be the ws object, not the wrapper's caller + expect(calls[0]?.self).toBe(ws) + }) + + test('binary payloads pass through untouched', () => { + const { ws, calls } = makeFakeWs() + applyDefaultWsCompression(ws) + + const buffer = new Uint8Array([1, 2, 3]) + const result = ws.send(buffer) + + expect(result).toBe(3) + expect(calls[0]?.data).toBe(buffer) + expect(calls[0]?.compress).toBe(true) + }) +}) diff --git a/hub/src/web/wsCompression.ts b/hub/src/web/wsCompression.ts new file mode 100644 index 00000000..aee1988e --- /dev/null +++ b/hub/src/web/wsCompression.ts @@ -0,0 +1,25 @@ +/** + * WebSocket per-message-deflate defaults. + * + * Bun negotiates the permessage-deflate extension when `perMessageDeflate: + * true` is set on the serve options, but actually compressing a frame is + * still opt-in per send() call — and @socket.io/bun-engine (plus the voice + * proxies) call `ws.send(data)` without the flag, so nothing would ever be + * compressed. Wrapping send() turns the negotiated extension into actual + * wire compression: terminal streams and CLI sync JSON shrink by 70-99% + * through the relay tunnel. + * + * Safe for every client: when the peer did not negotiate the extension Bun + * ignores the compress flag and sends plaintext (verified: RSV1 stays 0 and + * the payload arrives intact). Callers that pass an explicit flag keep it. + */ + +type CompressibleSend = ( + data: string | Bun.BufferSource, + compress?: boolean +) => number + +export function applyDefaultWsCompression(ws: { send: CompressibleSend }): void { + const original = ws.send.bind(ws) + ws.send = (data, compress) => original(data, compress ?? true) +} diff --git a/shared/src/schemas.ts b/shared/src/schemas.ts index f044c24c..4f1596f7 100644 --- a/shared/src/schemas.ts +++ b/shared/src/schemas.ts @@ -555,7 +555,14 @@ export const SyncEventSchema = z.discriminatedUnion('type', [ type: z.literal('connection-changed'), data: z.object({ status: z.string(), - subscriptionId: z.string().optional() + subscriptionId: z.string().optional(), + /** + * Reconnect verdict. 'ok' means the hub replayed every event the + * client missed (sent right after this one), so the client can skip + * its full refetch. 'gap' (or absence, on older hubs) means the + * client must resync from REST. + */ + resume: z.enum(['ok', 'gap']).optional() }).optional() }) ]) diff --git a/web/src/App.tsx b/web/src/App.tsx index 2bccf74b..3f9ea2c3 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -207,10 +207,18 @@ function AppInner() { void run() }, [api, isPushSupported, pushPermission, requestPermission, subscribe, token]) - const handleSseConnect = useCallback(() => { + const handleSseConnect = useCallback((info: { resumed: boolean }) => { // Clear disconnected state on successful connection reportSseConnect() + // The hub replayed every event missed during the gap, so the caches + // are already consistent - the full refetch below would only re-download + // what the replay just delivered. First connects and long gaps arrive + // with resumed=false and take the resync path. + if (info.resumed && !isFirstConnectRef.current) { + return + } + // Increment token to track this specific connection const token = ++syncTokenRef.current @@ -265,10 +273,15 @@ function AppInner() { void syncTailMessages(api, event.sessionId) }, [api, selectedSessionId]) - const handleSessionSseConnect = useCallback(() => { + const handleSessionSseConnect = useCallback((info: { resumed: boolean }) => { if (!api || !selectedSessionId) { return } + // A resumed connection replayed messages-consumed/message events for + // this session, so the queued-state snapshot cannot have drifted. + if (info.resumed) { + return + } void reconcileQueuedStateAfterConnect(api, selectedSessionId).catch((error) => { console.error('Failed to reconcile queued state after SSE connect:', error) }) diff --git a/web/src/hooks/queries/useSkills.ts b/web/src/hooks/queries/useSkills.ts index f25d5fb1..29d3ef1f 100644 --- a/web/src/hooks/queries/useSkills.ts +++ b/web/src/hooks/queries/useSkills.ts @@ -46,9 +46,12 @@ export function useSkills( return response }, enabled: Boolean(api && sessionId), - staleTime: 30_000, - refetchInterval: 30_000, - refetchOnWindowFocus: 'always', + // Skills change only when the user edits files on the machine, so + // polling them on a timer just burns relay bandwidth on an answer + // that is almost always identical. getSuggestions() refetches when + // the user actually types "$", which is the moment freshness matters. + staleTime: 5 * 60_000, + refetchOnWindowFocus: true, gcTime: 30 * 60 * 1000, retry: false, }) @@ -61,10 +64,12 @@ export function useSkills( }, [query.data]) const getSuggestions = useCallback(async (queryText: string): Promise => { - const refreshed = queryText === '$' ? await query.refetch() : null - const currentSkills = refreshed?.data?.success - ? (refreshed.data.skills ?? []) - : skills + // Fire-and-forget for the same reason as useSlashCommands: the RPC can + // stall behind a wedged CLI, and the menu must not block on it. + if (queryText === '$') { + void query.refetch() + } + const currentSkills = skills const recent = getRecentSkills() const getRecency = (name: string) => recent[name] ?? 0 const searchTerm = queryText.startsWith('$') diff --git a/web/src/hooks/queries/useSlashCommands.ts b/web/src/hooks/queries/useSlashCommands.ts index 593ae026..bb5a78f2 100644 --- a/web/src/hooks/queries/useSlashCommands.ts +++ b/web/src/hooks/queries/useSlashCommands.ts @@ -44,9 +44,11 @@ export function useSlashCommands( return await api.getSlashCommands(sessionId) }, enabled: Boolean(api && sessionId), - staleTime: 30_000, - refetchInterval: 30_000, - refetchOnWindowFocus: 'always', + // Same reasoning as useSkills: the command list is near-static, so it + // is refetched when the user opens the menu with "/" rather than on a + // 30s timer that runs for as long as the session is open. + staleTime: 5 * 60_000, + refetchOnWindowFocus: true, gcTime: 30 * 60 * 1000, retry: false, // Don't retry RPC failures }) @@ -66,6 +68,15 @@ export function useSlashCommands( }, [agentType, query.data]) const getSuggestions = useCallback(async (queryText: string): Promise => { + // Opening the menu is the one moment a stale list is visible, so + // refresh here instead of polling in the background. Fire-and-forget: + // the /slash-commands RPC can stall for its full 30s timeout when the + // CLI is wedged, and the menu must show cached + built-in commands + // immediately - the refreshed list flows in via query state. + if (queryText === '/') { + void query.refetch() + } + const searchTerm = queryText.startsWith('/') ? queryText.slice(1).toLowerCase() : queryText.toLowerCase() @@ -105,7 +116,7 @@ export function useSlashCommands( content: cmd.content, source: cmd.source })) - }, [commands]) + }, [commands, query.refetch]) return { commands, diff --git a/web/src/hooks/useAgentTerminalSocket.ts b/web/src/hooks/useAgentTerminalSocket.ts index 0a560fb4..d21b288b 100644 --- a/web/src/hooks/useAgentTerminalSocket.ts +++ b/web/src/hooks/useAgentTerminalSocket.ts @@ -92,6 +92,12 @@ export function useAgentTerminalSocket(options: UseAgentTerminalSocketOptions): reconnectionDelay: 1000, reconnectionDelayMax: 5000, transports: ['polling', 'websocket'], + // Once a websocket upgrade has succeeded, reconnects go straight + // to websocket instead of repeating the HTTP long-polling + // handshake — several fewer round-trips through the relay tunnel + // on every reconnect. First-ever connects keep the polling + // fallback for networks that block websockets. + rememberUpgrade: true, autoConnect: false }) const socket = manager.socket('/terminal', { diff --git a/web/src/hooks/useSSE.ts b/web/src/hooks/useSSE.ts index 009cdb44..9667197f 100644 --- a/web/src/hooks/useSSE.ts +++ b/web/src/hooks/useSSE.ts @@ -85,6 +85,11 @@ const HEARTBEAT_WATCHDOG_INTERVAL_MS = 10_000 const RECONNECT_BASE_DELAY_MS = 1_000 const RECONNECT_MAX_DELAY_MS = 30_000 const RECONNECT_JITTER_MS = 500 +// A hub that stays unreachable is usually offline for hours, not seconds. +// Retrying every 30s forever costs a full TLS handshake per attempt through +// the relay, so widen the ceiling once the fast retries have clearly failed. +const RECONNECT_SLOW_AFTER_ATTEMPTS = 8 +const RECONNECT_SLOW_MAX_DELAY_MS = 300_000 const INVALIDATION_BATCH_MS = 16 function sortSessionSummaries(left: SessionSummary, right: SessionSummary): number { @@ -288,7 +293,8 @@ function buildEventsUrl( baseUrl: string, token: string, subscription: SSESubscription, - visibility: VisibilityState + visibility: VisibilityState, + lastEventId: string | null ): string { const params = new URLSearchParams() params.set('token', token) @@ -302,6 +308,9 @@ function buildEventsUrl( if (subscription.machineId) { params.set('machineId', subscription.machineId) } + if (lastEventId) { + params.set('lastEventId', lastEventId) + } const path = `/api/events?${params.toString()}` try { @@ -318,7 +327,12 @@ export function useSSE(options: { subscription?: SSESubscription scope?: SSEScope onEvent: (event: SyncEvent) => void - onConnect?: () => void + /** + * Fires on the server's connection-changed handshake. `resumed` is true + * when the hub replayed everything missed since the last connection, so + * the caller can skip its full REST resync. + */ + onConnect?: (info: { resumed: boolean }) => void onDisconnect?: (reason: string) => void onError?: (error: unknown) => void onToast?: (event: ToastEvent) => void @@ -339,6 +353,15 @@ export function useSSE(options: { }>({ sessions: false, machines: false, sessionIds: new Set(), scratchlistSessionIds: new Set() }) const reconnectTimerRef = useRef | null>(null) const reconnectAttemptRef = useRef(0) + // Set when a reconnect was due while the tab was hidden. Hidden tabs do + // not schedule retries at all - they wait for the tab to come back. + const reconnectDeferredRef = useRef(false) + // Last SSE event id seen, keyed by the subscription it belongs to. Sent + // as ?lastEventId on reconnects of the SAME subscription so the hub can + // replay the gap instead of the client refetching everything. A cursor + // from a different subscription (e.g. after a session switch) would make + // the hub replay against the wrong filter set, hence the key check. + const lastEventCursorRef = useRef<{ key: string; id: string } | null>(null) const lastActivityAtRef = useRef(0) const [reconnectNonce, setReconnectNonce] = useState(0) const [subscriptionId, setSubscriptionId] = useState(null) @@ -387,15 +410,19 @@ export function useSSE(options: { reconnectTimerRef.current = null } reconnectAttemptRef.current = 0 + reconnectDeferredRef.current = false setSubscriptionId(null) return } setSubscriptionId(null) + const resumeCursor = lastEventCursorRef.current?.key === subscriptionKey + ? lastEventCursorRef.current.id + : null const url = buildEventsUrl(options.baseUrl, options.token, { ...subscription, sessionId: subscription.sessionId ?? undefined - }, getVisibilityState()) + }, getVisibilityState(), resumeCursor) const eventSource = new EventSource(url) let disconnectNotified = false let reconnectRequested = false @@ -403,8 +430,23 @@ export function useSSE(options: { lastActivityAtRef.current = Date.now() const scheduleReconnect = () => { + // Nobody is looking at a hidden tab, and every retry through the + // relay costs a TLS handshake. Defer until the tab is visible + // again; onVisibilityChange reconnects immediately at that point. + if (getVisibilityState() === 'hidden') { + reconnectDeferredRef.current = true + if (reconnectTimerRef.current) { + clearTimeout(reconnectTimerRef.current) + reconnectTimerRef.current = null + } + return + } + const attempt = reconnectAttemptRef.current - const exponentialDelay = Math.min(RECONNECT_MAX_DELAY_MS, RECONNECT_BASE_DELAY_MS * (2 ** attempt)) + const maxDelay = attempt >= RECONNECT_SLOW_AFTER_ATTEMPTS + ? RECONNECT_SLOW_MAX_DELAY_MS + : RECONNECT_MAX_DELAY_MS + const exponentialDelay = Math.min(maxDelay, RECONNECT_BASE_DELAY_MS * (2 ** attempt)) const jitter = Math.floor(Math.random() * (RECONNECT_JITTER_MS + 1)) reconnectAttemptRef.current = attempt + 1 if (reconnectTimerRef.current) { @@ -690,6 +732,13 @@ export function useSSE(options: { setSubscriptionId(nextId) } } + // The connect callback fires here, not on EventSource open: + // only the server handshake knows whether the gap was replayed + // (`resume: 'ok'`) or the caller must resync. Older hubs omit + // the field, which safely reads as a full resync. + const resumed = data && typeof data === 'object' + && (data as { resume?: unknown }).resume === 'ok' + onConnectRef.current?.({ resumed: Boolean(resumed) }) } if (event.type === 'toast') { @@ -815,6 +864,15 @@ export function useSSE(options: { } handleSyncEvent(parsed as SyncEvent) + + // Track the hub's replay cursor - after handling, so a throwing + // handler leaves the cursor behind the event and the hub replays + // it on the next reconnect (at-least-once). EventSource keeps + // lastEventId sticky across frames, so heartbeats (no id field) + // simply repeat the previous cursor. + if (message.lastEventId) { + lastEventCursorRef.current = { key: subscriptionKey, id: message.lastEventId } + } } eventSource.onmessage = handleMessage @@ -824,9 +882,11 @@ export function useSSE(options: { reconnectTimerRef.current = null } reconnectAttemptRef.current = 0 + reconnectDeferredRef.current = false disconnectNotified = false lastActivityAtRef.current = Date.now() - onConnectRef.current?.() + // onConnect intentionally does NOT fire here: it fires on the + // connection-changed handshake, which carries the resume verdict. } eventSource.onerror = (error) => { onErrorRef.current?.(error) @@ -834,7 +894,12 @@ export function useSSE(options: { requestReconnect('closed') return } - notifyDisconnect('error') + // CONNECTING means the browser would keep retrying natively - at + // its own few-second interval, with the original (stale) URL, and + // regardless of tab visibility. Take the source down and route the + // retry through our scheduler so hidden-tab deferral and the slow + // backoff actually govern every reconnect path. + requestReconnect('transport-error') } const watchdogTimer = setInterval(() => { @@ -856,6 +921,14 @@ export function useSSE(options: { // HEARTBEAT_WATCHDOG_INTERVAL_MS after switching back. const onVisibilityChange = () => { if (getVisibilityState() !== 'visible') return + // A retry fell due while the tab was hidden and was deliberately + // not scheduled. Run it now, before the identity guard below: + // requestReconnect has already cleared eventSourceRef. + if (reconnectDeferredRef.current) { + reconnectDeferredRef.current = false + setReconnectNonce((value) => value + 1) + return + } if (eventSourceRef.current !== eventSource) return if (Date.now() - lastActivityAtRef.current >= HEARTBEAT_STALE_MS) { requestReconnect('visibility-recovery') diff --git a/web/src/hooks/useTerminalSocket.ts b/web/src/hooks/useTerminalSocket.ts index 4789aeed..d96c452e 100644 --- a/web/src/hooks/useTerminalSocket.ts +++ b/web/src/hooks/useTerminalSocket.ts @@ -123,6 +123,9 @@ export function useTerminalSocket(options: UseTerminalSocketOptions): { reconnectionDelay: 1000, reconnectionDelayMax: 5000, transports: ['polling', 'websocket'], + // Skip the HTTP long-polling phase on reconnects once a websocket + // upgrade has succeeded — see useAgentTerminalSocket. + rememberUpgrade: true, autoConnect: false }) const socket = manager.socket('/terminal', {