When the computer sleeps and wakes, the SSE/WS event stream drops silently. Messages appeared sent (optimistic insert) but never reached the OpenCode server, and the user had no indication the system was disconnected. Three fixes: 1. Connection state tracking: add onDisconnect callback to the event pipeline. Stream failures set isConnected=false in useConfigStore; successful reconnect sets isConnected=true. 2. Send guard: optimisticSend, respondToPermission, and respondToQuestion now check isConnected before making API calls, throwing a clear error that surfaces as a toast to the user. The /compact command also checks connection with error feedback. 3. Faster server recovery: add triggerHealthCheck() to the server lifecycle and wire it into the WS event stream runtime. When the upstream OpenCode connection fails, the server immediately checks health and restarts if needed, instead of waiting up to 15s for the periodic health check.
568 lines
16 KiB
TypeScript
568 lines
16 KiB
TypeScript
/**
|
|
* Event Pipeline — transport connection, event coalescing, and batched flush.
|
|
*
|
|
* Plain closure API:
|
|
* const { cleanup } = createEventPipeline({ sdk, onEvent })
|
|
*
|
|
* No class, no start/stop lifecycle. One pipeline per mount.
|
|
* Abort controller created once at init, cleaned up via returned cleanup fn.
|
|
*/
|
|
|
|
import type { Event, OpencodeClient } from "@opencode-ai/sdk/v2/client"
|
|
import { opencodeClient } from "@/lib/opencode/client"
|
|
import { syncDebug } from "./debug"
|
|
|
|
export type QueuedEvent = {
|
|
directory: string
|
|
payload: Event
|
|
}
|
|
|
|
export type FlushHandler = (events: QueuedEvent[]) => void
|
|
|
|
const FLUSH_FRAME_MS = 16
|
|
const STREAM_YIELD_MS = 8
|
|
const RECONNECT_DELAY_MS = 250
|
|
const HEARTBEAT_TIMEOUT_MS = 15_000
|
|
const WS_FALLBACK_WINDOW_MS = 60_000
|
|
const WS_READY_TIMEOUT_MS = 2_000
|
|
const ABSOLUTE_URL_PATTERN = /^[a-zA-Z][a-zA-Z\d+\-.]*:\/\//
|
|
|
|
export type EventPipelineInput = {
|
|
sdk: OpencodeClient
|
|
onEvent: (directory: string, payload: Event) => void
|
|
routeDirectory?: (directory: string, payload: Event) => string
|
|
/** Called after stream reconnects (visibility restore or heartbeat timeout). */
|
|
onReconnect?: () => void
|
|
/** Called when the stream disconnects (heartbeat timeout, network error, or transport failure). */
|
|
onDisconnect?: () => void
|
|
transport?: "auto" | "ws" | "sse"
|
|
}
|
|
|
|
type MessageStreamWsFrame = {
|
|
type: "ready" | "event" | "error"
|
|
payload?: unknown
|
|
eventId?: string
|
|
directory?: string
|
|
message?: string
|
|
scope?: "global" | "directory"
|
|
}
|
|
|
|
const normalizeEventType = (payload: Event): Event => {
|
|
const type = (payload as { type?: unknown }).type
|
|
if (typeof type !== "string") {
|
|
return payload
|
|
}
|
|
|
|
const match = /^(.*)\.(\d+)$/.exec(type)
|
|
if (!match || !match[1]) {
|
|
return payload
|
|
}
|
|
|
|
return {
|
|
...payload,
|
|
type: match[1] as Event["type"],
|
|
} as unknown as Event
|
|
}
|
|
|
|
function resolveEventDirectory(event: unknown, payload: Event): string {
|
|
const directDirectory =
|
|
typeof event === "object" && event !== null && typeof (event as { directory?: unknown }).directory === "string"
|
|
? (event as { directory: string }).directory
|
|
: null
|
|
|
|
if (directDirectory && directDirectory.length > 0) {
|
|
return directDirectory
|
|
}
|
|
|
|
const properties =
|
|
typeof payload.properties === "object" && payload.properties !== null
|
|
? (payload.properties as Record<string, unknown>)
|
|
: null
|
|
const propertyDirectory = typeof properties?.directory === "string" ? properties.directory : null
|
|
|
|
return propertyDirectory && propertyDirectory.length > 0 ? propertyDirectory : "global"
|
|
}
|
|
|
|
function resolveEventPayload(payload: unknown): Event | null {
|
|
if (!payload || typeof payload !== "object") {
|
|
return null
|
|
}
|
|
|
|
const record = payload as { type?: unknown; payload?: unknown }
|
|
if (typeof record.type === "string") {
|
|
return payload as Event
|
|
}
|
|
|
|
if (record.payload && typeof record.payload === "object" && typeof (record.payload as { type?: unknown }).type === "string") {
|
|
return record.payload as Event
|
|
}
|
|
|
|
return null
|
|
}
|
|
|
|
function resolveAbsoluteUrl(candidate: string): string {
|
|
const normalized = typeof candidate === "string" && candidate.trim().length > 0 ? candidate.trim() : "/api"
|
|
if (ABSOLUTE_URL_PATTERN.test(normalized)) {
|
|
return normalized
|
|
}
|
|
|
|
if (typeof window === "undefined") {
|
|
return normalized
|
|
}
|
|
|
|
const baseReference = window.location?.href || window.location?.origin
|
|
if (!baseReference) {
|
|
return normalized
|
|
}
|
|
|
|
return new URL(normalized, baseReference).toString()
|
|
}
|
|
|
|
function toWebSocketUrl(candidate: string): string {
|
|
const url = new URL(resolveAbsoluteUrl(candidate))
|
|
url.protocol = url.protocol === "https:" ? "wss:" : "ws:"
|
|
return url.toString()
|
|
}
|
|
|
|
function buildGlobalEventWsUrl(lastEventId?: string): string {
|
|
const baseUrl = opencodeClient.getBaseUrl()
|
|
const normalizedBase = baseUrl.endsWith("/") ? baseUrl : `${baseUrl}/`
|
|
const httpUrl = new URL("global/event/ws", resolveAbsoluteUrl(normalizedBase))
|
|
if (lastEventId && lastEventId.length > 0) {
|
|
httpUrl.searchParams.set("lastEventId", lastEventId)
|
|
}
|
|
return toWebSocketUrl(httpUrl.toString())
|
|
}
|
|
|
|
type DirectoryQueue = {
|
|
queue: Event[]
|
|
buffer: Event[]
|
|
coalesced: Map<string, number>
|
|
staleDeltas: Set<string>
|
|
timer: ReturnType<typeof setTimeout> | undefined
|
|
last: number
|
|
}
|
|
|
|
export function createEventPipeline(input: EventPipelineInput) {
|
|
const { sdk, onEvent, onReconnect, onDisconnect, routeDirectory, transport = "auto" } = input
|
|
const abort = new AbortController()
|
|
let hasConnected = false
|
|
let disconnected = false
|
|
let lastEventId: string | undefined
|
|
let wsFallbackUntil = 0
|
|
|
|
const directories = new Map<string, DirectoryQueue>()
|
|
|
|
const getOrCreateDir = (directory: string): DirectoryQueue => {
|
|
let d = directories.get(directory)
|
|
if (d) return d
|
|
d = {
|
|
queue: [],
|
|
buffer: [],
|
|
coalesced: new Map(),
|
|
staleDeltas: new Set(),
|
|
timer: undefined,
|
|
last: 0,
|
|
}
|
|
directories.set(directory, d)
|
|
return d
|
|
}
|
|
|
|
const key = (payload: Event): string | undefined => {
|
|
if (payload.type === "session.status") {
|
|
const props = payload.properties as { sessionID: string }
|
|
return `session.status:${props.sessionID}`
|
|
}
|
|
if (payload.type === "lsp.updated") {
|
|
return "lsp.updated"
|
|
}
|
|
if (payload.type === "message.part.updated") {
|
|
const part = (payload.properties as { part: { messageID: string; id: string } }).part
|
|
return `message.part.updated:${part.messageID}:${part.id}`
|
|
}
|
|
if (payload.type === "message.part.delta") {
|
|
const props = payload.properties as { messageID: string; partID: string; field: string }
|
|
return `message.part.delta:${props.messageID}:${props.partID}:${props.field}`
|
|
}
|
|
return undefined
|
|
}
|
|
|
|
const deltaKey = (messageID: string, partID: string, field: string) => `${messageID}:${partID}:${field}`
|
|
|
|
const flushDir = (directory: string) => {
|
|
const d = directories.get(directory)
|
|
if (!d) return
|
|
if (d.timer) {
|
|
clearTimeout(d.timer)
|
|
d.timer = undefined
|
|
}
|
|
if (d.queue.length === 0) return
|
|
|
|
const events = d.queue
|
|
const staleDeltas = d.staleDeltas.size > 0 ? new Set(d.staleDeltas) : undefined
|
|
d.queue = d.buffer
|
|
d.buffer = events
|
|
d.queue.length = 0
|
|
d.coalesced.clear()
|
|
d.staleDeltas.clear()
|
|
|
|
d.last = Date.now()
|
|
syncDebug.pipeline.flush(events.length)
|
|
for (const payload of events) {
|
|
if (staleDeltas && payload.type === "message.part.delta") {
|
|
const props = payload.properties as { messageID: string; partID: string; field: string }
|
|
if (staleDeltas.has(deltaKey(props.messageID, props.partID, props.field))) {
|
|
continue
|
|
}
|
|
}
|
|
onEvent(directory, payload)
|
|
}
|
|
|
|
d.buffer.length = 0
|
|
}
|
|
|
|
const flushAll = () => {
|
|
for (const directory of directories.keys()) {
|
|
flushDir(directory)
|
|
}
|
|
}
|
|
|
|
const scheduleDir = (directory: string) => {
|
|
const d = getOrCreateDir(directory)
|
|
if (d.timer) return
|
|
const elapsed = Date.now() - d.last
|
|
d.timer = setTimeout(() => flushDir(directory), Math.max(0, FLUSH_FRAME_MS - elapsed))
|
|
}
|
|
|
|
const wait = (ms: number) => new Promise<void>((resolve) => setTimeout(resolve, ms))
|
|
const isAbortError = (error: unknown): boolean =>
|
|
error instanceof DOMException && error.name === "AbortError" ||
|
|
(typeof error === "object" && error !== null && (error as { name?: string }).name === "AbortError")
|
|
|
|
let streamErrorLogged = false
|
|
let attempt: AbortController | undefined
|
|
let lastEventAt = Date.now()
|
|
let heartbeat: ReturnType<typeof setTimeout> | undefined
|
|
|
|
const markConnected = () => {
|
|
disconnected = false
|
|
if (hasConnected) {
|
|
onReconnect?.()
|
|
return
|
|
}
|
|
hasConnected = true
|
|
}
|
|
|
|
const enqueueEvent = (directory: string, payload: Event) => {
|
|
const normalizedPayload = normalizeEventType(payload)
|
|
const routedDirectory = routeDirectory?.(directory, normalizedPayload) || directory
|
|
const d = getOrCreateDir(routedDirectory)
|
|
const k = key(normalizedPayload)
|
|
if (k) {
|
|
const i = d.coalesced.get(k)
|
|
if (i !== undefined) {
|
|
if (normalizedPayload.type === "message.part.delta") {
|
|
const prev = d.queue[i] as unknown as { properties: { delta: string } }
|
|
const inc = normalizedPayload.properties as { delta: string }
|
|
d.queue[i] = {
|
|
...normalizedPayload,
|
|
properties: {
|
|
...(normalizedPayload.properties as object),
|
|
delta: prev.properties.delta + inc.delta,
|
|
},
|
|
} as unknown as Event
|
|
} else {
|
|
d.queue[i] = normalizedPayload
|
|
if (normalizedPayload.type === "message.part.updated") {
|
|
const part = (normalizedPayload.properties as { part: { messageID: string; id: string } }).part
|
|
d.staleDeltas.add(deltaKey(part.messageID, part.id, "text"))
|
|
d.staleDeltas.add(deltaKey(part.messageID, part.id, "output"))
|
|
}
|
|
}
|
|
syncDebug.pipeline.coalesced(normalizedPayload.type, k)
|
|
return
|
|
}
|
|
d.coalesced.set(k, d.queue.length)
|
|
}
|
|
|
|
d.queue.push(normalizedPayload)
|
|
scheduleDir(routedDirectory)
|
|
}
|
|
|
|
const resetHeartbeat = () => {
|
|
lastEventAt = Date.now()
|
|
if (heartbeat) clearTimeout(heartbeat)
|
|
heartbeat = setTimeout(() => {
|
|
attempt?.abort()
|
|
}, HEARTBEAT_TIMEOUT_MS)
|
|
}
|
|
|
|
const clearHeartbeat = () => {
|
|
if (!heartbeat) return
|
|
clearTimeout(heartbeat)
|
|
heartbeat = undefined
|
|
}
|
|
|
|
const runSseAttempt = async (signal: AbortSignal) => {
|
|
const events = await sdk.global.event({
|
|
signal,
|
|
onSseError: (error: unknown) => {
|
|
if (isAbortError(error)) return
|
|
if (streamErrorLogged) return
|
|
streamErrorLogged = true
|
|
console.error("[event-pipeline] SSE stream error", error)
|
|
},
|
|
})
|
|
|
|
markConnected()
|
|
|
|
let yielded = Date.now()
|
|
resetHeartbeat()
|
|
|
|
for await (const event of events.stream) {
|
|
resetHeartbeat()
|
|
streamErrorLogged = false
|
|
const payload = resolveEventPayload((event as { payload?: Event }).payload ?? event)
|
|
if (!payload) {
|
|
continue
|
|
}
|
|
const directory = resolveEventDirectory(event, payload)
|
|
enqueueEvent(directory, payload)
|
|
|
|
if (Date.now() - yielded < STREAM_YIELD_MS) continue
|
|
yielded = Date.now()
|
|
await wait(0)
|
|
}
|
|
}
|
|
|
|
const runWsAttempt = async (signal: AbortSignal) => {
|
|
await new Promise<void>((resolve, reject) => {
|
|
let settled = false
|
|
let opened = false
|
|
const socket = new WebSocket(buildGlobalEventWsUrl(lastEventId))
|
|
const setFallbackCode = (error: Error) => {
|
|
if (!opened && transport === "auto") {
|
|
wsFallbackUntil = Date.now() + WS_FALLBACK_WINDOW_MS
|
|
;(error as Error & { code?: string }).code = "WS_FALLBACK"
|
|
}
|
|
}
|
|
|
|
let readyTimer: ReturnType<typeof setTimeout> | undefined = setTimeout(() => {
|
|
readyTimer = undefined
|
|
const error = new Error("Message stream WebSocket ready timeout")
|
|
setFallbackCode(error)
|
|
settleReject(error)
|
|
try {
|
|
socket.close()
|
|
} catch {
|
|
// ignore
|
|
}
|
|
}, WS_READY_TIMEOUT_MS)
|
|
|
|
const cleanup = () => {
|
|
if (readyTimer) {
|
|
clearTimeout(readyTimer)
|
|
readyTimer = undefined
|
|
}
|
|
socket.onopen = null
|
|
socket.onmessage = null
|
|
socket.onerror = null
|
|
socket.onclose = null
|
|
}
|
|
|
|
const settleResolve = () => {
|
|
if (settled) return
|
|
settled = true
|
|
signal.removeEventListener("abort", handleAbort)
|
|
cleanup()
|
|
resolve()
|
|
}
|
|
|
|
const settleReject = (error: unknown) => {
|
|
if (settled) return
|
|
settled = true
|
|
signal.removeEventListener("abort", handleAbort)
|
|
cleanup()
|
|
reject(error)
|
|
}
|
|
|
|
const handleAbort = () => {
|
|
try {
|
|
socket.close()
|
|
} catch {
|
|
// ignore close failures during abort
|
|
}
|
|
settleResolve()
|
|
}
|
|
|
|
signal.addEventListener("abort", handleAbort, { once: true })
|
|
|
|
socket.onopen = () => {
|
|
streamErrorLogged = false
|
|
}
|
|
|
|
socket.onmessage = (messageEvent) => {
|
|
resetHeartbeat()
|
|
streamErrorLogged = false
|
|
|
|
let frame: MessageStreamWsFrame | null = null
|
|
try {
|
|
frame = JSON.parse(String(messageEvent.data)) as MessageStreamWsFrame
|
|
} catch (error) {
|
|
console.warn("[event-pipeline] Failed to parse WS frame", error)
|
|
return
|
|
}
|
|
|
|
if (!frame || typeof frame.type !== "string") {
|
|
return
|
|
}
|
|
|
|
if (frame.type === "ready") {
|
|
opened = true
|
|
if (readyTimer) {
|
|
clearTimeout(readyTimer)
|
|
readyTimer = undefined
|
|
}
|
|
markConnected()
|
|
return
|
|
}
|
|
|
|
if (frame.type === "error") {
|
|
const error = new Error(frame.message || "Message stream WebSocket error")
|
|
setFallbackCode(error)
|
|
settleReject(error)
|
|
try {
|
|
socket.close()
|
|
} catch {
|
|
// ignore
|
|
}
|
|
return
|
|
}
|
|
|
|
if (frame.type !== "event") {
|
|
return
|
|
}
|
|
|
|
const payload = resolveEventPayload(frame.payload)
|
|
if (!payload) {
|
|
return
|
|
}
|
|
|
|
if (typeof frame.eventId === "string" && frame.eventId.length > 0) {
|
|
lastEventId = frame.eventId
|
|
}
|
|
|
|
const directory = resolveEventDirectory(
|
|
{ directory: frame.directory, payload },
|
|
payload,
|
|
)
|
|
enqueueEvent(directory, payload)
|
|
}
|
|
|
|
socket.onerror = () => {
|
|
void 0
|
|
}
|
|
|
|
socket.onclose = () => {
|
|
if (signal.aborted) {
|
|
settleResolve()
|
|
return
|
|
}
|
|
|
|
const error = new Error("Global message stream WebSocket closed")
|
|
setFallbackCode(error)
|
|
settleReject(error)
|
|
}
|
|
})
|
|
}
|
|
|
|
const resolveTransport = (): "ws" | "sse" => {
|
|
if (typeof WebSocket !== "function") {
|
|
return "sse"
|
|
}
|
|
if (transport === "ws") {
|
|
return "ws"
|
|
}
|
|
if (transport === "sse") {
|
|
return "sse"
|
|
}
|
|
return wsFallbackUntil > Date.now() ? "sse" : "ws"
|
|
}
|
|
|
|
void (async () => {
|
|
while (!abort.signal.aborted) {
|
|
attempt = new AbortController()
|
|
lastEventAt = Date.now()
|
|
let retryDelayMs = RECONNECT_DELAY_MS
|
|
const currentTransport = resolveTransport()
|
|
const onAbort = () => {
|
|
attempt?.abort()
|
|
}
|
|
abort.signal.addEventListener("abort", onAbort)
|
|
|
|
try {
|
|
if (currentTransport === "ws") {
|
|
await runWsAttempt(attempt.signal)
|
|
} else {
|
|
await runSseAttempt(attempt.signal)
|
|
}
|
|
} catch (error) {
|
|
const code = typeof error === "object" && error !== null ? (error as { code?: unknown }).code : undefined
|
|
if (currentTransport === "ws" && code === "WS_FALLBACK") {
|
|
retryDelayMs = 0
|
|
} else if (!isAbortError(error)) {
|
|
if (!streamErrorLogged) {
|
|
streamErrorLogged = true
|
|
console.error("[event-pipeline] stream failed", error)
|
|
}
|
|
// Notify consumer that the stream has disconnected, so it can
|
|
// update connection state (e.g. set isConnected = false).
|
|
// Guard: only fire once per disconnection cycle to avoid repeated
|
|
// setState calls on every failed retry attempt.
|
|
if (!disconnected) {
|
|
disconnected = true
|
|
onDisconnect?.()
|
|
}
|
|
}
|
|
} finally {
|
|
abort.signal.removeEventListener("abort", onAbort)
|
|
attempt = undefined
|
|
clearHeartbeat()
|
|
}
|
|
|
|
if (abort.signal.aborted) return
|
|
if (retryDelayMs > 0) {
|
|
await wait(retryDelayMs)
|
|
}
|
|
}
|
|
})().finally(flushAll)
|
|
|
|
const onVisibility = () => {
|
|
if (typeof document === "undefined") return
|
|
if (document.visibilityState !== "visible") return
|
|
if (Date.now() - lastEventAt < HEARTBEAT_TIMEOUT_MS) return
|
|
attempt?.abort()
|
|
}
|
|
|
|
const onPageShow = (event: PageTransitionEvent) => {
|
|
if (!event.persisted) return
|
|
attempt?.abort()
|
|
}
|
|
|
|
if (typeof document !== "undefined") {
|
|
document.addEventListener("visibilitychange", onVisibility)
|
|
window.addEventListener("pageshow", onPageShow)
|
|
}
|
|
|
|
const cleanup = () => {
|
|
if (typeof document !== "undefined") {
|
|
document.removeEventListener("visibilitychange", onVisibility)
|
|
window.removeEventListener("pageshow", onPageShow)
|
|
}
|
|
abort.abort()
|
|
flushAll()
|
|
}
|
|
|
|
return { cleanup }
|
|
}
|