From a47dabff22a4313f59420c4a4bc72a0cbaa3399f Mon Sep 17 00:00:00 2001 From: LukeParkerDev <10430890+Hona@users.noreply.github.com> Date: Tue, 4 Aug 2026 18:33:53 +1000 Subject: [PATCH] fix(app): parse session messages off thread --- packages/app/src/context/server-session.ts | 18 +++++++-- packages/app/src/context/server-sync.tsx | 2 + .../src/context/session-message-decoder.ts | 38 +++++++++++++++++++ .../context/session-message-decoder.worker.ts | 12 ++++++ 4 files changed, 67 insertions(+), 3 deletions(-) create mode 100644 packages/app/src/context/session-message-decoder.ts create mode 100644 packages/app/src/context/session-message-decoder.worker.ts diff --git a/packages/app/src/context/server-session.ts b/packages/app/src/context/server-session.ts index 967e76660fa..798f2952856 100644 --- a/packages/app/src/context/server-session.ts +++ b/packages/app/src/context/server-session.ts @@ -189,7 +189,11 @@ function reconcileFetched( return [...result.values()].sort((a, b) => cmp(a.id, b.id)) } -type ServerSessionOptions = { retry?: typeof retry; protocol?: Promise<"v1" | "v2"> } +type ServerSessionOptions = { + retry?: typeof retry + protocol?: Promise<"v1" | "v2"> + decodeMessages?: (buffer: ArrayBuffer) => Promise +} export function createServerSession( client: OpencodeClient, @@ -573,9 +577,17 @@ export function createServerSession( complete: response.data.length === 0, } } - const response = await (options?.retry ?? retry)(() => { + const response = await (options?.retry ?? retry)(async () => { onAttempt?.() - return client.session.messages({ sessionID, limit, before }) + 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>["data"]>>( + response.data, + ), + } }) await yieldToMain() const items = (response.data ?? []).filter((item) => !!item?.info?.id) diff --git a/packages/app/src/context/server-sync.tsx b/packages/app/src/context/server-sync.tsx index 13a0b74bc6f..82fa994eb28 100644 --- a/packages/app/src/context/server-sync.tsx +++ b/packages/app/src/context/server-sync.tsx @@ -59,6 +59,7 @@ import type { } from "@opencode-ai/client/promise" import { toggleMcp } from "./global-sync/mcp" import { createServerSession, type ServerSession } from "./server-session" +import { decodeSessionMessages } from "./session-message-decoder" type GlobalStore = { ready: boolean @@ -226,6 +227,7 @@ export function createServerSyncContextInner(serverSDK: ServerSDK) { const session = createServerSession(serverSDK.client, serverSDK.api.session, serverSDK.api.message, { protocol: serverSDK.protocol, + decodeMessages: decodeSessionMessages, }) const queryOptionsApi = makeQueryOptionsApi( serverSDK.scope, diff --git a/packages/app/src/context/session-message-decoder.ts b/packages/app/src/context/session-message-decoder.ts new file mode 100644 index 00000000000..691fa00e6fd --- /dev/null +++ b/packages/app/src/context/session-message-decoder.ts @@ -0,0 +1,38 @@ +import SessionMessageDecoderWorkerUrl from "./session-message-decoder.worker.ts?worker&url" + +type Response = { id: number; data?: unknown; error?: string } + +let worker: Worker | undefined +let nextID = 0 +const pending = new Map void; reject: (error: Error) => void }>() + +export function decodeSessionMessages(buffer: ArrayBuffer) { + const id = ++nextID + return new Promise((resolve, reject) => { + pending.set(id, { resolve: (value) => resolve(value as T), reject }) + getWorker().postMessage({ id, buffer }, [buffer]) + }) +} + +function getWorker() { + if (worker) return worker + worker = new Worker(SessionMessageDecoderWorkerUrl, { type: "module" }) + worker.onmessage = (event: MessageEvent) => { + 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 + } + 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 +} diff --git a/packages/app/src/context/session-message-decoder.worker.ts b/packages/app/src/context/session-message-decoder.worker.ts new file mode 100644 index 00000000000..024fee72d7d --- /dev/null +++ b/packages/app/src/context/session-message-decoder.worker.ts @@ -0,0 +1,12 @@ +type DecoderRequest = { id: number; buffer: ArrayBuffer } + +self.onmessage = (event: MessageEvent) => { + try { + const text = new TextDecoder().decode(event.data.buffer) + self.postMessage({ id: event.data.id, data: text ? JSON.parse(text) : {} }) + } catch (error) { + self.postMessage({ id: event.data.id, error: error instanceof Error ? error.message : String(error) }) + } +} + +export {}