mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-08 13:23:20 +00:00
fix(app): parse session messages off thread
This commit is contained in:
parent
1277ceb426
commit
a47dabff22
4 changed files with 67 additions and 3 deletions
|
|
@ -189,7 +189,11 @@ function reconcileFetched<T extends { id: string }>(
|
|||
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?: <T>(buffer: ArrayBuffer) => Promise<T>
|
||||
}
|
||||
|
||||
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<NonNullable<Awaited<ReturnType<typeof client.session.messages>>["data"]>>(
|
||||
response.data,
|
||||
),
|
||||
}
|
||||
})
|
||||
await yieldToMain()
|
||||
const items = (response.data ?? []).filter((item) => !!item?.info?.id)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
38
packages/app/src/context/session-message-decoder.ts
Normal file
38
packages/app/src/context/session-message-decoder.ts
Normal file
|
|
@ -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<number, { resolve: (value: unknown) => void; reject: (error: Error) => void }>()
|
||||
|
||||
export function decodeSessionMessages<T>(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])
|
||||
})
|
||||
}
|
||||
|
||||
function getWorker() {
|
||||
if (worker) return worker
|
||||
worker = new Worker(SessionMessageDecoderWorkerUrl, { 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
|
||||
}
|
||||
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
|
||||
}
|
||||
12
packages/app/src/context/session-message-decoder.worker.ts
Normal file
12
packages/app/src/context/session-message-decoder.worker.ts
Normal file
|
|
@ -0,0 +1,12 @@
|
|||
type DecoderRequest = { id: number; 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) : {} })
|
||||
} catch (error) {
|
||||
self.postMessage({ id: event.data.id, error: error instanceof Error ? error.message : String(error) })
|
||||
}
|
||||
}
|
||||
|
||||
export {}
|
||||
Loading…
Add table
Add a link
Reference in a new issue