Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 35 additions & 8 deletions packages/app/src/context/server-sdk.test.ts
Original file line number Diff line number Diff line change
@@ -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)
})
})

Expand Down
74 changes: 67 additions & 7 deletions packages/app/src/context/server-sdk.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof createGlobalEmitter<{ [key: string]: ServerEvent }>>
Expand Down Expand Up @@ -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[] = []
Expand Down Expand Up @@ -256,6 +268,20 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS
let run: Promise<void> | 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<typeof setTimeout> | 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
Expand All @@ -271,14 +297,25 @@ 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 =
kind === "v1"
? (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
Expand All @@ -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,
Expand All @@ -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
Expand All @@ -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(() => {
Expand Down
Loading