/* eslint-disable react-refresh/only-export-components */ import React, { createContext, useContext, useEffect, useRef, useCallback, useMemo } from "react" import type { Event, Message, Part } from "@opencode-ai/sdk/v2/client" import type { Session } from "@opencode-ai/sdk/v2" import type { StoreApi } from "zustand" import { useStore } from "zustand" import type { OpencodeClient } from "@opencode-ai/sdk/v2/client" import { createEventPipeline } from "./event-pipeline" import { reduceGlobalEvent, applyGlobalProject, applyDirectoryEvent } from "./event-reducer" import { useGlobalSyncStore, type GlobalSyncStore } from "./global-sync-store" import { ChildStoreManager, type DirectoryStore } from "./child-store" import { aggregateLiveSessions, aggregateLiveSessionStatuses, areSessionListsEquivalent, areStatusMapsEquivalent, findLiveSession, findLiveSessionStatus, } from "./live-aggregate" import { bootstrapGlobal, bootstrapDirectory } from "./bootstrap" import { retry } from "./retry" import { updateStreamingState } from "./streaming" 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" import { useTodosPersistStore } from "@/stores/useTodosPersistStore" import { toast } from "@/components/ui" import { appendNotification } from "./notification-store" import type { State } from "./types" import type { SessionStatus } from "@opencode-ai/sdk/v2/client" import type { PermissionRequest } from "@/types/permission" import type { QuestionRequest } from "@/types/question" import * as sessionActions from "./session-actions" import { getSessionMaterializationStatus, materializeSessionSnapshots } from "./materialization" // --------------------------------------------------------------------------- // Context // --------------------------------------------------------------------------- type SyncSystem = { childStores: ChildStoreManager sdk: OpencodeClient directory: string } const SYNC_CONTEXT_GLOBAL_KEY = "__openchamber_sync_context__" type SyncGlobal = typeof globalThis & { [SYNC_CONTEXT_GLOBAL_KEY]?: React.Context } const syncGlobal = globalThis as SyncGlobal const SyncContext = syncGlobal[SYNC_CONTEXT_GLOBAL_KEY] ?? createContext(null) syncGlobal[SYNC_CONTEXT_GLOBAL_KEY] = SyncContext function useSyncSystem() { const ctx = useContext(SyncContext) if (!ctx) throw new Error("useSyncSystem must be used within ") return ctx } function getLiveStates(childStores: ChildStoreManager): State[] { return Array.from(childStores.children.values(), (store) => store.getState()) } function useLiveSyncSelector(selector: (states: State[]) => T, isEqual: (left: T, right: T) => boolean = Object.is): T { const { childStores } = useSyncSystem() const cacheRef = useRef(undefined) const initializedRef = useRef(false) const getSnapshot = useCallback(() => { const next = selector(getLiveStates(childStores)) if (initializedRef.current && isEqual(cacheRef.current as T, next)) { return cacheRef.current as T } cacheRef.current = next initializedRef.current = true return next }, [childStores, isEqual, selector]) return React.useSyncExternalStore( useCallback((notify) => childStores.subscribeAll(notify), [childStores]), getSnapshot, getSnapshot, ) } // --------------------------------------------------------------------------- // Event handler — applies one SSE event at a time to the live store. // Each event reads live state, creates a shallow draft, applies, writes back. // React 18 batches synchronous setState calls automatically. // --------------------------------------------------------------------------- /** Read status for a session across all directories */ export function useGlobalSessionStatus(sessionId: string): SessionStatus | undefined { return useLiveSyncSelector( useCallback((states) => findLiveSessionStatus(states, sessionId), [sessionId]), ) } /** Read all session statuses (for sidebar) */ export function useAllSessionStatuses(): Record { return useLiveSyncSelector( useCallback((states) => aggregateLiveSessionStatuses(states), []), areStatusMapsEquivalent, ) } export function useAllLiveSessions(): Session[] { return useLiveSyncSelector( useCallback((states) => aggregateLiveSessions(states), []), areSessionListsEquivalent, ) } // Boot debounce — suppresses redundant refresh/re-bootstrap events during startup. let bootingRoot = false let bootedAt = 0 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 requestSignature = (items: Array<{ id: string }> | undefined): string => { if (!items || items.length === 0) return "" return items .map((item) => item.id) .sort(cmp) .join("|") } const cmp = (a: string, b: string) => (a < b ? -1 : a > b ? 1 : 0) const syncSnapshotSignature = (value: unknown): string => JSON.stringify(value) function haveEquivalentSyncSnapshots(left: unknown, right: unknown): boolean { return syncSnapshotSignature(left) === syncSnapshotSignature(right) } // --------------------------------------------------------------------------- // Session materialization scheduler — when local message/part state is incomplete, // fetch the canonical session snapshot and materialize messages and parts together. // Tracked per-directory, deduplicated, and auto-expiring. // --------------------------------------------------------------------------- type PendingSessionMaterialization = { sessionID: string directory: string enqueuedAt: number } const SESSION_MATERIALIZATION_COOLDOWN_MS = 5_000 const pendingSessionMaterializations = new Map() // key: directory:sessionID const materializationKey = (directory: string, sessionID: string) => `${directory}:${sessionID}` function enqueueSessionMaterialization(directory: string, sessionID: string, childStores: ChildStoreManager) { if (!directory || directory === "global" || !sessionID) return const k = materializationKey(directory, sessionID) const existing = pendingSessionMaterializations.get(k) if (existing && Date.now() - existing.enqueuedAt < SESSION_MATERIALIZATION_COOLDOWN_MS) return pendingSessionMaterializations.set(k, { sessionID, directory, enqueuedAt: Date.now() }) // Defer to next microtask so we don't hold up the current event batch void Promise.resolve().then(async () => { const store = childStores.getChild(directory) if (!store) { pendingSessionMaterializations.delete(k) return } try { await materializeSessionFromServer(directory, sessionID, store) } catch { // Transient failure — next SSE event or reconnect will catch up. } finally { pendingSessionMaterializations.delete(k) } }) } async function materializeSessionFromServer( directory: string, sessionID: string, store: StoreApi, ) { const scopedClient = opencodeClient.getScopedSdkClient(directory) const result = await retry(() => scopedClient.session.messages({ sessionID, limit: SESSION_MATERIALIZATION_MESSAGE_LIMIT }), ) const records = (result.data ?? []).filter((record: { info?: { id?: string } }) => !!record?.info?.id) if (records.length === 0) return store.setState((state: DirectoryStore) => { const materialized = materializeSessionSnapshots( state, sessionID, records.map((record: { info: Message; parts?: Part[] }) => ({ info: stripMessageDiffSnapshots(record.info), parts: record.parts ?? [], })), { skipPartTypes: RECONNECT_SKIP_PARTS }, ) return { message: materialized.message, part: materialized.part } }) } // Module-level refs for notification viewed check. // Used to determine if user is currently viewing the session when a notification arrives. let _activeDirectory = "" let _activeSession = "" const externallyViewedSessions = new Map() const EXTERNAL_VIEW_TTL_MS = 15_000 const viewedSessionKey = (directory: string, sessionId: string) => `${directory}\n${sessionId}` function pruneExternallyViewedSessions(now = Date.now()) { for (const [key, expiresAt] of externallyViewedSessions.entries()) { if (expiresAt <= now) { externallyViewedSessions.delete(key) } } } const pendingQuestionToastIds = new Set() const pendingPermissionToastIds = new Set() const getQuestionToastKey = (sessionID?: string, requestID?: string) => { if (!sessionID || !requestID) return null return `${sessionID}:${requestID}` } const getPermissionToastKey = (sessionID?: string, requestID?: string) => { if (!sessionID || !requestID) return null return `${sessionID}:${requestID}` } const openSessionFromToast = (sessionID: string, directory: string) => { void import("./session-ui-store") .then(({ useSessionUIStore }) => { useSessionUIStore.getState().setCurrentSession(sessionID, directory) }) .catch(() => undefined) } export function setActiveSession(directory: string, sessionId: string) { _activeDirectory = directory _activeSession = sessionId } export function setExternallyViewedSession(directory: string, sessionId: string, viewed: boolean) { if (!directory || !sessionId) return const key = viewedSessionKey(directory, sessionId) if (!viewed) { externallyViewedSessions.delete(key) return } externallyViewedSessions.set(key, Date.now() + EXTERNAL_VIEW_TTL_MS) } function isViewedInCurrentSession(directory: string, sessionId?: string): boolean { if (!sessionId) return false if (_activeDirectory && _activeSession && directory === _activeDirectory && sessionId === _activeSession) return true pruneExternallyViewedSessions() return externallyViewedSessions.has(viewedSessionKey(directory, sessionId)) } function isRecentBoot() { return bootingRoot || Date.now() - bootedAt < BOOT_DEBOUNCE_MS } function getViewedSessionMaterializationTarget(directory: string) { if (!_activeDirectory || !_activeSession) return null if (directory !== _activeDirectory) return null return { directory: _activeDirectory, sessionId: _activeSession, } } function toSessionStatus(status: Awaited>[string]): SessionStatus | undefined { if (!status) return undefined if (status.type === "idle" || status.type === "busy") { return { type: status.type } } if ( status.type === "retry" && typeof status.attempt === "number" && typeof status.message === "string" && typeof status.next === "number" ) { return { type: "retry", attempt: status.attempt, message: status.message, next: status.next, } } return undefined } type EventRoutingIndex = { sessionDirectoryById: Map messageSessionById: Map sessionMessageIdsById: Map> } const createEventRoutingIndex = (): EventRoutingIndex => ({ sessionDirectoryById: new Map(), messageSessionById: new Map(), sessionMessageIdsById: new Map(), }) const normalizeEventDirectory = (rawDirectory: string): string => { if (!rawDirectory || rawDirectory === "global") { return rawDirectory } const normalized = rawDirectory.replace(/\\/g, "/").replace(/^([a-z]):/, (_, l: string) => l.toUpperCase() + ":") // Strip trailing slashes to match child store keys (normalizeDirectoryPath in useDirectoryStore) return normalized.length > 1 ? normalized.replace(/\/+$/, "") : normalized } const getSessionIdFromPayload = (event: Event): string | null => { const properties = (event as { properties?: unknown }).properties if (!properties || typeof properties !== "object") { return null } const props = properties as Record if (event.type === "message.updated") { const info = props.info if (!info || typeof info !== "object") { return null } const sessionID = (info as { sessionID?: unknown }).sessionID return typeof sessionID === "string" && sessionID.length > 0 ? sessionID : null } if ( event.type === "message.removed" || event.type === "session.status" || event.type === "todo.updated" || event.type === "permission.asked" || event.type === "permission.replied" || event.type === "question.asked" || event.type === "question.replied" || event.type === "question.rejected" || event.type === "session.deleted" ) { const sessionID = props.sessionID return typeof sessionID === "string" && sessionID.length > 0 ? sessionID : null } if (event.type === "message.part.updated") { const part = props.part if (!part || typeof part !== "object") { return null } const sessionID = (part as { sessionID?: unknown }).sessionID return typeof sessionID === "string" && sessionID.length > 0 ? sessionID : null } if (event.type === "session.created" || event.type === "session.updated") { const info = props.info if (!info || typeof info !== "object") { return null } const id = (info as { id?: unknown }).id return typeof id === "string" && id.length > 0 ? id : null } return null } const getMessageIdFromPayload = (event: Event): string | null => { const properties = (event as { properties?: unknown }).properties if (!properties || typeof properties !== "object") { return null } const props = properties as Record if (event.type === "message.updated") { const info = props.info if (!info || typeof info !== "object") { return null } const id = (info as { id?: unknown }).id return typeof id === "string" && id.length > 0 ? id : null } if (event.type === "message.removed" || event.type === "message.part.delta" || event.type === "message.part.removed") { const messageID = props.messageID return typeof messageID === "string" && messageID.length > 0 ? messageID : null } if (event.type === "message.part.updated") { const part = props.part if (!part || typeof part !== "object") { return null } const messageID = (part as { messageID?: unknown }).messageID return typeof messageID === "string" && messageID.length > 0 ? messageID : null } return null } const setIndexedSessionDirectory = (routingIndex: EventRoutingIndex, sessionID: string, directory: string) => { if (!sessionID || !directory || directory === "global") { return } routingIndex.sessionDirectoryById.set(sessionID, directory) } const setIndexedSessionMessages = ( routingIndex: EventRoutingIndex, sessionID: string, directory: string, messages: Message[], ) => { if (!sessionID) { return } setIndexedSessionDirectory(routingIndex, sessionID, directory) const previous = routingIndex.sessionMessageIdsById.get(sessionID) const next = new Set() for (const message of messages) { if (!message?.id) { continue } next.add(message.id) routingIndex.messageSessionById.set(message.id, sessionID) } if (previous) { for (const previousMessageID of previous) { if (!next.has(previousMessageID)) { routingIndex.messageSessionById.delete(previousMessageID) } } } routingIndex.sessionMessageIdsById.set(sessionID, next) } const setIndexedMessage = ( routingIndex: EventRoutingIndex, sessionID: string, messageID: string, directory: string, ) => { if (!sessionID || !messageID) { return } setIndexedSessionDirectory(routingIndex, sessionID, directory) routingIndex.messageSessionById.set(messageID, sessionID) const existing = routingIndex.sessionMessageIdsById.get(sessionID) if (existing) { existing.add(messageID) } else { routingIndex.sessionMessageIdsById.set(sessionID, new Set([messageID])) } } const removeIndexedMessage = ( routingIndex: EventRoutingIndex, messageID: string, sessionHint?: string | null, ) => { if (!messageID) { return } const sessionID = sessionHint ?? routingIndex.messageSessionById.get(messageID) routingIndex.messageSessionById.delete(messageID) if (!sessionID) { return } const messageIds = routingIndex.sessionMessageIdsById.get(sessionID) if (!messageIds) { return } messageIds.delete(messageID) if (messageIds.size === 0) { routingIndex.sessionMessageIdsById.delete(sessionID) } } const removeIndexedSession = (routingIndex: EventRoutingIndex, sessionID: string) => { if (!sessionID) { return } routingIndex.sessionDirectoryById.delete(sessionID) const messageIds = routingIndex.sessionMessageIdsById.get(sessionID) if (messageIds) { for (const messageID of messageIds) { routingIndex.messageSessionById.delete(messageID) } } routingIndex.sessionMessageIdsById.delete(sessionID) } const ingestDirectoryStateIntoRoutingIndex = ( routingIndex: EventRoutingIndex, directory: string, state: State, ) => { const nextSessionIds = new Set() for (const session of state.session) { if (!session?.id) { continue } nextSessionIds.add(session.id) setIndexedSessionDirectory(routingIndex, session.id, directory) } for (const sessionID of Object.keys(state.message)) { nextSessionIds.add(sessionID) setIndexedSessionDirectory(routingIndex, sessionID, directory) setIndexedSessionMessages(routingIndex, sessionID, directory, state.message[sessionID] ?? EMPTY_MESSAGES) } for (const [indexedSessionID, indexedDirectory] of routingIndex.sessionDirectoryById) { if (indexedDirectory !== directory) { continue } if (!nextSessionIds.has(indexedSessionID)) { removeIndexedSession(routingIndex, indexedSessionID) } } } const findSessionInChildStores = ( sessionID: string, childStores: ChildStoreManager, routingIndex: EventRoutingIndex, ): string | null => { for (const [dir, store] of childStores.children) { const state = store.getState() if ( state.session.some((s) => s.id === sessionID) || Object.prototype.hasOwnProperty.call(state.message, sessionID) || Object.prototype.hasOwnProperty.call(state.session_status ?? {}, sessionID) ) { // Self-heal: populate the routing index so future events resolve instantly setIndexedSessionDirectory(routingIndex, sessionID, dir) return dir } } return null } const childStoreHasSessionState = ( childStores: ChildStoreManager, directory: string, sessionID: string, ): boolean => { const store = childStores.getChild(directory) if (!store) return false const state = store.getState() return state.session.some((session) => session.id === sessionID) || Object.prototype.hasOwnProperty.call(state.message, sessionID) || Object.prototype.hasOwnProperty.call(state.session_status ?? {}, sessionID) } const childStoreHasMessagePartState = ( childStores: ChildStoreManager, directory: string, messageID: string, ): boolean => { const store = childStores.getChild(directory) if (!store) return false return Object.prototype.hasOwnProperty.call(store.getState().part, messageID) } const resolveDirectoryFromRoutingIndex = ( routingIndex: EventRoutingIndex, rawDirectory: string, payload: Event, childStores: ChildStoreManager, ): string => { const normalizedDirectory = normalizeEventDirectory(rawDirectory) const sessionID = getSessionIdFromPayload(payload) if (sessionID) { if (normalizedDirectory && normalizedDirectory !== "global" && childStoreHasSessionState(childStores, normalizedDirectory, sessionID)) { setIndexedSessionDirectory(routingIndex, sessionID, normalizedDirectory) return normalizedDirectory } const indexedDirectory = routingIndex.sessionDirectoryById.get(sessionID) if (indexedDirectory && childStores.getChild(indexedDirectory)) { return indexedDirectory } // Routing index miss — scan child stores for this session. // Covers optimistic sessions not yet indexed and events with wrong/empty directory. const found = findSessionInChildStores(sessionID, childStores, routingIndex) if (found) { return found } } const messageID = getMessageIdFromPayload(payload) if (messageID) { if (normalizedDirectory && normalizedDirectory !== "global" && childStoreHasMessagePartState(childStores, normalizedDirectory, messageID)) { return normalizedDirectory } const sessionFromMessage = routingIndex.messageSessionById.get(messageID) if (sessionFromMessage) { const indexedDirectory = routingIndex.sessionDirectoryById.get(sessionFromMessage) if (indexedDirectory && childStores.getChild(indexedDirectory)) { return indexedDirectory } } // Scan child stores for a store that has parts for this message for (const [dir, store] of childStores.children) { if (Object.prototype.hasOwnProperty.call(store.getState().part, messageID)) { return dir } } } // Single-store fallback: if there's only one directory, use it if ( (sessionID || messageID) && (!normalizedDirectory || normalizedDirectory === "global") && childStores.children.size === 1 ) { const onlyDirectory = childStores.children.keys().next().value if (typeof onlyDirectory === "string" && onlyDirectory.length > 0) { return onlyDirectory } } return normalizedDirectory } const updateRoutingIndexFromEvent = ( routingIndex: EventRoutingIndex, directory: string, payload: Event, ) => { if (!directory || directory === "global") { return } const sessionID = getSessionIdFromPayload(payload) if (sessionID) { setIndexedSessionDirectory(routingIndex, sessionID, directory) } switch (payload.type) { case "session.created": case "session.updated": { const info = (payload.properties as { info?: Session }).info if (info?.id) { setIndexedSessionDirectory(routingIndex, info.id, directory) } return } case "session.deleted": { const deletedSessionID = (payload.properties as { sessionID?: string }).sessionID if (deletedSessionID) { removeIndexedSession(routingIndex, deletedSessionID) } return } case "message.updated": { const info = (payload.properties as { info?: Message }).info if (info?.id && info.sessionID) { setIndexedMessage(routingIndex, info.sessionID, info.id, directory) } return } case "message.removed": { const props = payload.properties as { sessionID?: string; messageID?: string } if (props.messageID) { removeIndexedMessage(routingIndex, props.messageID, props.sessionID) } return } case "message.part.updated": { const part = (payload.properties as { part?: Part }).part as (Part & { sessionID?: string; messageID?: string }) | undefined if (part?.messageID && part.sessionID) { setIndexedMessage(routingIndex, part.sessionID, part.messageID, directory) } return } default: return } } /** * Re-fetch pending questions and permissions for a directory and merge them * into the directory's child store, preserving any in-flight SSE updates that * arrived while the request was pending. Used by reconnect/materialization * recovery paths only; normal session switches rely on primary SSE reducer * state for `question.asked` / `permission.asked` events. When * `candidateSessionIds` is omitted, every session known to the directory store * is treated as a candidate. */ export async function resyncBlockingRequestsForDirectory( directory: string, store: StoreApi, candidateSessionIds?: string[], ) { const before = store.getState() const knownSessionIds = new Set([ ...before.session.map((session) => session.id), ...Object.keys(before.message ?? {}), ...Object.keys(before.session_status ?? {}), ...Object.keys(before.question ?? {}), ...Object.keys(before.permission ?? {}), ]) const candidates = candidateSessionIds ?? Array.from(knownSessionIds) if (candidates.length === 0) return // Re-fetch pending questions that may have been asked during an SSE gap, // reconnect window, or directory materialization gap. try { const beforeSignatures = new Map( candidates.map((sessionId) => [sessionId, requestSignature(before.question[sessionId])]), ) const pendingQuestions = await opencodeClient.listPendingQuestions({ directories: [directory] }) const grouped: Record = {} for (const q of pendingQuestions) { if (!q?.id || !q.sessionID) continue if (!knownSessionIds.has(q.sessionID)) continue const list = grouped[q.sessionID] if (list) list.push(q) else grouped[q.sessionID] = [q] } for (const sessionId of Object.keys(grouped)) { grouped[sessionId].sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)) } for (const [sessionId, questions] of Object.entries(grouped)) { const knownIds = new Set((before.question[sessionId] ?? []).map((item) => item.id)) const isViewed = isViewedInCurrentSession(directory, sessionId) if (isViewed) continue for (const question of questions) { if (knownIds.has(question.id)) continue const toastKey = getQuestionToastKey(sessionId, question.id) if (!toastKey || pendingQuestionToastIds.has(toastKey)) continue pendingQuestionToastIds.add(toastKey) const firstQuestion = question.questions?.[0] const title = firstQuestion?.header?.trim() || "Input needed" const description = firstQuestion?.question?.trim() || "Agent is waiting for your response" toast.info(title, { id: `question-${toastKey}`, description, action: { label: "Open session", onClick: () => openSessionFromToast(sessionId, directory), }, }) } } store.setState((state: DirectoryStore) => { const merged = { ...state.question } for (const [sessionId, questions] of Object.entries(grouped)) { merged[sessionId] = questions } for (const sessionId of candidates) { if (grouped[sessionId]) continue const beforeSignature = beforeSignatures.get(sessionId) ?? "" const currentSignature = requestSignature(state.question[sessionId]) if (currentSignature !== beforeSignature) continue delete merged[sessionId] } return { question: merged } }) } catch { // Non-fatal: question resync best-effort } // Re-fetch pending permissions — same rationale as questions. try { const beforeSignatures = new Map( candidates.map((sessionId) => [sessionId, requestSignature(before.permission[sessionId])]), ) const pendingPermissions = await opencodeClient.listPendingPermissions({ directories: [directory] }) const grouped: Record = {} for (const permission of pendingPermissions) { if (!permission?.id || !permission.sessionID) continue if (!knownSessionIds.has(permission.sessionID)) continue const list = grouped[permission.sessionID] if (list) list.push(permission) else grouped[permission.sessionID] = [permission] } for (const sessionId of Object.keys(grouped)) { grouped[sessionId].sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)) } const permissionStore = usePermissionStore.getState() const autoAcceptingSessionIds = Object.keys(grouped).filter((sessionId) => permissionStore.isSessionAutoAccepting(sessionId)) if (autoAcceptingSessionIds.length > 0) { await Promise.all( autoAcceptingSessionIds.flatMap((sessionId) => (grouped[sessionId] ?? []).map((permission) => sessionActions.respondToPermission(permission.sessionID, permission.id, "once").catch(() => undefined), ), ), ) for (const sessionId of autoAcceptingSessionIds) { delete grouped[sessionId] } } for (const [sessionId, permissions] of Object.entries(grouped)) { const knownIds = new Set((before.permission[sessionId] ?? []).map((item) => item.id)) const isViewed = isViewedInCurrentSession(directory, sessionId) if (isViewed) continue for (const permission of permissions) { if (knownIds.has(permission.id)) continue const toastKey = getPermissionToastKey(sessionId, permission.id) if (!toastKey || pendingPermissionToastIds.has(toastKey)) continue pendingPermissionToastIds.add(toastKey) const description = typeof permission.permission === "string" && permission.permission.trim().length > 0 ? permission.permission : "Agent needs your approval" toast.info("Permission needed", { id: `permission-${toastKey}`, description, action: { label: "Open session", onClick: () => openSessionFromToast(sessionId, directory), }, }) } } store.setState((state: DirectoryStore) => { const merged = { ...state.permission } for (const [sessionId, permissions] of Object.entries(grouped)) { merged[sessionId] = permissions } for (const sessionId of candidates) { if (grouped[sessionId]) continue const beforeSignature = beforeSignatures.get(sessionId) ?? "" const currentSignature = requestSignature(state.permission[sessionId]) if (currentSignature !== beforeSignature) continue delete merged[sessionId] } return { permission: merged } }) } catch { // Non-fatal: permission resync best-effort } } async function resyncDirectoryAfterReconnect( directory: string, store: StoreApi, routingIndex: EventRoutingIndex, ) { const current = store.getState() const candidateSessionIds = getReconnectCandidateSessionIds(current, { directory, viewedSession: getViewedSessionMaterializationTarget(directory), }) 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 }, } }) } } const scopedClient = opencodeClient.getScopedSdkClient(directory) await Promise.all(candidateSessionIds.map(async (sessionId) => { const [sessionResponse, messageResponse] = await Promise.all([ scopedClient.session.get({ sessionID: sessionId }).catch(() => null), scopedClient.session.messages({ sessionID: sessionId, limit: RECONNECT_MESSAGE_LIMIT }).catch(() => null), ]) const session = sessionResponse?.data const records = messageResponse?.data if (!session || !records) return const nextSession = stripSessionDiffSnapshots(session) const nextMessages = records .filter((record) => !!record?.info?.id) .map((record) => stripMessageDiffSnapshots(record.info)) .sort((a, b) => cmp(a.id, b.id)) store.setState((state: DirectoryStore) => { 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 (!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 sessionChanged = true } const materialized = materializeSessionSnapshots( state, sessionId, records.map((record) => ({ info: stripMessageDiffSnapshots(record.info), parts: record.parts ?? [], })), { skipPartTypes: RECONNECT_SKIP_PARTS }, ) const messagesChanged = materialized.messagesChanged const partsChanged = materialized.partsChanged if (!sessionChanged && !messagesChanged && !partsChanged) { return state } return { ...(sessionChanged ? { session: sessions, sessionTotal } : {}), ...(messagesChanged ? { message: materialized.message } : {}), ...(partsChanged ? { part: materialized.part } : {}), } }) setIndexedSessionDirectory(routingIndex, nextSession.id, directory) setIndexedSessionMessages(routingIndex, sessionId, directory, nextMessages) })) await resyncBlockingRequestsForDirectory(directory, store, candidateSessionIds) ingestDirectoryStateIntoRoutingIndex(routingIndex, directory, store.getState()) } function handleEvent( rawDirectory: string, payload: Event, childStores: ChildStoreManager, routingIndex: EventRoutingIndex, ) { const directory = resolveDirectoryFromRoutingIndex(routingIndex, rawDirectory, payload, childStores) // Global events if (directory === "global" || !directory) { const recent = isRecentBoot() const result = reduceGlobalEvent(payload) if (!result) return if (result.type === "refresh") { // Suppress refresh during/shortly after bootstrap if (!recent) { useGlobalSyncStore.setState({ reload: "pending" }) } } else if (result.type === "project") { const current = useGlobalSyncStore.getState() useGlobalSyncStore.setState({ projects: applyGlobalProject(current, result.project).projects, }) } // On server.connected / global.disposed, re-bootstrap all directories // but only if not during recent boot if (payload.type === "server.connected" || payload.type === "global.disposed") { if (!recent) { for (const dir of childStores.children.keys()) { const store = childStores.getChild(dir) if (store && store.getState().status !== "loading") { // Mark as loading to trigger re-bootstrap store.setState({ status: "loading" as const }) childStores.ensureChild(dir) } } } } return } // Directory events let store = childStores.getChild(directory) let resolvedDirectory = directory if (!store) { // Store not found for this directory — attempt recovery by scanning // child stores for the session. This handles directory mismatches // (trailing slashes, case differences, events with wrong directory). const sessionID = getSessionIdFromPayload(payload) if (sessionID) { const fallbackDir = findSessionInChildStores(sessionID, childStores, routingIndex) if (fallbackDir) { store = childStores.getChild(fallbackDir) resolvedDirectory = fallbackDir } } } if (!store) { // Try as global event for unknown directories const result = reduceGlobalEvent(payload) if (result?.type === "refresh") { useGlobalSyncStore.setState({ reload: "pending" }) } else if (result?.type === "project") { const current = useGlobalSyncStore.getState() useGlobalSyncStore.setState({ projects: applyGlobalProject(current, result.project).projects, }) } return } childStores.mark(resolvedDirectory) if (payload.type === "permission.asked") { const permission = payload.properties as PermissionRequest const permissionStore = usePermissionStore.getState() if (permissionStore.isSessionAutoAccepting(permission.sessionID)) { updateRoutingIndexFromEvent(routingIndex, resolvedDirectory, payload) void sessionActions.respondToPermission(permission.sessionID, permission.id, "once").catch(() => undefined) return } const toastKey = getPermissionToastKey(permission.sessionID, permission.id) const isViewed = isViewedInCurrentSession(resolvedDirectory, permission.sessionID) if (!isViewed && toastKey && !pendingPermissionToastIds.has(toastKey)) { pendingPermissionToastIds.add(toastKey) const description = typeof permission.permission === "string" && permission.permission.trim().length > 0 ? permission.permission : "Agent needs your approval" toast.info("Permission needed", { id: `permission-${toastKey}`, description, action: { label: "Open session", onClick: () => openSessionFromToast(permission.sessionID, resolvedDirectory), }, }) } } if (payload.type === "permission.replied") { const props = payload.properties as { sessionID?: string; requestID?: string } const toastKey = getPermissionToastKey(props.sessionID, props.requestID) if (toastKey) { pendingPermissionToastIds.delete(toastKey) toast.dismiss(`permission-${toastKey}`) } } if (payload.type === "question.asked") { const question = payload.properties as QuestionRequest const sessionID = question.sessionID const toastKey = getQuestionToastKey(sessionID, question.id) const isViewed = isViewedInCurrentSession(resolvedDirectory, sessionID) if (!isViewed && toastKey && !pendingQuestionToastIds.has(toastKey)) { pendingQuestionToastIds.add(toastKey) const firstQuestion = question.questions?.[0] const title = firstQuestion?.header?.trim() || "Input needed" const description = firstQuestion?.question?.trim() || "Agent is waiting for your response" toast.info(title, { id: `question-${toastKey}`, description, action: { label: "Open session", onClick: () => openSessionFromToast(sessionID, resolvedDirectory), }, }) } } if (payload.type === "question.replied" || payload.type === "question.rejected") { const props = payload.properties as { sessionID?: string; requestID?: string } const toastKey = getQuestionToastKey(props.sessionID, props.requestID) if (toastKey) { pendingQuestionToastIds.delete(toastKey) toast.dismiss(`question-${toastKey}`) } } // Notification dispatch for session turn-complete and error events. // These are NOT handled by the event reducer — only the notification store. if (payload.type === "session.idle" || payload.type === "session.error") { const props = payload.properties as { sessionID?: string; error?: { message?: string; code?: string } } const sessionID = props.sessionID // Skip subtask sessions — only top-level sessions generate notifications const storeState = store.getState() const session = storeState.session.find((s) => s.id === sessionID) if (session && (session as { parentID?: string }).parentID) { // subtask — skip notification } else if (sessionID) { appendNotification({ directory: resolvedDirectory, session: sessionID, time: Date.now(), viewed: isViewedInCurrentSession(resolvedDirectory, sessionID), ...(payload.type === "session.error" ? { type: "error" as const, error: props.error } : { type: "turn-complete" as const }), }) } } // Sync-layer parent resync: when a child session goes idle, recover // the parent session snapshot. This ensures the // parent's task tool part reflects the child's completion even when // no ToolPart component is mounted. if (payload.type === "session.idle") { const idleSessionId = getSessionIdFromPayload(payload) if (idleSessionId && resolvedDirectory && resolvedDirectory !== "global") { const sessionState = store.getState() const idleSession = sessionState.session.find((s) => s.id === idleSessionId) const parentID = idleSession ? (idleSession as Session & { parentID?: string | null }).parentID : null if (parentID) { enqueueSessionMaterialization(resolvedDirectory, parentID, childStores) } } } // Read live state, create targeted draft cloning ONLY fields that event // type will mutate. This preserves reference identity for untouched slices // so Zustand selectors skip re-renders for unrelated subscribers. const current = store.getState() const draft: State = { ...current } switch (payload.type) { case "session.created": case "session.updated": case "session.deleted": draft.session = [...current.session] draft.permission = { ...current.permission } draft.todo = { ...current.todo } draft.part = { ...current.part } break case "session.diff": draft.session_diff = { ...current.session_diff } break case "session.status": case "session.idle": case "session.error": draft.session_status = { ...(current.session_status ?? {}) } break case "todo.updated": draft.todo = { ...current.todo } break case "message.updated": draft.message = { ...current.message } break case "message.removed": draft.message = { ...current.message } draft.part = { ...current.part } break case "message.part.updated": case "message.part.removed": case "message.part.delta": draft.part = { ...current.part } break case "vcs.branch.updated": break case "permission.asked": case "permission.replied": draft.permission = { ...current.permission } break case "question.asked": case "question.replied": case "question.rejected": draft.question = { ...current.question } break case "lsp.updated": draft.lsp = [...current.lsp] break default: break } const reducerResult = applyDirectoryEvent(draft, payload, { onSetSessionTodo: (sessionID, todos) => { useTodosPersistStore.getState().setSessionTodos(sessionID, todos) }, }) const reducerChanged = typeof reducerResult === "boolean" ? reducerResult : reducerResult.changed const materializationResult = typeof reducerResult === "boolean" ? undefined : reducerResult.materialization if (reducerChanged) { store.setState(draft) const sessionID = getSessionIdFromPayload(payload) ?? undefined const messageID = getMessageIdFromPayload(payload) ?? undefined syncDebug.dispatch.eventApplied(payload.type, sessionID, messageID) // Snapshot materialization on message.updated: if the message was inserted or // replaced but draft.part[messageID] is empty, the parts were lost or // never arrived. Recover the session so the UI doesn't render a blank bubble. if (sessionID && messageID && payload.type === "message.updated") { const after = store.getState() const info = (payload.properties as { info: Message }).info if (info.role === "assistant" && (!after.part[messageID] || after.part[messageID].length === 0)) { enqueueSessionMaterialization(resolvedDirectory, sessionID, childStores) } } } else { const sessionID = getSessionIdFromPayload(payload) ?? undefined const messageID = getMessageIdFromPayload(payload) ?? undefined syncDebug.dispatch.eventNoChange(payload.type, sessionID, messageID) } // Snapshot materialization is driven by typed reducer outcomes, not by // inferring meaning from a generic false/no-change result. if (materializationResult) { const materializationSessionID = materializationResult.sessionID ?? getSessionIdFromPayload(payload) ?? undefined if (materializationSessionID) { enqueueSessionMaterialization(resolvedDirectory, materializationSessionID, childStores) } } updateRoutingIndexFromEvent(routingIndex, resolvedDirectory, payload) } // --------------------------------------------------------------------------- // Provider // --------------------------------------------------------------------------- const dispatchOpenCodeUpdateAvailable = (payload: { version: string }) => { if (typeof window === "undefined") return window.dispatchEvent(new CustomEvent("openchamber:opencode-update-available", { detail: payload })) } export function SyncProvider(props: { sdk: OpencodeClient directory: string children: React.ReactNode }) { const messageStreamTransport = useConfigStore((state) => state.settingsMessageStreamTransport) const childStoresRef = useRef(null) if (!childStoresRef.current) childStoresRef.current = new ChildStoreManager() const childStores = childStoresRef.current const routingIndexRef = useRef(null) if (!routingIndexRef.current) routingIndexRef.current = createEventRoutingIndex() const routingIndex = routingIndexRef.current const system = useMemo( () => ({ childStores, sdk: props.sdk, directory: props.directory, }), [childStores, props.sdk, props.directory], ) // Configure child store manager useEffect(() => { const bootingDirs = new Set() childStores.configure({ onBootstrap: (directory) => { if (bootingDirs.has(directory)) return bootingDirs.add(directory) const store = childStores.getChild(directory) if (!store) return const runBootstrap = async (attempt: number) => { const globalState = useGlobalSyncStore.getState() await bootstrapDirectory({ directory, sdk: props.sdk, getState: () => store.getState(), set: (patch) => { store.setState(patch) if (patch.session || patch.message) { ingestDirectoryStateIntoRoutingIndex(routingIndex, directory, store.getState()) } }, global: { config: globalState.config, projects: globalState.projects, providers: globalState.providers, }, loadSessions: (dir) => retry(async () => { const result = await props.sdk.session.list({ directory: dir, roots: true, limit: 50, }) // SDK returns { error } instead of { data } on non-ok responses (503). // Preserve HTTP status so retry()'s transient detection works. const rawError = (result as { error?: unknown }).error if (rawError) { const response = (result as { response?: { status?: number } }).response const status = response?.status const message = typeof rawError === "object" && rawError !== null && "message" in rawError ? String((rawError as { message?: unknown }).message) : String(rawError) const wrapped = new Error(`session.list failed${status ? ` (${status})` : ""}: ${message}`) if (status !== undefined) { ;(wrapped as Error & { status?: number }).status = status } throw wrapped } const sessions = (result.data ?? []) .filter((s) => !!s?.id) .sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)) // Race guard: if the list came back empty but event pipeline // already populated the store, don't clobber. OpenCode can // answer HTTP with empty sessions while WS delivers session // events for the same data (disk warmup race on app launch). const currentSessions = store.getState().session if (sessions.length === 0 && currentSessions.length > 0) { console.warn( `[bootstrap] session.list returned empty for ${dir}; preserving ${currentSessions.length} existing sessions`, ) return } store.setState({ session: sessions, sessionTotal: sessions.length, limit: Math.max(sessions.length, 50) }) ingestDirectoryStateIntoRoutingIndex(routingIndex, directory, store.getState()) }), }) // VS Code race: if sessions are still empty after bootstrap, OpenCode // wasn't ready yet (bridge returned 503). Retry a few times. const state = store.getState() if (state.session.length === 0 && attempt < 5) { console.warn(`[bootstrap] sessions empty for ${directory} after attempt ${attempt + 1}; retrying in 2s`) await new Promise((r) => setTimeout(r, 2000)) store.setState({ status: "loading" as const }) await runBootstrap(attempt + 1) } else if (state.session.length === 0) { console.warn(`[bootstrap] sessions empty for ${directory} after ${attempt + 1} attempts; giving up`) } } runBootstrap(0).finally(() => { bootingDirs.delete(directory) }) }, onDispose: (directory) => { bootingDirs.delete(directory) }, isBooting: (directory) => bootingDirs.has(directory), isLoadingSessions: () => false, }) }, [childStores, props.sdk, routingIndex]) // Bootstrap global state — set bootingRoot/bootedAt to suppress // redundant refresh events during startup useEffect(() => { bootingRoot = true const globalActions = useGlobalSyncStore.getState().actions bootstrapGlobal(props.sdk, globalActions.set) .then(() => { bootedAt = Date.now() }) .finally(() => { bootingRoot = false }) }, [props.sdk]) // 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({ sdk: props.sdk, transport: messageStreamTransport, routeDirectory: (directory, payload) => { return resolveDirectoryFromRoutingIndex(routingIndex, directory, payload, childStores) }, onEvent: (directory, payload) => { if (payload.type === "installation.update-available") { const version = typeof (payload.properties as { version?: unknown })?.version === "string" ? (payload.properties as { version: string }).version : "" if (version) { dispatchOpenCodeUpdateAvailable({ version }) } } handleEvent(directory, payload, childStores, routingIndex) }, onReconnect: () => { useConfigStore.setState({ isConnected: true, hasEverConnected: true, connectionPhase: "connected", }) for (const dir of childStores.children.keys()) { triggerReconnectMaterialization(dir) } }, onDisconnect: (reason) => { const { hasEverConnected } = useConfigStore.getState() useConfigStore.setState({ isConnected: false, connectionPhase: hasEverConnected ? "reconnecting" : "connecting", lastDisconnectReason: reason, }) }, 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. useConfigStore.setState({ isConnected: true, hasEverConnected: true, connectionPhase: "connected", }) if (_activeDirectory) { triggerReconnectMaterialization(_activeDirectory) } }, }) return cleanup }, [props.sdk, childStores, routingIndex, messageStreamTransport]) // Ensure current directory's child store exists useEffect(() => { if (props.directory) { const store = childStores.ensureChild(props.directory) ingestDirectoryStateIntoRoutingIndex(routingIndex, props.directory, store.getState()) } }, [props.directory, childStores, routingIndex]) // Set refs so non-React code (session-actions, session-ui-store) can access sync state useEffect(() => { setSyncRefs(props.sdk, childStores, props.directory, (sessionID, dir) => { setIndexedSessionDirectory(routingIndex, sessionID, dir) }) setActionRefs( props.sdk, childStores, () => opencodeClient.getDirectory() || props.directory, ) }, [props.sdk, props.directory, childStores, routingIndex]) // Subscribe to child store for streaming state derivation useEffect(() => { if (!props.directory) return const store = childStores.getChild(props.directory) if (!store) return const unsubscribe = store.subscribe((state) => { updateStreamingState(state) }) return unsubscribe }, [props.directory, childStores]) return {props.children} } // --------------------------------------------------------------------------- // Hooks // --------------------------------------------------------------------------- /** Access the global sync store */ export function useGlobalSync() { return useGlobalSyncStore() } /** Access the global sync store with a selector */ export function useGlobalSyncSelector(selector: (state: GlobalSyncStore) => T): T { return useGlobalSyncStore(selector) } /** Get the child store for a directory (defaults to current) */ export function useDirectoryStore(directory?: string): StoreApi { const system = useSyncSystem() const dir = directory ?? system.directory return system.childStores.ensureChild(dir) } /** Select from the current directory's store */ export function useDirectorySync(selector: (state: State) => T, directory?: string): T { const store = useDirectoryStore(directory) return useStore(store, selector) } /** Get the revert messageID for a session (if reverted) */ export function useSessionRevertMessageID(sessionID: string, directory?: string): string | undefined { return useDirectorySync( useCallback((state: State) => { const session = state.session.find((s) => s.id === sessionID) return (session as { revert?: { messageID?: string } } | undefined)?.revert?.messageID }, [sessionID]), directory, ) } /** Get session messages for a specific session */ export function useSessionMessages(sessionID: string, directory?: string) { return useDirectorySync( useCallback((state: State) => state.message[sessionID] ?? EMPTY_MESSAGES, [sessionID]), directory, ) } /** * Get visible session messages — filters out reverted messages. * Filters out reverted messages (id >= session.revert.messageID). */ export function useVisibleSessionMessages(sessionID: string, directory?: string) { const messages = useSessionMessages(sessionID, directory) const revertMessageID = useSessionRevertMessageID(sessionID, directory) return useMemo(() => { if (!revertMessageID) return messages return messages.filter((m) => m.id < revertMessageID) }, [messages, revertMessageID]) } /** Check whether the message list for a session has been loaded into sync state. */ export function useSessionMessagesResolved(sessionID: string, directory?: string): boolean { return useDirectorySync( useCallback((state: State) => { if (!sessionID) return false return Object.prototype.hasOwnProperty.call(state.message, sessionID) }, [sessionID]), directory, ) } /** Get parts for a specific message */ export function useSessionParts(messageID: string, directory?: string) { return useDirectorySync( useCallback((state: State) => state.part[messageID] ?? EMPTY_PARTS, [messageID]), directory, ) } /** Get status for a specific session */ export function useSessionStatus(sessionID: string, directory?: string) { return useDirectorySync( useCallback((state: State) => state.session_status?.[sessionID], [sessionID]), directory, ) } /** Get permissions for a specific session */ export function useSessionPermissions(sessionID: string, directory?: string) { return useDirectorySync( useCallback((state: State) => state.permission[sessionID] ?? EMPTY_PERMISSION_REQUESTS, [sessionID]), directory, ) } /** Get questions for a specific session */ export function useSessionQuestions(sessionID: string, directory?: string) { return useDirectorySync( useCallback((state: State) => state.question[sessionID] ?? EMPTY_QUESTION_REQUESTS, [sessionID]), directory, ) } /** Get sessions list for a directory */ export function useSessions(directory?: string) { return useDirectorySync( useCallback((state: State) => state.session, []), directory, ) } const getSidebarSessionSignature = (session: Session, stableUpdatedAt: number): string => { const directory = (session as Session & { directory?: string | null }).directory ?? '' const parentID = (session as Session & { parentID?: string | null }).parentID ?? '' const projectWorktree = (session as Session & { project?: { worktree?: string | null } | null }).project?.worktree ?? '' const shared = session.share?.url ?? '' return [ session.id, session.title ?? '', session.time?.created ?? 0, session.time?.archived ? 1 : 0, directory, parentID, projectWorktree, shared, stableUpdatedAt, ].join('|') } /** Get sessions stabilized for sidebar tree rendering */ export function useSidebarSessions(directory?: string): Session[] { const store = useDirectoryStore(directory) const cacheRef = React.useRef<{ source: Session[] streamingSignature: string array: Session[] signatures: Map sessionsById: Map stableUpdatedAtById: Map streamingById: Map } | null>(null) const getSnapshot = React.useCallback(() => { const state = store.getState() const source = state.session const cached = cacheRef.current const streamingSignature = source .map((session) => { const statusType = state.session_status?.[session.id]?.type const isStreaming = statusType === 'busy' || statusType === 'retry' return `${session.id}:${isStreaming ? 1 : 0}` }) .join('|') if (cached && cached.source === source && cached.streamingSignature === streamingSignature) { return cached.array } const signatures = new Map() const sessionsById = new Map() const stableUpdatedAtById = new Map() const streamingById = new Map() let changed = !cached || cached.array.length !== source.length const array = source.map((session) => { const rawUpdatedAt = Number(session.time?.updated ?? session.time?.created ?? 0) const statusType = state.session_status?.[session.id]?.type const isStreaming = statusType === 'busy' || statusType === 'retry' const cachedUpdatedAt = cached?.stableUpdatedAtById.get(session.id) ?? rawUpdatedAt const wasStreaming = cached?.streamingById.get(session.id) ?? false const stableUpdatedAt = isStreaming ? (wasStreaming ? cachedUpdatedAt : Math.max(rawUpdatedAt, cachedUpdatedAt, Date.now())) : Math.max(rawUpdatedAt, cachedUpdatedAt) const signature = getSidebarSessionSignature(session, stableUpdatedAt) signatures.set(session.id, signature) stableUpdatedAtById.set(session.id, stableUpdatedAt) streamingById.set(session.id, isStreaming) const cachedSession = cached?.sessionsById.get(session.id) if ( cachedSession && cached?.signatures.get(session.id) === signature ) { sessionsById.set(session.id, cachedSession) return cachedSession } changed = true const nextSession = stableUpdatedAt === rawUpdatedAt ? session : { ...session, time: { ...session.time, updated: stableUpdatedAt, }, } sessionsById.set(session.id, nextSession) return nextSession }) if (!changed && cached) { cacheRef.current = { source, streamingSignature, array: cached.array, signatures, sessionsById: cached.sessionsById, stableUpdatedAtById, streamingById, } return cached.array } cacheRef.current = { source, streamingSignature, array, signatures, sessionsById, stableUpdatedAtById, streamingById } return array }, [store]) return React.useSyncExternalStore(store.subscribe, getSnapshot, getSnapshot) } /** Get one session by id for a directory */ export function useSession(sessionID?: string | null, directory?: string) { const { childStores } = useSyncSystem() const getSnapshot = useCallback(() => { if (directory) { return childStores.getChild(directory)?.getState().session.find((session) => session.id === sessionID) } return findLiveSession(getLiveStates(childStores), sessionID) }, [childStores, directory, sessionID]) const subscribe = useCallback((notify: () => void) => { if (directory) { return childStores.ensureChild(directory).subscribe(notify) } return childStores.subscribeAll(notify) }, [childStores, directory]) return React.useSyncExternalStore(subscribe, getSnapshot, getSnapshot) } /** Get one session directory by id for a directory */ export function useSessionDirectory(sessionID?: string | null, directory?: string): string | undefined { const session = useSession(sessionID, directory) return (session as (typeof session & { directory?: string | null }) | undefined)?.directory ?? undefined } /** Get the SDK client */ export function useSyncSDK() { return useSyncSystem().sdk } /** Get the current directory */ export function useSyncDirectory() { return useSyncSystem().directory } /** Get the child store manager (for advanced operations) */ export function useChildStoreManager() { return useSyncSystem().childStores } export type SessionTextMessage = { id: string role: string | null text: string } const getPartText = (part: Part): string => { if (part?.type !== "text") return "" const text = (part as { text?: unknown }).text return typeof text === "string" ? text : "" } const getConcatenatedTextFromParts = (parts: Part[]): string => { let text = "" for (const part of parts) { text += getPartText(part) } return text } const getFirstTextFromParts = (parts: Part[]): string => { for (const part of parts) { const text = getPartText(part) if (text.length > 0) return text } return "" } type SessionMessageRecord = { info: Message; parts: Part[] } type SessionMessageRecordsSnapshot = { sessionID: string sourceMessages: Message[] visibleMessages: Message[] revertMessageID?: string list: SessionMessageRecord[] byId: Map } function getVisibleMessagesForSession(state: State, sessionID: string, previous?: SessionMessageRecordsSnapshot): { sourceMessages: Message[] visibleMessages: Message[] revertMessageID?: string } { const sourceMessages = state.message[sessionID] ?? EMPTY_MESSAGES const session = state.session.find((candidate) => candidate.id === sessionID) const revertMessageID = (session as { revert?: { messageID?: string } } | undefined)?.revert?.messageID if ( previous && previous.sourceMessages === sourceMessages && previous.revertMessageID === revertMessageID ) { return { sourceMessages, visibleMessages: previous.visibleMessages, revertMessageID, } } return { sourceMessages, visibleMessages: revertMessageID ? sourceMessages.filter((message) => message.id < revertMessageID) : sourceMessages, revertMessageID, } } export function buildSessionMessageRecordsSnapshot( state: State, sessionID: string, previous?: SessionMessageRecordsSnapshot, suspendPartUpdates = false, ): SessionMessageRecordsSnapshot { const { sourceMessages, visibleMessages, revertMessageID } = getVisibleMessagesForSession(state, sessionID, previous) const nextById = new Map() const nextList = visibleMessages.map((message) => { const previousRecord = previous?.byId.get(message.id) const parts = suspendPartUpdates && previousRecord ? previousRecord.parts : (state.part[message.id] ?? EMPTY_PARTS) const nextRecord = previousRecord && previousRecord.info === message && previousRecord.parts === parts ? previousRecord : { info: message, parts } nextById.set(message.id, nextRecord) return nextRecord }) const unchanged = Boolean(previous) && previous?.visibleMessages === visibleMessages && previous.list.length === nextList.length && previous.list.every((record, index) => record === nextList[index]) if (unchanged && previous) { return previous } return { sessionID, sourceMessages, visibleMessages, revertMessageID, list: nextList, byId: nextById, } } export function useSessionMessageCount(sessionID: string, directory?: string): number { return useDirectorySync( useCallback((state: State) => { if (!sessionID) return 0 return state.message[sessionID]?.length ?? 0 }, [sessionID]), directory, ) } export function useSessionTextMessages(sessionID: string, directory?: string): SessionTextMessage[] { const records = useSessionMessageRecords(sessionID, directory) return useMemo( () => records.map((record) => ({ id: record.info.id, role: typeof record.info.role === "string" ? record.info.role : null, text: getConcatenatedTextFromParts(record.parts), })), [records], ) } export function useUserMessageHistory(sessionID: string, directory?: string): string[] { const records = useSessionMessageRecords(sessionID, directory) const userMessages = useMemo(() => records.filter((record) => record.info.role === 'user'), [records]) return useMemo(() => { const history: string[] = [] for (let index = userMessages.length - 1; index >= 0; index -= 1) { const message = userMessages[index] const text = getFirstTextFromParts(message.parts) if (text.length > 0) { history.push(text) } } return history }, [userMessages]) } /** * Get messages for a session in the old {info, parts}[] format. * Uses visible messages (filtered by revert state). * * Uses a ref-stable parts lookup that only triggers re-renders when * a part array for one of our displayed messages actually changes. */ export function useSessionMessageRecords( sessionID: string, directory?: string, options?: { suspendPartUpdates?: boolean }, ) { const store = useDirectoryStore(directory) const snapshotRef = useRef({ sessionID, sourceMessages: EMPTY_MESSAGES, visibleMessages: EMPTY_MESSAGES, revertMessageID: undefined, list: [], byId: new Map(), }) const getSnapshot = useCallback(() => { const nextSnapshot = buildSessionMessageRecordsSnapshot( store.getState(), sessionID, snapshotRef.current.sessionID === sessionID ? snapshotRef.current : undefined, Boolean(options?.suspendPartUpdates), ) snapshotRef.current = nextSnapshot return nextSnapshot.list }, [options?.suspendPartUpdates, sessionID, store]) return React.useSyncExternalStore(store.subscribe, getSnapshot, getSnapshot) } /** * Ensures a session's messages are loaded into the sync store. * If the session exists in state.session but messages haven't been fetched * (state.message[sessionID] is absent), triggers a background API fetch. * * This covers the case where a user navigates to an old parent session * whose child session messages were never loaded — bootstrap only loads * session metadata, not messages. */ // Module-level in-flight tracking for useEnsureSessionMessages. // Prevents redundant parallel fetches when multiple component instances // (e.g. multiple ToolParts) request the same session's messages. const _ensureMessagesLoading = new Set() export function useEnsureSessionMessages(sessionID: string, directory?: string) { const store = useDirectoryStore(directory) React.useEffect(() => { if (!sessionID) return const state = store.getState() // Already loaded into a renderable message/part snapshot — nothing to do. if (getSessionMaterializationStatus(state, sessionID).renderable) return // Session doesn't exist — nothing to load if (!state.session.some((s) => s.id === sessionID)) return const dir = directory ?? opencodeClient.getDirectory() const loadingKey = `${dir ?? ""}:${sessionID}` // Already loading this session for this directory if (_ensureMessagesLoading.has(loadingKey)) return _ensureMessagesLoading.add(loadingKey) void (async () => { try { await materializeSessionFromServer(dir ?? "", sessionID, store) } catch { // Transient failure — next navigation or reconnect will retry } finally { _ensureMessagesLoading.delete(loadingKey) } })() }, [sessionID, store, directory]) } /** * Determines if a session is actively working. * Checks session_status and only falls back to incomplete assistant messages * when authoritative status is missing. * Returns false when permissions are pending (permission indicator takes priority). */ export function useIsSessionWorking(sessionID: string, directory?: string): boolean { const status = useSessionStatus(sessionID, directory) const permissions = useSessionPermissions(sessionID, directory) const messages = useSessionMessages(sessionID, directory) return useMemo(() => { // Permissions pending → not "working" (show permission indicator instead) if (permissions.length > 0) return false // Check session_status const hasAuthoritativeStatus = status !== undefined const statusWorking = hasAuthoritativeStatus && status.type !== "idle" // Check for incomplete assistant message (fallback if status event delayed) let hasPendingAssistant = false for (let i = messages.length - 1; i >= 0; i--) { const m = messages[i] if (m.role === "assistant" && typeof (m as { time?: { completed?: number } }).time?.completed !== "number") { hasPendingAssistant = true break } } if (hasAuthoritativeStatus) return statusWorking return hasPendingAssistant }, [status, permissions, messages]) } const EMPTY_MESSAGES: Message[] = [] const EMPTY_PARTS: Part[] = [] const EMPTY_PERMISSION_REQUESTS: PermissionRequest[] = [] const EMPTY_QUESTION_REQUESTS: QuestionRequest[] = []