mirror of
https://github.com/anomalyco/opencode.git
synced 2026-09-28 13:33:35 +08:00
refactor(desktop): send IPC messages by structured clone (#49701)
This commit is contained in:
@@ -456,7 +456,6 @@
|
||||
"@lydell/node-pty-linux-x64": "1.2.0-beta.12",
|
||||
"@lydell/node-pty-win32-arm64": "1.2.0-beta.12",
|
||||
"@lydell/node-pty-win32-x64": "1.2.0-beta.12",
|
||||
"msgpackr-extract": "3.0.4",
|
||||
},
|
||||
},
|
||||
"packages/enterprise": {
|
||||
|
||||
@@ -48,7 +48,8 @@ test("does not package external copies of bundled dependencies", () => {
|
||||
expect(pkg.devDependencies.effect).toBe("catalog:")
|
||||
expect(pkg.devDependencies["@effect/platform-node"]).toBe("catalog:")
|
||||
expect(pkg.devDependencies["drizzle-orm"]).toBe("catalog:")
|
||||
expect(pkg.optionalDependencies["msgpackr-extract"]).toBe("3.0.4")
|
||||
// IPC crosses the port by structured clone; no MessagePack runtime or native accelerator ships.
|
||||
expect(Object.keys(pkg.optionalDependencies)).not.toContain("msgpackr-extract")
|
||||
})
|
||||
|
||||
test("keeps PTY binaries without stale native packaging", () => {
|
||||
@@ -101,6 +102,6 @@ test("bundles one Effect runtime and Drizzle while keeping native dependencies e
|
||||
expect(imports).toContain("node:sqlite")
|
||||
expect(chunks.some((chunk) => chunk.dynamicImports.includes("@zip.js/zip.js"))).toBe(true)
|
||||
expect(imports).toContain(`@lydell/node-pty-${process.platform}-${process.arch}`)
|
||||
expect(modules.some((id) => id.includes("/node_modules/msgpackr-extract/"))).toBe(false)
|
||||
expect(chunks.some((chunk) => chunk.code.includes("msgpackr-extract"))).toBe(true)
|
||||
expect(modules.some((id) => id.includes("/node_modules/msgpackr"))).toBe(false)
|
||||
expect(chunks.some((chunk) => chunk.code.includes("msgpackr"))).toBe(false)
|
||||
}, 30_000)
|
||||
|
||||
@@ -58,9 +58,9 @@ const require = __cjs_mod__.createRequire(import.meta.url);
|
||||
},
|
||||
},
|
||||
externalizeDeps: {
|
||||
// Bundle the Effect family together; native MessagePack acceleration stays optional and external.
|
||||
// Bundle the Effect family together.
|
||||
exclude: ["effect", "@effect/platform-node", "@effect/platform-node-shared", "drizzle-orm"],
|
||||
include: [nodePtyPkg, "msgpackr-extract"],
|
||||
include: [nodePtyPkg],
|
||||
},
|
||||
},
|
||||
plugins: [
|
||||
|
||||
@@ -68,7 +68,6 @@
|
||||
"@lydell/node-pty-linux-arm64": "1.2.0-beta.12",
|
||||
"@lydell/node-pty-linux-x64": "1.2.0-beta.12",
|
||||
"@lydell/node-pty-win32-arm64": "1.2.0-beta.12",
|
||||
"@lydell/node-pty-win32-x64": "1.2.0-beta.12",
|
||||
"msgpackr-extract": "3.0.4"
|
||||
"@lydell/node-pty-win32-x64": "1.2.0-beta.12"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,17 +3,22 @@ import { EventEmitter } from "node:events"
|
||||
import { MessageChannel } from "node:worker_threads"
|
||||
import type { MessagePortMain, WebContents } from "electron"
|
||||
import { Context, Effect, Layer, ManagedRuntime, Option, Queue, Schema, Stream } from "effect"
|
||||
import { Rpc, RpcClient, RpcClientError, RpcGroup, RpcMessage, RpcSerialization, RpcServer } from "effect/unstable/rpc"
|
||||
import { Rpc, RpcClient, RpcClientError, RpcGroup, RpcMessage, RpcServer } from "effect/unstable/rpc"
|
||||
import { Transferable } from "effect/unstable/workers"
|
||||
import { IpcPortHandoff, IpcServerProtocolLive } from "./ipc-transport"
|
||||
|
||||
describe("desktop RPC transport", () => {
|
||||
test("keeps multiple renderer ports independent", async () => {
|
||||
let received: unknown
|
||||
const handlers = TestRpcs.toLayer(
|
||||
Effect.gen(function* () {
|
||||
const handoff = yield* IpcPortHandoff
|
||||
return TestRpcs.of({
|
||||
"test.focused": (_request, context) => Effect.succeed(handoff.sender(context.client.id)?.id === 1),
|
||||
"test.blob.put": ({ data }) => Effect.succeed([...data].join(",")),
|
||||
"test.blob.put": ({ data }) => {
|
||||
received = data
|
||||
return Effect.succeed([...data].join(","))
|
||||
},
|
||||
"test.blob.get": () => Effect.succeed(new Uint8Array([3, 1, 4])),
|
||||
"test.events": () => Stream.make(new TestEvent({ value: "session.new" })),
|
||||
})
|
||||
@@ -34,6 +39,8 @@ describe("desktop RPC transport", () => {
|
||||
expect(focused).toBe(true)
|
||||
expect(unfocused).toBe(false)
|
||||
expect(await putBlob(firstClient, new Uint8Array([2, 7, 1]))).toBe("2,7,1")
|
||||
// Binary payloads arrive as bytes, not as base64 text.
|
||||
expect(received).toBeInstanceOf(Uint8Array)
|
||||
expect(await getBlob(firstClient)).toEqual(new Uint8Array([3, 1, 4]))
|
||||
expect(await firstEvent(firstClient)).toEqual(new TestEvent({ value: "session.new" }))
|
||||
|
||||
@@ -54,8 +61,8 @@ describe("desktop RPC transport", () => {
|
||||
class TestEvent extends Schema.TaggedClass<TestEvent>()("TestEvent", { value: Schema.String }) {}
|
||||
const TestRpcs = RpcGroup.make(
|
||||
Rpc.make("test.focused", { success: Schema.Boolean }),
|
||||
Rpc.make("test.blob.put", { payload: { data: Schema.Uint8Array }, success: Schema.String }),
|
||||
Rpc.make("test.blob.get", { success: Schema.Uint8Array }),
|
||||
Rpc.make("test.blob.put", { payload: { data: Transferable.Uint8Array }, success: Schema.String }),
|
||||
Rpc.make("test.blob.get", { success: Transferable.Uint8Array }),
|
||||
Rpc.make("test.events", { success: TestEvent, stream: true }),
|
||||
)
|
||||
type TestRpcClient = RpcClient.FromGroup<typeof TestRpcs, RpcClientError.RpcClientError>
|
||||
@@ -109,13 +116,9 @@ function clientProtocol(port: MessagePort) {
|
||||
RpcClient.Protocol,
|
||||
RpcClient.Protocol.make(
|
||||
Effect.fnUntraced(function* (writeResponse, clientIds) {
|
||||
const serialization = yield* RpcSerialization.RpcSerialization
|
||||
const parser = serialization.makeUnsafe()
|
||||
const inbound = yield* Queue.unbounded<RpcMessage.FromServerEncoded>()
|
||||
const onMessage = (event: MessageEvent) =>
|
||||
parser
|
||||
.decode(event.data)
|
||||
.forEach((message) => Queue.offerUnsafe(inbound, message as RpcMessage.FromServerEncoded))
|
||||
Queue.offerUnsafe(inbound, event.data as RpcMessage.FromServerEncoded)
|
||||
port.addEventListener("message", onMessage)
|
||||
port.start()
|
||||
yield* Effect.addFinalizer(() =>
|
||||
@@ -131,18 +134,15 @@ function clientProtocol(port: MessagePort) {
|
||||
Effect.forkScoped,
|
||||
)
|
||||
return {
|
||||
codecFor: serialization.codecFor,
|
||||
codecFor: Schema.toCodecJson,
|
||||
send: (_clientId: number, request: RpcMessage.FromClientEncoded) =>
|
||||
Effect.sync(() => {
|
||||
const encoded = parser.encode(request)
|
||||
if (encoded !== undefined) port.postMessage(encoded)
|
||||
}),
|
||||
Effect.sync(() => port.postMessage(request)),
|
||||
supportsAck: true,
|
||||
supportsTransferables: false,
|
||||
}
|
||||
}),
|
||||
),
|
||||
).pipe(Layer.provide(RpcSerialization.layerMsgPack))
|
||||
)
|
||||
}
|
||||
|
||||
function sender(id: number) {
|
||||
|
||||
@@ -1,14 +1,12 @@
|
||||
import type { MessagePortMain, WebContents } from "electron"
|
||||
import { Context, Effect, Layer, Option, Queue, Stream } from "effect"
|
||||
import { RpcMessage, RpcSerialization, RpcServer } from "effect/unstable/rpc"
|
||||
import { createIpcCodec } from "../shared/ipc-codec"
|
||||
import { Context, Effect, Layer, Option, Queue, Schema, Stream } from "effect"
|
||||
import { RpcMessage, RpcServer } from "effect/unstable/rpc"
|
||||
import { bindIpcEvents } from "./ipc-events"
|
||||
|
||||
type PortBinding = {
|
||||
readonly id: number
|
||||
readonly sender: WebContents
|
||||
readonly port: MessagePortMain
|
||||
readonly parser: ReturnType<typeof createIpcCodec>
|
||||
readonly onMessage: (event: Electron.MessageEvent) => void
|
||||
readonly onClose: () => void
|
||||
readonly unbindEvents: Effect.Effect<void>
|
||||
@@ -21,6 +19,9 @@ type Handoff = {
|
||||
|
||||
export class IpcPortHandoff extends Context.Service<IpcPortHandoff, Handoff>()("opencode/desktop/IpcPortHandoff") {}
|
||||
|
||||
// Messages cross the port by structured clone, like Effect's worker protocol: no serialization
|
||||
// layer, so binary payloads stay binary and nothing is packed into a shared buffer. Electron's
|
||||
// MessagePortMain can only transfer ports, so byte payloads are cloned in both directions.
|
||||
export const IpcServerProtocolLive = Layer.unwrap(
|
||||
Effect.gen(function* () {
|
||||
const handoffs = yield* Queue.unbounded<readonly [WebContents, MessagePortMain]>()
|
||||
@@ -31,7 +32,6 @@ export const IpcServerProtocolLive = Layer.unwrap(
|
||||
RpcServer.Protocol,
|
||||
RpcServer.Protocol.make(
|
||||
Effect.fnUntraced(function* (writeRequest) {
|
||||
const serialization = yield* RpcSerialization.RpcSerialization
|
||||
const disconnects = yield* Queue.unbounded<number>()
|
||||
const inbound = yield* Queue.unbounded<readonly [number, RpcMessage.FromClientEncoded]>()
|
||||
const runFork = Effect.runForkWith(yield* Effect.context())
|
||||
@@ -59,21 +59,12 @@ export const IpcServerProtocolLive = Layer.unwrap(
|
||||
}
|
||||
|
||||
const id = nextClientId++
|
||||
const parser = createIpcCodec(serialization)
|
||||
const onMessage = (event: Electron.MessageEvent) => {
|
||||
try {
|
||||
parser
|
||||
.decode(event.data)
|
||||
.forEach((message) =>
|
||||
Queue.offerUnsafe(inbound, [id, message as RpcMessage.FromClientEncoded] as const),
|
||||
)
|
||||
} catch {
|
||||
return
|
||||
}
|
||||
Queue.offerUnsafe(inbound, [id, event.data as RpcMessage.FromClientEncoded] as const)
|
||||
}
|
||||
const onClose = () => runFork(disconnect(id))
|
||||
const unbindEvents = yield* bindIpcEvents(sender.id)
|
||||
const binding = { id, sender, port, parser, onMessage, onClose, unbindEvents }
|
||||
const binding = { id, sender, port, onMessage, onClose, unbindEvents }
|
||||
bindings.set(id, binding)
|
||||
senderBindings.set(sender.id, id)
|
||||
port.on("message", onMessage)
|
||||
@@ -93,14 +84,11 @@ export const IpcServerProtocolLive = Layer.unwrap(
|
||||
yield* Effect.addFinalizer(() => Effect.forEach([...bindings.keys()], disconnect, { discard: true }))
|
||||
|
||||
return {
|
||||
codecFor: serialization.codecFor,
|
||||
codecFor: Schema.toCodecJson,
|
||||
disconnects,
|
||||
send: (clientId, response) =>
|
||||
Effect.sync(() => {
|
||||
const binding = bindings.get(clientId)
|
||||
if (!binding) return
|
||||
const encoded = binding.parser.encode(response)
|
||||
if (encoded !== undefined) binding.port.postMessage(encoded)
|
||||
bindings.get(clientId)?.port.postMessage(response)
|
||||
}),
|
||||
end: disconnect,
|
||||
clientIds: Effect.sync(() => new Set(bindings.keys())),
|
||||
@@ -124,4 +112,4 @@ export const IpcServerProtocolLive = Layer.unwrap(
|
||||
}),
|
||||
)
|
||||
}),
|
||||
).pipe(Layer.provide(RpcSerialization.layerMsgPack))
|
||||
)
|
||||
|
||||
@@ -98,8 +98,10 @@ export function createDraftStore(
|
||||
orphans = true
|
||||
return id
|
||||
},
|
||||
getBlob(id: string): Uint8Array | null {
|
||||
return db.select({ data: blobs.data }).from(blobs).where(eq(blobs.id, id)).get()?.data ?? null
|
||||
getBlob(id: string): Uint8Array<ArrayBuffer> | null {
|
||||
const data = db.select({ data: blobs.data }).from(blobs).where(eq(blobs.id, id)).get()?.data
|
||||
// node:sqlite allocates a dedicated ArrayBuffer per BLOB column value.
|
||||
return data ? (data as Uint8Array<ArrayBuffer>) : null
|
||||
},
|
||||
flush: writer.flush,
|
||||
close: writer.close,
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
import { Context, Effect, Layer, ManagedRuntime, Queue, Stream } from "effect"
|
||||
import { RpcClient, RpcMessage, RpcSerialization } from "effect/unstable/rpc"
|
||||
import { createIpcCodec } from "../shared/ipc-codec"
|
||||
import { Context, Effect, Layer, ManagedRuntime, Queue, Schema, Stream } from "effect"
|
||||
import { RpcClient, RpcMessage } from "effect/unstable/rpc"
|
||||
import { DesktopRpcs, type DesktopRpcClient } from "../shared/ipc-rpc"
|
||||
import type { DesktopEvent } from "../shared/ipc-rpc/events"
|
||||
import { IpcTransportPort } from "../shared/ipc-transport"
|
||||
@@ -80,22 +79,17 @@ export function listen<Tag extends EventTag>(tag: Tag, listener: (value: EventVa
|
||||
}
|
||||
}
|
||||
|
||||
// Structured clone over the port, like Effect's worker protocol: no serialization layer, so binary
|
||||
// payloads stay binary. Buffers are cloned rather than transferred: Electron's MessagePortMain
|
||||
// drops transferred ArrayBuffers, so a request carrying one would never arrive.
|
||||
function clientProtocol(value: MessagePort) {
|
||||
return Layer.effect(
|
||||
RpcClient.Protocol,
|
||||
RpcClient.Protocol.make(
|
||||
Effect.fnUntraced(function* (writeResponse, clientIds) {
|
||||
const serialization = yield* RpcSerialization.RpcSerialization
|
||||
const parser = createIpcCodec(serialization)
|
||||
const inbound = yield* Queue.unbounded<RpcMessage.FromServerEncoded>()
|
||||
const onMessage = (event: MessageEvent) => {
|
||||
try {
|
||||
parser
|
||||
.decode(event.data)
|
||||
.forEach((message) => Queue.offerUnsafe(inbound, message as RpcMessage.FromServerEncoded))
|
||||
} catch {
|
||||
return
|
||||
}
|
||||
Queue.offerUnsafe(inbound, event.data as RpcMessage.FromServerEncoded)
|
||||
}
|
||||
value.addEventListener("message", onMessage)
|
||||
value.start()
|
||||
@@ -112,16 +106,15 @@ function clientProtocol(value: MessagePort) {
|
||||
Effect.forkScoped,
|
||||
)
|
||||
return {
|
||||
codecFor: serialization.codecFor,
|
||||
codecFor: Schema.toCodecJson,
|
||||
send: (_clientId, request) =>
|
||||
Effect.sync(() => {
|
||||
const encoded = parser.encode(request)
|
||||
if (encoded !== undefined) value.postMessage(encoded)
|
||||
value.postMessage(request)
|
||||
}),
|
||||
supportsAck: true,
|
||||
supportsTransferables: false,
|
||||
}
|
||||
}),
|
||||
),
|
||||
).pipe(Layer.provide(RpcSerialization.layerMsgPack))
|
||||
)
|
||||
}
|
||||
|
||||
@@ -1,61 +0,0 @@
|
||||
import { describe, expect, test } from "bun:test"
|
||||
import { RpcSerialization } from "effect/unstable/rpc"
|
||||
import { createIpcCodec, ipcLargeMessageBytes } from "./ipc-codec"
|
||||
|
||||
type Serialization = RpcSerialization.RpcSerialization["Service"]
|
||||
|
||||
const msgpack = RpcSerialization.makeMsgPack()
|
||||
|
||||
function counting(serialization: Serialization) {
|
||||
let created = 0
|
||||
const counted: Serialization = {
|
||||
...serialization,
|
||||
makeUnsafe: () => {
|
||||
created++
|
||||
return serialization.makeUnsafe()
|
||||
},
|
||||
}
|
||||
return { created: () => created, serialization: counted }
|
||||
}
|
||||
|
||||
describe("ipc codec", () => {
|
||||
test("posts an exact-size buffer instead of a view into the shared target", () => {
|
||||
const codec = createIpcCodec(msgpack)
|
||||
// The msgpack target starts at 8 KiB, so a view would drag a larger backing store along.
|
||||
const encoded = codec.encode({ _tag: "Request", id: "1", tag: "Ping", payload: {} })
|
||||
expect(encoded).toBeInstanceOf(Uint8Array)
|
||||
const bytes = encoded as Uint8Array
|
||||
expect(bytes.byteOffset).toBe(0)
|
||||
expect(bytes.buffer.byteLength).toBe(bytes.byteLength)
|
||||
})
|
||||
|
||||
test("copies Node Buffers, whose slice is only a view", () => {
|
||||
const backing = new ArrayBuffer(1024 * 1024)
|
||||
const view = Buffer.from(backing, 16, 32)
|
||||
view.fill(7)
|
||||
const stub: Serialization = { ...msgpack, makeUnsafe: () => ({ decode: () => [], encode: () => view }) }
|
||||
const encoded = createIpcCodec(stub).encode({}) as Uint8Array
|
||||
expect(encoded.buffer).not.toBe(backing)
|
||||
expect(encoded.buffer.byteLength).toBe(32)
|
||||
expect([...encoded]).toEqual(Array(32).fill(7))
|
||||
})
|
||||
|
||||
test("replaces the encoder after a large message and keeps the decoder", () => {
|
||||
const spy = counting(msgpack)
|
||||
const codec = createIpcCodec(spy.serialization)
|
||||
expect(spy.created()).toBe(2)
|
||||
codec.encode({ small: true })
|
||||
expect(spy.created()).toBe(2)
|
||||
codec.encode({ large: "x".repeat(ipcLargeMessageBytes) })
|
||||
expect(spy.created()).toBe(3)
|
||||
codec.decode(codec.encode({ after: 1 }) as Uint8Array)
|
||||
expect(spy.created()).toBe(3)
|
||||
})
|
||||
|
||||
test("round-trips messages through the wrapped serialization", () => {
|
||||
const client = createIpcCodec(msgpack)
|
||||
const server = createIpcCodec(msgpack)
|
||||
const message = { _tag: "Request", id: "7", tag: "DraftsSet", payload: { key: "k", value: "v" } }
|
||||
expect(server.decode(client.encode(message) as Uint8Array)).toEqual([message])
|
||||
})
|
||||
})
|
||||
@@ -1,25 +0,0 @@
|
||||
import type { RpcSerialization } from "effect/unstable/rpc"
|
||||
|
||||
// After a message this large the encoder is replaced so its grown target buffer can be collected.
|
||||
export const ipcLargeMessageBytes = 1024 * 1024
|
||||
|
||||
// The MessagePack parser packs into one shared, grow-only target buffer and returns a view into it.
|
||||
// Posting that view structured-clones the whole backing buffer, so after one large message every
|
||||
// later message (even a 55-byte ack) would copy the full grown buffer across processes on each
|
||||
// send. Encoding through this wrapper posts an exact-size copy instead. Decoding keeps a single
|
||||
// parser for the connection: record structures the peer defined inline must stay known.
|
||||
export function createIpcCodec(serialization: RpcSerialization.RpcSerialization["Service"]) {
|
||||
const decoder = serialization.makeUnsafe()
|
||||
let encoder = serialization.makeUnsafe()
|
||||
return {
|
||||
decode: (bytes: Uint8Array | string) => decoder.decode(bytes),
|
||||
encode(message: unknown) {
|
||||
const encoded = encoder.encode(message)
|
||||
if (!(encoded instanceof Uint8Array)) return encoded
|
||||
// Not `.slice()`: in the main process the packer hands out a Node Buffer, whose slice is a view.
|
||||
const copy = new Uint8Array(encoded)
|
||||
if (copy.byteLength > ipcLargeMessageBytes) encoder = serialization.makeUnsafe()
|
||||
return copy
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
import { Schema } from "effect"
|
||||
import { Rpc, RpcGroup } from "effect/unstable/rpc"
|
||||
import { Transferable } from "effect/unstable/workers"
|
||||
|
||||
const OptionalString = Schema.optional(Schema.String)
|
||||
const PickerOptions = Schema.Struct({
|
||||
@@ -18,7 +19,7 @@ const PickedFiles = Schema.Struct({
|
||||
token: Schema.String,
|
||||
files: Schema.Array(Schema.Struct({ path: Schema.String, name: Schema.String, size: Schema.Number })),
|
||||
})
|
||||
const ClipboardImage = Schema.Struct({ buffer: Schema.Uint8Array, width: Schema.Number, height: Schema.Number })
|
||||
const ClipboardImage = Schema.Struct({ buffer: Transferable.Uint8Array, width: Schema.Number, height: Schema.Number })
|
||||
|
||||
export const FilesOpenDirectoryPicker = Rpc.make("FilesOpenDirectoryPicker", {
|
||||
payload: { options: Schema.optional(PickerOptions) },
|
||||
@@ -30,7 +31,7 @@ export const FilesOpenFilePicker = Rpc.make("FilesOpenFilePicker", {
|
||||
})
|
||||
export const FilesReadPickedFile = Rpc.make("FilesReadPickedFile", {
|
||||
payload: { token: Schema.String, path: Schema.String },
|
||||
success: Schema.Uint8Array,
|
||||
success: Transferable.Uint8Array,
|
||||
})
|
||||
export const FilesReleasePickedFiles = Rpc.make("FilesReleasePickedFiles", {
|
||||
payload: { token: Schema.String },
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { Schema } from "effect"
|
||||
import { Rpc, RpcGroup } from "effect/unstable/rpc"
|
||||
import { Transferable } from "effect/unstable/workers"
|
||||
|
||||
export const StorageItems = Rpc.make("StorageItems", {
|
||||
payload: { name: Schema.String },
|
||||
@@ -24,12 +25,12 @@ export const DraftsSet = Rpc.make("DraftsSet", {
|
||||
})
|
||||
export const DraftsDelete = Rpc.make("DraftsDelete", { payload: { key: Schema.String } })
|
||||
export const DraftsPutBlob = Rpc.make("DraftsPutBlob", {
|
||||
payload: { data: Schema.Uint8Array },
|
||||
payload: { data: Transferable.Uint8Array },
|
||||
success: Schema.String,
|
||||
})
|
||||
export const DraftsGetBlob = Rpc.make("DraftsGetBlob", {
|
||||
payload: { id: Schema.String },
|
||||
success: Schema.NullOr(Schema.Uint8Array),
|
||||
success: Schema.NullOr(Transferable.Uint8Array),
|
||||
})
|
||||
|
||||
export const StorageRpcs = RpcGroup.make(
|
||||
|
||||
Reference in New Issue
Block a user