diff --git a/packages/ui/src/components/chat/MessageList.tsx b/packages/ui/src/components/chat/MessageList.tsx index 50a08020..da8d5c96 100644 --- a/packages/ui/src/components/chat/MessageList.tsx +++ b/packages/ui/src/components/chat/MessageList.tsx @@ -49,6 +49,87 @@ const resolveMessageRole = (message: ChatMessageEntry): string | null => { ?? null; }; +const hasCompactionPart = (message: ChatMessageEntry): boolean => { + return message.parts.some((part) => { + const type = (part as { type?: unknown }).type; + return type === 'compaction'; + }); +}; + +const getPartText = (part: Part): string => { + const text = (part as { text?: unknown }).text; + if (typeof text === 'string') { + return text; + } + const content = (part as { content?: unknown }).content; + if (typeof content === 'string') { + return content; + } + return ''; +}; + +const normalizeCompactionCommandMessage = (message: ChatMessageEntry): ChatMessageEntry => { + if (!hasCompactionPart(message)) { + return message; + } + + let changedParts = false; + const nextParts = message.parts.map((part) => { + const type = (part as { type?: unknown }).type; + if (type !== 'compaction') { + return part; + } + changedParts = true; + return { type: 'text', text: '/compact' } as Part; + }); + + const info = message.info as unknown as { clientRole?: string | null | undefined }; + const needsClientRole = info.clientRole !== 'user'; + + if (!changedParts && !needsClientRole) { + return message; + } + + return { + ...message, + info: needsClientRole + ? ({ + ...(message.info as unknown as Record), + clientRole: 'user', + } as unknown as typeof message.info) + : message.info, + parts: changedParts ? nextParts : message.parts, + }; +}; + +const normalizeCompactionSummaryMessage = ( + message: ChatMessageEntry, + compactionCommandIds: Set, +): ChatMessageEntry => { + const role = resolveMessageRole(message); + if (role !== 'system') { + return message; + } + + const parentID = getMessageParentId(message); + if (!parentID || !compactionCommandIds.has(parentID)) { + return message; + } + + const info = message.info as unknown as { clientRole?: string | null | undefined }; + if (info.clientRole === 'assistant') { + return message; + } + + return { + ...message, + info: ({ + ...(message.info as unknown as Record), + clientRole: 'assistant', + } as unknown as typeof message.info), + }; +}; + const isAssistantMessageCompleted = (message: ChatMessageEntry): boolean => { const info = message.info as { time?: { completed?: unknown }; status?: unknown }; const completed = info.time?.completed; @@ -279,11 +360,12 @@ const getNormalizedMessageForDisplay = (message: ChatMessageEntry): ChatMessageE return cached; } - const filteredParts = filterSyntheticParts(message.parts); - const normalized = filteredParts === message.parts - ? message + const normalizedCompactionMessage = normalizeCompactionCommandMessage(message); + const filteredParts = filterSyntheticParts(normalizedCompactionMessage.parts); + const normalized = filteredParts === normalizedCompactionMessage.parts + ? normalizedCompactionMessage : { - ...message, + ...normalizedCompactionMessage, parts: filteredParts, }; @@ -1242,12 +1324,17 @@ const MessageList = React.forwardRef(({ dedupedMessages.reverse(); const output: ChatMessageEntry[] = []; + const compactionCommandIds = new Set(); for (let index = 0; index < dedupedMessages.length; index += 1) { const current = dedupedMessages[index]; + const currentWithRole = normalizeCompactionSummaryMessage(current, compactionCommandIds); + if (hasCompactionPart(current) || current.parts.some((part) => part.type === 'text' && getPartText(part).trim() === '/compact')) { + compactionCommandIds.add(current.info.id); + } const previous = output.length > 0 ? output[output.length - 1] : undefined; if (isUserSubtaskMessage(previous)) { - const bridge = isSyntheticSubtaskBridgeAssistant(current); + const bridge = isSyntheticSubtaskBridgeAssistant(currentWithRole); if (bridge.hide) { output[output.length - 1] = withSubtaskSessionId(previous as ChatMessageEntry, bridge.taskSessionId); continue; @@ -1255,14 +1342,14 @@ const MessageList = React.forwardRef(({ } if (isUserShellMarkerMessage(previous)) { - const bridge = getShellBridgeAssistantDetails(current, getMessageId(previous)); + const bridge = getShellBridgeAssistantDetails(currentWithRole, getMessageId(previous)); if (bridge.hide) { output[output.length - 1] = withShellBridgeDetails(previous as ChatMessageEntry, bridge.details); continue; } } - output.push(current); + output.push(currentWithRole); } const outputIndexById = new Map(); diff --git a/packages/ui/src/sync/event-pipeline.ts b/packages/ui/src/sync/event-pipeline.ts index 977dee0c..efc32aca 100644 --- a/packages/ui/src/sync/event-pipeline.ts +++ b/packages/ui/src/sync/event-pipeline.ts @@ -41,6 +41,23 @@ export type EventPipelineInput = { onReconnect?: () => void } +const normalizeEventType = (payload: Event): Event => { + const type = (payload as { type?: unknown }).type + if (typeof type !== "string") { + return payload + } + + const match = /^(.*)\.(\d+)$/.exec(type) + if (!match || !match[1]) { + return payload + } + + return { + ...payload, + type: match[1] as Event["type"], + } as unknown as Event +} + function resolveEventDirectory(event: unknown, payload: Event): string { const directDirectory = typeof event === "object" && event !== null && typeof (event as { directory?: unknown }).directory === "string" @@ -190,21 +207,22 @@ export function createEventPipeline(input: EventPipelineInput) { if (!payload || typeof payload !== "object" || typeof (payload as { type?: unknown }).type !== "string") { continue } - const directory = resolveEventDirectory(event, payload) - const k = key(directory, payload) + const normalizedPayload = normalizeEventType(payload) + const directory = resolveEventDirectory(event, normalizedPayload) + const k = key(directory, normalizedPayload) if (k) { const i = coalesced.get(k) if (i !== undefined) { - queue[i] = { directory, payload } - if (payload.type === "message.part.updated") { - const part = (payload.properties as { part: { messageID: string; id: string } }).part + queue[i] = { directory, payload: normalizedPayload } + if (normalizedPayload.type === "message.part.updated") { + const part = (normalizedPayload.properties as { part: { messageID: string; id: string } }).part staleDeltas.add(deltaKey(directory, part.messageID, part.id)) } continue } coalesced.set(k, queue.length) } - queue.push({ directory, payload }) + queue.push({ directory, payload: normalizedPayload }) schedule() if (Date.now() - yielded < STREAM_YIELD_MS) continue diff --git a/packages/ui/src/sync/sync-context.tsx b/packages/ui/src/sync/sync-context.tsx index 5dc0846b..0a08d9af 100644 --- a/packages/ui/src/sync/sync-context.tsx +++ b/packages/ui/src/sync/sync-context.tsx @@ -178,7 +178,350 @@ function toSessionStatus(status: Awaited) { +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 + } + return rawDirectory.replace(/\\/g, "/").replace(/^([a-z]):/, (_, l: string) => l.toUpperCase() + ":") +} + +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 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 + } + } + + const messageID = getMessageIdFromPayload(payload) + if (messageID) { + const sessionFromMessage = routingIndex.messageSessionById.get(messageID) + if (sessionFromMessage) { + const indexedDirectory = routingIndex.sessionDirectoryById.get(sessionFromMessage) + if (indexedDirectory) { + return indexedDirectory + } + } + } + + if ((!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, + routingIndex: EventRoutingIndex, +) { const current = store.getState() const candidateSessionIds = getReconnectCandidateSessionIds(current) if (candidateSessionIds.length === 0) return @@ -256,20 +599,22 @@ async function resyncDirectoryAfterReconnect(directory: string, store: StoreApi< part: nextPartState, } }) + + setIndexedSessionDirectory(routingIndex, nextSession.id, directory) + setIndexedSessionMessages(routingIndex, sessionId, directory, nextMessages) })) + + ingestDirectoryStateIntoRoutingIndex(routingIndex, directory, store.getState()) } function handleEvent( rawDirectory: string, payload: Event, childStores: ChildStoreManager, + routingIndex: EventRoutingIndex, ) { - // Normalize directory path: SSE events from OpenCode use native OS separators - // (backslashes on Windows) and may differ in drive-letter case. - // Child stores are keyed with forward slashes and uppercase drive letters. - const directory = rawDirectory && rawDirectory !== "global" - ? rawDirectory.replace(/\\/g, "/").replace(/^([a-z]):/, (_, l: string) => l.toUpperCase() + ":") - : rawDirectory + const directory = resolveDirectoryFromRoutingIndex(routingIndex, rawDirectory, payload, childStores) + // Global events if (directory === "global" || !directory) { const recent = isRecentBoot() @@ -305,6 +650,7 @@ function handleEvent( // Directory events const store = childStores.getChild(directory) + if (!store) { // Try as global event for unknown directories const result = reduceGlobalEvent(payload) @@ -402,6 +748,8 @@ function handleEvent( store.setState(draft) } + updateRoutingIndexFromEvent(routingIndex, directory, payload) + // Update global session status for cross-directory sidebar visibility if (payload.type === "session.status") { const props = payload.properties as { sessionID: string; status: SessionStatus } @@ -435,6 +783,9 @@ export function SyncProvider(props: { 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( () => ({ @@ -465,6 +816,9 @@ export function SyncProvider(props: { getState: () => store.getState(), set: (patch) => { store.setState(patch) + if (patch.session || patch.message) { + ingestDirectoryStateIntoRoutingIndex(routingIndex, directory, store.getState()) + } if (patch.session_status) { const current = useGlobalSessionStatusStore.getState().statuses const merged = { ...current, ...patch.session_status } @@ -491,6 +845,7 @@ export function SyncProvider(props: { .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()) }), }) @@ -514,7 +869,7 @@ export function SyncProvider(props: { isBooting: (directory) => bootingDirs.has(directory), isLoadingSessions: () => false, }) - }, [childStores, props.sdk]) + }, [childStores, props.sdk, routingIndex]) // Bootstrap global state — set bootingRoot/bootedAt to suppress // redundant refresh events during startup @@ -538,7 +893,7 @@ export function SyncProvider(props: { const { cleanup } = createEventPipeline({ sdk: props.sdk, onEvent: (directory, payload) => { - handleEvent(directory, payload, childStores) + handleEvent(directory, payload, childStores, routingIndex) }, onReconnect: () => { for (const [dir, store] of childStores.children) { @@ -546,7 +901,7 @@ export function SyncProvider(props: { if (getReconnectCandidateSessionIds(store.getState()).length === 0) continue reconnectResyncing.add(dir) - void resyncDirectoryAfterReconnect(dir, store) + void resyncDirectoryAfterReconnect(dir, store, routingIndex) .catch(() => { // Transient failure during resync — next SSE event or reconnect will catch up. }) @@ -557,14 +912,15 @@ export function SyncProvider(props: { }, }) return cleanup - }, [props.sdk, childStores]) + }, [props.sdk, childStores, routingIndex]) // Ensure current directory's child store exists useEffect(() => { if (props.directory) { - childStores.ensureChild(props.directory) + const store = childStores.ensureChild(props.directory) + ingestDirectoryStateIntoRoutingIndex(routingIndex, props.directory, store.getState()) } - }, [props.directory, childStores]) + }, [props.directory, childStores, routingIndex]) // Set refs so non-React code (session-actions, session-ui-store) can access sync state useEffect(() => {