diff --git a/packages/app/src/context/server-sdk.test.ts b/packages/app/src/context/server-sdk.test.ts index 57e1cd86f3ac..faf4d1d3c90d 100644 --- a/packages/app/src/context/server-sdk.test.ts +++ b/packages/app/src/context/server-sdk.test.ts @@ -1,17 +1,44 @@ import { describe, expect, test } from "bun:test" -import { adaptServerEvent, coalesceServerEvents, enqueueServerEvent, resumeStreamAfterPageShow } from "./server-sdk" +import { + RECONNECT_DELAY_MAX_MS, + STREAM_STALL_TIMEOUT_MS, + adaptServerEvent, + coalesceServerEvents, + enqueueServerEvent, + reconnectDelay, +} from "./server-sdk" import type { OpenCodeEvent } from "@opencode-ai/client/promise" import type { Event } from "@opencode-ai/sdk/v2/client" -describe("resumeStreamAfterPageShow", () => { - test("restarts a stream only after a back-forward cache restore", () => { - let starts = 0 - const start = () => starts++ +describe("reconnectDelay", () => { + test("starts at the base delay for the first fast drop", () => { + expect(reconnectDelay(0, 0)).toBe(250) + }) + + test("grows exponentially while reconnects keep dropping fast", () => { + expect(reconnectDelay(1, 0)).toBe(500) + expect(reconnectDelay(2, 0)).toBe(1000) + expect(reconnectDelay(3, 0)).toBe(2000) + }) + + test("caps at the maximum delay", () => { + expect(reconnectDelay(20, 0)).toBe(RECONNECT_DELAY_MAX_MS) + }) - resumeStreamAfterPageShow({ persisted: false } as PageTransitionEvent, start) - resumeStreamAfterPageShow({ persisted: true } as PageTransitionEvent, start) + test("adds bounded jitter without exceeding backoff + base delay", () => { + for (const drops of [0, 1, 2, 10, 20]) { + const delay = reconnectDelay(drops, 1) + const backoff = Math.min(250 * 2 ** drops, RECONNECT_DELAY_MAX_MS) + expect(delay).toBeLessThanOrEqual(backoff + 250) + } + }) +}) - expect(starts).toBe(1) +describe("STREAM_STALL_TIMEOUT_MS", () => { + test("exceeds three server heartbeat intervals", () => { + // The server emits server.heartbeat every 10s; a live stream should never + // go this long without any event. + expect(STREAM_STALL_TIMEOUT_MS).toBeGreaterThan(30_000) }) }) diff --git a/packages/app/src/context/server-sdk.tsx b/packages/app/src/context/server-sdk.tsx index 7dd2a6e59edf..a43d76a5a250 100644 --- a/packages/app/src/context/server-sdk.tsx +++ b/packages/app/src/context/server-sdk.tsx @@ -159,9 +159,22 @@ function currentDeltaFragment(event: CurrentDelta) { return event.type === "session.compaction.delta" ? event.data.text : event.data.delta } -export function resumeStreamAfterPageShow(event: PageTransitionEvent, start: () => unknown) { - if (!event.persisted) return - start() +export const RECONNECT_DELAY_MAX_MS = 10_000 +export const STREAM_STALL_TIMEOUT_MS = 45_000 +const RECONNECT_DELAY_MS = 250 +// A stream that stayed up at least this long is treated as having connected for +// real, so the next attempt returns to the base delay instead of backing off. +const STABLE_CONNECTION_MS = 3_000 + +// Reconnects that keep dropping almost immediately indicate a flapping +// connection: each reconnect emits "server.connected", which the sync layer +// treats as "refresh everything", and without backoff each re-bootstrap gets +// cut off by the next reconnect before it can finish (surfacing as repeated +// HTTP 499 client-aborts). Exponential backoff turns that into a shrinking +// number of self-inflicted aborts instead of a sustained storm. +export function reconnectDelay(fastDrops: number, jitterSeed = Math.random()) { + const backoff = Math.min(RECONNECT_DELAY_MS * 2 ** fastDrops, RECONNECT_DELAY_MAX_MS) + return backoff + jitterSeed * Math.min(backoff, RECONNECT_DELAY_MS) } type ServerEventEmitter = ReturnType> @@ -217,7 +230,6 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS type Queued = QueuedServerEvent const FLUSH_FRAME_MS = 16 const STREAM_YIELD_MS = 8 - const RECONNECT_DELAY_MS = 250 let queue: Queued[] = [] let buffer: Queued[] = [] @@ -256,6 +268,20 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS let run: Promise | undefined let started = false let generation = 0 + // Consecutive attempts that dropped fast (under STABLE_CONNECTION_MS). + let consecutiveFastDrops = 0 + // The server emits a heartbeat every 10s; if nothing — heartbeat included — + // arrives for STREAM_STALL_TIMEOUT_MS the connection is dead in practice even + // though the streaming fetch never resolves or rejects (suspended tab, proxy + // or NAT drop without TCP RST). + let stallTimer: ReturnType | undefined + let stalled = false + + const resync = (reason: string) => { + if (!started || !attempt) return + console.debug("[global-sdk] event stream resync", { url: server.http.url, reason }) + attempt.abort() + } const start = () => { if (started) return run @@ -271,6 +297,15 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS attempt?.abort() } abort.signal.addEventListener("abort", onAbort) + const connectedAt = Date.now() + stalled = false + const armStallWatchdog = () => { + clearTimeout(stallTimer) + stallTimer = setTimeout(() => { + stalled = true + attempt?.abort() + }, STREAM_STALL_TIMEOUT_MS) + } try { const kind = await protocol const events = @@ -278,7 +313,9 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS ? (await eventSdk.global.event({ signal: attempt.signal })).stream : eventApi.event.subscribe({ signal: attempt.signal }) let yielded = Date.now() + armStallWatchdog() for await (const event of events) { + armStallWatchdog() streamErrorLogged = false const legacy = "payload" in event if (legacy && event.payload.type === "sync") continue @@ -291,7 +328,12 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS await wait(0) } } catch (error) { - if (!isStreamClosed(error, attempt?.signal) && !streamErrorLogged) { + if (stalled) { + if (!streamErrorLogged) { + streamErrorLogged = true + console.error("[global-sdk] event stream stalled", { url: server.http.url }) + } + } else if (!isStreamClosed(error, attempt?.signal) && !streamErrorLogged) { streamErrorLogged = true console.error("[global-sdk] event stream failed", { url: server.http.url, @@ -300,12 +342,16 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS }) } } finally { + clearTimeout(stallTimer) + stallTimer = undefined abort.signal.removeEventListener("abort", onAbort) attempt = undefined } if (abort.signal.aborted || !started || generation !== active) return - await wait(RECONNECT_DELAY_MS) + if (Date.now() - connectedAt >= STABLE_CONNECTION_MS) consecutiveFastDrops = 0 + else consecutiveFastDrops++ + await wait(reconnectDelay(consecutiveFastDrops)) } })().finally(() => { if (run !== current) return @@ -319,12 +365,26 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS const stop = () => { started = false generation++ + clearTimeout(stallTimer) + stallTimer = undefined attempt?.abort() } onMount(() => { makeEventListener(window, "pagehide", stop) - makeEventListener(window, "pageshow", (event) => resumeStreamAfterPageShow(event, start)) + // Always restart on pageshow: an open streaming fetch makes the page + // ineligible for bfcache, so persisted is almost always false here — the + // previous check left the stream stopped after ordinary tab returns. + makeEventListener(window, "pageshow", () => start()) + // A suspended or backgrounded tab can leave the TCP connection half-open: + // no events arrive, but the fetch neither resolves nor rejects. When the + // page becomes visible again, force a resync instead of waiting for the + // stall watchdog to notice. + makeEventListener(document, "visibilitychange", () => { + if (document.visibilityState !== "visible") return + resync("Page returned to the foreground") + }) + makeEventListener(window, "online", () => resync("Network connection restored")) }) onCleanup(() => {