fix(app): move hot response processing off thread

This commit is contained in:
LukeParkerDev 2026-08-04 20:07:35 +10:00
parent a47dabff22
commit ccc11dc92d
12 changed files with 290 additions and 66 deletions

View file

@ -13,6 +13,8 @@ import { useGlobal } from "./global"
import { ServerScope } from "@/utils/server-scope"
import { detectServerProtocol, type ServerProtocol } from "@/utils/server-protocol"
import { createCompatibleApi, type CompatibleApi } from "@/utils/server-compat"
import { decodeVcsDiff } from "@/utils/vcs-diff-decoder"
import { decodeSessionList } from "./session-message-decoder"
const isAbortError = (error: unknown) =>
error !== null && typeof error === "object" && "name" in error && error.name === "AbortError"
@ -346,7 +348,7 @@ function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerS
throwOnError: true,
directory,
})
const api = createCompatibleApi({ protocol, current: currentApi, legacy })
const api = createCompatibleApi({ protocol, current: currentApi, legacy, decodeVcsDiff, decodeSessionList })
return {
server,
@ -432,6 +434,8 @@ function createDirSdkContext(directory: string, serverSDK: ServerSDKBase) {
current: serverSDK.currentApi,
legacy: (next) => serverSDK.createClient({ directory: next ?? directory, throwOnError: true }),
directory,
decodeVcsDiff,
decodeSessionList,
}),
event: emitter,
get url() {

View file

@ -22,6 +22,7 @@ import { normalizeSessionMessages } from "@/utils/session-message"
import { dropSessionCaches, pickSessionCacheEvictions, SESSION_CACHE_LIMIT } from "./global-sync/session-cache"
import { createV2SessionReducer, type V2SessionReduction } from "./server-session-v2-reducer"
import type { ServerApi } from "@/utils/server"
import type { DecodedLegacyMessagePage } from "./session-message-decode"
type MessageApi = ServerApi["message"]
@ -192,7 +193,7 @@ function reconcileFetched<T extends { id: string }>(
type ServerSessionOptions = {
retry?: typeof retry
protocol?: Promise<"v1" | "v2">
decodeMessages?: <T>(buffer: ArrayBuffer) => Promise<T>
decodeMessages?: (buffer: ArrayBuffer) => Promise<DecodedLegacyMessagePage>
}
export function createServerSession(
@ -582,14 +583,16 @@ export function createServerSession(
if (!options?.decodeMessages) return client.session.messages({ sessionID, limit, before })
const response = await client.session.messages({ sessionID, limit, before }, { parseAs: "arrayBuffer" })
if (!(response.data instanceof ArrayBuffer)) throw new Error("Session messages response is not an ArrayBuffer")
return {
...response,
data: await options.decodeMessages<NonNullable<Awaited<ReturnType<typeof client.session.messages>>["data"]>>(
response.data,
),
}
return { response, decoded: await options.decodeMessages(response.data) }
})
await yieldToMain()
if ("decoded" in response)
return {
...response.decoded,
sourceMode: before ? ("older" as const) : ("latest" as const),
cursor: response.response.response.headers.get("x-next-cursor") ?? undefined,
complete: !response.response.response.headers.get("x-next-cursor"),
}
const items = (response.data ?? []).filter((item) => !!item?.info?.id)
return {
session: items.map((item) => cleanMessage(item.info)).sort((a, b) => cmp(a.id, b.id)),

View file

@ -0,0 +1,42 @@
import { expect, test } from "bun:test"
import type { Message, Part, Session } from "@opencode-ai/sdk/v2/client"
import { decodeLegacyMessagePage, decodeLegacySessionList } from "./session-message-decode"
test("decodes and projects a legacy message page", () => {
const info = {
id: "message",
sessionID: "session",
role: "user",
time: { created: 1 },
agent: "build",
model: { providerID: "provider", modelID: "model" },
} as Message
const part = {
id: "part",
sessionID: "session",
messageID: info.id,
type: "text",
text: "hello",
} as Part
const result = decodeLegacyMessagePage(new TextEncoder().encode(JSON.stringify([{ info, parts: [part] }])).buffer)
expect(result.session).toEqual([info])
expect(result.part).toEqual([{ id: info.id, part: [part] }])
expect(result.source).toEqual([{ id: info.id, type: "user", text: "hello", time: info.time }])
})
test("decodes and projects a legacy session list", () => {
const session = {
id: "session",
projectID: "project",
directory: "/repo",
title: "Session",
version: "1",
time: { created: 1, updated: 1 },
} as Session
const result = decodeLegacySessionList(new TextEncoder().encode(JSON.stringify([session])).buffer)
expect(result).toEqual([
expect.objectContaining({ id: session.id, title: session.title, location: { directory: "/repo" } }),
])
})

View file

@ -0,0 +1,77 @@
import type { SessionInfo, SessionMessageInfo } from "@opencode-ai/client/promise"
import type { Message, Part, Session } from "@opencode-ai/sdk/v2/client"
import { message as cleanMessage } from "@/utils/diffs"
export type DecodedLegacyMessagePage = {
session: Message[]
part: { id: string; part: Part[] }[]
source: SessionMessageInfo[]
}
export function decodeLegacyMessagePage(buffer: ArrayBuffer): DecodedLegacyMessagePage {
const text = new TextDecoder().decode(buffer)
const items = (text ? (JSON.parse(text) as { info?: Message; parts?: Part[] }[]) : []).filter(
(item): item is { info: Message; parts: Part[] } => !!item.info?.id && Array.isArray(item.parts),
)
return {
session: items.map((item) => cleanMessage(item.info)).sort((a, b) => compare(a.id, b.id)),
part: items.map((item) => ({
id: item.info.id,
part: item.parts.filter((part) => !!part?.id).sort((a, b) => compare(a.id, b.id)),
})),
source: items
.slice()
.sort((a, b) => compare(a.info.id, b.info.id))
.map((item) =>
item.info.role === "user"
? {
id: item.info.id,
type: "user" as const,
text: item.parts.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("\n"),
time: item.info.time,
}
: {
id: item.info.id,
type: "assistant" as const,
agent: item.info.agent ?? item.info.mode,
model: { id: item.info.modelID, providerID: item.info.providerID, variant: item.info.variant },
content: [],
time: item.info.time,
},
),
}
}
export function decodeLegacySessionList(buffer: ArrayBuffer) {
const text = new TextDecoder().decode(buffer)
return (text ? (JSON.parse(text) as Session[]) : []).map(legacySessionInfo)
}
export function legacySessionInfo(session: Session): SessionInfo {
return {
id: session.id,
parentID: session.parentID,
projectID: session.projectID,
agent: session.agent,
model: session.model && {
id: session.model.id,
providerID: session.model.providerID,
variant: session.model.variant,
},
cost: session.cost ?? 0,
tokens: session.tokens ?? { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
time: session.time,
title: session.title,
location: { directory: session.directory, workspaceID: session.workspaceID },
subpath: session.path,
revert: session.revert && {
messageID: session.revert.messageID,
partID: session.revert.partID,
snapshot: session.revert.snapshot,
},
}
}
function compare(a: string, b: string) {
return a < b ? -1 : a > b ? 1 : 0
}

View file

@ -1,4 +1,6 @@
import SessionMessageDecoderWorkerUrl from "./session-message-decoder.worker.ts?worker&url"
import type { DecodedLegacyMessagePage } from "./session-message-decode"
import type { SessionInfo } from "@opencode-ai/client/promise"
type Response = { id: number; data?: unknown; error?: string }
@ -6,11 +8,19 @@ let worker: Worker | undefined
let nextID = 0
const pending = new Map<number, { resolve: (value: unknown) => void; reject: (error: Error) => void }>()
export function decodeSessionMessages<T>(buffer: ArrayBuffer) {
export function decodeSessionMessages(buffer: ArrayBuffer) {
return decode<DecodedLegacyMessagePage>("messages", buffer)
}
export function decodeSessionList(buffer: ArrayBuffer) {
return decode<SessionInfo[]>("sessions", buffer)
}
function decode<T>(type: "messages" | "sessions", buffer: ArrayBuffer) {
const id = ++nextID
return new Promise<T>((resolve, reject) => {
pending.set(id, { resolve: (value) => resolve(value as T), reject })
getWorker().postMessage({ id, buffer }, [buffer])
getWorker().postMessage({ id, type, buffer }, [buffer])
})
}

View file

@ -1,9 +1,16 @@
type DecoderRequest = { id: number; buffer: ArrayBuffer }
import { decodeLegacyMessagePage, decodeLegacySessionList } from "./session-message-decode"
type DecoderRequest = { id: number; type: "messages" | "sessions"; buffer: ArrayBuffer }
self.onmessage = (event: MessageEvent<DecoderRequest>) => {
try {
const text = new TextDecoder().decode(event.data.buffer)
self.postMessage({ id: event.data.id, data: text ? JSON.parse(text) : {} })
self.postMessage({
id: event.data.id,
data:
event.data.type === "messages"
? decodeLegacyMessagePage(event.data.buffer)
: decodeLegacySessionList(event.data.buffer),
})
} catch (error) {
self.postMessage({ id: event.data.id, error: error instanceof Error ? error.message : String(error) })
}

View file

@ -690,8 +690,12 @@ export default function Page() {
queryFn: mode
? () =>
sdk()
.api.vcs.diff({ location: { directory: sdk().directory }, mode: mode === "git" ? "working" : mode })
.then((result) => result.data)
.api.vcs.diff({
location: { directory: sdk().directory },
mode: mode === "git" ? "working" : mode,
context: 0,
})
.then((result) => result.data.map((diff) => ({ ...diff, patch: "" })))
.catch((error) => {
console.debug("[session-review] failed to load vcs diff", { mode, error })
return []

View file

@ -1,6 +1,8 @@
import { describe, expect, test } from "bun:test"
import { createApiForServer, createSdkForServer } from "./server"
import { createCompatibleApi } from "./server-compat"
import { decodeVcsDiffData } from "./vcs-diff-data"
import { decodeLegacySessionList } from "@/context/session-message-decode"
function setup(
protocol: "v1" | "v2" | Promise<"v1" | "v2">,
@ -48,6 +50,8 @@ function setup(
current: createApiForServer({ server, fetch: fetcher }),
legacy: (directory) => createSdkForServer({ server, fetch: fetcher, directory, throwOnError: true }),
directory: "/repo",
decodeVcsDiff: async (buffer) => decodeVcsDiffData(buffer),
decodeSessionList: async (buffer) => decodeLegacySessionList(buffer),
})
return { api, requests }
}

View file

@ -1,7 +1,8 @@
import type { ServerApi } from "./server"
import type { ServerProtocol } from "./server-protocol"
import type { AgentPartInput, FilePartInput, OpencodeClient, Session, TextPartInput } from "@opencode-ai/sdk/v2/client"
import type { AgentPartInput, FilePartInput, OpencodeClient, TextPartInput } from "@opencode-ai/sdk/v2/client"
import type {
FileDiffInfo,
Project,
ProjectCurrent,
SessionApi,
@ -15,6 +16,7 @@ import type {
SessionShellInput,
SessionShellOutput,
} from "@opencode-ai/client/promise"
import { legacySessionInfo } from "@/context/session-message-decode"
type LegacyClient = OpencodeClient
type LegacyFor = (directory?: string) => LegacyClient
@ -51,6 +53,8 @@ type CompatibleInput = {
current: ServerApi
legacy: LegacyFor
directory?: string
decodeVcsDiff: (buffer: ArrayBuffer) => Promise<FileDiffInfo[]>
decodeSessionList: (buffer: ArrayBuffer) => Promise<SessionInfo[]>
}
function mime(uri: string) {
@ -58,31 +62,6 @@ function mime(uri: string) {
return match?.[1] ?? "application/octet-stream"
}
function sessionInfo(session: Session): SessionInfo {
return {
id: session.id,
parentID: session.parentID,
projectID: session.projectID,
agent: session.agent,
model: session.model && {
id: session.model.id,
providerID: session.model.providerID,
variant: session.model.variant,
},
cost: session.cost ?? 0,
tokens: session.tokens ?? { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
time: session.time,
title: session.title,
location: { directory: session.directory, workspaceID: session.workspaceID },
subpath: session.path,
revert: session.revert && {
messageID: session.revert.messageID,
partID: session.revert.partID,
snapshot: session.revert.snapshot,
},
}
}
export function createCompatibleApi(input: CompatibleInput): CompatibleApi {
const v1 = createV1Api(input)
return lazyApi(
@ -148,29 +127,34 @@ function createV1Api(input: CompatibleInput): CompatibleApi {
search: value.search,
limit: value.limit,
},
options,
{ ...options, parseAs: "arrayBuffer" },
)
return { data: (result.data ?? []).map(sessionInfo), cursor: {} }
if (!(result.data instanceof ArrayBuffer)) throw new Error("Session list response is not an ArrayBuffer")
return { data: await input.decodeSessionList(result.data), cursor: {} }
}
const result = await legacy({ directory: value?.directory }).session.list({
directory: value?.directory,
roots: value?.parentID === null ? true : undefined,
search: value?.search,
limit: value?.limit,
})
return { data: (result.data ?? []).map(sessionInfo), cursor: {} }
const result = await legacy({ directory: value?.directory }).session.list(
{
directory: value?.directory,
roots: value?.parentID === null ? true : undefined,
search: value?.search,
limit: value?.limit,
},
{ parseAs: "arrayBuffer" },
)
if (!(result.data instanceof ArrayBuffer)) throw new Error("Session list response is not an ArrayBuffer")
return { data: await input.decodeSessionList(result.data), cursor: {} }
},
async create(value?: Parameters<ServerApi["session"]["create"]>[0]) {
const result = await legacy(value?.location ?? undefined).session.create({
directory: directory(value?.location ?? undefined),
})
if (!result.data) throw new Error("Failed to create session")
return sessionInfo(result.data)
return legacySessionInfo(result.data)
},
async get(value: Parameters<ServerApi["session"]["get"]>[0]) {
const result = await legacy().session.get(value)
if (!result.data) throw new Error(`Session not found: ${value.sessionID}`)
return sessionInfo(result.data)
return legacySessionInfo(result.data)
},
async active() {
const result = await legacy().session.status()
@ -192,7 +176,7 @@ function createV1Api(input: CompatibleInput): CompatibleApi {
async fork(value: Parameters<ServerApi["session"]["fork"]>[0]) {
const result = await legacy().session.fork(value)
if (!result.data) throw new Error("Failed to fork session")
return sessionInfo(result.data)
return legacySessionInfo(result.data)
},
async interrupt(value: Parameters<ServerApi["session"]["interrupt"]>[0]) {
await legacy().session.abort(value)
@ -341,20 +325,15 @@ function createV1Api(input: CompatibleInput): CompatibleApi {
return located(result.data ?? [], value?.location)
},
async diff(value: Parameters<ServerApi["vcs"]["diff"]>[0]) {
const result = await legacy(value.location).vcs.diff({
mode: value.mode === "working" ? "git" : value.mode,
context: value.context,
})
return located(
(result.data ?? []).map((file) => ({
file: file.file,
patch: file.patch ?? "",
additions: file.additions,
deletions: file.deletions,
status: file.status ?? "modified",
})),
value.location,
const result = await legacy(value.location).vcs.diff(
{
mode: value.mode === "working" ? "git" : value.mode,
context: value.context,
},
{ parseAs: "arrayBuffer" },
)
if (!(result.data instanceof ArrayBuffer)) throw new Error("VCS diff response is not an ArrayBuffer")
return located(await input.decodeVcsDiff(result.data), value.location)
},
},
file: {

View file

@ -0,0 +1,20 @@
import type { FileDiffInfo } from "@opencode-ai/client/promise"
export function decodeVcsDiffData(buffer: ArrayBuffer): FileDiffInfo[] {
const text = new TextDecoder().decode(buffer)
return (text ? JSON.parse(text) : []).map(
(file: {
file: string
patch?: string
additions: number
deletions: number
status?: "added" | "deleted" | "modified"
}) => ({
file: file.file,
patch: file.patch ?? "",
additions: file.additions,
deletions: file.deletions,
status: file.status ?? "modified",
}),
)
}

View file

@ -0,0 +1,61 @@
import type { FileDiffInfo } from "@opencode-ai/client/promise"
import VcsDiffDecoderWorkerUrl from "./vcs-diff-decoder.worker.ts?worker&url"
type Response = { id: number; data?: FileDiffInfo[]; error?: string }
let worker: Worker | undefined
let nextID = 0
const pending = new Map<number, { resolve: (value: FileDiffInfo[]) => void; reject: (error: Error) => void }>()
let lastInput = 0
document.addEventListener(
"beforeinput",
() => {
lastInput = performance.now()
},
{ capture: true },
)
export function decodeVcsDiff(buffer: ArrayBuffer) {
const id = ++nextID
return new Promise<FileDiffInfo[]>((resolve, reject) => {
pending.set(id, { resolve, reject })
getWorker().postMessage({ id, buffer }, [buffer])
})
}
function getWorker() {
if (worker) return worker
worker = new Worker(VcsDiffDecoderWorkerUrl, { type: "module" })
worker.onmessage = (event: MessageEvent<Response>) => {
const request = pending.get(event.data.id)
if (!request) return
pending.delete(event.data.id)
if (event.data.error) {
request.reject(new Error(event.data.error))
return
}
resolveWhenInputIdle(request.resolve, event.data.data ?? [])
}
worker.onerror = (event) => {
const error = new Error(event.message)
pending.forEach((request) => request.reject(error))
pending.clear()
worker?.terminate()
worker = undefined
}
return worker
}
function resolveWhenInputIdle(resolve: (value: FileDiffInfo[]) => void, value: FileDiffInfo[], initial = true) {
const active = document.activeElement
const editing =
active instanceof HTMLInputElement ||
active instanceof HTMLTextAreaElement ||
(active instanceof HTMLElement && active.isContentEditable)
const delay = Math.max(lastInput + 100 - performance.now(), initial && editing ? 100 : 0)
if (delay <= 0) {
resolve(value)
return
}
setTimeout(() => resolveWhenInputIdle(resolve, value, false), delay)
}

View file

@ -0,0 +1,13 @@
import { decodeVcsDiffData } from "./vcs-diff-data"
type Request = { id: number; buffer: ArrayBuffer }
self.onmessage = (event: MessageEvent<Request>) => {
try {
self.postMessage({ id: event.data.id, data: decodeVcsDiffData(event.data.buffer) })
} catch (error) {
self.postMessage({ id: event.data.id, error: error instanceof Error ? error.message : String(error) })
}
}
export {}