From 0c8301d579c74f6bb7ffb83e222f3c57e664842b Mon Sep 17 00:00:00 2001 From: Bohdan Triapitsyn Date: Sun, 24 May 2026 19:17:04 +0300 Subject: [PATCH] fix: recover stuck live session updates Reconnects and resyncs active sessions when live updates stall Normalizes synthetic session status events Uses authoritative status snapshots to clear stale busy states --- packages/ui/src/sync/event-pipeline.test.ts | 39 +++ packages/ui/src/sync/event-pipeline.ts | 90 +++++- packages/ui/src/sync/sync-context.tsx | 268 ++++++++++++++---- packages/web/server/index.js | 2 +- .../server/lib/opencode/session-runtime.js | 4 +- .../lib/opencode/session-runtime.test.js | 6 +- 6 files changed, 333 insertions(+), 76 deletions(-) diff --git a/packages/ui/src/sync/event-pipeline.test.ts b/packages/ui/src/sync/event-pipeline.test.ts index 9a4ea81f..6ebcc41d 100644 --- a/packages/ui/src/sync/event-pipeline.test.ts +++ b/packages/ui/src/sync/event-pipeline.test.ts @@ -96,4 +96,43 @@ describe("createEventPipeline", () => { return `updated:${((event.properties as { part: { text: string } }).part).text}` })).toEqual(["updated:a", "delta:b", "updated:ab"]) }) + + test("normalizes openchamber session status events", async () => { + let resolveStreamFinished!: () => void + const streamFinished = new Promise((resolve) => { + resolveStreamFinished = resolve + }) + let resolveDelivered!: (event: Event) => void + const deliveredEvent = new Promise((resolve) => { + resolveDelivered = resolve + }) + const pipeline = createEventPipeline({ + sdk: createSdk([ + { + type: "openchamber:session-status", + properties: { + sessionID: "ses_1", + status: "idle", + }, + } as unknown as Event, + ], resolveStreamFinished), + onEvent: (_directory, payload) => { + resolveDelivered(payload) + }, + transport: "sse", + heartbeatTimeoutMs: 1_000, + }) + + try { + await streamFinished + const delivered = await Promise.race([deliveredEvent, failAfter(500)]) + expect(delivered.type).toBe("session.status") + expect(delivered.properties).toEqual({ + sessionID: "ses_1", + status: { type: "idle" }, + }) + } finally { + pipeline.cleanup() + } + }) }) diff --git a/packages/ui/src/sync/event-pipeline.ts b/packages/ui/src/sync/event-pipeline.ts index cef03ece..c1a681bb 100644 --- a/packages/ui/src/sync/event-pipeline.ts +++ b/packages/ui/src/sync/event-pipeline.ts @@ -12,7 +12,7 @@ * Abort controller created once at init, cleaned up via returned cleanup fn. */ -import type { Event, OpencodeClient } from "@opencode-ai/sdk/v2/client" +import type { Event, OpencodeClient, SessionStatus } from "@opencode-ai/sdk/v2/client" import { opencodeClient } from "@/lib/opencode/client" import { syncDebug } from "./debug" @@ -59,6 +59,11 @@ export type EventPipelineInput = { wsReadyTimeoutMs?: number } +export type EventPipeline = { + cleanup: () => void + reconnect: (reason?: string) => void +} + type MessageStreamWsFrame = { type: "ready" | "event" | "error" | "backpressure" payload?: unknown @@ -68,7 +73,70 @@ type MessageStreamWsFrame = { scope?: "global" | "directory" } +const normalizeOpenChamberSessionStatus = (payload: Event): Event | null => { + const record = payload as unknown as { + id?: unknown + type?: unknown + properties?: { + sessionID?: unknown + sessionId?: unknown + status?: unknown + metadata?: { + attempt?: unknown + message?: unknown + next?: unknown + } + } + } + + if (record.type !== "openchamber:session-status") return null + + const sessionID = typeof record.properties?.sessionID === "string" && record.properties.sessionID.length > 0 + ? record.properties.sessionID + : typeof record.properties?.sessionId === "string" && record.properties.sessionId.length > 0 + ? record.properties.sessionId + : "" + const rawStatus = typeof record.properties?.status === "string" ? record.properties.status : "" + if (!sessionID || !rawStatus) return null + + let status: SessionStatus | null = null + if (rawStatus === "idle" || rawStatus === "busy") { + status = { type: rawStatus } + } else if (rawStatus === "retry") { + const metadata = record.properties?.metadata + if ( + typeof metadata?.attempt === "number" + && typeof metadata.message === "string" + && typeof metadata.next === "number" + ) { + status = { + type: "retry", + attempt: metadata.attempt, + message: metadata.message, + next: metadata.next, + } + } + } + if (!status) return null + + return { + id: typeof record.id === "string" && record.id.length > 0 + ? record.id + : `openchamber-status-${sessionID}-${Date.now()}`, + type: "session.status", + properties: { + sessionID, + status, + }, + } as Event +} + const normalizeEventType = (payload: Event): Event => { + const normalizedOpenChamberStatus = normalizeOpenChamberSessionStatus(payload) + if (normalizedOpenChamberStatus) { + return normalizedOpenChamberStatus + } + const type = (payload as { type?: unknown }).type if (typeof type !== "string") { return payload @@ -173,13 +241,10 @@ type DirectoryQueue = { type AttemptAbortReason = | "pipeline_stopped" - | "ws_heartbeat_timeout" - | "sse_heartbeat_timeout" - | "ws_system_resume" - | "sse_system_resume" + | `${"ws" | "sse"}_${string}` | null -export function createEventPipeline(input: EventPipelineInput) { +export function createEventPipeline(input: EventPipelineInput): EventPipeline { const { sdk, onEvent, @@ -669,10 +734,8 @@ export function createEventPipeline(input: EventPipelineInput) { if (currentTransport === "ws" && code === "WS_FALLBACK") { retryDelayMs = 0 // Transport switch (WS → SSE fallback), not a real disconnection. - // No events were lost — the next attempt will use SSE and carry - // lastEventId for gapless replay. Notify consumer so it can set - // isConnected, but do NOT treat this as a disconnection requiring - // a full directory resync. + // The consumer still gets a hook so it can resync authoritative + // state; real networks can lose/buffer events around transport flips. onTransportSwitch?.() } else if (!isAbortError(error)) { consecutiveFailures += 1 @@ -772,6 +835,11 @@ export function createEventPipeline(input: EventPipelineInput) { attempt?.abort() } + const reconnect = (reason = "manual") => { + attemptAbortReason = `${activeTransport}_${reason}` + attempt?.abort() + } + if (typeof document !== "undefined") { document.addEventListener("visibilitychange", onVisibility) window.addEventListener("pageshow", onPageShow) @@ -799,5 +867,5 @@ export function createEventPipeline(input: EventPipelineInput) { flushAll() } - return { cleanup } + return { cleanup, reconnect } } diff --git a/packages/ui/src/sync/sync-context.tsx b/packages/ui/src/sync/sync-context.tsx index ebf77aab..dc9ce48b 100644 --- a/packages/ui/src/sync/sync-context.tsx +++ b/packages/ui/src/sync/sync-context.tsx @@ -127,6 +127,10 @@ const BOOT_DEBOUNCE_MS = 1500 const RECONNECT_MESSAGE_LIMIT = 30 const SESSION_MATERIALIZATION_MESSAGE_LIMIT = 30 const RECONNECT_SKIP_PARTS = new Set(["patch", "step-start", "step-finish"]) +const ACTIVE_SESSION_WATCHDOG_INTERVAL_MS = 5_000 +const ACTIVE_SESSION_STATUS_POLL_INTERVAL_MS = 5_000 +const ACTIVE_SESSION_STALE_EVENT_MS = 20_000 +const ACTIVE_SESSION_FULL_RESYNC_COOLDOWN_MS = 15_000 const requestSignature = (items: Array<{ id: string }> | undefined): string => { if (!items || items.length === 0) return "" return items @@ -312,6 +316,87 @@ function toSessionStatus(status: Awaited>, + candidateSessionIds: string[], +): Record | null { + if (nextStatuses === null) return null + const relevantStatuses: Record = {} + for (const sessionId of candidateSessionIds) { + relevantStatuses[sessionId] = toSessionStatus(nextStatuses[sessionId]) ?? { type: "idle" } + } + return relevantStatuses +} + +function applySessionStatusSnapshot( + store: StoreApi, + relevantStatuses: Record, +): boolean { + if (Object.keys(relevantStatuses).length === 0) return false + + let changed = false + store.setState((state: DirectoryStore) => { + for (const [sessionId, nextStatus] of Object.entries(relevantStatuses)) { + if (!haveEquivalentSyncSnapshots(state.session_status?.[sessionId], nextStatus)) { + changed = true + break + } + } + + if (!changed) { + return state + } + + return { + session_status: { ...state.session_status, ...relevantStatuses }, + } + }) + + return changed +} + +async function resyncDirectorySessionStatuses( + directory: string, + store: StoreApi, + candidateSessionIds: string[], +): Promise | null> { + const nextStatuses = await opencodeClient.getSessionStatusForDirectory(directory) + // null = fetch failed; preserve existing state. {} or populated = authoritative + // snapshot of active sessions — candidates not listed are idle now. + const relevantStatuses = buildRelevantSessionStatuses(nextStatuses, candidateSessionIds) + if (relevantStatuses === null) return null + applySessionStatusSnapshot(store, relevantStatuses) + return relevantStatuses +} + +function needsSnapshotAfterStatusPoll( + state: DirectoryStore, + sessionId: string, + nextStatus: SessionStatus | undefined, +): boolean { + if (nextStatus?.type !== "idle") return false + const currentStatus = state.session_status?.[sessionId] + if (currentStatus && currentStatus.type !== "idle") return true + + const messages = state.message[sessionId] + const lastMessage = messages?.[messages.length - 1] + return !!lastMessage + && lastMessage.role === "assistant" + && typeof (lastMessage as { time?: { completed?: number } }).time?.completed !== "number" +} + type EventRoutingIndex = { sessionDirectoryById: Map messageSessionById: Map @@ -901,41 +986,10 @@ async function resyncDirectoryAfterReconnect( routingIndex: EventRoutingIndex, ) { const current = store.getState() - const candidateSessionIds = getReconnectCandidateSessionIds(current, { - directory, - viewedSession: getViewedSessionMaterializationTarget(directory), - }) + const candidateSessionIds = getActiveSessionCandidateIds(directory, current) if (candidateSessionIds.length === 0) return - const nextStatuses = await opencodeClient.getSessionStatusForDirectory(directory) - // null = fetch failed; preserve existing state. {} or populated = authoritative - // snapshot of active sessions — candidates not listed are idle now. - if (nextStatuses !== null) { - const relevantStatuses: Record = {} - for (const sessionId of candidateSessionIds) { - relevantStatuses[sessionId] = toSessionStatus(nextStatuses[sessionId]) ?? { type: "idle" } - } - - if (Object.keys(relevantStatuses).length > 0) { - store.setState((state: DirectoryStore) => { - let changed = false - for (const [sessionId, nextStatus] of Object.entries(relevantStatuses)) { - if (!haveEquivalentSyncSnapshots(state.session_status?.[sessionId], nextStatus)) { - changed = true - break - } - } - - if (!changed) { - return state - } - - return { - session_status: { ...state.session_status, ...relevantStatuses }, - } - }) - } - } + await resyncDirectorySessionStatuses(directory, store, candidateSessionIds) const scopedClient = opencodeClient.getScopedSdkClient(directory) await Promise.all(candidateSessionIds.map(async (sessionId) => { @@ -1314,6 +1368,12 @@ export function SyncProvider(props: { const routingIndexRef = useRef(null) if (!routingIndexRef.current) routingIndexRef.current = createEventRoutingIndex() const routingIndex = routingIndexRef.current + const lastActiveEventAtByDirectoryRef = useRef(new Map()) + const lastStatusPollAtByDirectoryRef = useRef(new Map()) + const lastFullResyncAtByDirectoryRef = useRef(new Map()) + const resyncingDirectoriesRef = useRef(new Set()) + const statusPollingDirectoriesRef = useRef(new Set()) + const pipelineReconnectRef = useRef<((reason?: string) => void) | null>(null) const system = useMemo( () => ({ @@ -1324,6 +1384,23 @@ export function SyncProvider(props: { [childStores, props.sdk, props.directory], ) + const triggerDirectoryResync = useCallback((directory: string) => { + const store = childStores.children.get(directory) + if (!store) return + const resyncing = resyncingDirectoriesRef.current + if (resyncing.has(directory)) return + + lastFullResyncAtByDirectoryRef.current.set(directory, Date.now()) + resyncing.add(directory) + void resyncDirectoryAfterReconnect(directory, store, routingIndex) + .catch(() => { + // Transient failure — the watchdog, next SSE event, or reconnect will catch up. + }) + .finally(() => { + resyncing.delete(directory) + }) + }, [childStores, routingIndex]) + // Configure child store manager useEffect(() => { const bootingDirs = new Set() @@ -1435,30 +1512,16 @@ export function SyncProvider(props: { // Event pipeline — created once per mount. No class, no start/stop. // Abort controller owned by the pipeline closure. Cleanup aborts + flushes. useEffect(() => { - const reconnectMaterializing = new Set() - const triggerReconnectMaterialization = (directory: string) => { - const store = childStores.children.get(directory) - if (!store) return - if (reconnectMaterializing.has(directory)) return - - reconnectMaterializing.add(directory) - void resyncDirectoryAfterReconnect(directory, store, routingIndex) - .catch(() => { - // Transient failure during materialization — next SSE event, transport switch, - // or reconnect will catch up. - }) - .finally(() => { - reconnectMaterializing.delete(directory) - }) - } - - const { cleanup } = createEventPipeline({ + const pipeline = createEventPipeline({ sdk: props.sdk, transport: messageStreamTransport, routeDirectory: (directory, payload) => { return resolveDirectoryFromRoutingIndex(routingIndex, directory, payload, childStores) }, onEvent: (directory, payload) => { + if (!isStreamHeartbeatEvent(payload)) { + lastActiveEventAtByDirectoryRef.current.set(directory, Date.now()) + } dispatchVSCodeRuntimeNotificationEvent(directory, payload) if (payload.type === "installation.update-available") { const version = typeof (payload.properties as { version?: unknown })?.version === "string" @@ -1477,7 +1540,7 @@ export function SyncProvider(props: { connectionPhase: "connected", }) for (const dir of childStores.children.keys()) { - triggerReconnectMaterialization(dir) + triggerDirectoryResync(dir) } }, onDisconnect: (reason) => { @@ -1489,21 +1552,108 @@ export function SyncProvider(props: { }) }, onTransportSwitch: () => { - // Transport switched (e.g. WS timeout → SSE fallback) without a full - // disconnect. If the active session missed the transition into a busy - // turn, force a targeted resync for the viewed directory. + // Transport changes are gap-prone in real networks. Treat them like a + // reconnect and refresh active session snapshots from HTTP. useConfigStore.setState({ isConnected: true, hasEverConnected: true, connectionPhase: "connected", }) - if (_activeDirectory) { - triggerReconnectMaterialization(_activeDirectory) + for (const dir of childStores.children.keys()) { + triggerDirectoryResync(dir) } }, }) - return cleanup - }, [props.sdk, childStores, routingIndex, messageStreamTransport]) + pipelineReconnectRef.current = pipeline.reconnect + return () => { + if (pipelineReconnectRef.current === pipeline.reconnect) { + pipelineReconnectRef.current = null + } + pipeline.cleanup() + } + }, [props.sdk, childStores, routingIndex, messageStreamTransport, triggerDirectoryResync]) + + useEffect(() => { + let stopped = false + let running = false + + const pollDirectoryStatuses = async ( + directory: string, + store: StoreApi, + candidateSessionIds: string[], + ) => { + const polling = statusPollingDirectoriesRef.current + if (polling.has(directory)) return + polling.add(directory) + try { + const before = store.getState() + const statuses = await resyncDirectorySessionStatuses(directory, store, candidateSessionIds) + if (!statuses) return + const needsSnapshot = candidateSessionIds.some((sessionId) => ( + needsSnapshotAfterStatusPoll(before, sessionId, statuses[sessionId]) + )) + if (needsSnapshot) { + triggerDirectoryResync(directory) + } + } finally { + polling.delete(directory) + } + } + + const tick = () => { + if (running || stopped) return + running = true + void Promise.resolve() + .then(() => { + if (stopped) return + const now = Date.now() + for (const [directory, store] of childStores.children.entries()) { + const state = store.getState() + const candidateSessionIds = getActiveSessionCandidateIds(directory, state) + if (candidateSessionIds.length === 0) { + lastActiveEventAtByDirectoryRef.current.delete(directory) + lastStatusPollAtByDirectoryRef.current.delete(directory) + lastFullResyncAtByDirectoryRef.current.delete(directory) + continue + } + + if (!lastActiveEventAtByDirectoryRef.current.has(directory)) { + lastActiveEventAtByDirectoryRef.current.set(directory, now) + } + + const lastStatusPollAt = lastStatusPollAtByDirectoryRef.current.get(directory) ?? 0 + if (now - lastStatusPollAt >= ACTIVE_SESSION_STATUS_POLL_INTERVAL_MS) { + lastStatusPollAtByDirectoryRef.current.set(directory, now) + void pollDirectoryStatuses(directory, store, candidateSessionIds).catch(() => undefined) + } + + const lastActiveEventAt = lastActiveEventAtByDirectoryRef.current.get(directory) ?? now + const lastFullResyncAt = lastFullResyncAtByDirectoryRef.current.get(directory) ?? 0 + if ( + now - lastActiveEventAt >= ACTIVE_SESSION_STALE_EVENT_MS + && now - lastFullResyncAt >= ACTIVE_SESSION_FULL_RESYNC_COOLDOWN_MS + ) { + pipelineReconnectRef.current?.("active_stream_stale") + triggerDirectoryResync(directory) + } + } + }) + .finally(() => { + running = false + if (stopped) { + statusPollingDirectoriesRef.current.clear() + } + }) + } + + const interval = setInterval(tick, ACTIVE_SESSION_WATCHDOG_INTERVAL_MS) + tick() + + return () => { + stopped = true + clearInterval(interval) + } + }, [childStores, triggerDirectoryResync]) // Ensure current directory's child store exists useEffect(() => { diff --git a/packages/web/server/index.js b/packages/web/server/index.js index b6f9b6cf..d1fbfc24 100644 --- a/packages/web/server/index.js +++ b/packages/web/server/index.js @@ -726,7 +726,7 @@ const processForwardedEventPayload = (payload, emitSyntheticEvent) => { emitSyntheticEvent({ type: 'openchamber:session-status', properties: { - sessionId, + sessionID: sessionId, status, timestamp: Date.now(), metadata: { diff --git a/packages/web/server/lib/opencode/session-runtime.js b/packages/web/server/lib/opencode/session-runtime.js index 4aea8cf6..680d51ea 100644 --- a/packages/web/server/lib/opencode/session-runtime.js +++ b/packages/web/server/lib/opencode/session-runtime.js @@ -158,7 +158,7 @@ export const createSessionRuntime = ({ writeSseEvent, getNotificationClients, br const syntheticPayload = { type: 'openchamber:session-status', properties: { - sessionId, + sessionID: sessionId, status: state.status, timestamp: state.lastUpdateAt, metadata: state.metadata, @@ -216,7 +216,7 @@ export const createSessionRuntime = ({ writeSseEvent, getNotificationClients, br const syntheticPayload = { type: 'openchamber:session-status', properties: { - sessionId, + sessionID: sessionId, status: state.status, timestamp: Date.now(), metadata: {}, diff --git a/packages/web/server/lib/opencode/session-runtime.test.js b/packages/web/server/lib/opencode/session-runtime.test.js index 1d4807d2..11fbcbf7 100644 --- a/packages/web/server/lib/opencode/session-runtime.test.js +++ b/packages/web/server/lib/opencode/session-runtime.test.js @@ -49,7 +49,7 @@ describe('session runtime', () => { expect(events).toContainEqual({ type: 'openchamber:session-status', properties: expect.objectContaining({ - sessionId: 'session-1', + sessionID: 'session-1', status: 'idle', needsAttention: true, }), @@ -57,7 +57,7 @@ describe('session runtime', () => { expect(events.at(-1)).toEqual({ type: 'openchamber:session-status', properties: { - sessionId: 'session-1', + sessionID: 'session-1', status: 'idle', timestamp: expect.any(Number), metadata: {}, @@ -92,7 +92,7 @@ describe('session runtime', () => { expect(events).toContainEqual({ type: 'openchamber:session-status', properties: expect.objectContaining({ - sessionId: 'legacy-session-1', + sessionID: 'legacy-session-1', status: 'busy', }), });