Files
openchamber/packages/ui/src/sync/sync-context.tsx
T
Bohdan Triapitsyn 767f423500 fix: suppress auto-approved permission notifications
Skip permission prompts when auto-approve is enabled
Keep request handling reliable after reconnect
Remove unused icon import to keep lint clean
2026-04-17 13:22:14 +03:00

1856 lines
63 KiB
TypeScript

/* 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 { opencodeClient } from "@/lib/opencode/client"
import { usePermissionStore } from "@/stores/permissionStore"
import { useConfigStore } from "@/stores/useConfigStore"
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"
// ---------------------------------------------------------------------------
// 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<SyncSystem | null>
}
const syncGlobal = globalThis as SyncGlobal
const SyncContext = syncGlobal[SYNC_CONTEXT_GLOBAL_KEY] ?? createContext<SyncSystem | null>(null)
syncGlobal[SYNC_CONTEXT_GLOBAL_KEY] = SyncContext
function useSyncSystem() {
const ctx = useContext(SyncContext)
if (!ctx) throw new Error("useSyncSystem must be used within <SyncProvider>")
return ctx
}
function getLiveStates(childStores: ChildStoreManager): State[] {
return Array.from(childStores.children.values(), (store) => store.getState())
}
function useLiveSyncSelector<T>(selector: (states: State[]) => T, isEqual: (left: T, right: T) => boolean = Object.is): T {
const { childStores } = useSyncSystem()
const cacheRef = useRef<T | undefined>(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<string, SessionStatus> {
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 = 200
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 partRepairSignature = (part: Part): string => JSON.stringify(part)
function haveEquivalentPartSnapshots(left: Part[] | undefined, right: Part[]): boolean {
if (!left) {
return right.length === 0
}
if (left.length !== right.length) {
return false
}
for (let index = 0; index < left.length; index += 1) {
const leftPart = left[index]
const rightPart = right[index]
if (!leftPart || !rightPart) {
return false
}
if (leftPart.id !== rightPart.id) {
return false
}
if (partRepairSignature(leftPart) !== partRepairSignature(rightPart)) {
return false
}
}
return true
}
// ---------------------------------------------------------------------------
// Parts-gap recovery — when SSE events arrive but parts are missing,
// trigger a targeted re-fetch for the affected sessions.
// Tracked per-directory, deduplicated, and auto-expiring.
// ---------------------------------------------------------------------------
type PendingRepair = {
sessionID: string
directory: string
enqueuedAt: number
}
const REPAIR_COOLDOWN_MS = 5_000
const pendingRepairs = new Map<string, PendingRepair>() // key: directory:sessionID
const repairKey = (directory: string, sessionID: string) => `${directory}:${sessionID}`
function enqueuePartsRepair(directory: string, sessionID: string, childStores: ChildStoreManager) {
if (!directory || directory === "global" || !sessionID) return
const k = repairKey(directory, sessionID)
const existing = pendingRepairs.get(k)
if (existing && Date.now() - existing.enqueuedAt < REPAIR_COOLDOWN_MS) return
pendingRepairs.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) {
pendingRepairs.delete(k)
return
}
try {
await repairSessionParts(directory, sessionID, store)
} catch {
// Transient failure — next SSE event or reconnect will catch up.
} finally {
pendingRepairs.delete(k)
}
})
}
async function repairSessionParts(
directory: string,
sessionID: string,
store: StoreApi<DirectoryStore>,
) {
const scopedClient = opencodeClient.getScopedSdkClient(directory)
const result = await retry(() =>
scopedClient.session.messages({ sessionID, limit: RECONNECT_MESSAGE_LIMIT }),
)
const records = (result.data ?? []).filter((record: { info?: { id?: string } }) => !!record?.info?.id)
if (records.length === 0) return
store.setState((state: DirectoryStore) => {
const nextPartState = { ...state.part }
for (const record of records) {
const messageId = record?.info?.id
if (!messageId) continue
const newParts = (record.parts ?? [])
.filter((part: Part) => !!part?.id && !RECONNECT_SKIP_PARTS.has(part.type))
.sort((a: Part, b: Part) => cmp(a.id, b.id))
const existing = nextPartState[messageId]
// Repair when parts are missing, truncated, or stale-but-same-length.
if (!haveEquivalentPartSnapshots(existing, newParts)) {
nextPartState[messageId] = newParts
}
}
return { part: nextPartState }
})
}
// 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 pendingQuestionToastIds = new Set<string>()
const pendingPermissionToastIds = new Set<string>()
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
}
function isViewedInCurrentSession(directory: string, sessionId?: string): boolean {
if (!_activeDirectory || !_activeSession || !sessionId) return false
if (directory !== _activeDirectory) return false
return sessionId === _activeSession
}
function isRecentBoot() {
return bootingRoot || Date.now() - bootedAt < BOOT_DEBOUNCE_MS
}
function getReconnectCandidateSessionIds(state: State) {
const ids = new Set<string>()
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)
}
}
return Array.from(ids)
}
function toSessionStatus(status: Awaited<ReturnType<typeof opencodeClient.getSessionStatus>>[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<string, string>
messageSessionById: Map<string, string>
sessionMessageIdsById: Map<string, Set<string>>
}
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<string, unknown>
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<string, unknown>
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<string>()
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<string>()
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 resolveDirectoryFromRoutingIndex = (
routingIndex: EventRoutingIndex,
rawDirectory: string,
payload: Event,
childStores: ChildStoreManager,
): string => {
const normalizedDirectory = normalizeEventDirectory(rawDirectory)
const sessionID = getSessionIdFromPayload(payload)
if (sessionID) {
const indexedDirectory = routingIndex.sessionDirectoryById.get(sessionID)
if (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) {
const sessionFromMessage = routingIndex.messageSessionById.get(messageID)
if (sessionFromMessage) {
const indexedDirectory = routingIndex.sessionDirectoryById.get(sessionFromMessage)
if (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
}
}
async function resyncDirectoryAfterReconnect(
directory: string,
store: StoreApi<DirectoryStore>,
routingIndex: EventRoutingIndex,
) {
const current = store.getState()
const candidateSessionIds = getReconnectCandidateSessionIds(current)
if (candidateSessionIds.length === 0) return
const nextStatuses = await opencodeClient.getSessionStatusForDirectory(directory)
const relevantStatuses: Record<string, SessionStatus> = {}
for (const sessionId of candidateSessionIds) {
const nextStatus = toSessionStatus(nextStatuses[sessionId])
if (nextStatus) {
relevantStatuses[sessionId] = nextStatus
}
}
if (Object.keys(relevantStatuses).length > 0) {
store.setState((state: DirectoryStore) => ({
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))
const nextMessageIds = new Set(nextMessages.map((message) => message.id))
store.setState((state: DirectoryStore) => {
const sessions = [...state.session]
const sessionIndex = sessions.findIndex((item) => item.id === nextSession.id)
let sessionChanged = false
let sessionTotal = state.sessionTotal
if (sessionIndex >= 0) {
if (sessions[sessionIndex] !== nextSession) {
sessions[sessionIndex] = nextSession
sessionChanged = true
}
} else {
sessions.push(nextSession)
sessions.sort((a, b) => cmp(a.id, b.id))
if (!nextSession.parentID) sessionTotal += 1
sessionChanged = true
}
const nextPartState = { ...state.part }
const previousMessages = state.message[sessionId] ?? []
for (const message of previousMessages) {
if (!nextMessageIds.has(message.id)) {
delete nextPartState[message.id]
}
}
for (const record of records) {
const messageId = record?.info?.id
if (!messageId) continue
nextPartState[messageId] = (record.parts ?? [])
.filter((part) => !!part?.id && !RECONNECT_SKIP_PARTS.has(part.type))
.sort((a, b) => cmp(a.id, b.id))
}
return {
...(sessionChanged ? { session: sessions, sessionTotal } : {}),
message: { ...state.message, [sessionId]: nextMessages },
part: nextPartState,
}
})
setIndexedSessionDirectory(routingIndex, nextSession.id, directory)
setIndexedSessionMessages(routingIndex, sessionId, directory, nextMessages)
}))
// Re-fetch pending questions on reconnect — they may have been asked
// during the SSE disconnection window and will not arrive via SSE events.
// Overwrite sessions covered by API response, and clear reconnect candidates
// that remain unchanged during the request but are absent from the response.
// If SSE changed a session while the request was in-flight, keep that data.
try {
const before = store.getState()
const knownSessionIds = new Set<string>([
...before.session.map((session) => session.id),
...Object.keys(before.message ?? {}),
...Object.keys(before.session_status ?? {}),
...Object.keys(before.question ?? {}),
...Object.keys(before.permission ?? {}),
])
const beforeSignatures = new Map(
candidateSessionIds.map((sessionId) => [sessionId, requestSignature(before.question[sessionId])]),
)
const pendingQuestions = await opencodeClient.listPendingQuestions({ directories: [directory] })
const grouped: Record<string, QuestionRequest[]> = {}
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]
}
// Sort each group by id for binary-search compatibility
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 candidateSessionIds) {
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 on reconnect — same rationale as questions.
try {
const before = store.getState()
const knownSessionIds = new Set<string>([
...before.session.map((session) => session.id),
...Object.keys(before.message ?? {}),
...Object.keys(before.session_status ?? {}),
...Object.keys(before.question ?? {}),
...Object.keys(before.permission ?? {}),
])
const beforeSignatures = new Map(
candidateSessionIds.map((sessionId) => [sessionId, requestSignature(before.permission[sessionId])]),
)
const pendingPermissions = await opencodeClient.listPendingPermissions({ directories: [directory] })
const grouped: Record<string, PermissionRequest[]> = {}
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 candidateSessionIds) {
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
}
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 }),
})
}
}
// Read live state, create targeted draft cloning ONLY fields the 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
}
if (applyDirectoryEvent(draft, payload)) {
store.setState(draft)
const sessionID = getSessionIdFromPayload(payload) ?? undefined
const messageID = getMessageIdFromPayload(payload) ?? undefined
syncDebug.dispatch.eventApplied(payload.type, sessionID, messageID)
// Parts-gap recovery on message.updated: if the message was inserted or
// replaced but draft.part[messageID] is empty, the parts were lost or
// never arrived. Trigger repair 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)) {
enqueuePartsRepair(resolvedDirectory, sessionID, childStores)
}
}
} else {
const sessionID = getSessionIdFromPayload(payload) ?? undefined
const messageID = getMessageIdFromPayload(payload) ?? undefined
syncDebug.dispatch.eventNoChange(payload.type, sessionID, messageID)
// Parts-gap recovery: if a part event was dropped because the parts array
// was missing (message not yet inserted or parts lost), trigger a repair
// fetch for the session.
if (sessionID && messageID && (
payload.type === "message.part.delta" || payload.type === "message.part.updated"
)) {
enqueuePartsRepair(resolvedDirectory, sessionID, childStores)
}
}
updateRoutingIndexFromEvent(routingIndex, resolvedDirectory, payload)
}
// ---------------------------------------------------------------------------
// Provider
// ---------------------------------------------------------------------------
export function SyncProvider(props: {
sdk: OpencodeClient
directory: string
children: React.ReactNode
}) {
const messageStreamTransport = useConfigStore((state) => state.settingsMessageStreamTransport)
const childStoresRef = useRef<ChildStoreManager | null>(null)
if (!childStoresRef.current) childStoresRef.current = new ChildStoreManager()
const childStores = childStoresRef.current
const routingIndexRef = useRef<EventRoutingIndex | null>(null)
if (!routingIndexRef.current) routingIndexRef.current = createEventRoutingIndex()
const routingIndex = routingIndexRef.current
const system = useMemo<SyncSystem>(
() => ({
childStores,
sdk: props.sdk,
directory: props.directory,
}),
[childStores, props.sdk, props.directory],
)
// Configure child store manager
useEffect(() => {
const bootingDirs = new Set<string>()
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).
// Throw so retry() retries and allSettled marks it as rejected.
if ((result as { error?: unknown }).error) {
throw new Error("session.list failed: " + String((result as { error?: unknown }).error))
}
const sessions = (result.data ?? [])
.filter((s) => !!s?.id)
.sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0))
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) {
await new Promise((r) => setTimeout(r, 2000))
store.setState({ status: "loading" as const })
await runBootstrap(attempt + 1)
}
}
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 reconnectResyncing = new Set<string>()
const { cleanup } = createEventPipeline({
sdk: props.sdk,
transport: messageStreamTransport,
routeDirectory: (directory, payload) => {
return resolveDirectoryFromRoutingIndex(routingIndex, directory, payload, childStores)
},
onEvent: (directory, payload) => {
handleEvent(directory, payload, childStores, routingIndex)
},
onReconnect: () => {
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)
})
}
},
})
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 <SyncContext.Provider value={system}>{props.children}</SyncContext.Provider>
}
// ---------------------------------------------------------------------------
// Hooks
// ---------------------------------------------------------------------------
/** Access the global sync store */
export function useGlobalSync() {
return useGlobalSyncStore()
}
/** Access the global sync store with a selector */
export function useGlobalSyncSelector<T>(selector: (state: GlobalSyncStore) => T): T {
return useGlobalSyncStore(selector)
}
/** Get the child store for a directory (defaults to current) */
export function useDirectoryStore(directory?: string): StoreApi<DirectoryStore> {
const system = useSyncSystem()
const dir = directory ?? system.directory
return system.childStores.ensureChild(dir)
}
/** Select from the current directory's store */
export function useDirectorySync<T>(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<string, string>
sessionsById: Map<string, Session>
stableUpdatedAtById: Map<string, number>
streamingById: Map<string, boolean>
} | 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<string, string>()
const sessionsById = new Map<string, Session>()
const stableUpdatedAtById = new Map<string, number>()
const streamingById = new Map<string, boolean>()
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()))
: 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<string, SessionMessageRecord>
}
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,
}
}
function buildSessionMessageRecordsSnapshot(
state: State,
sessionID: string,
previous?: SessionMessageRecordsSnapshot,
suspendPartUpdates = false,
): SessionMessageRecordsSnapshot {
const { sourceMessages, visibleMessages, revertMessageID } = getVisibleMessagesForSession(state, sessionID, previous)
const nextById = new Map<string, SessionMessageRecord>()
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<SessionMessageRecordsSnapshot>({
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)
}
/**
* 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[] = []