From 5c7f5aa4f7c24ec8aee390b6a774c560243e6804 Mon Sep 17 00:00:00 2001 From: jwcrystal <121911854+jwcrystal@users.noreply.github.com> Date: Fri, 24 Apr 2026 17:39:58 +0800 Subject: [PATCH] fix(sync): resync the active session after reconnect transitions (#1011) * fix(sync): resync the active session after reconnect transitions Include the viewed session in reconnect recovery and trigger a targeted resync on transport switches so the active chat catches up after missed live events. Constraint: Keep the fix in the sync layer instead of adding ChatContainer-only recovery Rejected: Widen reconnect heuristics for every cached session | broader recovery scope than needed Confidence: high Scope-risk: narrow * Avoid no-op reconnect resync store writes --------- Co-authored-by: Bohdan Triapitsyn --- .../ui/src/sync/reconnect-recovery.test.ts | 69 ++++++++ packages/ui/src/sync/reconnect-recovery.ts | 62 +++++++ packages/ui/src/sync/sync-context.tsx | 163 +++++++++++------- 3 files changed, 235 insertions(+), 59 deletions(-) create mode 100644 packages/ui/src/sync/reconnect-recovery.test.ts create mode 100644 packages/ui/src/sync/reconnect-recovery.ts diff --git a/packages/ui/src/sync/reconnect-recovery.test.ts b/packages/ui/src/sync/reconnect-recovery.test.ts new file mode 100644 index 00000000..e4e8e016 --- /dev/null +++ b/packages/ui/src/sync/reconnect-recovery.test.ts @@ -0,0 +1,69 @@ +import { describe, expect, test } from "bun:test" +import type { Message, SessionStatus } from "@opencode-ai/sdk/v2/client" +import type { Session } from "@opencode-ai/sdk/v2" +import { getReconnectCandidateSessionIds } from "./reconnect-recovery" + +function createSession(id: string, overrides: Partial = {}): Session { + return { + id, + title: id, + time: { created: 1, updated: 1 }, + version: "1", + ...overrides, + } as Session +} + +function createAssistantMessage(id: string, sessionID: string, completed?: number): Message { + return { + id, + sessionID, + role: "assistant", + time: completed ? { created: 1, updated: 1, completed } : { created: 1, updated: 1 }, + parts: [], + } as unknown as Message +} + +describe("getReconnectCandidateSessionIds", () => { + test("includes non-idle, incomplete assistant, and parent sessions", () => { + const busyStatus = { type: "busy" } as SessionStatus + + expect(getReconnectCandidateSessionIds({ + session: [ + createSession("busy"), + createSession("child", { parentID: "parent" }), + createSession("parent"), + createSession("incomplete"), + ], + session_status: { busy: busyStatus }, + message: { + incomplete: [createAssistantMessage("m-1", "incomplete")], + }, + }).sort()).toEqual(["busy", "incomplete", "parent"]) + }) + + test("includes the currently viewed session even when it looks idle and complete", () => { + expect(getReconnectCandidateSessionIds({ + session: [createSession("active")], + session_status: { active: { type: "idle" } as SessionStatus }, + message: { + active: [createAssistantMessage("m-1", "active", 1)], + }, + }, { + directory: "/repo", + viewedSession: { directory: "/repo", sessionId: "active" }, + }).sort()).toContain("active") + }) + + test("does not include a viewed session from another directory", () => { + expect(getReconnectCandidateSessionIds({ + session: [createSession("active")], + session_status: { active: { type: "idle" } as SessionStatus }, + message: { + active: [createAssistantMessage("m-1", "active", 1)], + }, + }, { + directory: "/repo-a", + viewedSession: { directory: "/repo-b", sessionId: "active" }, + }).sort()).not.toContain("active") + }) +}) diff --git a/packages/ui/src/sync/reconnect-recovery.ts b/packages/ui/src/sync/reconnect-recovery.ts new file mode 100644 index 00000000..7a4a44f0 --- /dev/null +++ b/packages/ui/src/sync/reconnect-recovery.ts @@ -0,0 +1,62 @@ +import type { SessionStatus, Message } from "@opencode-ai/sdk/v2/client" +import type { Session } from "@opencode-ai/sdk/v2" + +type ReconnectRecoveryState = { + session: Session[] + session_status?: Record + message?: Record +} + +export type ViewedSessionRecoveryTarget = { + directory: string + sessionId: string +} + +type ReconnectCandidateOptions = { + directory?: string + viewedSession?: ViewedSessionRecoveryTarget | null +} + +export function getReconnectCandidateSessionIds(state: ReconnectRecoveryState, options?: ReconnectCandidateOptions) { + const ids = new Set() + + for (const [sessionId, status] of Object.entries(state.session_status ?? {})) { + if (status && status.type !== "idle") ids.add(sessionId) + } + + for (const [sessionId, messages] of Object.entries(state.message ?? {})) { + const lastMessage = messages[messages.length - 1] + if ( + lastMessage + && lastMessage.role === "assistant" + && typeof (lastMessage as { time?: { completed?: number } }).time?.completed !== "number" + ) { + ids.add(sessionId) + } + } + + const parentIds = new Set() + for (const session of state.session) { + const parentId = (session as Session & { parentID?: string | null }).parentID + if (parentId) { + parentIds.add(parentId) + } + } + for (const pid of parentIds) { + ids.add(pid) + } + + const viewedSession = options?.viewedSession + if (viewedSession?.sessionId && viewedSession.directory === options?.directory) { + const sessionId = viewedSession.sessionId + const sessionExists = state.session.some((session) => session.id === sessionId) + || Object.hasOwn(state.session_status ?? {}, sessionId) + || Object.hasOwn(state.message ?? {}, sessionId) + + if (sessionExists) { + ids.add(sessionId) + } + } + + return Array.from(ids) +} diff --git a/packages/ui/src/sync/sync-context.tsx b/packages/ui/src/sync/sync-context.tsx index 50531cf3..4f1eca59 100644 --- a/packages/ui/src/sync/sync-context.tsx +++ b/packages/ui/src/sync/sync-context.tsx @@ -24,6 +24,7 @@ import { setActionRefs } from "./session-actions" import { setSyncRefs } from "./sync-refs" import { stripMessageDiffSnapshots, stripSessionDiffSnapshots } from "./sanitize" import { syncDebug } from "./debug" +import { getReconnectCandidateSessionIds } from "./reconnect-recovery" import { opencodeClient } from "@/lib/opencode/client" import { usePermissionStore } from "@/stores/permissionStore" import { useConfigStore } from "@/stores/useConfigStore" @@ -133,6 +134,7 @@ const requestSignature = (items: Array<{ id: string }> | undefined): string => { const cmp = (a: string, b: string) => (a < b ? -1 : a > b ? 1 : 0) const partRepairSignature = (part: Part): string => JSON.stringify(part) +const syncSnapshotSignature = (value: unknown): string => JSON.stringify(value) function haveEquivalentPartSnapshots(left: Part[] | undefined, right: Part[]): boolean { if (!left) { @@ -160,6 +162,36 @@ function haveEquivalentPartSnapshots(left: Part[] | undefined, right: Part[]): b return true } +function haveEquivalentMessageSnapshots(left: Message[] | undefined, right: Message[]): boolean { + if (!left) { + return right.length === 0 + } + + if (left.length !== right.length) { + return false + } + + for (let index = 0; index < left.length; index += 1) { + const leftMessage = left[index] + const rightMessage = right[index] + if (!leftMessage || !rightMessage) { + return false + } + if (leftMessage.id !== rightMessage.id) { + return false + } + if (syncSnapshotSignature(leftMessage) !== syncSnapshotSignature(rightMessage)) { + return false + } + } + + return true +} + +function haveEquivalentSyncSnapshots(left: unknown, right: unknown): boolean { + return syncSnapshotSignature(left) === syncSnapshotSignature(right) +} + // --------------------------------------------------------------------------- // Parts-gap recovery — when SSE events arrive but parts are missing, // trigger a targeted re-fetch for the affected sessions. @@ -273,40 +305,13 @@ function isRecentBoot() { return bootingRoot || Date.now() - bootedAt < BOOT_DEBOUNCE_MS } -function getReconnectCandidateSessionIds(state: State) { - const ids = new Set() - - for (const [sessionId, status] of Object.entries(state.session_status ?? {})) { - if (status && status.type !== "idle") ids.add(sessionId) +function getViewedSessionRecoveryTarget(directory: string) { + if (!_activeDirectory || !_activeSession) return null + if (directory !== _activeDirectory) return null + return { + directory: _activeDirectory, + sessionId: _activeSession, } - - for (const [sessionId, messages] of Object.entries(state.message ?? {})) { - const lastMessage = messages[messages.length - 1] - if ( - lastMessage - && lastMessage.role === "assistant" - && typeof (lastMessage as { time?: { completed?: number } }).time?.completed !== "number" - ) { - ids.add(sessionId) - } - } - - // Ensure parent sessions of any child sessions are also resynced. - // A child may have completed during the disconnect gap without the - // store knowing (SSE events lost). The parent's task tool part needs - // to reflect the child's final state. - const parentIds = new Set() - for (const session of state.session) { - const parentId = (session as Session & { parentID?: string | null }).parentID - if (parentId) { - parentIds.add(parentId) - } - } - for (const pid of parentIds) { - ids.add(pid) - } - - return Array.from(ids) } function toSessionStatus(status: Awaited>[string]): SessionStatus | undefined { @@ -716,7 +721,10 @@ async function resyncDirectoryAfterReconnect( routingIndex: EventRoutingIndex, ) { const current = store.getState() - const candidateSessionIds = getReconnectCandidateSessionIds(current) + const candidateSessionIds = getReconnectCandidateSessionIds(current, { + directory, + viewedSession: getViewedSessionRecoveryTarget(directory), + }) if (candidateSessionIds.length === 0) return const nextStatuses = await opencodeClient.getSessionStatusForDirectory(directory) @@ -730,9 +738,23 @@ async function resyncDirectoryAfterReconnect( } if (Object.keys(relevantStatuses).length > 0) { - store.setState((state: DirectoryStore) => ({ - session_status: { ...state.session_status, ...relevantStatuses }, - })) + 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 }, + } + }) } const scopedClient = opencodeClient.getScopedSdkClient(directory) @@ -752,17 +774,19 @@ async function resyncDirectoryAfterReconnect( .sort((a, b) => cmp(a.id, b.id)) store.setState((state: DirectoryStore) => { - const sessions = [...state.session] - const sessionIndex = sessions.findIndex((item) => item.id === nextSession.id) + const sessionIndex = state.session.findIndex((item) => item.id === nextSession.id) + let sessions = state.session let sessionChanged = false let sessionTotal = state.sessionTotal if (sessionIndex >= 0) { - if (sessions[sessionIndex] !== nextSession) { + if (!haveEquivalentSyncSnapshots(sessions[sessionIndex], nextSession)) { + sessions = [...state.session] sessions[sessionIndex] = nextSession sessionChanged = true } } else { + sessions = [...state.session] sessions.push(nextSession) sessions.sort((a, b) => cmp(a.id, b.id)) if (!nextSession.parentID) sessionTotal += 1 @@ -772,19 +796,32 @@ async function resyncDirectoryAfterReconnect( // Merge parts: overwrite only messages present in the fetch snapshot. // Do NOT delete parts for messages that may have been added by SSE // events arriving between the fetch and the setState — those are more recent. - const nextPartState = { ...state.part } + let nextPartState = state.part + let partsChanged = false for (const record of records) { const messageId = record?.info?.id if (!messageId) continue - nextPartState[messageId] = (record.parts ?? []) + const nextParts = (record.parts ?? []) .filter((part) => !!part?.id && !RECONNECT_SKIP_PARTS.has(part.type)) .sort((a, b) => cmp(a.id, b.id)) + if (!haveEquivalentPartSnapshots(state.part[messageId], nextParts)) { + if (!partsChanged) { + nextPartState = { ...state.part } + partsChanged = true + } + nextPartState[messageId] = nextParts + } + } + + const messagesChanged = !haveEquivalentMessageSnapshots(state.message[sessionId], nextMessages) + if (!sessionChanged && !messagesChanged && !partsChanged) { + return state } return { ...(sessionChanged ? { session: sessions, sessionTotal } : {}), - message: { ...state.message, [sessionId]: nextMessages }, - part: nextPartState, + ...(messagesChanged ? { message: { ...state.message, [sessionId]: nextMessages } } : {}), + ...(partsChanged ? { part: nextPartState } : {}), } }) @@ -1365,6 +1402,21 @@ export function SyncProvider(props: { // Abort controller owned by the pipeline closure. Cleanup aborts + flushes. useEffect(() => { const reconnectResyncing = new Set() + const triggerRecoveryResync = (directory: string) => { + const store = childStores.children.get(directory) + if (!store) return + if (reconnectResyncing.has(directory)) return + + reconnectResyncing.add(directory) + void resyncDirectoryAfterReconnect(directory, store, routingIndex) + .catch(() => { + // Transient failure during resync — next SSE event, transport switch, + // or reconnect will catch up. + }) + .finally(() => { + reconnectResyncing.delete(directory) + }) + } const { cleanup } = createEventPipeline({ sdk: props.sdk, @@ -1381,18 +1433,8 @@ export function SyncProvider(props: { hasEverConnected: true, connectionPhase: "connected", }) - for (const [dir, store] of childStores.children) { - if (reconnectResyncing.has(dir)) continue - if (getReconnectCandidateSessionIds(store.getState()).length === 0) continue - - reconnectResyncing.add(dir) - void resyncDirectoryAfterReconnect(dir, store, routingIndex) - .catch(() => { - // Transient failure during resync — next SSE event or reconnect will catch up. - }) - .finally(() => { - reconnectResyncing.delete(dir) - }) + for (const dir of childStores.children.keys()) { + triggerRecoveryResync(dir) } }, onDisconnect: (reason) => { @@ -1404,14 +1446,17 @@ export function SyncProvider(props: { }) }, onTransportSwitch: () => { - // Transport switched (e.g. WS timeout → SSE fallback) without - // actual disconnection. No events lost — just update connection - // state without triggering a full directory resync. + // 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. useConfigStore.setState({ isConnected: true, hasEverConnected: true, connectionPhase: "connected", }) + if (_activeDirectory) { + triggerRecoveryResync(_activeDirectory) + } }, }) return cleanup