mirror of
https://github.com/wu736139669/hapi.git
synced 2026-10-06 18:39:47 +00:00
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
This commit is contained in:
@@ -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' })
|
||||
})
|
||||
})
|
||||
|
||||
+159
-6
@@ -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<void>
|
||||
send: (event: SyncEvent, eventId?: string) => void | Promise<void>
|
||||
sendHeartbeat: () => void | Promise<void>
|
||||
/**
|
||||
* 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<string, SSEConnection> = 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<string, string>()
|
||||
|
||||
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<void>
|
||||
/** Last event id the client saw; enables replay instead of resync. */
|
||||
resumeFrom?: string | null
|
||||
send: (event: SyncEvent, eventId?: string) => void | Promise<void>
|
||||
sendHeartbeat: () => void | Promise<void>
|
||||
}): 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<void> {
|
||||
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<SyncEvent, { type: 'toast' }>): Promise<number> {
|
||||
const deliveries: Array<Promise<{ id: string; ok: boolean }>> = []
|
||||
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)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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<string, unknown> }
|
||||
|
||||
/**
|
||||
* 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<Frame[]> {
|
||||
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<string, unknown> })
|
||||
}
|
||||
}
|
||||
}
|
||||
} 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<WebAppEnv> {
|
||||
const app = new Hono<WebAppEnv>()
|
||||
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<WebAppEnv>, query: string): Promise<Response> {
|
||||
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)
|
||||
})
|
||||
})
|
||||
@@ -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<void>((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<void>((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'))
|
||||
|
||||
+33
-1
@@ -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<unknown>)
|
||||
const wsAny = ws as ServerWebSocket<{ _qwenProxy?: boolean; _geminiProxy?: boolean }>
|
||||
if (wsAny.data?._geminiProxy) {
|
||||
geminiProxyHandler.open(wsAny)
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
})
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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()
|
||||
})
|
||||
])
|
||||
|
||||
+15
-2
@@ -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)
|
||||
})
|
||||
|
||||
@@ -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<Suggestion[]> => {
|
||||
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('$')
|
||||
|
||||
@@ -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<Suggestion[]> => {
|
||||
// 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,
|
||||
|
||||
@@ -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', {
|
||||
|
||||
+79
-6
@@ -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<ReturnType<typeof setTimeout> | 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<string | null>(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')
|
||||
|
||||
@@ -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', {
|
||||
|
||||
Reference in New Issue
Block a user