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
This commit is contained in:
Bohdan Triapitsyn
2026-05-24 19:17:04 +03:00
parent c4956af565
commit 0c8301d579
6 changed files with 333 additions and 76 deletions
@@ -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<void>((resolve) => {
resolveStreamFinished = resolve
})
let resolveDelivered!: (event: Event) => void
const deliveredEvent = new Promise<Event>((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()
}
})
})
+79 -11
View File
@@ -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 }
}
+209 -59
View File
@@ -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<ReturnType<typeof opencodeClient.getSes
return undefined
}
function isStreamHeartbeatEvent(payload: Event): boolean {
const type = (payload as { type?: unknown }).type
return type === "server.heartbeat" || type === "openchamber:heartbeat"
}
function getActiveSessionCandidateIds(directory: string, state: DirectoryStore): string[] {
return getReconnectCandidateSessionIds(state, {
directory,
viewedSession: getViewedSessionMaterializationTarget(directory),
})
}
function buildRelevantSessionStatuses(
nextStatuses: Awaited<ReturnType<typeof opencodeClient.getSessionStatusForDirectory>>,
candidateSessionIds: string[],
): Record<string, SessionStatus> | null {
if (nextStatuses === null) return null
const relevantStatuses: Record<string, SessionStatus> = {}
for (const sessionId of candidateSessionIds) {
relevantStatuses[sessionId] = toSessionStatus(nextStatuses[sessionId]) ?? { type: "idle" }
}
return relevantStatuses
}
function applySessionStatusSnapshot(
store: StoreApi<DirectoryStore>,
relevantStatuses: Record<string, SessionStatus>,
): 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<DirectoryStore>,
candidateSessionIds: string[],
): Promise<Record<string, SessionStatus> | 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<string, string>
messageSessionById: Map<string, string>
@@ -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<string, SessionStatus> = {}
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<EventRoutingIndex | null>(null)
if (!routingIndexRef.current) routingIndexRef.current = createEventRoutingIndex()
const routingIndex = routingIndexRef.current
const lastActiveEventAtByDirectoryRef = useRef(new Map<string, number>())
const lastStatusPollAtByDirectoryRef = useRef(new Map<string, number>())
const lastFullResyncAtByDirectoryRef = useRef(new Map<string, number>())
const resyncingDirectoriesRef = useRef(new Set<string>())
const statusPollingDirectoriesRef = useRef(new Set<string>())
const pipelineReconnectRef = useRef<((reason?: string) => void) | null>(null)
const system = useMemo<SyncSystem>(
() => ({
@@ -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<string>()
@@ -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<string>()
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<DirectoryStore>,
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(() => {
+1 -1
View File
@@ -726,7 +726,7 @@ const processForwardedEventPayload = (payload, emitSyntheticEvent) => {
emitSyntheticEvent({
type: 'openchamber:session-status',
properties: {
sessionId,
sessionID: sessionId,
status,
timestamp: Date.now(),
metadata: {
@@ -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: {},
@@ -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',
}),
});