diff --git a/packages/client/README.md b/packages/client/README.md index 8c53e47c7e..06d0dd17ea 100644 --- a/packages/client/README.md +++ b/packages/client/README.md @@ -1,17 +1,57 @@ # @opencode-ai/client -Private generation target for clients derived directly from OpenCode's authoritative Effect `HttpApi`. +Promise and Effect clients derived from OpenCode's authoritative Effect `HttpApi`, plus handwritten Node transports. ## Entrypoints - `@opencode-ai/client`: zero-Effect Promise client using `fetch`. +- `@opencode-ai/client/node`: Promise client plus Node-hosted browser attachments. - `@opencode-ai/client/effect`: rich Effect network client using an environment-provided `HttpClient`. The generated surface includes every standard HTTP group from Server's concrete API. The build compiler reads `@opencode-ai/server/api`; the generated Effect runtime imports a client-local projection built from Protocol, with a generation-equivalence test preventing transport drift. Custom transports such as the PTY WebSocket connection remain outside the generic HTTP client. Run `bun run generate` after changing the contract and `bun run check:generated` to detect committed-output drift. The Effect entrypoint uses canonical decoded values such as `Session.ID`, `Location.Ref`, and `Prompt`. These datatypes come from the lightweight `@opencode-ai/schema` package and are re-exported so callers depend only on the client surface. Protocol owns endpoint construction and middleware placement; Server supplies the concrete middleware keys used by the build-time API. -The Promise root remains structural and has no Core or Effect runtime dependency. `/effect` depends only on Effect, Schema, and Protocol and is browser-bundle safe. Bundle-boundary tests enforce both import graphs. +The Promise root remains structural and has no Core, Effect, Schema, Protocol, or WebSocket runtime dependency. `/node` adds Effect, Schema, Protocol, and `ws`, but never Core or Server. `/effect` depends only on Effect, Schema, and Protocol and is browser-bundle safe. Bundle-boundary tests enforce these import graphs. + +## Node browser attachments + +The Node entrypoint owns the control connection, Session lease, authenticated proxy, and network tunnels. Consumers provide a browser adapter once with `BrowserDriver.define`; normal attachment calls only provide a Session ID and that descriptor. + +```ts +import { BrowserDriver, BrowserDriverError, OpenCode } from "@opencode-ai/client/node" + +const driver = BrowserDriver.define(async ({ proxy, signal }) => { + const browser = await launchBrowser({ proxy, signal }) + return { + resource: browser, + state: () => browser.state(), + subscribe: (listener) => browser.subscribe(listener), + execute: async (command, options) => { + try { + return await browser.execute(command, options) + } catch (cause) { + throw new BrowserDriverError("internal", "Browser command failed", { cause }) + } + }, + dispose: () => browser.close(), + } +}) + +const client = OpenCode.make({ + baseUrl: "https://opencode.example", + headers: { authorization: `Basic ${credentials}` }, +}) +const attachment = await client.browser.attach({ sessionID, driver }) + +await attachment.close() +``` + +`attach` resolves only after the server acknowledges the exact Session lease. Each attachment has its own proxy and driver resource; `close()` and `Symbol.asyncDispose` are idempotent. A single Node client multiplexes up to 16 distinct Sessions over one lazily opened control WebSocket. + +Driver factories should return after configuring their resource rather than await a proxied navigation: tunnel dialing is deliberately held behind the first lease acknowledgement, which is published after the driver supplies its initial state. + +`BrowserDriver` descriptors are structural factory functions, so adapters remain compatible across duplicate client package instances. The Node entrypoint also re-exports canonical `Browser` contracts. Throw `BrowserDriverError` for typed command failures; structurally equivalent errors are accepted only when their `code` is a valid `Browser.ErrorCode`. Effect consumers construct canonical decoded inputs: diff --git a/packages/client/package.json b/packages/client/package.json index fe379b5457..a62614e3a9 100644 --- a/packages/client/package.json +++ b/packages/client/package.json @@ -17,6 +17,7 @@ ], "exports": { ".": "./src/promise/index.ts", + "./node": "./src/node/index.ts", "./promise": "./src/promise/index.ts", "./promise/api": "./src/promise/api.ts", "./service": "./src/promise/service.ts", @@ -28,12 +29,14 @@ "build": "bun run script/build-package.ts", "generate": "bun run script/build.ts", "check:generated": "bun run generate && git diff --exit-code -- src/promise/generated src/effect/generated src/effect/api", - "test": "bun test --timeout 5000", - "typecheck": "tsgo --noEmit" + "test": "bun test --timeout 5000 && bun run test:node-package", + "test:node-package": "bun test ./test/node/package-smoke.ts --timeout 60000", + "typecheck": "tsgo --noEmit && tsgo -p test/types/tsconfig.json --noEmit && tsgo -p test/types/tsconfig.nodenext.json --noEmit" }, "dependencies": { "@opencode-ai/schema": "workspace:*", - "@opencode-ai/protocol": "workspace:*" + "@opencode-ai/protocol": "workspace:*", + "ws": "8.21.0" }, "peerDependencies": { "effect": "4.0.0-beta.101" @@ -48,6 +51,7 @@ "@opencode-ai/httpapi-codegen": "workspace:*", "@tsconfig/bun": "catalog:", "@types/bun": "catalog:", + "@types/ws": "8.18.1", "@typescript/native-preview": "catalog:", "effect": "catalog:" } diff --git a/packages/client/script/build-package.ts b/packages/client/script/build-package.ts index 323a63ddf9..c26896428c 100644 --- a/packages/client/script/build-package.ts +++ b/packages/client/script/build-package.ts @@ -7,3 +7,4 @@ process.chdir(fileURLToPath(new URL("..", import.meta.url))) await $`rm -rf dist` await $`bun tsc -p tsconfig.build.json` +await $`bun build src/node/index.ts --outfile dist/node/index.js --target=node --format=esm --packages=external` diff --git a/packages/client/src/effect/generated/client.ts b/packages/client/src/effect/generated/client.ts index 8b457db56f..6341c6bf87 100644 --- a/packages/client/src/effect/generated/client.ts +++ b/packages/client/src/effect/generated/client.ts @@ -3,7 +3,7 @@ import { Effect, Stream, Schema } from "effect" import { Sse } from "effect/unstable/encoding" import { HttpClientError } from "effect/unstable/http" import { HttpApiClient } from "effect/unstable/httpapi" -import { ClientApi } from "../../contract" +import { ClientApi } from "../../contract.js" import type { Endpoint0_0Output, Endpoint0_1Input, @@ -218,7 +218,7 @@ import type { Endpoint27_1Input, Endpoint27_1Output, } from "../api/api.js" -import { ClientError } from "./client-error" +import { ClientError } from "./client-error.js" type RawClient = HttpApiClient.ForApi diff --git a/packages/client/src/effect/generated/index.ts b/packages/client/src/effect/generated/index.ts index bc0dbc9fa4..0589939b46 100644 --- a/packages/client/src/effect/generated/index.ts +++ b/packages/client/src/effect/generated/index.ts @@ -1,2 +1,2 @@ -export { ClientError } from "./client-error" -export * as OpenCode from "./client" +export { ClientError } from "./client-error.js" +export * as OpenCode from "./client.js" diff --git a/packages/client/src/node/browser/client.ts b/packages/client/src/node/browser/client.ts new file mode 100644 index 0000000000..6e76de55a7 --- /dev/null +++ b/packages/client/src/node/browser/client.ts @@ -0,0 +1,639 @@ +import { BrowserControlProtocol } from "@opencode-ai/protocol/browser-control" +import { BROWSER_CONTROL_PROTOCOL } from "@opencode-ai/protocol/groups/browser" +import { Browser } from "@opencode-ai/schema/browser" +import { BrowserControl } from "@opencode-ai/schema/browser-control" +import { Session } from "@opencode-ai/schema/session" +import { Effect, Schema } from "effect" +import WebSocket from "ws" +import type { ClientOptions } from "../../promise/generated/client.js" +import { + browserDriverFactory, + type BrowserDriver, + type BrowserDriverContext, + type BrowserDriverInstance, + type BrowserProxy, +} from "./driver.js" +import { createBrowserProxy } from "./proxy.js" +import { openBrowserTunnel, type BrowserTunnelEndpoint } from "./tunnel.js" + +export interface BrowserAttachOptions { + readonly sessionID: string + readonly driver: BrowserDriver + readonly signal?: AbortSignal +} + +export interface BrowserAttachment extends AsyncDisposable { + readonly resource: Resource + readonly close: () => Promise +} + +export interface BrowserClient { + readonly attach: (options: BrowserAttachOptions) => Promise> +} + +type ProxyServer = Awaited> + +type AttachmentRecord = { + readonly sessionID: Session.ID + readonly leaseID: Browser.LeaseID + readonly abort: AbortController + readonly setup: Promise + readonly finishSetup: () => void + readonly externalSignal?: AbortSignal + externalAbort?: () => void + active: boolean + closed: boolean + state?: Browser.State + execute?: BrowserDriverInstance["execute"] + dispose?: BrowserDriverInstance["dispose"] + unsubscribe?: () => void + proxy?: ProxyServer + closing?: Promise +} + +type ActiveAttachment = AttachmentRecord & { + readonly active: true + readonly state: Browser.State + readonly execute: BrowserDriverInstance["execute"] +} + +type Waiter = { + readonly key: string + readonly signal: AbortSignal + readonly resolve: () => void + readonly reject: (error: Error) => void + readonly abort: () => void +} + +/** Creates the Node-only browser host attached to one immutable server endpoint. */ +export function createBrowserClient(options: ClientOptions): BrowserClient { + const control = new BrowserClientControl(endpoint(options)) + return { attach: (input) => control.attach(input) } +} + +class BrowserClientControl { + private readonly records = new Map() + private readonly requests = new Map< + BrowserControl.RequestID, + { readonly leaseID: Browser.LeaseID; readonly abort: AbortController } + >() + private readonly waiters = new Set() + private socket?: WebSocket + private retry?: ReturnType + private retryAttempt = 0 + private revision = 0 + private ready = false + private opening?: Promise + private openingAbort?: AbortController + private synced?: string + private syncing?: { readonly revision: number; readonly snapshot: string; readonly leases: Set } + private syncedLeases = new Set() + private inbound = Promise.resolve() + private failing = false + + constructor(private readonly server: BrowserTunnelEndpoint) {} + + async attach(input: BrowserAttachOptions): Promise> { + if (!Schema.is(Session.ID)(input.sessionID)) throw new TypeError("Browser attachment requires a valid Session ID") + if (input.signal?.aborted) throw abortError(input.signal, "Browser attachment was aborted") + const sessionID = Session.ID.make(input.sessionID) + if (this.records.has(sessionID)) throw new Error(`Browser is already attached to Session ${sessionID}`) + if (this.records.size >= 16) throw new Error("A Node client cannot attach more than 16 browser Sessions") + + let finishSetup!: () => void + const setup = new Promise((resolve) => { + finishSetup = resolve + }) + const record: AttachmentRecord = { + sessionID, + leaseID: Browser.LeaseID.create(), + abort: new AbortController(), + setup, + finishSetup, + externalSignal: input.signal, + active: false, + closed: false, + } + this.records.set(sessionID, record) + record.externalAbort = () => { + void this.close(record, abortError(input.signal, "Browser attachment was aborted")).catch(() => undefined) + } + input.signal?.addEventListener("abort", record.externalAbort, { once: true }) + + try { + record.proxy = await createBrowserProxy({ + connect: async (target, signal) => { + await this.waitForAttachment(record, signal) + const tunnelSignal = AbortSignal.any([signal, record.abort.signal]) + if (tunnelSignal.aborted) throw abortError(tunnelSignal, "Browser tunnel was aborted") + return openBrowserTunnel({ + endpoint: this.server, + sessionID: record.sessionID, + leaseID: record.leaseID, + target, + signal: tunnelSignal, + }) + }, + }) + this.requireOpen(record) + + const instance = await createDriver( + input.driver, + { proxy: exposedProxy(record.proxy), signal: record.abort.signal }, + record.abort.signal, + ) + if (instance !== null && typeof instance === "object" && typeof instance.dispose === "function") { + record.dispose = () => instance.dispose() + } + if ( + instance === null || + typeof instance !== "object" || + typeof instance.state !== "function" || + typeof instance.subscribe !== "function" || + typeof instance.execute !== "function" || + typeof instance.dispose !== "function" + ) { + throw new TypeError("Browser driver factory returned an invalid driver instance") + } + record.execute = (command, options) => instance.execute(command, options) + + let receivedState = false + record.unsubscribe = instance.subscribe((state) => { + if (record.closed) return + const next = contractState(state) + if (!next) { + void this.close(record, new TypeError("Browser driver published an invalid state")).catch(() => undefined) + return + } + receivedState = true + record.state = next + if (record.active) this.changed() + }) + if (typeof record.unsubscribe !== "function") { + throw new TypeError("Browser driver subscribe must return an unsubscribe function") + } + const initial = contractState(instance.state()) + if (!initial) throw new TypeError("Browser driver returned an invalid initial state") + if (!receivedState) record.state = initial + this.requireOpen(record) + + record.active = true + this.changed() + record.finishSetup() + await this.waitForAttachment(record, record.abort.signal) + this.requireOpen(record) + + const close = () => this.close(record) + return Object.freeze({ + resource: instance.resource, + close, + [Symbol.asyncDispose]: close, + }) + } catch (error) { + record.finishSetup() + await this.close(record).catch(() => undefined) + throw error + } + } + + private close(record: AttachmentRecord, reason = new Error("Browser attachment was closed")) { + if (!record.closed) { + record.closed = true + record.active = false + if (this.records.get(record.sessionID) === record) this.records.delete(record.sessionID) + record.abort.abort(reason) + this.syncedLeases.delete(attachmentKey(record.sessionID, record.leaseID)) + this.rejectAttachmentWaiters(record, reason) + this.abortRequests(record.leaseID) + this.changed() + } + if (record.closing) return record.closing + record.closing = record.setup.then(async () => { + if (record.externalAbort) record.externalSignal?.removeEventListener("abort", record.externalAbort) + const unsubscribe = await cleanup(() => record.unsubscribe?.()) + const dispose = await cleanup(() => record.dispose?.()) + const proxy = await cleanup(() => record.proxy?.close()) + const failure = [unsubscribe, dispose, proxy].find( + (result): result is { readonly ok: false; readonly error: unknown } => !result.ok, + ) + if (failure) throw failure.error + }) + return record.closing + } + + private requireOpen(record: AttachmentRecord) { + if (!record.closed && !record.abort.signal.aborted) return + throw abortError(record.abort.signal, "Browser attachment was closed") + } + + private changed() { + if (this.failing) return + if (this.attachments().length === 0) { + this.stopControl() + return + } + this.connect() + this.publish() + } + + private attachments() { + return [...this.records.values()].filter( + (record): record is ActiveAttachment => + record.active && record.state !== undefined && record.execute !== undefined && !record.closed, + ) + } + + private connect() { + if (this.socket || this.retry || this.opening || this.attachments().length === 0) return + const abort = new AbortController() + this.openingAbort = abort + this.opening = this.open(abort.signal) + .catch((error) => this.failControl(error instanceof Error ? error : new Error(String(error)))) + .finally(() => { + if (this.openingAbort !== abort) return + this.openingAbort = undefined + this.opening = undefined + }) + } + + private async open(lifetime: AbortSignal) { + if (process.versions.bun) { + // TODO: Remove the HTTP auth probe once Bun exposes ws upgrade responses. + // https://github.com/oven-sh/bun/issues/5951 + const signal = AbortSignal.any([lifetime, AbortSignal.timeout(10_000)]) + const response = await (this.server.fetch ?? globalThis.fetch)(new URL("/api/health", this.server.url), { + headers: this.server.authorization ? { Authorization: this.server.authorization } : undefined, + signal, + }).catch(() => undefined) + if (lifetime.aborted) return + if (!response) { + this.scheduleReconnect() + return + } + if (response.status === 401 || response.status === 403) { + this.failControl(new Error(`Browser control connection was rejected with HTTP ${response.status}`)) + return + } + } + if (this.socket || this.retry || this.attachments().length === 0) return + const socket = new WebSocket(controlURL(this.server), BROWSER_CONTROL_PROTOCOL, { + ...(this.server.authorization ? { headers: { Authorization: this.server.authorization } } : {}), + handshakeTimeout: 10_000, + maxPayload: BrowserControlProtocol.MaxMessageBytes, + perMessageDeflate: false, + followRedirects: false, + }) + this.socket = socket + socket.on("message", (data, binary) => { + if (socket !== this.socket) return + this.inbound = this.inbound.then( + () => this.receive(socket, data, binary), + () => this.receive(socket, data, binary), + ) + }) + socket.on("error", (error) => { + if (socket !== this.socket) return + const status = /^Unexpected server response: (401|403|426)$/.exec(error.message)?.[1] + if (status) this.failControl(new Error(`Browser control connection was rejected with HTTP ${status}`)) + }) + if (!process.versions.bun) { + socket.on("unexpected-response", (_request, response) => { + response.resume() + if (socket !== this.socket) return + if (response.statusCode === 401 || response.statusCode === 403 || response.statusCode === 426) { + this.failControl(new Error(`Browser control connection was rejected with HTTP ${response.statusCode}`)) + return + } + socket.terminate() + }) + } + socket.on("close", (code, reason) => { + if (socket !== this.socket) return + if (code === 1002 || code === 1007 || code === 1009) { + this.failControl(controlCloseError(code, reason)) + return + } + this.socket = undefined + this.resetConnection() + this.scheduleReconnect() + }) + } + + private async receive(socket: WebSocket, data: WebSocket.RawData, binary: boolean) { + if (binary) return this.protocolError(socket) + const decoded = await Effect.runPromise(BrowserControlProtocol.decodeFromServer(rawData(data))).catch( + () => undefined, + ) + if (!decoded || socket !== this.socket) return this.protocolError(socket) + if (decoded.type === "browser.control.ready") { + if (this.ready) return this.protocolError(socket) + this.ready = true + this.publish() + return + } + if (!this.ready) return this.protocolError(socket) + if (decoded.type === "browser.control.synced") { + if (!this.syncing || this.syncing.revision !== decoded.revision) return this.protocolError(socket) + this.synced = this.syncing.snapshot + this.syncedLeases = this.syncing.leases + this.syncing = undefined + this.retryAttempt = 0 + this.resolveWaiters() + this.publish() + return + } + if (decoded.type === "browser.control.cancel") { + const request = this.requests.get(decoded.requestID) + if (!request) return + if (request.leaseID !== decoded.leaseID) return this.protocolError(socket) + this.requests.delete(decoded.requestID) + request.abort.abort(new Error("Browser command was cancelled")) + return + } + void this.request(socket, decoded) + } + + private publish() { + const socket = this.socket + if (!socket || socket.readyState !== WebSocket.OPEN || !this.ready || this.syncing) return + const attachments = this.attachments().map((record) => ({ + sessionID: record.sessionID, + leaseID: record.leaseID, + state: record.state, + })) + const snapshot = JSON.stringify(attachments) + if (snapshot === this.synced) return + const revision = this.revision++ + this.syncing = { + revision, + snapshot, + leases: new Set(attachments.map((attachment) => attachmentKey(attachment.sessionID, attachment.leaseID))), + } + this.send(socket, { type: "browser.control.sync", revision, attachments }) + } + + private async request(socket: WebSocket, message: BrowserControl.Request) { + if (this.requests.has(message.requestID)) return this.protocolError(socket) + const record = this.records.get(message.sessionID) + if (!record?.active || record.closed || record.leaseID !== message.leaseID || !record.execute) { + this.send(socket, { + type: "browser.control.response", + requestID: message.requestID, + leaseID: message.leaseID, + outcome: { type: "failure", code: "not_attached", message: "The browser attachment is no longer available." }, + }) + return + } + + const abort = new AbortController() + this.requests.set(message.requestID, { leaseID: message.leaseID, abort }) + const outcome = await Promise.resolve() + .then(() => record.execute?.(message.command, { signal: abort.signal })) + .then( + (result): Browser.Outcome => { + if (!Schema.is(Browser.Result)(result) || result.type !== message.command.type) { + return { type: "failure", code: "protocol", message: "Browser driver returned an invalid command result." } + } + return { type: "success", result } + }, + (error): Browser.Outcome => driverFailure(error), + ) + if (this.requests.get(message.requestID)?.abort !== abort) return + this.requests.delete(message.requestID) + if (socket !== this.socket) return + this.send(socket, { + type: "browser.control.response", + requestID: message.requestID, + leaseID: message.leaseID, + outcome, + }) + if (outcome.type === "success") { + const state = contractState(outcome.result.state) + if (state) { + record.state = state + this.changed() + } + } + } + + private send(socket: WebSocket, message: BrowserControl.FromDesktop) { + if (socket !== this.socket || socket.readyState !== WebSocket.OPEN) return + socket.send(BrowserControlProtocol.encodeFromDesktop(message), (error) => { + if (error && socket === this.socket) socket.terminate() + }) + } + + private protocolError(socket: WebSocket) { + if (socket.readyState === WebSocket.OPEN) socket.close(1002, "Invalid browser control message") + else socket.terminate() + } + + private waitForAttachment(record: AttachmentRecord, signal: AbortSignal) { + const key = attachmentKey(record.sessionID, record.leaseID) + if (record.closed || this.records.get(record.sessionID) !== record) { + return Promise.reject(new Error("Browser attachment is no longer available")) + } + if (record.active && this.syncedLeases.has(key)) return Promise.resolve() + if (signal.aborted) return Promise.reject(abortError(signal, "Browser attachment wait was aborted")) + return new Promise((resolve, reject) => { + const waiter: Waiter = { + key, + signal, + resolve, + reject, + abort: () => { + this.waiters.delete(waiter) + reject(abortError(signal, "Browser attachment wait was aborted")) + }, + } + signal.addEventListener("abort", waiter.abort, { once: true }) + this.waiters.add(waiter) + }) + } + + private resolveWaiters() { + for (const waiter of this.waiters) { + if (!this.syncedLeases.has(waiter.key)) continue + const record = [...this.records.values()].find( + (item) => !item.closed && item.active && attachmentKey(item.sessionID, item.leaseID) === waiter.key, + ) + if (!record) continue + this.waiters.delete(waiter) + waiter.signal.removeEventListener("abort", waiter.abort) + waiter.resolve() + } + } + + private rejectAttachmentWaiters(record: AttachmentRecord, error: Error) { + const key = attachmentKey(record.sessionID, record.leaseID) + for (const waiter of this.waiters) { + if (waiter.key !== key) continue + this.waiters.delete(waiter) + waiter.signal.removeEventListener("abort", waiter.abort) + waiter.reject(error) + } + } + + private abortRequests(leaseID?: Browser.LeaseID) { + for (const [requestID, request] of this.requests) { + if (leaseID && request.leaseID !== leaseID) continue + this.requests.delete(requestID) + request.abort.abort(new Error("Browser command was aborted")) + } + } + + private stopControl() { + this.openingAbort?.abort() + this.openingAbort = undefined + this.opening = undefined + if (this.retry) clearTimeout(this.retry) + this.retry = undefined + const socket = this.socket + this.socket = undefined + this.resetConnection() + socket?.terminate() + } + + private failControl(error: Error) { + this.failing = true + this.stopControl() + this.retryAttempt = 0 + const records = [...this.records.values()] + records.forEach((record) => { + void this.close(record, error).catch(() => undefined) + }) + this.failing = false + } + + private resetConnection() { + this.ready = false + this.synced = undefined + this.syncedLeases.clear() + this.syncing = undefined + this.revision = 0 + this.abortRequests() + } + + private scheduleReconnect() { + if (this.retry || this.attachments().length === 0) return + const delay = Math.min(5_000, 100 * 2 ** this.retryAttempt++) + this.retry = setTimeout(() => { + this.retry = undefined + this.connect() + }, delay) + this.retry.unref() + } +} + +function createDriver(driver: BrowserDriver, context: BrowserDriverContext, signal: AbortSignal) { + const creating = Promise.resolve().then(() => browserDriverFactory(driver)(context)) + return new Promise>((resolve, reject) => { + let settled = false + const abort = () => { + if (settled) return + settled = true + signal.removeEventListener("abort", abort) + reject(abortError(signal, "Browser driver creation was aborted")) + } + signal.addEventListener("abort", abort, { once: true }) + creating.then( + (instance) => { + if (settled) { + void disposeLateDriver(instance) + return + } + settled = true + signal.removeEventListener("abort", abort) + resolve(instance) + }, + (error) => { + if (settled) return + settled = true + signal.removeEventListener("abort", abort) + reject(error) + }, + ) + if (signal.aborted) abort() + }) +} + +async function disposeLateDriver(instance: unknown) { + if (instance === null || typeof instance !== "object" || !("dispose" in instance)) return + const dispose = instance.dispose + if (typeof dispose !== "function") return + await Promise.resolve() + .then(() => dispose.call(instance)) + .catch(() => undefined) +} + +function cleanup(task: () => unknown) { + return Promise.resolve() + .then(task) + .then( + () => ({ ok: true as const }), + (error: unknown) => ({ ok: false as const, error }), + ) +} + +function exposedProxy(proxy: ProxyServer): BrowserProxy { + return Object.freeze({ + url: proxy.url, + host: proxy.host, + port: proxy.port, + credentials: Object.freeze({ ...proxy.credentials }), + certificateFingerprint: proxy.certificateFingerprint, + }) +} + +function contractState(state: Browser.State) { + if (!Schema.is(Browser.State)(state)) return undefined + return Object.freeze({ ...state }) +} + +function driverFailure(error: unknown): Browser.Failure { + return { + type: "failure", + code: + error !== null && typeof error === "object" && "code" in error && Schema.is(Browser.ErrorCode)(error.code) + ? error.code + : "internal", + message: (error instanceof Error ? error.message : String(error)).slice(0, 1_024), + } +} + +function endpoint(options: ClientOptions): BrowserTunnelEndpoint { + const url = new URL(options.baseUrl) + if ((url.protocol !== "http:" && url.protocol !== "https:") || url.username || url.password) { + throw new TypeError("Browser server endpoint must be an HTTP URL without embedded credentials") + } + const authorization = new Headers(options.headers).get("authorization") ?? undefined + return Object.freeze({ url: url.href, ...(authorization ? { authorization } : {}), fetch: options.fetch }) +} + +function controlURL(endpoint: BrowserTunnelEndpoint) { + const url = new URL(endpoint.url) + url.protocol = url.protocol === "https:" ? "wss:" : "ws:" + url.pathname = "/api/browser/control" + url.search = "" + url.hash = "" + return url +} + +function attachmentKey(sessionID: string, leaseID: Browser.LeaseID) { + return `${sessionID}\0${leaseID}` +} + +function abortError(signal: AbortSignal | undefined, message: string) { + return signal?.reason instanceof Error ? signal.reason : new Error(message) +} + +function controlCloseError(code: number, reason: Buffer) { + const detail = reason.toString("utf8").slice(0, 256) + return new Error(`Browser control connection closed with fatal code ${code}${detail ? `: ${detail}` : ""}`) +} + +function rawData(data: WebSocket.RawData) { + if (data instanceof ArrayBuffer) return new Uint8Array(data) + if (Array.isArray(data)) return new Uint8Array(Buffer.concat(data)) + return new Uint8Array(data.buffer, data.byteOffset, data.byteLength) +} diff --git a/packages/client/src/node/browser/driver.ts b/packages/client/src/node/browser/driver.ts new file mode 100644 index 0000000000..dda0103c91 --- /dev/null +++ b/packages/client/src/node/browser/driver.ts @@ -0,0 +1,58 @@ +import type { Browser } from "@opencode-ai/schema/browser" + +/** Connection details for the attachment's private authenticated proxy. */ +export interface BrowserProxy { + readonly url: string + readonly host: string + readonly port: number + readonly credentials: { + readonly username: string + readonly password: string + } + readonly certificateFingerprint: string +} + +export interface BrowserDriverContext { + readonly proxy: BrowserProxy + readonly signal: AbortSignal +} + +export interface BrowserDriverInstance { + readonly resource: Resource + readonly state: () => Browser.State + readonly subscribe: (listener: (state: Browser.State) => void) => () => void + readonly execute: (command: Browser.Command, options: { readonly signal: AbortSignal }) => Promise + readonly dispose: () => Promise | void +} + +export type BrowserDriverFactory = ( + context: BrowserDriverContext, +) => Promise> | BrowserDriverInstance + +/** Error returned to the server when a browser adapter cannot execute a command. */ +export class BrowserDriverError extends Error { + override readonly name = "BrowserDriverError" + + constructor( + readonly code: Browser.ErrorCode, + message: string, + options?: ErrorOptions, + ) { + super(message, options) + } +} + +/** Structural adapter descriptor, represented by its factory function. */ +export type BrowserDriver = BrowserDriverFactory + +export const BrowserDriver = { + define(create: BrowserDriverFactory): BrowserDriver { + if (typeof create !== "function") throw new TypeError("Browser driver factory must be a function") + return Object.freeze(create) + }, +} + +export function browserDriverFactory(driver: BrowserDriver) { + if (typeof driver !== "function") throw new TypeError("Browser driver must be a factory function") + return driver +} diff --git a/packages/client/src/node/browser/proxy-certificate.ts b/packages/client/src/node/browser/proxy-certificate.ts new file mode 100644 index 0000000000..3e2a022c3a --- /dev/null +++ b/packages/client/src/node/browser/proxy-certificate.ts @@ -0,0 +1,118 @@ +import { createHash, createSign, generateKeyPairSync, randomBytes } from "node:crypto" + +// Certificate construction follows VS Code's MIT-licensed tunnel proxy implementation. +// Copyright (c) Microsoft Corporation. + +/** Creates the ephemeral certificate pinned by the browser proxy's host. */ +export function createBrowserProxyCertificate(hostname: string) { + const keys = generateKeyPairSync("ec", { + namedCurve: "prime256v1", + publicKeyEncoding: { type: "spki", format: "pem" }, + privateKeyEncoding: { type: "pkcs8", format: "pem" }, + }) + const certificate = createCertificate(keys.privateKey, keys.publicKey, hostname) + const fingerprint = `sha256/${createHash("sha256").update(pemToDer(certificate)).digest("base64")}` + return { key: keys.privateKey, certificate, fingerprint } +} + +function createCertificate(privateKey: string, publicKey: string, hostname: string) { + const serial = randomBytes(8) + serial[0] &= 0x7f + if (serial.every((byte) => byte === 0)) serial[serial.length - 1] = 1 + const now = new Date(Date.now() - 60_000) + const expires = new Date(now) + expires.setUTCFullYear(expires.getUTCFullYear() + 1) + const name = sequence([set([sequence([oid(Buffer.from([0x55, 0x04, 0x03])), utf8("OpenCode Browser Proxy")])])]) + const signatureAlgorithm = sequence([oid(Buffer.from([0x2a, 0x86, 0x48, 0xce, 0x3d, 0x04, 0x03, 0x02]))]) + const extensions = Buffer.concat([ + Buffer.from([0xa3]), + lengthPrefix( + sequence([ + sequence([ + oid(Buffer.from([0x55, 0x1d, 0x11])), + octetString( + sequence([tagged(0x82, Buffer.from(hostname, "ascii")), tagged(0x87, Buffer.from([127, 0, 0, 1]))]), + ), + ]), + ]), + ), + ]) + const body = sequence([ + Buffer.from([0xa0, 0x03, 0x02, 0x01, 0x02]), + integer(serial), + signatureAlgorithm, + name, + sequence([time(now), time(expires)]), + name, + pemToDer(publicKey), + extensions, + ]) + const signer = createSign("SHA256") + signer.update(body) + const signature = signer.sign(privateKey) + const certificate = sequence([ + body, + signatureAlgorithm, + Buffer.concat([Buffer.from([0x03]), length(signature.length + 1), Buffer.from([0]), signature]), + ]) + const encoded = certificate + .toString("base64") + .match(/.{1,64}/g) + ?.join("\n") + if (!encoded) throw new Error("Failed to encode browser proxy certificate") + return `-----BEGIN CERTIFICATE-----\n${encoded}\n-----END CERTIFICATE-----\n` +} + +function pemToDer(pem: string) { + return Buffer.from(pem.replace(/-----[A-Z ]+-----/g, "").replaceAll(/\s/g, ""), "base64") +} + +function length(value: number) { + if (value < 0x80) return Buffer.from([value]) + if (value < 0x100) return Buffer.from([0x81, value]) + if (value < 0x10000) return Buffer.from([0x82, value >> 8, value & 0xff]) + if (value < 0x1000000) return Buffer.from([0x83, value >> 16, (value >> 8) & 0xff, value & 0xff]) + throw new RangeError(`ASN.1 value is too large: ${value}`) +} + +function lengthPrefix(value: Buffer) { + return Buffer.concat([length(value.length), value]) +} + +function tagged(tag: number, value: Buffer) { + return Buffer.concat([Buffer.from([tag]), lengthPrefix(value)]) +} + +function sequence(items: ReadonlyArray) { + return tagged(0x30, Buffer.concat(items)) +} + +function set(items: ReadonlyArray) { + return tagged(0x31, Buffer.concat(items)) +} + +function integer(value: Buffer) { + const first = value.findIndex( + (byte, index) => byte !== 0 || index === value.length - 1 || (value[index + 1] & 0x80) !== 0, + ) + const canonical = value.subarray(first < 0 ? value.length - 1 : first) + return tagged(0x02, canonical[0] & 0x80 ? Buffer.concat([Buffer.from([0]), canonical]) : canonical) +} + +function oid(value: Buffer) { + return tagged(0x06, value) +} + +function utf8(value: string) { + return tagged(0x0c, Buffer.from(value, "utf8")) +} + +function octetString(value: Buffer) { + return tagged(0x04, value) +} + +function time(value: Date) { + const encoded = value.toISOString().replaceAll(/[-:T]/g, "") + if (value.getUTCFullYear() < 2050) return tagged(0x17, Buffer.from(encoded.slice(2, 14) + "Z", "ascii")) + return tagged(0x18, Buffer.from(encoded.slice(0, 14) + "Z", "ascii")) +} diff --git a/packages/client/src/node/browser/proxy.ts b/packages/client/src/node/browser/proxy.ts new file mode 100644 index 0000000000..0653f23c5a --- /dev/null +++ b/packages/client/src/node/browser/proxy.ts @@ -0,0 +1,410 @@ +import { BrowserTunnel } from "@opencode-ai/schema/browser-tunnel" +import { randomBytes, timingSafeEqual } from "node:crypto" +import { once } from "node:events" +import { Agent, request, type IncomingHttpHeaders, type IncomingMessage, type ServerResponse } from "node:http" +import { createServer } from "node:https" +import { Duplex } from "node:stream" +import { createBrowserProxyCertificate } from "./proxy-certificate.js" + +export type BrowserProxyConnector = (target: BrowserTunnel.Target, signal: AbortSignal) => Promise + +/** Starts an authenticated TLS forward proxy whose target sockets can only come from the tunnel connector. */ +export async function createBrowserProxy(input: { readonly connect: BrowserProxyConnector }) { + const candidate = `opencode-${randomBytes(16).toString("hex")}.localhost` + const certificate = createBrowserProxyCertificate(candidate) + const username = randomBytes(16).toString("hex") + const password = randomBytes(32).toString("hex") + const expectedAuthorization = Buffer.from(`Basic ${Buffer.from(`${username}:${password}`).toString("base64")}`) + const clients = new Set() + const tunnels = new Set() + const pending = new Set() + const limiter = connectionLimiter(6, 64) + let closed = false + + const connect = async (target: BrowserTunnel.Target, signal?: AbortSignal) => { + if (closed) throw new Error("Browser proxy is closed") + const abort = new AbortController() + const cancel = () => abort.abort(signal?.reason) + const done = () => { + pending.delete(abort) + signal?.removeEventListener("abort", cancel) + } + signal?.addEventListener("abort", cancel, { once: true }) + if (signal?.aborted) cancel() + pending.add(abort) + return limiter + .run(() => input.connect(target, abort.signal), abort.signal) + .then( + (tunnel) => { + done() + if (closed || abort.signal.aborted) { + tunnel.destroy() + throw abort.signal.reason ?? new Error("Browser proxy is closed") + } + tunnels.add(tunnel) + tunnel.once("close", () => tunnels.delete(tunnel)) + tunnel.on("error", () => tunnel.destroy()) + return tunnel + }, + (error) => { + done() + throw error + }, + ) + } + + const authorized = (header: string | undefined) => { + if (!header) return false + const actual = Buffer.from(header) + return actual.length === expectedAuthorization.length && timingSafeEqual(actual, expectedAuthorization) + } + + const server = createServer( + { key: certificate.key, cert: certificate.certificate, allowHalfOpen: true, maxHeaderSize: 64 * 1_024 }, + (incoming, response) => { + void forward(incoming, response, connect, authorized).catch(() => response.destroy()) + }, + ) + server.requestTimeout = 30_000 + server.headersTimeout = 10_000 + server.keepAliveTimeout = 5_000 + server.maxRequestsPerSocket = 1_000 + server.on("connection", (socket) => { + clients.add(socket) + socket.once("close", () => clients.delete(socket)) + }) + server.on("connect", (incoming, socket, head) => { + clients.add(socket) + socket.once("close", () => clients.delete(socket)) + let established = false + void (async () => { + if (!authorized(singleHeader(incoming.headers["proxy-authorization"]))) { + socket.end( + 'HTTP/1.1 407 Proxy Authentication Required\r\nProxy-Authenticate: Basic realm="OpenCode Browser Proxy"\r\nContent-Length: 0\r\nConnection: close\r\n\r\n', + ) + return + } + const parsed = authority(incoming.url ?? "", 443) + if (!parsed) { + socket.end("HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") + return + } + const abort = new AbortController() + let tunnel: Duplex | undefined + const buffered: Uint8Array[] = head.byteLength ? [head] : [] + let bufferedBytes = head.byteLength + const cancel = () => { + abort.abort(new Error("Browser proxy client closed during tunnel setup")) + tunnel?.destroy() + } + const onReadable = () => { + while (true) { + const data: unknown = socket.read() + if (data === null) break + if (!(data instanceof Uint8Array)) { + cancel() + socket.destroy() + return + } + bufferedBytes += data.byteLength + if (bufferedBytes > 256 * 1_024) { + cancel() + socket.destroy() + return + } + buffered.push(data) + } + if (socket.readableEnded) cancel() + } + incoming.once("aborted", cancel) + incoming.once("error", cancel) + socket.once("close", cancel) + socket.once("end", cancel) + socket.once("error", cancel) + socket.on("readable", onReadable) + socket.pause() + tunnel = await connect(target(parsed.host, parsed.port), abort.signal) + onReadable() + if (socket.destroyed || socket.readableEnded || abort.signal.aborted) { + socket.off("readable", onReadable) + tunnel.destroy() + return + } + socket.write("HTTP/1.1 200 Connection Established\r\n\r\n") + established = true + for (const data of buffered) { + if (!tunnel.write(data)) await once(tunnel, "drain", { signal: abort.signal }) + } + socket.off("readable", onReadable) + bridge(socket, tunnel) + incoming.off("aborted", cancel) + incoming.off("error", cancel) + socket.off("close", cancel) + socket.off("end", cancel) + socket.off("error", cancel) + socket.resume() + })().catch(() => { + if (socket.destroyed) return + if (established) { + socket.destroy() + return + } + socket.end("HTTP/1.1 502 Bad Gateway\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") + }) + }) + server.on("error", () => undefined) + server.on("clientError", (_error, socket) => { + if (!socket.destroyed) socket.end("HTTP/1.1 400 Bad Request\r\nConnection: close\r\n\r\n") + }) + + await new Promise((resolve, reject) => { + const onError = (error: Error) => { + server.off("listening", onListening) + reject(error) + } + const onListening = () => { + server.off("error", onError) + resolve() + } + server.once("error", onError) + server.once("listening", onListening) + server.listen(0, "127.0.0.1") + }) + const address = server.address() + if (!address || typeof address === "string") { + server.close() + throw new Error("Browser proxy did not bind a TCP address") + } + const host = candidate + let closing: Promise | undefined + + return { + url: `https://${host}:${address.port}`, + host, + port: address.port, + credentials: { username, password }, + certificateFingerprint: certificate.fingerprint, + close() { + if (closing) return closing + closed = true + pending.forEach((abort) => abort.abort()) + pending.clear() + limiter.close() + tunnels.forEach((tunnel) => tunnel.destroy()) + tunnels.clear() + clients.forEach((client) => client.destroy()) + clients.clear() + closing = new Promise((resolve) => server.close(() => resolve())) + return closing + }, + } +} + +async function forward( + incoming: IncomingMessage, + response: ServerResponse, + connect: (target: BrowserTunnel.Target, signal?: AbortSignal) => Promise, + authorized: (header: string | undefined) => boolean, +) { + if (!authorized(singleHeader(incoming.headers["proxy-authorization"]))) { + response.writeHead(407, { "Proxy-Authenticate": 'Basic realm="OpenCode Browser Proxy"' }) + response.end() + return + } + const url = (() => { + try { + return new URL(incoming.url ?? "") + } catch { + return undefined + } + })() + if (!url || url.protocol !== "http:" || !url.hostname || url.username || url.password) { + response.writeHead(400) + response.end() + return + } + const port = url.port ? Number(url.port) : 80 + const abort = new AbortController() + let tunnel: Duplex | undefined + let agent: Agent | undefined + const cancel = () => { + abort.abort(new Error("Browser proxy downstream closed")) + tunnel?.destroy() + } + incoming.once("aborted", cancel) + incoming.once("error", cancel) + incoming.socket.once("close", cancel) + incoming.socket.once("end", cancel) + incoming.socket.once("error", cancel) + response.once("close", cancel) + response.once("error", cancel) + try { + tunnel = await connect(target(normalizeHostname(url.hostname), port), abort.signal) + const headers = forwardedHeaders(incoming.headers) + headers.host = url.host + headers.connection = "close" + const connection = tunnel + agent = new Agent({ keepAlive: false, maxSockets: 1, noDelay: false }) + agent.createConnection = () => connection + await new Promise((resolve, reject) => { + const upstream = request( + { + agent, + hostname: url.hostname, + port, + path: `${url.pathname}${url.search}`, + method: incoming.method, + headers, + maxHeaderSize: 64 * 1_024, + signal: abort.signal, + }, + (result) => { + const headers = forwardedHeaders(result.headers) + headers.connection = "close" + response.writeHead(result.statusCode ?? 502, result.statusMessage, headers) + result.once("error", reject) + result.once("aborted", () => reject(new Error("Browser proxy target response was aborted"))) + result.once("end", resolve) + result.pipe(response) + }, + ) + upstream.once("continue", () => response.writeContinue()) + upstream.once("error", reject) + incoming.pipe(upstream) + }) + } finally { + incoming.off("aborted", cancel) + incoming.off("error", cancel) + incoming.socket.off("close", cancel) + incoming.socket.off("end", cancel) + incoming.socket.off("error", cancel) + response.off("close", cancel) + response.off("error", cancel) + agent?.destroy() + tunnel?.destroy() + } +} + +function forwardedHeaders(input: IncomingHttpHeaders) { + const headers = { ...input } + const named = singleHeader(headers.connection) + ?.split(",") + .map((value) => value.trim().toLowerCase()) + .filter(Boolean) + named?.forEach((name) => delete headers[name]) + ;[ + "connection", + "keep-alive", + "proxy-authenticate", + "proxy-authorization", + "proxy-connection", + "te", + "trailer", + "transfer-encoding", + "upgrade", + ].forEach((name) => delete headers[name]) + return headers +} + +function bridge(client: Duplex, tunnel: Duplex) { + client.on("error", () => tunnel.destroy()) + tunnel.on("error", () => client.destroy()) + client.on("close", () => { + if (!tunnel.destroyed) tunnel.destroy() + }) + tunnel.on("close", () => { + if (!client.destroyed) client.destroy() + }) + client.pipe(tunnel) + tunnel.pipe(client) +} + +function authority(value: string, defaultPort: number) { + const bracket = /^\[([^\]]+)](?::([0-9]+))?$/.exec(value) + if (bracket) return validAuthority(bracket[1], bracket[2], defaultPort) + if (value.includes("[")) return undefined + const separator = value.lastIndexOf(":") + if (separator < 0) return validAuthority(value, undefined, defaultPort) + if (value.slice(0, separator).includes(":")) return undefined + return validAuthority(value.slice(0, separator), value.slice(separator + 1), defaultPort) +} + +function validAuthority(host: string, value: string | undefined, defaultPort: number) { + if (!host || (value !== undefined && !/^[0-9]+$/.test(value))) return undefined + const port = value === undefined ? defaultPort : Number(value) + if (!Number.isSafeInteger(port) || port < 1 || port > 65_535) return undefined + try { + return { host: BrowserTunnel.Host.make(host), port: BrowserTunnel.Port.make(port) } + } catch { + return undefined + } +} + +function target(host: string, port: number): BrowserTunnel.Target { + return { host: BrowserTunnel.Host.make(host), port: BrowserTunnel.Port.make(port) } +} + +function normalizeHostname(hostname: string) { + return hostname.startsWith("[") && hostname.endsWith("]") ? hostname.slice(1, -1) : hostname +} + +function singleHeader(value: string | ReadonlyArray | undefined) { + return typeof value === "string" ? value : undefined +} + +function connectionLimiter(limit: number, capacity: number) { + const queue: Array<() => void> = [] + let active = 0 + let closed = false + const drain = (): void => { + if (active >= limit) return + const start = queue.shift() + if (!start) return + start() + drain() + } + return { + run(task: () => Promise, signal: AbortSignal) { + return new Promise((resolve, reject) => { + if (closed || signal.aborted) { + reject(signal.reason ?? new Error("Browser proxy connection was cancelled")) + return + } + if (active >= limit && queue.length >= capacity) { + reject(new Error("Browser proxy connection queue is full")) + return + } + const start = () => { + signal.removeEventListener("abort", cancel) + if (closed || signal.aborted) { + reject(signal.reason ?? new Error("Browser proxy connection was cancelled")) + return + } + active++ + void Promise.resolve() + .then(task) + .then(resolve, reject) + .finally(() => { + active-- + drain() + }) + } + const cancel = () => { + const index = queue.indexOf(start) + if (index === -1) return + queue.splice(index, 1) + reject(signal.reason ?? new Error("Browser proxy connection was cancelled")) + } + if (active < limit) start() + else { + signal.addEventListener("abort", cancel, { once: true }) + queue.push(start) + } + }) + }, + close() { + closed = true + queue.splice(0).forEach((start) => start()) + }, + } +} diff --git a/packages/client/src/node/browser/tunnel.ts b/packages/client/src/node/browser/tunnel.ts new file mode 100644 index 0000000000..20bcb656d4 --- /dev/null +++ b/packages/client/src/node/browser/tunnel.ts @@ -0,0 +1,389 @@ +import { BrowserTunnelProtocol } from "@opencode-ai/protocol/browser-tunnel" +import { BROWSER_TUNNEL_PROTOCOL } from "@opencode-ai/protocol/groups/browser" +import { Browser } from "@opencode-ai/schema/browser" +import { BrowserTunnel } from "@opencode-ai/schema/browser-tunnel" +import { Session } from "@opencode-ai/schema/session" +import { Effect, Schema } from "effect" +import { Duplex } from "node:stream" +import WebSocket from "ws" + +export type BrowserTunnelEndpoint = { + readonly url: string + readonly authorization?: string + readonly fetch?: typeof globalThis.fetch +} + +export type BrowserTunnelOpen = { + readonly endpoint: BrowserTunnelEndpoint + readonly sessionID: Session.ID + readonly leaseID: Browser.LeaseID + readonly target: BrowserTunnel.Target + readonly signal?: AbortSignal +} + +export class BrowserTunnelError extends Error { + constructor( + readonly code: BrowserTunnel.OpenErrorCode | BrowserTunnel.ResetCode | "transport", + message: string, + ) { + super(message) + this.name = "BrowserTunnelError" + } +} + +/** Opens one flow-controlled WebSocket tunnel and exposes it as a single TCP-like stream. */ +export async function openBrowserTunnel(input: BrowserTunnelOpen): Promise { + const tunnel = new BrowserTunnelStream(input) + await tunnel.opened + return tunnel +} + +class BrowserTunnelStream extends Duplex { + readonly opened: Promise + readonly connecting = false + private resolveOpened!: () => void + private rejectOpened!: (error: Error) => void + private readonly socket: WebSocket + private readonly signal?: AbortSignal + private state: "ready" | "opening" | "open" | "closed" = "ready" + private inbound = Promise.resolve() + private timer?: ReturnType + private sendWindowBytes = 0 + private sendWindowFrames = 0 + private readonly outstanding: number[] = [] + private activeWrite?: { readonly data: Buffer; offset: number; readonly callback: (error?: Error | null) => void } + private sending = false + private receiveWindowBytes = BrowserTunnelProtocol.InitialWindowBytes + private receiveWindowFrames = BrowserTunnelProtocol.InitialFrameWindow + private heldWindowBytes = 0 + private heldWindowFrames = 0 + private paused = false + private localEnded = false + private remoteEnded = false + private remoteReset = false + private resetCode: BrowserTunnel.ResetCode = "cancelled" + + constructor(input: BrowserTunnelOpen) { + super({ + allowHalfOpen: true, + readableHighWaterMark: BrowserTunnelProtocol.InitialWindowBytes, + writableHighWaterMark: BrowserTunnelProtocol.InitialWindowBytes, + }) + this.opened = new Promise((resolve, reject) => { + this.resolveOpened = resolve + this.rejectOpened = reject + }) + this.on("error", () => undefined) + this.signal = input.signal + this.socket = new WebSocket(endpointURL(input.endpoint), BROWSER_TUNNEL_PROTOCOL, { + ...(input.endpoint.authorization ? { headers: { Authorization: input.endpoint.authorization } } : {}), + handshakeTimeout: 10_000, + maxPayload: Math.max(BrowserTunnelProtocol.MaxDataBytes, BrowserTunnelProtocol.MaxControlBytes) + 1, + perMessageDeflate: false, + followRedirects: false, + }) + this.timer = setTimeout( + () => this.fail(new BrowserTunnelError("transport", "Browser tunnel handshake timed out.")), + 15_000, + ) + this.timer.unref() + this.socket.on("message", (data, binary) => { + this.socket.pause() + const processing = this.inbound.then( + () => this.receive(input, data, binary), + () => this.receive(input, data, binary), + ) + this.inbound = processing + void processing.then( + () => { + if (this.inbound === processing && !this.paused && this.socket.readyState === WebSocket.OPEN) { + this.socket.resume() + } + }, + (error) => + this.fail(new BrowserTunnelError("transport", error instanceof Error ? error.message : String(error))), + ) + }) + this.socket.on("error", (error) => this.fail(new BrowserTunnelError("transport", error.message))) + this.socket.on("close", (code, reason) => { + if (this.state === "closed") return + if (this.localEnded && this.remoteEnded) { + this.state = "closed" + this.destroy() + return + } + this.fail(new BrowserTunnelError("transport", `Browser tunnel closed (${code}): ${reason.toString()}`)) + }) + this.signal?.addEventListener("abort", this.onAbort, { once: true }) + if (this.signal?.aborted) this.onAbort() + } + + override _read() { + if (!this.paused) return + this.paused = false + this.releaseReceiveWindow() + if (this.socket.readyState === WebSocket.OPEN) this.socket.resume() + } + + override _write(chunk: Buffer | string, encoding: BufferEncoding, callback: (error?: Error | null) => void) { + const data = typeof chunk === "string" ? Buffer.from(chunk, encoding) : chunk + if (this.state !== "open" || this.localEnded) { + callback(new BrowserTunnelError("transport", "Browser tunnel is not writable.")) + return + } + if (data.byteLength === 0) { + callback() + return + } + this.activeWrite = { data, offset: 0, callback } + this.pumpWrite() + } + + override _final(callback: (error?: Error | null) => void) { + if (this.state !== "open" || this.localEnded) { + callback(new BrowserTunnelError("transport", "Browser tunnel is not writable.")) + return + } + this.send(BrowserTunnelProtocol.encodeFromDesktop({ type: "browser.tunnel.end" }), (error) => { + if (error) { + callback(error) + this.fail(new BrowserTunnelError("transport", error.message)) + return + } + this.localEnded = true + callback() + this.finish() + }) + } + + override _destroy(error: Error | null, callback: (error?: Error | null) => void) { + if (this.timer) clearTimeout(this.timer) + this.signal?.removeEventListener("abort", this.onAbort) + if (this.state === "closed") { + callback(error) + return + } + const opened = this.state === "open" + this.state = "closed" + const pending = this.activeWrite + this.activeWrite = undefined + pending?.callback(error ?? new BrowserTunnelError("cancelled", "Browser tunnel was closed.")) + if ( + !opened || + this.remoteReset || + this.socket.readyState !== WebSocket.OPEN || + (this.localEnded && this.remoteEnded) + ) { + this.socket.terminate() + callback(error) + return + } + this.socket.send( + BrowserTunnelProtocol.encodeFromDesktop({ type: "browser.tunnel.reset", code: this.resetCode }), + { binary: true }, + () => { + this.socket.close(1000) + }, + ) + const terminate = setTimeout(() => this.socket.terminate(), 1_000) + terminate.unref() + callback(error) + } + + setKeepAlive() { + return this + } + + setNoDelay() { + return this + } + + setTimeout(_timeout: number, callback?: () => void) { + if (callback) this.once("timeout", callback) + return this + } + + ref() { + return this + } + + unref() { + return this + } + + private async receive(input: BrowserTunnelOpen, data: WebSocket.RawData, binary: boolean) { + if (!binary) return this.protocolFailure("Browser tunnel frames must be binary.") + const frame = await Effect.runPromise(BrowserTunnelProtocol.decodeFromServer(rawData(data))).catch(() => undefined) + if (!frame || this.state === "closed") return this.protocolFailure("Browser tunnel frame is invalid.") + if (frame.type === "data") { + if (this.state !== "open" || this.remoteEnded) return this.protocolFailure("Unexpected browser tunnel data.") + if (frame.data.byteLength > this.receiveWindowBytes || this.receiveWindowFrames === 0) { + return this.protocolFailure("Browser tunnel receive window was exceeded.") + } + this.receiveWindowBytes -= frame.data.byteLength + this.receiveWindowFrames-- + this.heldWindowBytes += frame.data.byteLength + this.heldWindowFrames++ + if (!this.push(frame.data)) { + this.paused = true + return + } + this.releaseReceiveWindow() + return + } + const message = frame.message + if (message.type === "browser.tunnel.ready") { + if (this.state !== "ready") return this.protocolFailure("Browser tunnel sent duplicate ready.") + this.state = "opening" + this.send( + BrowserTunnelProtocol.encodeFromDesktop({ + type: "browser.tunnel.open", + sessionID: input.sessionID, + leaseID: input.leaseID, + target: input.target, + receiveWindow: BrowserTunnel.WindowSize.make(BrowserTunnelProtocol.InitialWindowBytes), + receiveFrames: BrowserTunnel.FrameWindow.make(BrowserTunnelProtocol.InitialFrameWindow), + }), + (error) => { + if (error) this.fail(new BrowserTunnelError("transport", error.message)) + }, + ) + return + } + if (message.type === "browser.tunnel.opened") { + if (this.state !== "opening") return this.protocolFailure("Unexpected browser tunnel opened message.") + this.state = "open" + this.sendWindowBytes = message.receiveWindow + this.sendWindowFrames = message.receiveFrames + if (this.timer) clearTimeout(this.timer) + this.resolveOpened() + this.pumpWrite() + return + } + if (message.type === "browser.tunnel.rejected") { + if (this.state !== "opening") return this.protocolFailure("Unexpected browser tunnel rejection.") + this.fail(new BrowserTunnelError(message.code, message.message)) + return + } + if (message.type === "browser.tunnel.window") { + if (this.state !== "open" || message.frames > this.outstanding.length) { + return this.protocolFailure("Unexpected browser tunnel window update.") + } + const bytes = this.outstanding.slice(0, message.frames).reduce((total, size) => total + size, 0) + if (bytes !== message.bytes) + return this.protocolFailure("Browser tunnel window update does not match sent frames.") + this.outstanding.splice(0, message.frames) + this.sendWindowBytes += message.bytes + this.sendWindowFrames += message.frames + this.pumpWrite() + return + } + if (message.type === "browser.tunnel.end") { + if (this.state !== "open" || this.remoteEnded) return this.protocolFailure("Unexpected browser tunnel end.") + this.remoteEnded = true + this.push(null) + this.finish() + return + } + this.remoteReset = true + this.fail(new BrowserTunnelError(message.code, `Browser tunnel was reset: ${message.code}`)) + } + + private pumpWrite() { + const write = this.activeWrite + if (!write || this.sending || this.state !== "open") return + if (this.sendWindowBytes === 0 || this.sendWindowFrames === 0) return + const size = Math.min( + BrowserTunnelProtocol.MaxDataBytes, + this.sendWindowBytes, + write.data.byteLength - write.offset, + ) + this.sendWindowBytes -= size + this.sendWindowFrames-- + this.outstanding.push(size) + this.sending = true + this.send(BrowserTunnelProtocol.data(write.data.subarray(write.offset, write.offset + size)), (error) => { + this.sending = false + if (this.activeWrite !== write) return + if (error) { + this.activeWrite = undefined + write.callback(error) + this.fail(new BrowserTunnelError("transport", error.message)) + return + } + write.offset += size + if (write.offset === write.data.byteLength) { + this.activeWrite = undefined + write.callback() + return + } + this.pumpWrite() + }) + } + + private releaseReceiveWindow() { + if (this.heldWindowFrames === 0 || this.state !== "open") return + const bytes = this.heldWindowBytes + const frames = this.heldWindowFrames + this.heldWindowBytes = 0 + this.heldWindowFrames = 0 + this.receiveWindowBytes += bytes + this.receiveWindowFrames += frames + this.send( + BrowserTunnelProtocol.encodeFromDesktop({ + type: "browser.tunnel.window", + bytes: BrowserTunnel.WindowBytes.make(bytes), + frames: BrowserTunnel.FrameWindow.make(frames), + }), + (error) => { + if (error) this.fail(new BrowserTunnelError("transport", error.message)) + }, + ) + } + + private send(frame: Uint8Array, callback: (error?: Error) => void) { + if (this.socket.readyState !== WebSocket.OPEN) { + callback(new Error("Browser tunnel WebSocket is not open.")) + return + } + this.socket.send(frame, { binary: true }, callback) + } + + private protocolFailure(message: string) { + this.fail(new BrowserTunnelError("protocol_error", message)) + } + + private fail(error: BrowserTunnelError) { + if (this.state === "closed") return + if (Schema.is(BrowserTunnel.ResetCode)(error.code)) this.resetCode = error.code + if (this.state !== "open") this.rejectOpened(error) + this.destroy(error) + } + + private finish() { + if (!this.localEnded || !this.remoteEnded || this.socket.readyState !== WebSocket.OPEN) return + this.socket.close(1000) + } + + private readonly onAbort = () => { + this.fail(new BrowserTunnelError("cancelled", "Browser tunnel was cancelled.")) + } +} + +function endpointURL(endpoint: BrowserTunnelEndpoint) { + const url = new URL(endpoint.url) + if ((url.protocol !== "http:" && url.protocol !== "https:") || url.username || url.password) { + throw new TypeError("Browser server endpoint must be an HTTP URL without embedded credentials") + } + url.protocol = url.protocol === "https:" ? "wss:" : "ws:" + url.pathname = "/api/browser/tunnel" + url.search = "" + url.hash = "" + return url +} + +function rawData(data: WebSocket.RawData) { + if (data instanceof ArrayBuffer) return new Uint8Array(data) + if (Array.isArray(data)) return new Uint8Array(Buffer.concat(data)) + return new Uint8Array(data.buffer, data.byteOffset, data.byteLength) +} diff --git a/packages/client/src/node/client.ts b/packages/client/src/node/client.ts new file mode 100644 index 0000000000..2869537cac --- /dev/null +++ b/packages/client/src/node/client.ts @@ -0,0 +1,13 @@ +import { OpenCode } from "../promise/generated/index.js" +import { createBrowserClient } from "./browser/client.js" + +export type ClientOptions = OpenCode.ClientOptions +export type RequestOptions = OpenCode.RequestOptions + +/** Creates the Promise client with Node-only browser attachment support. */ +export function make(options: ClientOptions) { + return { + ...OpenCode.make(options), + browser: createBrowserClient(options), + } +} diff --git a/packages/client/src/node/index.ts b/packages/client/src/node/index.ts new file mode 100644 index 0000000000..df778c406d --- /dev/null +++ b/packages/client/src/node/index.ts @@ -0,0 +1,28 @@ +export { ClientError, type ClientErrorReason } from "../promise/generated/client-error.js" +export * from "../promise/generated/types.js" +export type { + AgentApi, + CatalogApi, + CommandApi, + EventApi, + IntegrationApi, + ModelApi, + PluginApi, + ProviderApi, + ReferenceApi, + WebSearchApi, + SessionApi, + SkillApi, +} from "../promise/api.js" +export * as OpenCode from "./client.js" +export { Browser } from "@opencode-ai/schema/browser" +export { BrowserDriver, BrowserDriverError } from "./browser/driver.js" +export type { + BrowserDriverContext, + BrowserDriverFactory, + BrowserDriverInstance, + BrowserProxy, +} from "./browser/driver.js" +export type { BrowserAttachment, BrowserAttachOptions, BrowserClient } from "./browser/client.js" +export type { EventSubscribeOutput as OpenCodeEvent } from "../promise/generated/types.js" +export type OpenCodeClient = ReturnType diff --git a/packages/client/src/promise/generated/client.ts b/packages/client/src/promise/generated/client.ts index e6edc8a856..88ade3f7d3 100644 --- a/packages/client/src/promise/generated/client.ts +++ b/packages/client/src/promise/generated/client.ts @@ -213,8 +213,8 @@ import type { WebsearchProvidersOutput, WebsearchQueryInput, WebsearchQueryOutput, -} from "./types" -import { ClientError } from "./client-error" +} from "./types.js" +import { ClientError } from "./client-error.js" export interface ClientOptions { readonly baseUrl: string diff --git a/packages/client/src/promise/generated/index.ts b/packages/client/src/promise/generated/index.ts index 2570372cf8..1a4e5ea3fe 100644 --- a/packages/client/src/promise/generated/index.ts +++ b/packages/client/src/promise/generated/index.ts @@ -1,3 +1,3 @@ -export { ClientError, type ClientErrorReason } from "./client-error" -export * as OpenCode from "./client" -export * from "./types" +export { ClientError, type ClientErrorReason } from "./client-error.js" +export * as OpenCode from "./client.js" +export * from "./types.js" diff --git a/packages/client/test/import-boundaries.test.ts b/packages/client/test/import-boundaries.test.ts index 6a979b00a7..8e09588b56 100644 --- a/packages/client/test/import-boundaries.test.ts +++ b/packages/client/test/import-boundaries.test.ts @@ -5,6 +5,7 @@ import { join, resolve, sep } from "node:path" const directory = resolve(import.meta.dir, "..") const effect = realpathSync(resolve(import.meta.dir, "../node_modules/effect")) +const ws = realpathSync(resolve(import.meta.dir, "../node_modules/ws")) const schema = resolve(import.meta.dir, "../../schema") const protocol = resolve(import.meta.dir, "../../protocol") const core = resolve(import.meta.dir, "../../core") @@ -17,6 +18,7 @@ describe("public import boundaries", () => { expect(within(root, effect)).toEqual([]) expect(within(root, schema)).toEqual([]) expect(within(root, protocol)).toEqual([]) + expect(within(root, ws)).toEqual([]) expect(within(root, core)).toEqual([]) expect(within(root, server)).toEqual([]) @@ -28,6 +30,15 @@ describe("public import boundaries", () => { expect(within(network, core)).toEqual([]) expect(within(network, server)).toEqual([]) + const node = await bundleInputs("@opencode-ai/client/node", "node") + + expect(within(node, effect).length).toBeGreaterThan(0) + expect(within(node, schema).length).toBeGreaterThan(0) + expect(within(node, protocol).length).toBeGreaterThan(0) + expect(within(node, ws).length).toBeGreaterThan(0) + expect(within(node, core)).toEqual([]) + expect(within(node, server)).toEqual([]) + const promiseService = await bundleInputs("@opencode-ai/client/service", "bun") expect(within(promiseService, effect)).toEqual([]) @@ -45,7 +56,7 @@ describe("public import boundaries", () => { }) }) -async function bundleInputs(specifier: string, target: "browser" | "bun") { +async function bundleInputs(specifier: string, target: "browser" | "bun" | "node") { const temporary = await mkdtemp(join(import.meta.dir, ".import-boundary-")) const entrypoint = join(temporary, "index.ts") const metafile = join(temporary, "meta.json") diff --git a/packages/client/test/node/browser-client.test.ts b/packages/client/test/node/browser-client.test.ts new file mode 100644 index 0000000000..7d58e661b2 --- /dev/null +++ b/packages/client/test/node/browser-client.test.ts @@ -0,0 +1,856 @@ +import { BrowserControlProtocol } from "@opencode-ai/protocol/browser-control" +import { BrowserTunnelProtocol } from "@opencode-ai/protocol/browser-tunnel" +import { BROWSER_CONTROL_PROTOCOL, BROWSER_TUNNEL_PROTOCOL } from "@opencode-ai/protocol/groups/browser" +import { BrowserControl } from "@opencode-ai/schema/browser-control" +import { BrowserTunnel } from "@opencode-ai/schema/browser-tunnel" +import { Session } from "@opencode-ai/schema/session" +import { describe, expect, test } from "bun:test" +import { Effect } from "effect" +import { once } from "node:events" +import { createServer } from "node:http" +import { connect, type TLSSocket } from "node:tls" +import WebSocket, { WebSocketServer } from "ws" +import { + Browser, + BrowserDriver, + BrowserDriverError, + OpenCode, + type BrowserDriverContext, + type BrowserDriverInstance, + type BrowserProxy, +} from "@opencode-ai/client/node" + +const initialState: Browser.State = { + url: "https://example.com/", + title: "Example", + loading: false, + canGoBack: false, + canGoForward: false, + generation: 1, +} + +describe("Node browser client", () => { + test("rejects authentication failures on Bun without unsupported ws response events", async () => { + const server = await controlServer("Basic expected") + const adapter = driver("unauthorized") + const client = OpenCode.make({ baseUrl: server.url }) + + try { + const result = await Promise.race([ + rejected(client.browser.attach({ sessionID: "ses_node_unauthorized", driver: adapter.descriptor })), + Bun.sleep(1_000).then(() => new Error("Authentication rejection timed out")), + ]) + expect(result.message).toContain("HTTP 401") + expect(server.authorizations).toEqual([undefined]) + expect(adapter.disposeCount).toBe(1) + } finally { + await server.close() + } + }) + + test("cancels a pending Bun authentication probe when attachment setup is aborted", async () => { + const abort = new AbortController() + const adapter = driver("preflight") + let resolveStarted!: () => void + const started = new Promise((resolve) => { + resolveStarted = resolve + }) + let cancelled = false + const fetch: typeof globalThis.fetch = (_input, init) => + new Promise((_resolve, reject) => { + resolveStarted() + const signal = init?.signal + if (!signal) return reject(new Error("Missing preflight signal")) + const stop = () => { + cancelled = true + reject(signal.reason) + } + signal.addEventListener("abort", stop, { once: true }) + if (signal.aborted) stop() + }) + const client = OpenCode.make({ baseUrl: "http://127.0.0.1:1", fetch }) + const attaching = client.browser.attach({ + sessionID: "ses_node_preflight_abort", + driver: adapter.descriptor, + signal: abort.signal, + }) + + await started + const reason = new Error("stop preflight") + abort.abort(reason) + expect(await rejected(attaching)).toBe(reason) + expect(cancelled).toBe(true) + expect(adapter.disposeCount).toBe(1) + }) + + test("keeps constructor authentication immutable and waits for the lease sync acknowledgement", async () => { + const authorization = `Basic ${Buffer.from("opencode:secret").toString("base64")}` + const server = await controlServer(authorization) + const headers = new Headers({ authorization }) + const adapter = driver("first") + const client = OpenCode.make({ baseUrl: server.url, headers }) + headers.set("authorization", "Basic changed") + + try { + expect(client.session).toBeDefined() + expect(client.browser.attach).toBeFunction() + let settled = false + expect(adapter.descriptor).toBe(adapter.factory) + const attaching = client.browser.attach({ sessionID: "ses_node_barrier", driver: adapter.factory }) + void attaching.then(() => { + settled = true + }) + const socket = await server.nextConnection() + const next = controlReader(socket) + + expect(server.authorizations).toEqual([authorization]) + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + const sync = await next() + if (sync.type !== "browser.control.sync") throw new Error("Expected browser control sync") + expect(sync.attachments).toHaveLength(1) + expect(sync.attachments[0]?.sessionID).toBe("ses_node_barrier") + expect(sync.attachments[0]?.leaseID.startsWith("brl_")).toBe(true) + expect(adapter.proxy?.url.startsWith("https://")).toBe(true) + expect(adapter.proxy?.credentials.username).toBeTruthy() + await Bun.sleep(10) + expect(settled).toBe(false) + + if (!adapter.proxy) throw new Error("Browser proxy is unavailable") + const proxyConnection = connectProxy(adapter.proxy, "target.example:443") + await Bun.sleep(10) + expect(server.tunnelConnections).toBe(0) + + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: sync.revision })) + const attachment = await attaching + expect(attachment.resource).toEqual({ name: "first" }) + + const tunnel = await server.nextTunnelConnection() + const nextTunnel = tunnelReader(tunnel) + tunnel.send(BrowserTunnelProtocol.encodeFromServer({ type: "browser.tunnel.ready" }), { binary: true }) + const open = await nextTunnel() + expect(open).toMatchObject({ + type: "control", + message: { + type: "browser.tunnel.open", + sessionID: "ses_node_barrier", + leaseID: sync.attachments[0]?.leaseID, + target: { host: "target.example", port: 443 }, + }, + }) + tunnel.send( + BrowserTunnelProtocol.encodeFromServer({ + type: "browser.tunnel.opened", + receiveWindow: BrowserTunnel.WindowSize.make(BrowserTunnelProtocol.InitialWindowBytes), + receiveFrames: BrowserTunnel.FrameWindow.make(BrowserTunnelProtocol.InitialFrameWindow), + }), + { binary: true }, + ) + const connected = await proxyConnection + expect(connected.response.startsWith("HTTP/1.1 200 Connection Established")).toBe(true) + connected.socket.destroy() + + const closed = once(socket, "close") + await attachment.close() + await closed + expect(adapter.disposeCount).toBe(1) + expect(adapter.unsubscribeCount).toBe(1) + } finally { + await server.close() + } + }) + + test("multiplexes Sessions, coalesces state, routes requests and cancellation, and preserves siblings", async () => { + const server = await controlServer() + const client = OpenCode.make({ baseUrl: server.url }) + const first = driver("first") + const second = driver("second") + + try { + const attachingFirst = client.browser.attach({ sessionID: "ses_node_first", driver: first.descriptor }) + const socket = await server.nextConnection() + const next = controlReader(socket) + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + const firstSync = await nextSync(next) + socket.send( + BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: firstSync.revision }), + ) + const firstAttachment = await attachingFirst + + const attachingSecond = client.browser.attach({ sessionID: "ses_node_second", driver: second.descriptor }) + const secondSync = await nextSync(next) + expect(secondSync.attachments.map((attachment) => attachment.sessionID)).toEqual([ + "ses_node_first", + "ses_node_second", + ]) + socket.send( + BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: secondSync.revision }), + ) + const secondAttachment = await attachingSecond + expect(server.connections).toBe(1) + + second.setState({ ...initialState, title: "Intermediate" }) + const pendingState = await nextSync(next) + second.setState({ ...initialState, title: "Latest" }) + socket.send( + BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: pendingState.revision }), + ) + const latestState = await nextSync(next) + expect( + latestState.attachments.find((attachment) => attachment.sessionID === "ses_node_second")?.state.title, + ).toBe("Latest") + socket.send( + BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: latestState.revision }), + ) + + const secondLease = secondSync.attachments.find( + (attachment) => attachment.sessionID === "ses_node_second", + )?.leaseID + if (!secondLease) throw new Error("Second browser lease is unavailable") + const requestID = BrowserControl.RequestID.create() + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.request", + requestID, + sessionID: Session.ID.make("ses_node_second"), + leaseID: secondLease, + command: { type: "snapshot", generation: 1 }, + }), + ) + expect(await next()).toMatchObject({ + type: "browser.control.response", + requestID, + outcome: { type: "success", result: { type: "snapshot", content: "second" } }, + }) + + const commandState = { ...initialState, title: "Updated by command", generation: 2 } + second.setExecute(async () => ({ type: "navigate", state: commandState })) + const stateRequestID = BrowserControl.RequestID.create() + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.request", + requestID: stateRequestID, + sessionID: Session.ID.make("ses_node_second"), + leaseID: secondLease, + command: { type: "navigate", url: "https://updated.example/", generation: 1 }, + }), + ) + expect(await next()).toMatchObject({ + type: "browser.control.response", + requestID: stateRequestID, + outcome: { type: "success", result: { state: { title: "Updated by command" } } }, + }) + const commandStateSync = await nextSync(next) + expect( + commandStateSync.attachments.find((attachment) => attachment.sessionID === "ses_node_second")?.state, + ).toEqual(commandState) + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.synced", + revision: commandStateSync.revision, + }), + ) + + second.setExecute(async () => ({ type: "click", state: initialState })) + const invalidID = BrowserControl.RequestID.create() + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.request", + requestID: invalidID, + sessionID: Session.ID.make("ses_node_second"), + leaseID: secondLease, + command: { type: "snapshot", generation: 1 }, + }), + ) + expect(await next()).toMatchObject({ + type: "browser.control.response", + requestID: invalidID, + outcome: { type: "failure", code: "protocol" }, + }) + + second.setExecute(async () => { + throw new BrowserDriverError("stale_ref", "The reference is stale") + }) + const errorID = BrowserControl.RequestID.create() + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.request", + requestID: errorID, + sessionID: Session.ID.make("ses_node_second"), + leaseID: secondLease, + command: { type: "click", ref: Browser.Ref.make("e1"), generation: 1 }, + }), + ) + expect(await next()).toMatchObject({ + type: "browser.control.response", + requestID: errorID, + outcome: { type: "failure", code: "stale_ref", message: "The reference is stale" }, + }) + + second.setExecute(async () => { + throw Object.assign(new Error("Invalid adapter code"), { code: "not_a_browser_error" }) + }) + const unsafeErrorID = BrowserControl.RequestID.create() + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.request", + requestID: unsafeErrorID, + sessionID: Session.ID.make("ses_node_second"), + leaseID: secondLease, + command: { type: "snapshot", generation: 2 }, + }), + ) + expect(await next()).toMatchObject({ + type: "browser.control.response", + requestID: unsafeErrorID, + outcome: { type: "failure", code: "internal", message: "Invalid adapter code" }, + }) + + let resolveStarted!: () => void + let resolveCancelled!: () => void + const started = new Promise((resolve) => { + resolveStarted = resolve + }) + const cancelled = new Promise((resolve) => { + resolveCancelled = resolve + }) + second.setExecute( + (_command, options) => + new Promise((_, reject) => { + resolveStarted() + options.signal.addEventListener( + "abort", + () => { + resolveCancelled() + reject(Object.assign(new Error("cancelled"), { code: "aborted" })) + }, + { once: true }, + ) + }), + ) + const cancelID = BrowserControl.RequestID.create() + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.request", + requestID: cancelID, + sessionID: Session.ID.make("ses_node_second"), + leaseID: secondLease, + command: { type: "click", ref: Browser.Ref.make("e1"), generation: 1 }, + }), + ) + await started + socket.send( + BrowserControlProtocol.encodeFromServer({ + type: "browser.control.cancel", + requestID: cancelID, + leaseID: secondLease, + }), + ) + await cancelled + + await firstAttachment.close() + const detached = await nextSync(next) + expect(detached.attachments).toHaveLength(1) + expect(detached.attachments[0]?.sessionID).toBe("ses_node_second") + socket.send( + BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: detached.revision }), + ) + expect(server.connections).toBe(1) + expect(second.disposeCount).toBe(0) + + const closed = once(socket, "close") + await secondAttachment.close() + await closed + expect(first.disposeCount).toBe(1) + expect(second.disposeCount).toBe(1) + } finally { + await server.close() + } + }) + + test("reconnects and republishes the complete attachment snapshot", async () => { + const server = await controlServer() + const client = OpenCode.make({ baseUrl: server.url }) + const adapter = driver("reconnect") + + try { + const attaching = client.browser.attach({ sessionID: "ses_node_reconnect", driver: adapter.descriptor }) + const original = await server.nextConnection() + const nextOriginal = controlReader(original) + original.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + const initial = await nextSync(nextOriginal) + original.send( + BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: initial.revision }), + ) + const attachment = await attaching + original.terminate() + + const replacement = await server.nextConnection() + const nextReplacement = controlReader(replacement) + replacement.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + const republished = await nextSync(nextReplacement) + expect(republished.attachments).toEqual(initial.attachments) + replacement.send( + BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: republished.revision }), + ) + expect(server.connections).toBe(2) + + const closed = once(replacement, "close") + await attachment.close() + await closed + } finally { + await server.close() + } + }) + + test("rejects duplicate Sessions and aborts setup while close remains idempotent", async () => { + const server = await controlServer() + const client = OpenCode.make({ baseUrl: server.url }) + const adapter = driver("owner") + + try { + const invalid = driver("invalid") + expect( + (await rejected(client.browser.attach({ sessionID: "invalid", driver: invalid.descriptor }))).message, + ).toContain("valid Session ID") + expect(invalid.createCount).toBe(0) + + const attaching = client.browser.attach({ sessionID: "ses_node_owner", driver: adapter.descriptor }) + const socket = await server.nextConnection() + const next = controlReader(socket) + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + const sync = await nextSync(next) + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: sync.revision })) + const attachment = await attaching + + const duplicate = driver("duplicate") + expect( + (await rejected(client.browser.attach({ sessionID: "ses_node_owner", driver: duplicate.descriptor }))).message, + ).toContain("already attached") + expect(duplicate.createCount).toBe(0) + + const closed = once(socket, "close") + await Promise.all([attachment.close(), attachment.close(), attachment[Symbol.asyncDispose]()]) + await closed + expect(adapter.disposeCount).toBe(1) + expect(adapter.unsubscribeCount).toBe(1) + + const abort = new AbortController() + const abortedAdapter = driver("aborted") + const aborted = client.browser.attach({ + sessionID: "ses_node_aborted", + driver: abortedAdapter.descriptor, + signal: abort.signal, + }) + const abortSocket = await server.nextConnection() + const nextAbort = controlReader(abortSocket) + abortSocket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + await nextSync(nextAbort) + const reason = new Error("stop attaching") + const abortClosed = once(abortSocket, "close") + abort.abort(reason) + expect(await rejected(aborted)).toBe(reason) + await abortClosed + expect(abortedAdapter.disposeCount).toBe(1) + expect(abortedAdapter.signal?.aborted).toBe(true) + + const alreadyAborted = new AbortController() + alreadyAborted.abort(reason) + const unused = driver("unused") + expect( + await rejected( + client.browser.attach({ + sessionID: "ses_node_unused", + driver: unused.descriptor, + signal: alreadyAborted.signal, + }), + ), + ).toBe(reason) + expect(unused.createCount).toBe(0) + } finally { + await server.close() + } + }) + + test("fails pending attachments on fatal control closes without reconnecting", async () => { + const server = await controlServer() + const client = OpenCode.make({ baseUrl: server.url }) + const adapter = driver("fatal") + + try { + const attaching = client.browser.attach({ sessionID: "ses_node_fatal", driver: adapter.descriptor }) + const socket = await server.nextConnection() + const next = controlReader(socket) + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + await nextSync(next) + const closed = once(socket, "close") + socket.close(1002, "invalid control state") + + expect((await rejected(attaching)).message).toContain("fatal code 1002") + await closed + await Bun.sleep(250) + expect(server.connections).toBe(1) + expect(adapter.disposeCount).toBe(1) + expect(adapter.signal?.aborted).toBe(true) + } finally { + await server.close() + } + }) + + test("aborts non-cooperative driver creation and disposes a late instance", async () => { + const server = await controlServer() + const client = OpenCode.make({ baseUrl: server.url }) + const abort = new AbortController() + let resolveStarted!: () => void + let resolveFactory!: (instance: BrowserDriverInstance<{ readonly name: string }>) => void + let resolveDisposed!: () => void + const started = new Promise((resolve) => { + resolveStarted = resolve + }) + const disposed = new Promise((resolve) => { + resolveDisposed = resolve + }) + const factory = () => { + resolveStarted() + return new Promise>((resolve) => { + resolveFactory = resolve + }) + } + + try { + const attaching = client.browser.attach({ + sessionID: "ses_node_blocked_factory", + driver: factory, + signal: abort.signal, + }) + await started + const reason = new Error("abort blocked factory") + abort.abort(reason) + expect( + await rejected( + Promise.race([ + attaching, + Bun.sleep(250).then(() => { + throw new Error("Browser attachment did not abort promptly") + }), + ]), + ), + ).toBe(reason) + + resolveFactory({ + resource: { name: "late" }, + state: () => initialState, + subscribe: () => () => undefined, + execute: async () => ({ + type: "snapshot", + state: initialState, + format: "opencode.semantic.v1", + content: "late", + }), + dispose: () => resolveDisposed(), + }) + await disposed + expect(server.authorizations).toEqual([]) + } finally { + await server.close() + } + }) + + test("cleans up in reverse acquisition order and surfaces the first error", async () => { + const server = await controlServer() + const client = OpenCode.make({ baseUrl: server.url }) + const first = new Error("unsubscribe failed") + const second = new Error("dispose failed") + const order: string[] = [] + let proxy: BrowserProxy | undefined + let proxyOpenDuringDispose = false + const descriptor = BrowserDriver.define((context) => { + proxy = context.proxy + return { + resource: context.proxy, + state: () => initialState, + subscribe: () => () => { + order.push("unsubscribe") + throw first + }, + execute: async () => ({ + type: "snapshot", + state: initialState, + format: "opencode.semantic.v1", + content: "cleanup", + }), + dispose: async () => { + order.push("dispose") + proxyOpenDuringDispose = await proxyAvailable(context.proxy) + throw second + }, + } + }) + + try { + const attaching = client.browser.attach({ sessionID: "ses_node_cleanup", driver: descriptor }) + const socket = await server.nextConnection() + const next = controlReader(socket) + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.ready" })) + const sync = await nextSync(next) + socket.send(BrowserControlProtocol.encodeFromServer({ type: "browser.control.synced", revision: sync.revision })) + const attachment = await attaching + + expect(await rejected(attachment.close())).toBe(first) + expect(await rejected(attachment.close())).toBe(first) + expect(order).toEqual(["unsubscribe", "dispose"]) + expect(proxyOpenDuringDispose).toBe(true) + if (!proxy) throw new Error("Browser proxy is unavailable") + expect(await proxyAvailable(proxy)).toBe(false) + } finally { + await server.close() + } + }) +}) + +type Execute = BrowserDriverInstance["execute"] + +function driver(name: string) { + let current = initialState + let execute: Execute = async () => ({ + type: "snapshot", + state: current, + format: "opencode.semantic.v1", + content: name, + }) + let proxy: BrowserProxy | undefined + let signal: AbortSignal | undefined + let createCount = 0 + let disposeCount = 0 + let unsubscribeCount = 0 + const listeners = new Set<(state: Browser.State) => void>() + const factory = ({ proxy: nextProxy, signal: nextSignal }: BrowserDriverContext) => { + createCount++ + proxy = nextProxy + signal = nextSignal + return { + resource: { name }, + state: () => current, + subscribe(listener) { + listeners.add(listener) + return () => { + if (listeners.delete(listener)) unsubscribeCount++ + } + }, + execute: (command, options) => execute(command, options), + dispose() { + disposeCount++ + }, + } + } + const descriptor = BrowserDriver.define(factory) + return { + factory, + descriptor, + get proxy() { + return proxy + }, + get signal() { + return signal + }, + get createCount() { + return createCount + }, + get disposeCount() { + return disposeCount + }, + get unsubscribeCount() { + return unsubscribeCount + }, + setState(state: Browser.State) { + current = state + listeners.forEach((listener) => listener(state)) + }, + setExecute(next: Execute) { + execute = next + }, + } +} + +async function rejected(promise: Promise) { + return promise.then( + () => new Error("Expected operation to reject"), + (error) => (error instanceof Error ? error : new Error(String(error))), + ) +} + +async function nextSync(next: () => Promise) { + const message = await next() + if (message.type !== "browser.control.sync") throw new Error("Expected browser control sync") + return message +} + +async function controlServer(authorization?: string) { + const http = createServer((request, response) => { + response.statusCode = request.headers.authorization === authorization ? 200 : 401 + response.end() + }) + const webSockets = new WebSocketServer({ noServer: true }) + const tunnels = new WebSocketServer({ noServer: true }) + const waiting: Array<(socket: WebSocket) => void> = [] + const queued: WebSocket[] = [] + const tunnelWaiting: Array<(socket: WebSocket) => void> = [] + const tunnelQueued: WebSocket[] = [] + const authorizations: Array = [] + let connections = 0 + let tunnelConnections = 0 + webSockets.on("connection", (socket) => { + connections++ + const resolve = waiting.shift() + if (resolve) resolve(socket) + else queued.push(socket) + }) + tunnels.on("connection", (socket) => { + tunnelConnections++ + const resolve = tunnelWaiting.shift() + if (resolve) resolve(socket) + else tunnelQueued.push(socket) + }) + http.on("upgrade", (request, socket, head) => { + authorizations.push(request.headers.authorization) + if (request.headers.authorization !== authorization) { + socket.end("HTTP/1.1 401 Unauthorized\r\n\r\n") + return + } + if ( + request.url === "/api/browser/control" && + request.headers["sec-websocket-protocol"] === BROWSER_CONTROL_PROTOCOL + ) { + webSockets.handleUpgrade(request, socket, head, (webSocket) => webSockets.emit("connection", webSocket, request)) + return + } + if ( + request.url === "/api/browser/tunnel" && + request.headers["sec-websocket-protocol"] === BROWSER_TUNNEL_PROTOCOL + ) { + tunnels.handleUpgrade(request, socket, head, (webSocket) => tunnels.emit("connection", webSocket, request)) + return + } + socket.end("HTTP/1.1 401 Unauthorized\r\n\r\n") + }) + await new Promise((resolve) => http.listen(0, "127.0.0.1", resolve)) + const address = http.address() + if (address === null || typeof address === "string") throw new Error("Browser control server did not bind TCP") + return { + authorizations, + get connections() { + return connections + }, + get tunnelConnections() { + return tunnelConnections + }, + url: `http://127.0.0.1:${address.port}`, + nextConnection() { + const socket = queued.shift() + if (socket) return Promise.resolve(socket) + return new Promise((resolve) => waiting.push(resolve)) + }, + nextTunnelConnection() { + const socket = tunnelQueued.shift() + if (socket) return Promise.resolve(socket) + return new Promise((resolve) => tunnelWaiting.push(resolve)) + }, + async close() { + webSockets.clients.forEach((socket) => socket.terminate()) + tunnels.clients.forEach((socket) => socket.terminate()) + webSockets.close() + tunnels.close() + http.closeAllConnections() + http.close() + await Bun.sleep(10) + }, + } +} + +function proxyAvailable(proxy: BrowserProxy) { + return new Promise((resolve) => { + const socket = connect({ + host: "127.0.0.1", + port: proxy.port, + servername: proxy.host, + rejectUnauthorized: false, + }) + let settled = false + const finish = (available: boolean) => { + if (settled) return + settled = true + clearTimeout(timeout) + socket.destroy() + resolve(available) + } + const timeout = setTimeout(() => finish(false), 250) + timeout.unref() + socket.once("secureConnect", () => finish(true)) + socket.once("error", () => finish(false)) + }) +} + +function connectProxy(proxy: BrowserProxy, authority: string) { + return new Promise<{ readonly socket: TLSSocket; readonly response: string }>((resolve, reject) => { + const socket = connect( + { + host: "127.0.0.1", + port: proxy.port, + servername: proxy.host, + rejectUnauthorized: false, + }, + () => { + socket.write( + `CONNECT ${authority} HTTP/1.1\r\nHost: ${authority}\r\nProxy-Authorization: Basic ${Buffer.from(`${proxy.credentials.username}:${proxy.credentials.password}`).toString("base64")}\r\n\r\n`, + ) + }, + ) + let response = "" + socket.on("error", () => undefined) + socket.once("error", reject) + socket.on("data", (data) => { + response += data.toString() + if (!response.includes("\r\n\r\n")) return + resolve({ socket, response }) + }) + }) +} + +function tunnelReader(socket: WebSocket) { + const queued: WebSocket.RawData[] = [] + const waiting: Array<(data: WebSocket.RawData) => void> = [] + socket.on("message", (data, binary) => { + if (!binary) throw new Error("Expected a binary browser tunnel frame") + const resolve = waiting.shift() + if (resolve) resolve(data) + else queued.push(data) + }) + return async () => { + const data = queued.shift() ?? (await new Promise((resolve) => waiting.push(resolve))) + return Effect.runPromise(BrowserTunnelProtocol.decodeFromDesktop(rawData(data))) + } +} + +function controlReader(socket: WebSocket) { + const queued: Array<{ readonly data: WebSocket.RawData; readonly binary: boolean }> = [] + const waiting: Array<(message: { readonly data: WebSocket.RawData; readonly binary: boolean }) => void> = [] + socket.on("message", (data, binary) => { + const resolve = waiting.shift() + if (resolve) resolve({ data, binary }) + else queued.push({ data, binary }) + }) + return async () => { + const message = + queued.shift() ?? + (await new Promise<{ readonly data: WebSocket.RawData; readonly binary: boolean }>((resolve) => + waiting.push(resolve), + )) + if (message.binary) throw new Error("Expected a text browser control message") + return Effect.runPromise(BrowserControlProtocol.decodeFromDesktop(rawData(message.data))) + } +} + +function rawData(data: WebSocket.RawData) { + if (data instanceof ArrayBuffer) return new Uint8Array(data) + if (Array.isArray(data)) return new Uint8Array(Buffer.concat(data)) + return new Uint8Array(data.buffer, data.byteOffset, data.byteLength) +} diff --git a/packages/client/test/node/package-smoke.ts b/packages/client/test/node/package-smoke.ts new file mode 100644 index 0000000000..33049538bb --- /dev/null +++ b/packages/client/test/node/package-smoke.ts @@ -0,0 +1,164 @@ +import { expect, test } from "bun:test" +import { mkdir, mkdtemp, rm } from "node:fs/promises" +import { join, relative, resolve } from "node:path" +import { pathToFileURL } from "node:url" + +const directory = resolve(import.meta.dir, "../..") + +test("built Node entrypoint imports in Node and treats HTTP authentication rejection as fatal", async () => { + await buildClient() + const output = await Bun.file(join(directory, "dist/node/index.js")).text() + expect(output).not.toMatch(/(?:from\s+|import\s*)["']\.\.?\//) + + const temporary = await mkdtemp(join(import.meta.dir, ".node-package-")) + try { + await Bun.write(join(temporary, "index.mjs"), output) + await stageWorkspaceDependencies(temporary) + const child = Bun.spawn( + ["node", "--input-type=module", "-e", nodeScenario(pathToFileURL(join(temporary, "index.mjs")).href)], + { cwd: temporary, stdout: "pipe", stderr: "pipe" }, + ) + const [exitCode, stdout, stderr] = await Promise.all([ + child.exited, + new Response(child.stdout).text(), + new Response(child.stderr).text(), + ]) + if (exitCode !== 0) throw new Error(stderr || stdout) + expect(stdout.trim()).toBe("ok") + } finally { + await rm(temporary, { recursive: true, force: true }) + } +}, 60_000) + +async function buildClient() { + const child = Bun.spawn([process.execPath, "run", "build"], { + cwd: directory, + stdout: "pipe", + stderr: "pipe", + }) + const [exitCode, stdout, stderr] = await Promise.all([ + child.exited, + new Response(child.stdout).text(), + new Response(child.stderr).text(), + ]) + if (exitCode !== 0) throw new Error(stdout + stderr) +} + +async function stageWorkspaceDependencies(temporary: string) { + const schema = join(temporary, "node_modules/@opencode-ai/schema") + const protocol = join(temporary, "node_modules/@opencode-ai/protocol") + await Promise.all([mkdir(schema, { recursive: true }), mkdir(protocol, { recursive: true })]) + + const schemaEntry = join(temporary, "schema.ts") + const protocolEntry = join(temporary, "protocol.ts") + await Promise.all([ + Bun.write( + schemaEntry, + [ + `export { Browser } from ${JSON.stringify(importPath(temporary, resolve(directory, "../schema/src/browser.ts")))}`, + `export { BrowserControl } from ${JSON.stringify(importPath(temporary, resolve(directory, "../schema/src/browser-control.ts")))}`, + `export { BrowserTunnel } from ${JSON.stringify(importPath(temporary, resolve(directory, "../schema/src/browser-tunnel.ts")))}`, + `export { Session } from ${JSON.stringify(importPath(temporary, resolve(directory, "../schema/src/session.ts")))}`, + ].join("\n"), + ), + Bun.write( + protocolEntry, + [ + `export { BrowserControlProtocol } from ${JSON.stringify(importPath(temporary, resolve(directory, "../protocol/src/browser-control.ts")))}`, + `export { BrowserTunnelProtocol } from ${JSON.stringify(importPath(temporary, resolve(directory, "../protocol/src/browser-tunnel.ts")))}`, + `export { BROWSER_CONTROL_PROTOCOL, BROWSER_TUNNEL_PROTOCOL } from ${JSON.stringify(importPath(temporary, resolve(directory, "../protocol/src/groups/browser.ts")))}`, + ].join("\n"), + ), + ]) + const [schemaBuild, protocolBuild] = await Promise.all([ + Bun.build({ + entrypoints: [schemaEntry], + outdir: schema, + naming: "index.js", + target: "node", + format: "esm", + packages: "bundle", + }), + Bun.build({ + entrypoints: [protocolEntry], + outdir: protocol, + naming: "index.js", + target: "node", + format: "esm", + packages: "bundle", + }), + ]) + if (!schemaBuild.success) throw new Error(schemaBuild.logs.map((log) => log.message).join("\n")) + if (!protocolBuild.success) throw new Error(protocolBuild.logs.map((log) => log.message).join("\n")) + await Promise.all([ + Bun.write( + join(schema, "package.json"), + JSON.stringify({ + type: "module", + exports: { + "./browser": "./index.js", + "./browser-control": "./index.js", + "./browser-tunnel": "./index.js", + "./session": "./index.js", + }, + }), + ), + Bun.write( + join(protocol, "package.json"), + JSON.stringify({ + type: "module", + exports: { + "./browser-control": "./index.js", + "./browser-tunnel": "./index.js", + "./groups/browser": "./index.js", + }, + }), + ), + ]) +} + +function importPath(from: string, to: string) { + const path = relative(from, to).replaceAll("\\", "/") + return path.startsWith(".") ? path : `./${path}` +} + +function nodeScenario(moduleURL: string) { + return `import { createServer } from "node:http" +const sdk = await import(${JSON.stringify(moduleURL)}) +if (typeof sdk.OpenCode.make !== "function") throw new Error("Missing OpenCode.make") +if (typeof sdk.BrowserDriver.define !== "function") throw new Error("Missing BrowserDriver.define") +if (typeof sdk.BrowserDriverError !== "function") throw new Error("Missing BrowserDriverError") +if (!sdk.Browser.State) throw new Error("Missing canonical Browser export") + +let upgrades = 0 +const server = createServer() +server.on("upgrade", (_request, socket) => { + upgrades++ + socket.end("HTTP/1.1 401 Unauthorized\\r\\nContent-Length: 0\\r\\nConnection: close\\r\\n\\r\\n") +}) +await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)) +const address = server.address() +if (!address || typeof address === "string") throw new Error("Server did not bind") +let disposed = 0 +const driver = () => ({ + resource: undefined, + state: () => ({ url: "about:blank", title: "", loading: false, canGoBack: false, canGoForward: false, generation: 0 }), + subscribe: () => () => undefined, + execute: async () => { throw new sdk.BrowserDriverError("internal", "unavailable") }, + dispose: () => { disposed++ }, +}) +const client = sdk.OpenCode.make({ baseUrl: \`http://127.0.0.1:\${address.port}\` }) +const error = await Promise.race([ + client.browser.attach({ sessionID: "ses_node_auth_rejection", driver }).then( + () => new Error("Expected authentication rejection"), + (cause) => cause, + ), + new Promise((resolve) => setTimeout(() => resolve(new Error("Authentication rejection timed out")), 2_000)), +]) +await new Promise((resolve) => setTimeout(resolve, 250)) +await new Promise((resolve) => server.close(resolve)) +if (!(error instanceof Error) || !error.message.includes("HTTP 401")) throw error +if (upgrades !== 1) throw new Error(\`Expected one upgrade, received \${upgrades}\`) +if (disposed !== 1) throw new Error(\`Expected one driver disposal, received \${disposed}\`) +console.log("ok")` +} diff --git a/packages/client/test/node/proxy-certificate.test.ts b/packages/client/test/node/proxy-certificate.test.ts new file mode 100644 index 0000000000..e0b1030644 --- /dev/null +++ b/packages/client/test/node/proxy-certificate.test.ts @@ -0,0 +1,22 @@ +import { createHash, X509Certificate } from "node:crypto" +import { createSecureContext } from "node:tls" +import { describe, expect, test } from "bun:test" +import { createBrowserProxyCertificate } from "../../src/node/browser/proxy-certificate" + +describe("browser proxy certificate", () => { + test("creates a unique valid certificate pinned to its loopback names", () => { + const hostname = "opencode-test.localhost" + const first = createBrowserProxyCertificate(hostname) + const second = createBrowserProxyCertificate(hostname) + const certificate = new X509Certificate(first.certificate) + const fingerprint = createHash("sha256").update(certificate.raw).digest("base64") + + expect(first.key.startsWith("-----BEGIN PRIVATE KEY-----")).toBe(true) + expect(first.certificate.startsWith("-----BEGIN CERTIFICATE-----")).toBe(true) + expect(first.fingerprint).toBe(`sha256/${fingerprint}`) + expect(certificate.subjectAltName).toContain(`DNS:${hostname}`) + expect(certificate.subjectAltName).toContain("IP Address:127.0.0.1") + expect(first.fingerprint).not.toBe(second.fingerprint) + expect(() => createSecureContext({ key: first.key, cert: first.certificate })).not.toThrow() + }) +}) diff --git a/packages/client/test/node/proxy.test.ts b/packages/client/test/node/proxy.test.ts new file mode 100644 index 0000000000..0d7d7dbd3a --- /dev/null +++ b/packages/client/test/node/proxy.test.ts @@ -0,0 +1,452 @@ +import { createHash } from "node:crypto" +import { once } from "node:events" +import { mkdtemp, rm } from "node:fs/promises" +import { createServer } from "node:http" +import { Socket } from "node:net" +import { tmpdir } from "node:os" +import { join } from "node:path" +import { Duplex } from "node:stream" +import { connect, type TLSSocket } from "node:tls" +import { pathToFileURL } from "node:url" +import { describe, expect, test } from "bun:test" +import { createBrowserProxy } from "../../src/node/browser/proxy" + +describe("browser proxy", () => { + test("forwards authenticated CONNECT streams without directly dialing the target", async () => { + let directConnections = 0 + const decoy = createServer((_incoming, response) => response.end("DIRECT DIAL")) + decoy.on("connection", () => directConnections++) + await new Promise((resolve) => decoy.listen(0, "127.0.0.1", resolve)) + const decoyAddress = decoy.address() + if (decoyAddress === null || typeof decoyAddress === "string") throw new Error("decoy server did not bind TCP") + const target = createServer((incoming, response) => response.end(`TARGET ${incoming.method} ${incoming.url}`)) + await new Promise((resolve) => target.listen(0, "127.0.0.1", resolve)) + const targetAddress = target.address() + if (targetAddress === null || typeof targetAddress === "string") throw new Error("target server did not bind TCP") + const opened: Array<{ host: string; port: number }> = [] + const proxy = await createBrowserProxy({ + connect: async (destination, signal) => { + opened.push(destination) + const socket = new Socket({ allowHalfOpen: true }) + const abort = () => socket.destroy(new Error("cancelled")) + signal.addEventListener("abort", abort, { once: true }) + socket.once("close", () => signal.removeEventListener("abort", abort)) + // Emulate the server-side tunnel dial. The local decoy must never receive a connection. + socket.connect(targetAddress.port, "127.0.0.1") + await once(socket, "connect") + return socket + }, + }) + try { + const tunnel = await proxyConnect(proxy, `127.0.0.1:${decoyAddress.port}`) + expect(tunnel.status).toBe(200) + tunnel.socket.write( + `GET /connected HTTP/1.1\r\nHost: 127.0.0.1:${decoyAddress.port}\r\nConnection: close\r\n\r\n`, + ) + const response = await readAll(tunnel.socket, tunnel.leftover) + expect(response).toContain("TARGET GET /connected") + expect(opened).toEqual([{ host: "127.0.0.1", port: decoyAddress.port }]) + + const certificate = await peerCertificate(proxy) + expect(`sha256/${createHash("sha256").update(certificate).digest("base64")}`).toBe(proxy.certificateFingerprint) + expect(directConnections).toBe(0) + } finally { + await proxy.close() + await Promise.all([ + new Promise((resolve) => target.close(() => resolve())), + new Promise((resolve) => decoy.close(() => resolve())), + ]) + } + }) + + test("forwards absolute-form HTTP with Node framing and closes held-open targets", async () => { + expect(await nodeAbsoluteFormScenario()).toEqual({ + unauthorized: { status: 407, body: "" }, + forwarded: { status: 200, body: "TARGET GET /absolute?q=1" }, + posted: { status: 200, body: "TARGET POST /post request body" }, + chunked: { status: 200, body: "TARGET CHUNKED" }, + held: { status: 200, body: "ok" }, + opened: 4, + destinationsCorrect: true, + directConnections: 0, + heldDestroyed: true, + }) + }) + + test("limits concurrent tunnel setup to six connections", async () => { + let active = 0 + let maximum = 0 + let started = 0 + const releases: Array<() => void> = [] + const proxy = await createBrowserProxy({ + connect: (_destination, signal) => + new Promise((resolve, reject) => { + active++ + started++ + maximum = Math.max(maximum, active) + const finish = (result: () => void) => { + signal.removeEventListener("abort", abort) + active-- + result() + } + const abort = () => finish(() => reject(new Error("cancelled"))) + signal.addEventListener("abort", abort, { once: true }) + releases.push(() => + finish(() => + resolve( + new Duplex({ + read() {}, + write(_chunk, _encoding, callback) { + callback() + }, + }), + ), + ), + ) + }), + }) + const connections = Array.from({ length: 7 }, () => + proxyConnect(proxy, "target.example:443").then( + (result) => result.socket.destroy(), + () => undefined, + ), + ) + try { + await waitFor(() => started === 6) + expect(maximum).toBe(6) + expect(started).toBe(6) + + releases.shift()?.() + await waitFor(() => started === 7) + expect(maximum).toBe(6) + } finally { + await proxy.close() + await Promise.all(connections) + } + }) + + test("aborts pending CONNECT setup when the downstream client closes", async () => { + let aborted = false + let resolveStarted!: () => void + const started = new Promise((resolve) => { + resolveStarted = resolve + }) + const proxy = await createBrowserProxy({ + connect: (_destination, signal) => + new Promise((_resolve, reject) => { + resolveStarted() + const abort = () => { + aborted = true + reject(new Error("cancelled")) + } + signal.addEventListener("abort", abort, { once: true }) + if (signal.aborted) abort() + }), + }) + const socket = proxyClient( + proxy, + `CONNECT target.example:443 HTTP/1.1\r\nHost: target.example:443\r\nProxy-Authorization: ${authorization(proxy)}\r\n\r\n`, + ) + try { + await started + socket.destroy() + await waitFor(() => aborted) + } finally { + socket.destroy() + await proxy.close() + } + }) + + test("aborts pending absolute-form setup when the downstream client closes", async () => { + let aborted = false + let resolveStarted!: () => void + const started = new Promise((resolve) => { + resolveStarted = resolve + }) + const proxy = await createBrowserProxy({ + connect: (_destination, signal) => + new Promise((_resolve, reject) => { + resolveStarted() + const abort = () => { + aborted = true + reject(new Error("cancelled")) + } + signal.addEventListener("abort", abort, { once: true }) + if (signal.aborted) abort() + }), + }) + const socket = proxyClient( + proxy, + `GET http://target.example/pending HTTP/1.1\r\nHost: target.example\r\nProxy-Authorization: ${authorization(proxy)}\r\n\r\n`, + ) + try { + await started + socket.destroy() + await waitFor(() => aborted) + } finally { + socket.destroy() + await proxy.close() + } + }) + + test("closes absolute-form requests when the target tunnel cannot connect", async () => { + const proxy = await createBrowserProxy({ connect: () => Promise.reject(new Error("target unavailable")) }) + const socket = proxyClient( + proxy, + `GET http://target.example/unavailable HTTP/1.1\r\nHost: target.example\r\nProxy-Authorization: ${authorization(proxy)}\r\n\r\n`, + ) + const chunks: Buffer[] = [] + socket.on("data", (chunk) => chunks.push(chunk)) + try { + await once(socket, "close") + expect(Buffer.concat(chunks).toString()).not.toContain("502 Bad Gateway") + } finally { + socket.destroy() + await proxy.close() + } + }) +}) + +type ProxyInfo = Awaited> + +function authorization(proxy: ProxyInfo) { + return `Basic ${Buffer.from(`${proxy.credentials.username}:${proxy.credentials.password}`).toString("base64")}` +} + +function proxyClient(proxy: ProxyInfo, request: string) { + const socket = connect( + { + host: "127.0.0.1", + port: proxy.port, + servername: proxy.host.endsWith(".localhost") ? proxy.host : undefined, + rejectUnauthorized: false, + }, + () => socket.write(request), + ) + socket.on("error", () => undefined) + return socket +} + +async function nodeAbsoluteFormScenario() { + const directory = await mkdtemp(join(tmpdir(), "opencode-browser-proxy-")) + try { + const built = await Bun.build({ + entrypoints: [join(import.meta.dir, "../../src/node/browser/proxy.ts")], + outdir: directory, + naming: "browser-proxy.mjs", + target: "node", + format: "esm", + }) + if (!built.success) throw new Error(built.logs.map((log) => log.message).join("\n")) + const output = built.outputs[0] + if (!output) throw new Error("Browser proxy Node bundle was not emitted") + const child = Bun.spawn( + ["node", "--input-type=module", "-e", absoluteFormScenario(pathToFileURL(output.path).href)], + { stdout: "pipe", stderr: "pipe" }, + ) + const [code, stdout, stderr] = await Promise.all([ + child.exited, + new Response(child.stdout).text(), + new Response(child.stderr).text(), + ]) + if (code !== 0) throw new Error(stderr || stdout) + const result: unknown = JSON.parse(stdout) + return result + } finally { + await rm(directory, { recursive: true, force: true }) + } +} + +function absoluteFormScenario(moduleURL: string) { + return `import { createBrowserProxy } from ${JSON.stringify(moduleURL)} +import { once } from "node:events" +import { Agent, createServer, request } from "node:http" +import { Socket } from "node:net" +import { Duplex } from "node:stream" +import { connect } from "node:tls" + +const proxyRequest = async (proxy, path, authenticated, options = {}) => { + const socket = await new Promise((resolve, reject) => { + const socket = connect({ host: "127.0.0.1", port: proxy.port, servername: proxy.host, rejectUnauthorized: false }, () => resolve(socket)) + socket.once("error", reject) + }) + const agent = new Agent({ keepAlive: false, maxSockets: 1, noDelay: false }) + agent.createConnection = () => socket + const body = Buffer.from(options.body ?? "") + try { + return await new Promise((resolve, reject) => { + const incoming = request({ + agent, + hostname: proxy.host, + path, + method: options.method ?? "GET", + headers: { + connection: "close", + ...(authenticated ? { "proxy-authorization": "Basic " + Buffer.from(proxy.credentials.username + ":" + proxy.credentials.password).toString("base64") } : {}), + ...(options.chunked ? { "transfer-encoding": "chunked" } : { "content-length": body.byteLength.toString() }), + }, + }, (response) => { + const chunks = [] + response.on("data", (chunk) => chunks.push(chunk)) + response.once("end", () => resolve({ status: response.statusCode ?? 0, body: Buffer.concat(chunks).toString() })) + }) + incoming.once("error", reject) + incoming.end(body) + }) + } finally { + agent.destroy() + socket.destroy() + } +} + +let directConnections = 0 +const decoy = createServer((_incoming, response) => response.end("DIRECT DIAL")) +decoy.on("connection", () => directConnections++) +decoy.listen(0, "127.0.0.1") +await once(decoy, "listening") +const decoyAddress = decoy.address() +if (!decoyAddress || typeof decoyAddress === "string") throw new Error("decoy server did not bind TCP") + +const targetServer = createServer((incoming, response) => { + const chunks = [] + incoming.on("data", (chunk) => chunks.push(chunk)) + incoming.on("end", () => { + response.writeHead(200, { "content-type": "text/plain" }) + if (incoming.url === "/chunked") { + response.write("TARGET ") + response.end("CHUNKED") + return + } + const body = Buffer.concat(chunks).toString() + response.end("TARGET " + incoming.method + " " + incoming.url + (body ? " " + body : "")) + }) +}) +targetServer.listen(0, "127.0.0.1") +await once(targetServer, "listening") +const targetAddress = targetServer.address() +if (!targetAddress || typeof targetAddress === "string") throw new Error("target server did not bind TCP") + +const opened = [] +let held +const proxy = await createBrowserProxy({ + connect: async (destination, signal) => { + opened.push(destination) + if (destination.host === "held.example") { + let responded = false + held = new Duplex({ + read() {}, + write(chunk, _encoding, callback) { + if (!responded && Buffer.from(chunk).includes(Buffer.from("\\r\\n\\r\\n"))) { + responded = true + this.push(Buffer.from("HTTP/1.1 200 OK\\r\\nContent-Length: 2\\r\\n\\r\\nok")) + } + callback() + }, + }) + return held + } + const socket = new Socket({ allowHalfOpen: true }) + const abort = () => socket.destroy(new Error("cancelled")) + signal.addEventListener("abort", abort, { once: true }) + socket.once("close", () => signal.removeEventListener("abort", abort)) + socket.connect(targetAddress.port, "127.0.0.1") + await once(socket, "connect") + return socket + }, +}) + +try { + const base = "http://127.0.0.1:" + decoyAddress.port + const unauthorized = await proxyRequest(proxy, base + "/private", false) + const forwarded = await proxyRequest(proxy, base + "/absolute?q=1", true) + const posted = await proxyRequest(proxy, base + "/post", true, { method: "POST", body: "request body", chunked: true }) + const chunked = await proxyRequest(proxy, base + "/chunked", true) + const heldResult = await proxyRequest(proxy, "http://held.example/held-open", true) + await new Promise((resolve) => setImmediate(resolve)) + console.log(JSON.stringify({ + unauthorized, + forwarded, + posted, + chunked, + held: heldResult, + opened: opened.length, + destinationsCorrect: opened.slice(0, 3).every((item) => item.host === "127.0.0.1" && item.port === decoyAddress.port) && opened[3]?.host === "held.example" && opened[3]?.port === 80, + directConnections, + heldDestroyed: held?.destroyed === true, + })) +} finally { + await proxy.close() + targetServer.closeAllConnections() + decoy.closeAllConnections() + await Promise.all([new Promise((resolve) => targetServer.close(resolve)), new Promise((resolve) => decoy.close(resolve))]) +}` +} + +function proxyConnect(proxy: ProxyInfo, destination: string) { + return new Promise<{ status: number; socket: TLSSocket; leftover: Buffer }>((resolve, reject) => { + const socket = connect( + { + host: "127.0.0.1", + port: proxy.port, + servername: proxy.host.endsWith(".localhost") ? proxy.host : undefined, + rejectUnauthorized: false, + }, + () => { + socket.write( + `CONNECT ${destination} HTTP/1.1\r\nHost: ${destination}\r\nProxy-Authorization: ${authorization(proxy)}\r\n\r\n`, + ) + }, + ) + let buffered = Buffer.alloc(0) + let settled = false + const onData = (chunk: Buffer) => { + buffered = Buffer.concat([buffered, chunk]) + const end = buffered.indexOf("\r\n\r\n") + if (end < 0) return + socket.off("data", onData) + const status = Number(buffered.toString("ascii", 0, end).split(" ", 2)[1]) + settled = true + resolve({ status, socket, leftover: buffered.subarray(end + 4) }) + } + socket.on("data", onData) + socket.once("error", reject) + socket.once("close", () => { + if (!settled) reject(new Error("Browser proxy connection closed before responding")) + }) + }) +} + +function readAll(socket: TLSSocket, initial: Buffer) { + return new Promise((resolve, reject) => { + const chunks = [initial] + socket.on("data", (chunk) => chunks.push(chunk)) + socket.on("end", () => resolve(Buffer.concat(chunks).toString())) + socket.once("error", reject) + }) +} + +function peerCertificate(proxy: ProxyInfo) { + return new Promise((resolve, reject) => { + const socket = connect( + { + host: "127.0.0.1", + port: proxy.port, + servername: proxy.host.endsWith(".localhost") ? proxy.host : undefined, + rejectUnauthorized: false, + }, + () => { + resolve(socket.getPeerCertificate().raw) + socket.end() + }, + ) + socket.once("error", reject) + }) +} + +async function waitFor(check: () => boolean) { + for (let attempt = 0; attempt < 100; attempt++) { + if (check()) return + await Bun.sleep(5) + } + throw new Error("Timed out waiting for browser proxy test condition") +} diff --git a/packages/client/test/node/tunnel.test.ts b/packages/client/test/node/tunnel.test.ts new file mode 100644 index 0000000000..d7336b3235 --- /dev/null +++ b/packages/client/test/node/tunnel.test.ts @@ -0,0 +1,208 @@ +import { BrowserTunnelProtocol } from "@opencode-ai/protocol/browser-tunnel" +import { BROWSER_TUNNEL_PROTOCOL } from "@opencode-ai/protocol/groups/browser" +import { Browser } from "@opencode-ai/schema/browser" +import { BrowserTunnel } from "@opencode-ai/schema/browser-tunnel" +import { Session } from "@opencode-ai/schema/session" +import { describe, expect, test } from "bun:test" +import { Effect } from "effect" +import { once } from "node:events" +import { createServer } from "node:http" +import WebSocket, { WebSocketServer } from "ws" +import { BrowserTunnelError, openBrowserTunnel } from "../../src/node/browser/tunnel" + +const sessionID = Session.ID.make("ses_client_tunnel") +const leaseID = Browser.LeaseID.make("brl_clienttunnel") + +describe("browser tunnel", () => { + test("enforces byte and frame windows and preserves both half-closes", async () => { + const server = await tunnelServer() + try { + const opening = openBrowserTunnel({ + endpoint: server.endpoint, + sessionID, + leaseID, + target: { host: BrowserTunnel.Host.make("target.example"), port: BrowserTunnel.Port.make(443) }, + }) + const socket = await server.connected + const next = frameReader(socket) + socket.send(BrowserTunnelProtocol.encodeFromServer({ type: "browser.tunnel.ready" }), { binary: true }) + const open = await next() + expect(open).toMatchObject({ + type: "control", + message: { + type: "browser.tunnel.open", + sessionID, + leaseID, + target: { host: "target.example", port: 443 }, + }, + }) + socket.send( + BrowserTunnelProtocol.encodeFromServer({ + type: "browser.tunnel.opened", + receiveWindow: BrowserTunnel.WindowSize.make(BrowserTunnelProtocol.InitialWindowBytes), + receiveFrames: BrowserTunnel.FrameWindow.make(BrowserTunnelProtocol.InitialFrameWindow), + }), + { binary: true }, + ) + const tunnel = await opening + + let writeSettled = false + const writing = new Promise((resolve, reject) => { + tunnel.write(Buffer.alloc(BrowserTunnelProtocol.InitialWindowBytes + 1, 7), (error) => { + writeSettled = true + if (error) reject(error) + else resolve() + }) + }) + const sent = [] + for (let index = 0; index < 4; index++) sent.push(await next()) + expect(sent.every((frame) => frame.type === "data" && frame.data.byteLength === 64 * 1_024)).toBe(true) + await Bun.sleep(25) + expect(writeSettled).toBe(false) + + socket.send( + BrowserTunnelProtocol.encodeFromServer({ + type: "browser.tunnel.window", + bytes: BrowserTunnel.WindowBytes.make(64 * 1_024), + frames: BrowserTunnel.FrameWindow.make(1), + }), + { binary: true }, + ) + const final = await next() + expect(final.type === "data" && final.data.byteLength).toBe(1) + await writing + + const received = once(tunnel, "data") + socket.send(BrowserTunnelProtocol.data(Buffer.from("from server")), { binary: true }) + expect(Buffer.from((await received)[0]).toString()).toBe("from server") + expect(await next()).toEqual({ + type: "control", + message: { + type: "browser.tunnel.window", + bytes: Buffer.byteLength("from server"), + frames: 1, + }, + }) + + const localEnd = next() + tunnel.end() + expect(await localEnd).toEqual({ type: "control", message: { type: "browser.tunnel.end" } }) + const remoteEnd = once(tunnel, "end") + socket.send(BrowserTunnelProtocol.encodeFromServer({ type: "browser.tunnel.end" }), { binary: true }) + await remoteEnd + } finally { + await server.close() + } + }) + + test("rejects window credit that does not exactly match complete frames", async () => { + const server = await tunnelServer() + try { + const opening = openBrowserTunnel({ + endpoint: server.endpoint, + sessionID, + leaseID, + target: { host: BrowserTunnel.Host.make("target.example"), port: BrowserTunnel.Port.make(80) }, + }) + const socket = await server.connected + const next = frameReader(socket) + socket.send(BrowserTunnelProtocol.encodeFromServer({ type: "browser.tunnel.ready" }), { binary: true }) + await next() + socket.send( + BrowserTunnelProtocol.encodeFromServer({ + type: "browser.tunnel.opened", + receiveWindow: BrowserTunnel.WindowSize.make(BrowserTunnelProtocol.InitialWindowBytes), + receiveFrames: BrowserTunnel.FrameWindow.make(BrowserTunnelProtocol.InitialFrameWindow), + }), + { binary: true }, + ) + const tunnel = await opening + await new Promise((resolve, reject) => + tunnel.write(Buffer.from("frame"), (error) => (error ? reject(error) : resolve())), + ) + await next() + const failed = once(tunnel, "error") + socket.send( + BrowserTunnelProtocol.encodeFromServer({ + type: "browser.tunnel.window", + bytes: BrowserTunnel.WindowBytes.make(4), + frames: BrowserTunnel.FrameWindow.make(1), + }), + { binary: true }, + ) + const error = (await failed)[0] + if (!(error instanceof BrowserTunnelError)) throw new Error("expected browser tunnel error") + expect(error.code).toBe("protocol_error") + expect(await next()).toEqual({ + type: "control", + message: { type: "browser.tunnel.reset", code: "protocol_error" }, + }) + } finally { + await server.close() + } + }) +}) + +async function tunnelServer() { + const authorization = `Basic ${Buffer.from("opencode:secret").toString("base64")}` + const http = createServer() + const webSockets = new WebSocketServer({ noServer: true }) + let resolveConnected!: (socket: WebSocket) => void + const connected = new Promise((resolve) => { + resolveConnected = resolve + }) + webSockets.once("connection", resolveConnected) + http.on("upgrade", (request, socket, head) => { + if ( + request.url !== "/api/browser/tunnel" || + request.headers.authorization !== authorization || + request.headers["sec-websocket-protocol"] !== BROWSER_TUNNEL_PROTOCOL + ) { + socket.end("HTTP/1.1 401 Unauthorized\r\n\r\n") + return + } + webSockets.handleUpgrade(request, socket, head, (webSocket) => webSockets.emit("connection", webSocket, request)) + }) + await new Promise((resolve) => http.listen(0, "127.0.0.1", resolve)) + const address = http.address() + if (address === null || typeof address === "string") throw new Error("browser tunnel server did not bind TCP") + return { + connected, + endpoint: { + url: `http://127.0.0.1:${address.port}`, + authorization, + }, + async close() { + webSockets.clients.forEach((socket) => socket.terminate()) + webSockets.close() + http.closeAllConnections() + http.close() + await Bun.sleep(10) + }, + } +} + +function frameReader(socket: WebSocket) { + const queued: Array<{ readonly data: WebSocket.RawData; readonly binary: boolean }> = [] + const waiting: Array<(message: { readonly data: WebSocket.RawData; readonly binary: boolean }) => void> = [] + socket.on("message", (data, binary) => { + const resolve = waiting.shift() + if (resolve) resolve({ data, binary }) + else queued.push({ data, binary }) + }) + return async () => { + const message = + queued.shift() ?? + (await new Promise<{ readonly data: WebSocket.RawData; readonly binary: boolean }>((resolve) => + waiting.push(resolve), + )) + if (!message.binary) throw new Error("Expected a binary browser tunnel frame") + return Effect.runPromise(BrowserTunnelProtocol.decodeFromDesktop(rawData(message.data))) + } +} + +function rawData(data: WebSocket.RawData) { + if (data instanceof ArrayBuffer) return new Uint8Array(data) + if (Array.isArray(data)) return new Uint8Array(Buffer.concat(data)) + return new Uint8Array(data.buffer, data.byteOffset, data.byteLength) +} diff --git a/packages/client/test/types/node-consumer.ts b/packages/client/test/types/node-consumer.ts new file mode 100644 index 0000000000..5ce08473c2 --- /dev/null +++ b/packages/client/test/types/node-consumer.ts @@ -0,0 +1,29 @@ +import { Browser, BrowserDriver, BrowserDriverError, OpenCode, type BrowserAttachment } from "@opencode-ai/client/node" + +const state: Browser.State = { + url: "about:blank", + title: "", + loading: false, + canGoBack: false, + canGoForward: false, + generation: 0, +} + +const factory: BrowserDriver<{ readonly proxyURL: string }> = (context) => ({ + resource: { proxyURL: context.proxy.url }, + state: () => state, + subscribe: () => () => undefined, + execute: async (_command, options) => { + throw new BrowserDriverError(options.signal.aborted ? "aborted" : "internal", "Command unavailable") + }, + dispose: () => undefined, +}) +const driver = BrowserDriver.define(factory) + +declare const client: ReturnType +const attachment: Promise> = client.browser.attach({ + sessionID: "ses_type_fixture", + driver, +}) + +void attachment diff --git a/packages/client/test/types/tsconfig.json b/packages/client/test/types/tsconfig.json new file mode 100644 index 0000000000..96a0b35cfa --- /dev/null +++ b/packages/client/test/types/tsconfig.json @@ -0,0 +1,8 @@ +{ + "$schema": "https://json.schemastore.org/tsconfig", + "extends": "../../tsconfig.json", + "compilerOptions": { + "noEmit": true + }, + "include": ["node-consumer.ts"] +} diff --git a/packages/client/test/types/tsconfig.nodenext.json b/packages/client/test/types/tsconfig.nodenext.json new file mode 100644 index 0000000000..aaa0942d6c --- /dev/null +++ b/packages/client/test/types/tsconfig.nodenext.json @@ -0,0 +1,10 @@ +{ + "$schema": "https://json.schemastore.org/tsconfig", + "extends": "../../tsconfig.json", + "compilerOptions": { + "module": "NodeNext", + "moduleResolution": "NodeNext", + "noEmit": true + }, + "include": ["node-consumer.ts"] +} diff --git a/packages/httpapi-codegen/src/index.ts b/packages/httpapi-codegen/src/index.ts index 90ce130d1d..d98b194233 100644 --- a/packages/httpapi-codegen/src/index.ts +++ b/packages/httpapi-codegen/src/index.ts @@ -351,7 +351,7 @@ export function emitPromise( { path: "index.ts", content: - 'export { ClientError, type ClientErrorReason } from "./client-error"\nexport * as OpenCode from "./client"\nexport * from "./types"\n', + 'export { ClientError, type ClientErrorReason } from "./client-error.js"\nexport * as OpenCode from "./client.js"\nexport * from "./types.js"\n', }, ], } @@ -371,7 +371,9 @@ function renderEffectShape( .map((field) => { const schema = effectInputSchema(endpoint, field) if (schema === undefined) { - throw new GenerationError({ reason: `Missing Effect input schema: ${endpoint.group}.${endpoint.endpoint.identifier}.${field.name}` }) + throw new GenerationError({ + reason: `Missing Effect input schema: ${endpoint.group}.${endpoint.endpoint.identifier}.${field.name}`, + }) } return `readonly ${JSON.stringify(field.name)}${field.optional ? "?" : ""}: ${effectType(schema, references, imports)}` }) @@ -453,11 +455,7 @@ function effectTypeReferences(input: ReadonlyArray) { return { names, asts, brands } } -function effectType( - schema: Schema.Top, - references: ReturnType, - imports: Set, -) { +function effectType(schema: Schema.Top, references: ReturnType, imports: Set) { const projected = Schema.toType(schema) const direct = references.asts.get(schema.ast) ?? references.asts.get(projected.ast) if (direct !== undefined) { @@ -595,7 +593,7 @@ function renderEffectFiles(groups: ReadonlyArray): Output["files"] { { path: "client.ts", content: renderClient(groups) }, { path: "index.ts", - content: 'export { ClientError } from "./client-error"\nexport * as OpenCode from "./client"\n', + content: 'export { ClientError } from "./client-error.js"\nexport * as OpenCode from "./client.js"\n', }, ] } @@ -713,7 +711,7 @@ function renderImportedEffectFiles( { path: "client.ts", content: client }, { path: "index.ts", - content: 'export { ClientError } from "./client-error"\nexport * as OpenCode from "./client"\n', + content: 'export { ClientError } from "./client-error.js"\nexport * as OpenCode from "./client.js"\n', }, ] } @@ -1313,7 +1311,13 @@ export function write( output.files, (file) => Effect.tryPromise({ - try: () => format(file.content, { filepath: file.path, parser: "typescript", semi: false, printWidth: 120 }), + try: () => + format(withJsExtensions(file.content), { + filepath: file.path, + parser: "typescript", + semi: false, + printWidth: 120, + }), catch: (error) => new GenerationError({ reason: `Failed to format ${file.path}: ${String(error)}` }), }).pipe(Effect.flatMap((content) => fs.writeFileString(join(directory, file.path), content))), { concurrency: 8, discard: true }, @@ -1322,6 +1326,16 @@ export function write( }) } +function withJsExtensions(content: string) { + return content + .replaceAll(/(from ")(\.\.?\/[^"\n]+)(")/g, (_match, start: string, path: string, end: string) => + /\.[cm]?[jt]sx?$/.test(path) ? `${start}${path}${end}` : `${start}${path}.js${end}`, + ) + .replaceAll(/(import\(")(\.\.?\/[^"\n]+)("\))/g, (_match, start: string, path: string, end: string) => + /\.[cm]?[jt]sx?$/.test(path) ? `${start}${path}${end}` : `${start}${path}.js${end}`, + ) +} + function isSafeOutputPath(path: string) { return path !== manifestName && !isAbsolute(path) && path !== "." && path !== ".." && !/[\\/]/.test(path) } diff --git a/packages/httpapi-codegen/test/generated/index.ts b/packages/httpapi-codegen/test/generated/index.ts index bc0dbc9fa4..0589939b46 100644 --- a/packages/httpapi-codegen/test/generated/index.ts +++ b/packages/httpapi-codegen/test/generated/index.ts @@ -1,2 +1,2 @@ -export { ClientError } from "./client-error" -export * as OpenCode from "./client" +export { ClientError } from "./client-error.js" +export * as OpenCode from "./client.js" diff --git a/packages/httpapi-codegen/test/write.test.ts b/packages/httpapi-codegen/test/write.test.ts index f0704abea2..26d9811228 100644 --- a/packages/httpapi-codegen/test/write.test.ts +++ b/packages/httpapi-codegen/test/write.test.ts @@ -8,14 +8,14 @@ describe("HttpApiCodegen.write", () => { const writes: Array<{ readonly path: string; readonly content: string }> = [] const output: Output = { operations: [], - files: [{ path: "session.ts", content: "export const session = {}" }], + files: [{ path: "session.ts", content: 'export { session } from "../../session-data"' }], } return Effect.gen(function* () { yield* write(output, "/generated") expect(writes).toEqual([ - { path: "/generated/session.ts", content: "export const session = {}\n" }, + { path: "/generated/session.ts", content: 'export { session } from "../../session-data.js"\n' }, { path: "/generated/.httpapi-codegen.json", content: '[\n "session.ts"\n]\n' }, ]) }).pipe(