diff --git a/packages/protocol/src/api.ts b/packages/protocol/src/api.ts index 2249797df4..7cf396ea7f 100644 --- a/packages/protocol/src/api.ts +++ b/packages/protocol/src/api.ts @@ -14,6 +14,7 @@ import { SkillGroup } from "./groups/skill.js" import { EventGroup, makeEventGroup } from "./groups/event.js" import type { Definition } from "@opencode-ai/schema/event" import { AgentGroup } from "./groups/agent.js" +import { BrowserGroup } from "./groups/browser.js" import { PluginGroup } from "./groups/plugin.js" import { HealthGroup } from "./groups/health.js" import { ServerGroup } from "./groups/server.js" @@ -85,6 +86,7 @@ type ApiGroups< | typeof HealthGroup | typeof ServerGroup | typeof DebugGroup + | typeof BrowserGroup | LocationGroups | FormGroups | SessionGroups @@ -146,6 +148,7 @@ const makeApiFromGroup = < HttpApi.make("server") .add(HealthGroup) .add(ServerGroup) + .add(BrowserGroup) .add(LocationGroup.middleware(locationMiddleware)) .add(AgentGroup.middleware(locationMiddleware)) .add(PluginGroup.middleware(locationMiddleware)) diff --git a/packages/protocol/src/client.ts b/packages/protocol/src/client.ts index 24e68e7a51..ea04d01185 100644 --- a/packages/protocol/src/client.ts +++ b/packages/protocol/src/client.ts @@ -38,6 +38,7 @@ export const groupNames = { "server.debug": "debug", "server.location": "location", "server.agent": "agent", + "server.browser": "browser", "server.plugin": "plugin", "server.session": "session", "server.message": "message", @@ -63,5 +64,16 @@ export const groupNames = { "server.vcs": "vcs", } as const -export const promiseOmitEndpoints = new Set(["pty.connect", "pty.connectToken"]) -export const effectOmitEndpoints = new Set(["fs.read", "pty.connect", "pty.connectToken"]) +export const promiseOmitEndpoints = new Set([ + "browser.control.connect", + "browser.tunnel.connect", + "pty.connect", + "pty.connectToken", +]) +export const effectOmitEndpoints = new Set([ + "browser.control.connect", + "browser.tunnel.connect", + "fs.read", + "pty.connect", + "pty.connectToken", +]) diff --git a/packages/server/src/browser-control-connection.ts b/packages/server/src/browser-control-connection.ts new file mode 100644 index 0000000000..fcf8cf1f43 --- /dev/null +++ b/packages/server/src/browser-control-connection.ts @@ -0,0 +1,115 @@ +export * as BrowserControlConnection from "./browser-control-connection" + +import { BrowserHost } from "@opencode-ai/core/browser-host" +import { BrowserControlProtocol } from "@opencode-ai/protocol/browser-control" +import { Browser } from "@opencode-ai/schema/browser" +import { BrowserControl } from "@opencode-ai/schema/browser-control" +import { Session } from "@opencode-ai/schema/session" +import { Deferred, Effect } from "effect" +import { Socket } from "effect/unstable/socket" + +const registrations = new Map() + +export function isAttached(sessionID: Session.ID, leaseID: Browser.LeaseID) { + return registrations.get(sessionID)?.leaseID === leaseID +} + +export const run = Effect.fn("BrowserControlConnection.run")(function* ( + socket: Socket.Socket, + opened: Effect.Effect = Effect.void, +) { + const browser = yield* BrowserHost.Service + const write = yield* socket.writer + const pending = new Map>() + const token = {} + let sessionID: Session.ID | undefined + let controller: BrowserHost.Controller | undefined + + const send = (message: BrowserControl.FromServer) => + Effect.try({ + try: () => BrowserControlProtocol.encodeFromServer(message), + catch: () => + new BrowserHost.RequestError({ code: "protocol", message: "Failed to encode browser control message." }), + }).pipe( + Effect.flatMap(write), + Effect.mapError( + () => new BrowserHost.RequestError({ code: "internal", message: "Browser control connection failed." }), + ), + ) + + const peer: BrowserHost.Peer = { + open: send({ type: "browser.control.open" }), + request: (command, leaseID) => + Effect.gen(function* () { + const requestID = BrowserControl.RequestID.create() + const done = yield* Deferred.make() + pending.set(requestID, done) + yield* send({ type: "browser.control.request", requestID, leaseID, command }) + const outcome = yield* Deferred.await(done).pipe( + Effect.onInterrupt(() => send({ type: "browser.control.cancel", requestID, leaseID }).pipe(Effect.ignore)), + Effect.ensuring(Effect.sync(() => pending.delete(requestID))), + ) + if (outcome.type === "failure") return yield* new BrowserHost.RequestError(outcome) + return outcome.result + }), + } + + yield* Effect.addFinalizer(() => + Effect.sync(() => { + if (sessionID && registrations.get(sessionID)?.token === token) registrations.delete(sessionID) + pending.forEach((done) => + Deferred.doneUnsafe( + done, + Effect.succeed({ type: "failure", code: "not_attached", message: "Browser control connection closed." }), + ), + ) + pending.clear() + }), + ) + + const receive = Effect.fnUntraced(function* (raw: string | Uint8Array) { + const message = yield* BrowserControlProtocol.decodeFromClient(raw) + if (!controller) { + if (message.type !== "browser.control.register") + return yield* Effect.fail(new Error("Expected browser registration.")) + sessionID = message.sessionID + controller = yield* browser.register(message.sessionID, peer) + registrations.set(message.sessionID, { token }) + yield* send({ type: "browser.control.registered" }) + return + } + if (!sessionID || message.type === "browser.control.register") { + return yield* Effect.fail(new Error("Browser control connection is already registered.")) + } + if (message.type === "browser.control.attach") { + yield* controller.attach(message.leaseID, message.state) + registrations.set(sessionID, { token, leaseID: message.leaseID }) + yield* send({ type: "browser.control.attached", leaseID: message.leaseID }) + return + } + if (message.type === "browser.control.state") { + yield* controller.state(message.leaseID, message.state) + return + } + if (message.type === "browser.control.detach") { + yield* controller.detach(message.leaseID) + registrations.set(sessionID, { token }) + return + } + const done = pending.get(message.requestID) + if (!done || registrations.get(sessionID)?.leaseID !== message.leaseID) { + return yield* Effect.fail(new Error("Browser response does not match a pending request.")) + } + Deferred.doneUnsafe(done, Effect.succeed(message.outcome)) + }) + + yield* socket.runRaw(receive, { onOpen: opened }).pipe( + Effect.catchCause((cause) => + write(new Socket.CloseEvent(1002, "Invalid browser control message")).pipe( + Effect.timeoutOrElse({ duration: "1 second", orElse: () => Effect.void }), + Effect.catch(() => Effect.void), + Effect.andThen(Effect.logDebug("Browser control connection closed", { cause })), + ), + ), + ) +}) diff --git a/packages/server/src/browser-tunnel.ts b/packages/server/src/browser-tunnel.ts new file mode 100644 index 0000000000..f470c90cb2 --- /dev/null +++ b/packages/server/src/browser-tunnel.ts @@ -0,0 +1,292 @@ +export * as BrowserTunnelServer from "./browser-tunnel" + +import { BrowserHost } from "@opencode-ai/core/browser-host" +import { BrowserTunnelProtocol } from "@opencode-ai/protocol/browser-tunnel" +import { BrowserTunnel } from "@opencode-ai/schema/browser-tunnel" +import { + Cause, + Context, + Effect, + Fiber, + Layer, + Option, + Queue, + Ref, + Result, + Schema, + Scope, + SynchronizedRef, +} from "effect" +import { Socket } from "effect/unstable/socket" +import { BrowserControlConnection } from "./browser-control-connection" + +const ActiveLimit = 64 + +export class CapacityError extends Schema.TaggedErrorClass()("BrowserTunnel.CapacityError", { + limit: Schema.Int, + message: Schema.String, +}) {} + +class TunnelError extends Schema.TaggedErrorClass()("BrowserTunnel.TunnelError", { + kind: Schema.Literals(["closed", "protocol", "target", "revoked"]), + message: Schema.String, + cause: Schema.optional(Schema.Defect()), +}) {} + +class ConnectError extends Schema.TaggedErrorClass()("BrowserTunnel.ConnectError", { + kind: Schema.Literals(["failed", "timeout"]), + message: Schema.String, + cause: Schema.optional(Schema.Defect()), +}) {} + +type Dial = (host: string, port: number) => Effect.Effect +type State = { readonly active: number; readonly shutdown: boolean } + +export interface Connection { + readonly run: (socket: Socket.Socket, opened?: Effect.Effect) => Effect.Effect +} + +export interface Interface { + readonly acquire: Effect.Effect + readonly shutdown: Effect.Effect +} + +export class Service extends Context.Service()("@opencode/server/BrowserTunnel") {} + +export function make(dial: Dial = connect) { + return Effect.gen(function* () { + const browser = yield* BrowserHost.Service + const state = yield* SynchronizedRef.make({ active: 0, shutdown: false }) + const connections = new Set>() + const shutdown = Effect.fn("BrowserTunnel.shutdown")(function* () { + const close = yield* SynchronizedRef.modify(state, (current) => [ + !current.shutdown, + { ...current, shutdown: true }, + ]) + if (close) yield* Effect.all(connections, { concurrency: "unbounded", discard: true }) + }) + yield* Effect.addFinalizer(shutdown) + + const acquire: Interface["acquire"] = Effect.acquireRelease( + SynchronizedRef.modifyEffect( + state, + Effect.fnUntraced(function* (current) { + if (current.shutdown || current.active >= ActiveLimit) { + return yield* new CapacityError({ limit: ActiveLimit, message: "Browser tunnel capacity is unavailable." }) + } + return [undefined, { ...current, active: current.active + 1 }] as const + }), + ), + () => SynchronizedRef.update(state, (current) => ({ ...current, active: Math.max(0, current.active - 1) })), + ).pipe( + Effect.andThen(Ref.make(false)), + Effect.map((started) => ({ + run: (socket: Socket.Socket, opened = Effect.void) => + Effect.gen(function* () { + const write = yield* socket.writer + if (yield* Ref.getAndSet(started, true)) return + const restart = close(write, 1012, "Server restarting") + connections.add(restart) + yield* serve(browser, socket, write, dial, opened).pipe( + Effect.catch(() => Effect.void), + Effect.ensuring(Effect.sync(() => connections.delete(restart))), + ) + }), + })), + ) + return Service.of({ acquire, shutdown: shutdown() }) + }) +} + +export const layer = Layer.effect(Service, make()) + +const serve = Effect.fn("BrowserTunnel.serve")(function* ( + browser: BrowserHost.Interface, + socket: Socket.Socket, + writeSocket: (data: string | Uint8Array | Socket.CloseEvent) => Effect.Effect, + dial: Dial, + opened: Effect.Effect, +) { + const inbound = yield* Queue.bounded(16) + const reader = yield* socket + .runRaw( + (message) => { + if (typeof message !== "string" && message.byteLength > BrowserTunnelProtocol.MaxFrameBytes) { + return fail(inbound, new TunnelError({ kind: "protocol", message: "Browser tunnel frame is too large." })) + } + return Queue.offer(inbound, message).pipe(Effect.asVoid) + }, + { onOpen: opened }, + ) + .pipe( + Effect.onExit(() => fail(inbound, new TunnelError({ kind: "closed", message: "Browser tunnel closed." }))), + Effect.forkScoped, + ) + + const first = yield* Queue.take(inbound).pipe( + Effect.timeoutOrElse({ + duration: "5 seconds", + orElse: () => Effect.fail(new TunnelError({ kind: "protocol", message: "Browser tunnel open timed out." })), + }), + Effect.flatMap(BrowserTunnelProtocol.decodeFromClient), + Effect.mapError(() => new TunnelError({ kind: "protocol", message: "Browser tunnel open message is invalid." })), + Effect.result, + ) + if (Result.isFailure(first)) { + yield* reject(writeSocket, "invalid_open", first.failure.message) + return + } + const input = first.success + const capability = yield* browser.get(input.sessionID) + if (Option.isNone(capability) || capability.value.type !== "attached") { + yield* reject(writeSocket, "not_attached", "No browser is attached to this Session.") + return + } + if (!BrowserControlConnection.isAttached(input.sessionID, input.leaseID)) { + yield* reject(writeSocket, "stale_lease", "The browser attachment lease is stale.") + return + } + + const target = yield* Effect.result( + Effect.raceFirst( + dial(input.target.host, input.target.port), + Effect.raceFirst( + Fiber.join(reader).pipe(Effect.andThen(new TunnelError({ kind: "closed", message: "Browser tunnel closed." }))), + capability.value.revoked.pipe( + Effect.andThen(new TunnelError({ kind: "revoked", message: "Browser lease was revoked." })), + ), + ), + ), + ) + if (Result.isFailure(target)) { + if (target.failure instanceof ConnectError) { + yield* reject( + writeSocket, + target.failure.kind === "timeout" ? "connect_timeout" : "connect_failed", + target.failure.message, + ) + } + return + } + const tcp = target.success + yield* Effect.addFinalizer(() => Effect.sync(() => tcp.destroy())) + yield* writeSocket(BrowserTunnelProtocol.encodeFromServer({ type: "browser.tunnel.opened" })) + + const output = yield* Queue.bounded(1) + const onData = (data: Buffer) => { + tcp.pause() + Queue.offerUnsafe(output, data) + } + const onClose = () => + Queue.failCauseUnsafe(output, Cause.fail(new TunnelError({ kind: "closed", message: "Target closed." }))) + const onError = (cause: Error) => + Queue.failCauseUnsafe(output, Cause.fail(new TunnelError({ kind: "target", message: "Target failed.", cause }))) + tcp.on("data", onData) + tcp.once("close", onClose) + tcp.once("error", onError) + yield* Effect.addFinalizer(() => + Effect.sync(() => { + tcp.off("data", onData) + tcp.off("close", onClose) + tcp.off("error", onError) + }).pipe(Effect.andThen(Queue.shutdown(output))), + ) + + const fromClient = Effect.forever( + Queue.take(inbound).pipe( + Effect.flatMap((message) => + typeof message === "string" + ? new TunnelError({ kind: "protocol", message: "Tunnel payloads must be binary." }) + : writeTarget(tcp, message), + ), + ), + ) + const fromTarget = Effect.forever( + Queue.take(output).pipe( + Effect.flatMap((data) => + Effect.forEach( + Array.from({ length: Math.ceil(data.byteLength / BrowserTunnelProtocol.MaxFrameBytes) }, (_, index) => + data.subarray( + index * BrowserTunnelProtocol.MaxFrameBytes, + (index + 1) * BrowserTunnelProtocol.MaxFrameBytes, + ), + ), + writeSocket, + { discard: true }, + ), + ), + Effect.ensuring(Effect.sync(() => tcp.resume())), + ), + ) + yield* Effect.raceFirst( + Effect.all([fromClient, fromTarget], { concurrency: "unbounded", discard: true }), + Effect.raceFirst(Fiber.join(reader), capability.value.revoked), + ).pipe(Effect.ensuring(close(writeSocket, 1000, "Browser tunnel closed"))) +}) + +function connect(host: string, port: number) { + return Effect.gen(function* () { + const net = yield* Effect.promise(() => import("node:net")) + return yield* Effect.acquireRelease( + Effect.callback((resume) => { + const socket = new net.Socket() + const onError = (cause: Error) => + resume( + Effect.fail( + new ConnectError({ kind: "failed", message: "Failed to connect browser tunnel target.", cause }), + ), + ) + socket.once("error", onError) + socket.connect(port, host, () => { + socket.off("error", onError) + socket.setNoDelay(true) + resume(Effect.succeed(socket)) + }) + return Effect.sync(() => socket.destroy()) + }).pipe( + Effect.timeoutOrElse({ + duration: "10 seconds", + orElse: () => + Effect.fail(new ConnectError({ kind: "timeout", message: "Browser tunnel target connection timed out." })), + }), + ), + (socket) => Effect.sync(() => socket.destroy()), + ) + }) +} + +function writeTarget(socket: import("node:net").Socket, data: Uint8Array) { + return Effect.callback((resume) => { + socket.write(data, (cause) => + resume( + cause ? Effect.fail(new TunnelError({ kind: "target", message: "Target write failed.", cause })) : Effect.void, + ), + ) + }) +} + +function reject( + write: (data: string | Uint8Array | Socket.CloseEvent) => Effect.Effect, + code: BrowserTunnel.OpenErrorCode, + message: string, +) { + return write(BrowserTunnelProtocol.encodeFromServer({ type: "browser.tunnel.rejected", code, message })).pipe( + Effect.catch(() => Effect.void), + Effect.andThen(close(write, 1000, message)), + ) +} + +function close( + write: (data: string | Uint8Array | Socket.CloseEvent) => Effect.Effect, + code: number, + reason: string, +) { + return write(new Socket.CloseEvent(code, reason.slice(0, 123))).pipe( + Effect.timeoutOrElse({ duration: "1 second", orElse: () => Effect.void }), + Effect.catch(() => Effect.void), + ) +} + +function fail(queue: Queue.Queue, error: TunnelError) { + return Effect.sync(() => Queue.failCauseUnsafe(queue, Cause.fail(error))) +} diff --git a/packages/server/src/handlers.ts b/packages/server/src/handlers.ts index a0126ea89e..eb2a6f2532 100644 --- a/packages/server/src/handlers.ts +++ b/packages/server/src/handlers.ts @@ -11,6 +11,7 @@ import { CommandHandler } from "./handlers/command" import { SkillHandler } from "./handlers/skill" import { EventHandler } from "./handlers/event" import { AgentHandler } from "./handlers/agent" +import { BrowserHandler } from "./handlers/browser" import { PluginHandler } from "./handlers/plugin" import { HealthHandler } from "./handlers/health" import { ServerHandler } from "./handlers/server" @@ -35,6 +36,7 @@ export const handlers = Layer.mergeAll( DebugHandler, LocationHandler, AgentHandler, + BrowserHandler, PluginHandler, SessionHandler, MessageHandler, diff --git a/packages/server/src/handlers/browser.ts b/packages/server/src/handlers/browser.ts new file mode 100644 index 0000000000..3280a81094 --- /dev/null +++ b/packages/server/src/handlers/browser.ts @@ -0,0 +1,72 @@ +import { NodeHttpServerRequest } from "@effect/platform-node" +import { BrowserControlProtocol } from "@opencode-ai/protocol/browser-control" +import { BrowserTunnelProtocol } from "@opencode-ai/protocol/browser-tunnel" +import { ServiceUnavailableError } from "@opencode-ai/protocol/errors" +import { Effect } from "effect" +import { HttpServerRequest, HttpServerResponse } from "effect/unstable/http" +import { HttpApiBuilder } from "effect/unstable/httpapi" +import { ServerResponse } from "node:http" +import { Api } from "../api" +import { BrowserControlConnection } from "../browser-control-connection" +import { BrowserTunnelServer } from "../browser-tunnel" +import { CorsConfig, isAllowedRequestOrigin, type CorsOptions } from "../cors" + +export const BrowserHandler = HttpApiBuilder.group(Api, "server.browser", (handlers) => + Effect.gen(function* () { + const tunnels = yield* BrowserTunnelServer.Service + const cors = yield* CorsConfig + + return handlers + .handleRaw( + "browser.control.connect", + Effect.fn("BrowserHandler.control")(function* (ctx) { + const rejected = rejectUpgrade(ctx.request.headers, BrowserControlProtocol.Subprotocol, cors) + if (rejected) return rejected + const socket = yield* Effect.orDie(ctx.request.upgrade) + yield* BrowserControlConnection.run( + socket, + Effect.sync(() => markUpgraded(ctx.request)), + ) + return HttpServerResponse.empty() + }), + ) + .handleRaw( + "browser.tunnel.connect", + Effect.fn("BrowserHandler.tunnel")(function* (ctx) { + const rejected = rejectUpgrade(ctx.request.headers, BrowserTunnelProtocol.Subprotocol, cors) + if (rejected) return rejected + const connection = yield* tunnels.acquire.pipe( + Effect.mapError((error) => new ServiceUnavailableError({ service: "browser", message: error.message })), + ) + const socket = yield* Effect.orDie(ctx.request.upgrade) + yield* connection.run( + socket, + Effect.sync(() => markUpgraded(ctx.request)), + ) + return HttpServerResponse.empty() + }), + ) + }), +) + +function markUpgraded(request: HttpServerRequest.HttpServerRequest) { + const socket = NodeHttpServerRequest.toIncomingMessage(request).socket + // Bun leaves its HTTP handshake response assigned after ws takes ownership. Detaching + // matches Node's post-upgrade socket state and lets Effect complete the raw handler normally. + const response = Reflect.get(socket, "_httpMessage") + if (response instanceof ServerResponse) response.detachSocket(socket) +} + +function rejectUpgrade( + headers: Readonly>, + protocol: string, + cors: CorsOptions | undefined, +) { + if (!isAllowedRequestOrigin(headers.origin, headers.host, cors)) { + return HttpServerResponse.empty({ status: 403 }) + } + if (headers["sec-websocket-protocol"]?.split(",", 1)[0]?.trim() !== protocol) { + return HttpServerResponse.empty({ status: 426, headers: { "sec-websocket-protocol": protocol } }) + } + return undefined +} diff --git a/packages/server/src/middleware/authorization.ts b/packages/server/src/middleware/authorization.ts index a004fa973d..16c4910af5 100644 --- a/packages/server/src/middleware/authorization.ts +++ b/packages/server/src/middleware/authorization.ts @@ -1,9 +1,9 @@ import { ServerAuth } from "../auth" import { UnauthorizedError } from "@opencode-ai/protocol/errors" -import { Authorization } from "@opencode-ai/protocol/middleware/authorization" +import { Authorization, HeaderOnlyAuthorization } from "@opencode-ai/protocol/middleware/authorization" export { Authorization } from "@opencode-ai/protocol/middleware/authorization" import { hasPtyConnectTicketURL } from "@opencode-ai/protocol/groups/pty" -import { Effect, Encoding, Layer, Redacted } from "effect" +import { Context, Effect, Encoding, Layer, Redacted } from "effect" import { HttpEffect, HttpServerRequest, HttpServerResponse } from "effect/unstable/http" const AUTH_TOKEN_QUERY = "auth_token" @@ -26,9 +26,9 @@ function decodeCredential(input: string) { ) } -function credentialFromRequest(request: HttpServerRequest.HttpServerRequest) { - const url = new URL(request.url, "http://localhost") - const token = url.searchParams.get(AUTH_TOKEN_QUERY) +function credentialFromRequest(request: HttpServerRequest.HttpServerRequest, headerOnly = false) { + const url = new URL(request.url, "http://opencode.invalid") + const token = headerOnly ? undefined : url.searchParams.get(AUTH_TOKEN_QUERY) if (token) return decodeCredential(token) const match = /^Basic\s+(.+)$/i.exec(request.headers.authorization ?? "") if (match) return decodeCredential(match[1]) @@ -44,13 +44,17 @@ export const authorizationLayer = Layer.effect( Effect.gen(function* () { const config = yield* ServerAuth.Config if (!ServerAuth.required(config)) return Authorization.of((effect) => effect) - return Authorization.of((effect) => + return Authorization.of((effect, options) => Effect.gen(function* () { const request = yield* HttpServerRequest.HttpServerRequest // Browsers cannot set headers on WebSocket upgrades, so a ticketed PTY connect skips // credential checks here; the connect handler consumes and validates the ticket. - if (hasPtyConnectTicketURL(new URL(request.url, "http://localhost"))) return yield* effect - if (yield* authorizedRequest(request, config)) return yield* effect + if (hasPtyConnectTicketURL(new URL(request.url, "http://opencode.invalid"))) return yield* effect + const headerOnly = Context.get(options.endpoint.annotations, HeaderOnlyAuthorization) + const authorized = yield* credentialFromRequest(request, headerOnly).pipe( + Effect.map((credential) => ServerAuth.authorized(credential, config)), + ) + if (authorized) return yield* effect yield* HttpEffect.appendPreResponseHandler((_request, response) => Effect.succeed(HttpServerResponse.setHeader(response, "www-authenticate", WWW_AUTHENTICATE)), ) diff --git a/packages/server/src/process.ts b/packages/server/src/process.ts index 483745b404..7d0dc0a15f 100644 --- a/packages/server/src/process.ts +++ b/packages/server/src/process.ts @@ -130,9 +130,39 @@ function bind(hostname: string, port: number) { const parentScope = yield* Scope.Scope const serverScope = yield* Scope.fork(parentScope) const server = createServer() + const sockets = new Set() + const onConnection = (socket: import("node:net").Socket) => { + sockets.add(socket) + socket.once("close", () => sockets.delete(socket)) + } + const onUpgrade = (_request: unknown, socket: import("node:net").Socket) => sockets.add(socket) + server.on("connection", onConnection) + server.on("upgrade", onUpgrade) return yield* Effect.gen(function* () { const http = yield* NodeHttpServer.make(() => server, { port, host: hostname }) yield* Effect.addFinalizer(() => Effect.sync(() => server.closeAllConnections())) + // Node's closeAllConnections deliberately excludes upgraded sockets. + yield* Effect.addFinalizer(() => + Effect.sync(() => { + server.off("connection", onConnection) + server.off("upgrade", onUpgrade) + server.closeAllConnections() + }).pipe( + Effect.andThen( + Effect.suspend(() => { + if (sockets.size === 0) return Effect.void + return Effect.sleep("1 second").pipe( + Effect.andThen( + Effect.sync(() => { + for (const socket of sockets) socket.destroy() + sockets.clear() + }), + ), + ) + }), + ), + ), + ) return { http, server, scope: serverScope } }).pipe( Effect.provideService(Scope.Scope, serverScope), @@ -241,7 +271,8 @@ function unavailable(status: Status.State) { /** * The managed server owns restart continuity: it resumes Sessions the previous server suspended and * suspends its own active Sessions on graceful shutdown. Suspension runs while the drains are still - * alive: connections close first, this finalizer runs next, and Session execution teardown follows. + * alive: request admission stops first, application-owned transports receive their shutdown signal, + * listener connections close, and this finalizer runs during application teardown. */ const installRestartContinuity = Effect.fnUntraced(function* (restart: SessionRestart.Interface) { yield* Effect.forkScoped(restart.resumeSuspendedSessions) diff --git a/packages/server/src/routes.ts b/packages/server/src/routes.ts index f2d0d68df7..76202037ed 100644 --- a/packages/server/src/routes.ts +++ b/packages/server/src/routes.ts @@ -4,6 +4,7 @@ import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { httpClient } from "@opencode-ai/util/effect/app-node-platform" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" import { Bus } from "@opencode-ai/core/bus" +import { BrowserHost } from "@opencode-ai/core/browser-host" import { EventLogger } from "@opencode-ai/core/event-logger" import { FileSystemSearch } from "@opencode-ai/core/filesystem/search" import { Observability } from "@opencode-ai/util/observability" @@ -40,11 +41,13 @@ import { layer } from "./location" import { formLocationLayer } from "./middleware/form-location" import { sessionLocationLayer } from "./middleware/session-location" import { ServerInfo } from "./server-info" +import { BrowserTunnelServer } from "./browser-tunnel" import type { ServerOptions } from "./options" const applicationServices = LayerNode.group([ Database.node, Bus.node, + BrowserHost.node, EventLogger.node, httpClient, Job.node, @@ -131,8 +134,11 @@ function makeRoutes( return serviceLayer.pipe( Layer.flatMap((context) => { const services = Layer.succeedContext(context) + const browserTunnel = BrowserTunnelServer.layer.pipe(Layer.provide(services)) const requestServices = Layer.merge( - Layer.succeedContext(Context.pick(PermissionSaved.Service, Project.Service, WellKnown.Service)(context)), + Layer.succeedContext( + Context.pick(BrowserHost.Service, PermissionSaved.Service, Project.Service, WellKnown.Service)(context), + ), ServerInfo.layer(serviceURLs, options.app), ) return HttpApiBuilder.layer(Api, { openapiPath: "/openapi.json" }).pipe( @@ -144,6 +150,7 @@ function makeRoutes( Layer.provide(schemaErrorLayer), Layer.provide(auth), HttpRouter.provideRequest(requestServices), + Layer.provideMerge(browserTunnel), Layer.provideMerge(services), Layer.provideMerge(HttpRouter.layer), ) diff --git a/packages/server/test/browser.test.ts b/packages/server/test/browser.test.ts new file mode 100644 index 0000000000..a66c8831f9 --- /dev/null +++ b/packages/server/test/browser.test.ts @@ -0,0 +1,120 @@ +import { BrowserHost } from "@opencode-ai/core/browser-host" +import { BrowserControlProtocol } from "@opencode-ai/protocol/browser-control" +import { BrowserTunnelProtocol } from "@opencode-ai/protocol/browser-tunnel" +import { Browser } from "@opencode-ai/schema/browser" +import { BrowserTunnel } from "@opencode-ai/schema/browser-tunnel" +import { Session } from "@opencode-ai/schema/session" +import { expect } from "bun:test" +import { Effect, Fiber, Queue } from "effect" +import { Socket } from "effect/unstable/socket" +import { createServer } from "node:net" +import { it } from "../../core/test/lib/effect" +import { BrowserControlConnection } from "../src/browser-control-connection" +import { BrowserTunnelServer } from "../src/browser-tunnel" + +const sessionID = Session.ID.make("ses_browser_server") +const leaseID = Browser.LeaseID.make("brl_browserserver") +const state: Browser.State = { + url: "http://localhost/", + title: "Local", + loading: false, + canGoBack: false, + canGoForward: false, + generation: 1, +} +const end = Symbol("end") + +it.live("registers and attaches with the real host before dialing remote TCP", () => + Effect.scoped( + Effect.gen(function* () { + const browser = yield* BrowserHost.make(() => Effect.succeed(true)) + const control = yield* makeSocket + const controlFiber = yield* BrowserControlConnection.run(control.socket).pipe( + Effect.provideService(BrowserHost.Service, browser), + Effect.forkChild, + ) + yield* Queue.offer( + control.inbound, + BrowserControlProtocol.encodeFromClient({ type: "browser.control.register", sessionID }), + ) + expect(yield* controlMessage(control)).toEqual({ type: "browser.control.registered" }) + yield* Queue.offer( + control.inbound, + BrowserControlProtocol.encodeFromClient({ type: "browser.control.attach", leaseID, state }), + ) + expect(yield* controlMessage(control)).toEqual({ type: "browser.control.attached", leaseID }) + + const target = yield* echoServer + const address = target.address() + if (!address || typeof address === "string") throw new Error("echo server did not bind") + const tunnels = yield* BrowserTunnelServer.make().pipe(Effect.provideService(BrowserHost.Service, browser)) + const connection = yield* tunnels.acquire + const transport = yield* makeSocket + const running = yield* connection.run(transport.socket).pipe(Effect.forkChild) + yield* Queue.offer( + transport.inbound, + BrowserTunnelProtocol.encodeFromClient({ + type: "browser.tunnel.open", + sessionID, + leaseID, + target: { host: BrowserTunnel.Host.make("127.0.0.1"), port: BrowserTunnel.Port.make(address.port) }, + }), + ) + const opened = yield* Queue.take(transport.outbound) + if (typeof opened !== "string") throw new Error("expected text tunnel handshake") + expect(yield* BrowserTunnelProtocol.decodeFromServer(opened)).toEqual({ type: "browser.tunnel.opened" }) + + yield* Queue.offer(transport.inbound, Buffer.from("through server")) + const echoed = yield* Queue.take(transport.outbound) + if (!(echoed instanceof Uint8Array)) throw new Error("expected raw tunnel bytes") + expect(Buffer.from(echoed).toString()).toBe("through server") + + yield* Queue.offer(transport.inbound, end) + yield* Fiber.join(running) + yield* Queue.offer(control.inbound, end) + yield* Fiber.join(controlFiber) + }), + ), +) + +const makeSocket = Effect.gen(function* () { + const inbound = yield* Queue.unbounded() + const outbound = yield* Queue.unbounded() + return { + inbound, + outbound, + socket: Socket.make({ + runRaw: (handler, options) => + Effect.gen(function* () { + if (options?.onOpen) yield* options.onOpen + while (true) { + const message = yield* Queue.take(inbound) + if (message === end) return + const handled = handler(message) + if (Effect.isEffect(handled)) yield* Effect.asVoid(handled) + } + }), + writer: Effect.succeed((message) => Queue.offer(outbound, message).pipe(Effect.asVoid)), + }), + } +}) + +function controlMessage(transport: Effect.Success) { + return Queue.take(transport.outbound).pipe( + Effect.flatMap((message) => + typeof message === "string" + ? BrowserControlProtocol.decodeFromServer(message) + : Effect.fail(new Error("expected text control message")), + ), + ) +} + +const echoServer = Effect.acquireRelease( + Effect.callback, Error>((resume) => { + const server = createServer((socket) => socket.pipe(socket)) + server.once("error", (error) => resume(Effect.fail(error))) + server.listen(0, "127.0.0.1", () => resume(Effect.succeed(server))) + return Effect.sync(() => server.close()) + }), + (server) => Effect.sync(() => server.close()), +)