feat(server): expose browser transport
This commit is contained in:
parent
883176b218
commit
d90fc005c5
10 changed files with 670 additions and 12 deletions
|
|
@ -14,6 +14,7 @@ import { SkillGroup } from "./groups/skill.js"
|
||||||
import { EventGroup, makeEventGroup } from "./groups/event.js"
|
import { EventGroup, makeEventGroup } from "./groups/event.js"
|
||||||
import type { Definition } from "@opencode-ai/schema/event"
|
import type { Definition } from "@opencode-ai/schema/event"
|
||||||
import { AgentGroup } from "./groups/agent.js"
|
import { AgentGroup } from "./groups/agent.js"
|
||||||
|
import { BrowserGroup } from "./groups/browser.js"
|
||||||
import { PluginGroup } from "./groups/plugin.js"
|
import { PluginGroup } from "./groups/plugin.js"
|
||||||
import { HealthGroup } from "./groups/health.js"
|
import { HealthGroup } from "./groups/health.js"
|
||||||
import { ServerGroup } from "./groups/server.js"
|
import { ServerGroup } from "./groups/server.js"
|
||||||
|
|
@ -85,6 +86,7 @@ type ApiGroups<
|
||||||
| typeof HealthGroup
|
| typeof HealthGroup
|
||||||
| typeof ServerGroup
|
| typeof ServerGroup
|
||||||
| typeof DebugGroup
|
| typeof DebugGroup
|
||||||
|
| typeof BrowserGroup
|
||||||
| LocationGroups<LocationId>
|
| LocationGroups<LocationId>
|
||||||
| FormGroups<LocationId, LocationService, FormLocationId, FormLocationService>
|
| FormGroups<LocationId, LocationService, FormLocationId, FormLocationService>
|
||||||
| SessionGroups<SessionLocationId, SessionLocationService>
|
| SessionGroups<SessionLocationId, SessionLocationService>
|
||||||
|
|
@ -146,6 +148,7 @@ const makeApiFromGroup = <
|
||||||
HttpApi.make("server")
|
HttpApi.make("server")
|
||||||
.add(HealthGroup)
|
.add(HealthGroup)
|
||||||
.add(ServerGroup)
|
.add(ServerGroup)
|
||||||
|
.add(BrowserGroup)
|
||||||
.add(LocationGroup.middleware(locationMiddleware))
|
.add(LocationGroup.middleware(locationMiddleware))
|
||||||
.add(AgentGroup.middleware(locationMiddleware))
|
.add(AgentGroup.middleware(locationMiddleware))
|
||||||
.add(PluginGroup.middleware(locationMiddleware))
|
.add(PluginGroup.middleware(locationMiddleware))
|
||||||
|
|
|
||||||
|
|
@ -38,6 +38,7 @@ export const groupNames = {
|
||||||
"server.debug": "debug",
|
"server.debug": "debug",
|
||||||
"server.location": "location",
|
"server.location": "location",
|
||||||
"server.agent": "agent",
|
"server.agent": "agent",
|
||||||
|
"server.browser": "browser",
|
||||||
"server.plugin": "plugin",
|
"server.plugin": "plugin",
|
||||||
"server.session": "session",
|
"server.session": "session",
|
||||||
"server.message": "message",
|
"server.message": "message",
|
||||||
|
|
@ -63,5 +64,16 @@ export const groupNames = {
|
||||||
"server.vcs": "vcs",
|
"server.vcs": "vcs",
|
||||||
} as const
|
} as const
|
||||||
|
|
||||||
export const promiseOmitEndpoints = new Set(["pty.connect", "pty.connectToken"])
|
export const promiseOmitEndpoints = new Set([
|
||||||
export const effectOmitEndpoints = new Set(["fs.read", "pty.connect", "pty.connectToken"])
|
"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",
|
||||||
|
])
|
||||||
|
|
|
||||||
115
packages/server/src/browser-control-connection.ts
Normal file
115
packages/server/src/browser-control-connection.ts
Normal file
|
|
@ -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<Session.ID, { readonly token: object; readonly leaseID?: Browser.LeaseID }>()
|
||||||
|
|
||||||
|
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<void> = Effect.void,
|
||||||
|
) {
|
||||||
|
const browser = yield* BrowserHost.Service
|
||||||
|
const write = yield* socket.writer
|
||||||
|
const pending = new Map<BrowserControl.RequestID, Deferred.Deferred<Browser.Outcome>>()
|
||||||
|
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<Browser.Outcome>()
|
||||||
|
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 })),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
})
|
||||||
292
packages/server/src/browser-tunnel.ts
Normal file
292
packages/server/src/browser-tunnel.ts
Normal file
|
|
@ -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<CapacityError>()("BrowserTunnel.CapacityError", {
|
||||||
|
limit: Schema.Int,
|
||||||
|
message: Schema.String,
|
||||||
|
}) {}
|
||||||
|
|
||||||
|
class TunnelError extends Schema.TaggedErrorClass<TunnelError>()("BrowserTunnel.TunnelError", {
|
||||||
|
kind: Schema.Literals(["closed", "protocol", "target", "revoked"]),
|
||||||
|
message: Schema.String,
|
||||||
|
cause: Schema.optional(Schema.Defect()),
|
||||||
|
}) {}
|
||||||
|
|
||||||
|
class ConnectError extends Schema.TaggedErrorClass<ConnectError>()("BrowserTunnel.ConnectError", {
|
||||||
|
kind: Schema.Literals(["failed", "timeout"]),
|
||||||
|
message: Schema.String,
|
||||||
|
cause: Schema.optional(Schema.Defect()),
|
||||||
|
}) {}
|
||||||
|
|
||||||
|
type Dial = (host: string, port: number) => Effect.Effect<import("node:net").Socket, ConnectError, Scope.Scope>
|
||||||
|
type State = { readonly active: number; readonly shutdown: boolean }
|
||||||
|
|
||||||
|
export interface Connection {
|
||||||
|
readonly run: (socket: Socket.Socket, opened?: Effect.Effect<void>) => Effect.Effect<void, never, Scope.Scope>
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface Interface {
|
||||||
|
readonly acquire: Effect.Effect<Connection, CapacityError, Scope.Scope>
|
||||||
|
readonly shutdown: Effect.Effect<void>
|
||||||
|
}
|
||||||
|
|
||||||
|
export class Service extends Context.Service<Service, Interface>()("@opencode/server/BrowserTunnel") {}
|
||||||
|
|
||||||
|
export function make(dial: Dial = connect) {
|
||||||
|
return Effect.gen(function* () {
|
||||||
|
const browser = yield* BrowserHost.Service
|
||||||
|
const state = yield* SynchronizedRef.make<State>({ active: 0, shutdown: false })
|
||||||
|
const connections = new Set<Effect.Effect<void>>()
|
||||||
|
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<void, Socket.SocketError>,
|
||||||
|
dial: Dial,
|
||||||
|
opened: Effect.Effect<void>,
|
||||||
|
) {
|
||||||
|
const inbound = yield* Queue.bounded<string | Uint8Array, TunnelError>(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<Uint8Array, TunnelError>(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<import("node:net").Socket, ConnectError>((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<void, TunnelError>((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<void, Socket.SocketError>,
|
||||||
|
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<void, Socket.SocketError>,
|
||||||
|
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<string | Uint8Array, TunnelError>, error: TunnelError) {
|
||||||
|
return Effect.sync(() => Queue.failCauseUnsafe(queue, Cause.fail(error)))
|
||||||
|
}
|
||||||
|
|
@ -11,6 +11,7 @@ import { CommandHandler } from "./handlers/command"
|
||||||
import { SkillHandler } from "./handlers/skill"
|
import { SkillHandler } from "./handlers/skill"
|
||||||
import { EventHandler } from "./handlers/event"
|
import { EventHandler } from "./handlers/event"
|
||||||
import { AgentHandler } from "./handlers/agent"
|
import { AgentHandler } from "./handlers/agent"
|
||||||
|
import { BrowserHandler } from "./handlers/browser"
|
||||||
import { PluginHandler } from "./handlers/plugin"
|
import { PluginHandler } from "./handlers/plugin"
|
||||||
import { HealthHandler } from "./handlers/health"
|
import { HealthHandler } from "./handlers/health"
|
||||||
import { ServerHandler } from "./handlers/server"
|
import { ServerHandler } from "./handlers/server"
|
||||||
|
|
@ -35,6 +36,7 @@ export const handlers = Layer.mergeAll(
|
||||||
DebugHandler,
|
DebugHandler,
|
||||||
LocationHandler,
|
LocationHandler,
|
||||||
AgentHandler,
|
AgentHandler,
|
||||||
|
BrowserHandler,
|
||||||
PluginHandler,
|
PluginHandler,
|
||||||
SessionHandler,
|
SessionHandler,
|
||||||
MessageHandler,
|
MessageHandler,
|
||||||
|
|
|
||||||
72
packages/server/src/handlers/browser.ts
Normal file
72
packages/server/src/handlers/browser.ts
Normal file
|
|
@ -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<Record<string, string | undefined>>,
|
||||||
|
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
|
||||||
|
}
|
||||||
|
|
@ -1,9 +1,9 @@
|
||||||
import { ServerAuth } from "../auth"
|
import { ServerAuth } from "../auth"
|
||||||
import { UnauthorizedError } from "@opencode-ai/protocol/errors"
|
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"
|
export { Authorization } from "@opencode-ai/protocol/middleware/authorization"
|
||||||
import { hasPtyConnectTicketURL } from "@opencode-ai/protocol/groups/pty"
|
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"
|
import { HttpEffect, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
|
||||||
|
|
||||||
const AUTH_TOKEN_QUERY = "auth_token"
|
const AUTH_TOKEN_QUERY = "auth_token"
|
||||||
|
|
@ -26,9 +26,9 @@ function decodeCredential(input: string) {
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
function credentialFromRequest(request: HttpServerRequest.HttpServerRequest) {
|
function credentialFromRequest(request: HttpServerRequest.HttpServerRequest, headerOnly = false) {
|
||||||
const url = new URL(request.url, "http://localhost")
|
const url = new URL(request.url, "http://opencode.invalid")
|
||||||
const token = url.searchParams.get(AUTH_TOKEN_QUERY)
|
const token = headerOnly ? undefined : url.searchParams.get(AUTH_TOKEN_QUERY)
|
||||||
if (token) return decodeCredential(token)
|
if (token) return decodeCredential(token)
|
||||||
const match = /^Basic\s+(.+)$/i.exec(request.headers.authorization ?? "")
|
const match = /^Basic\s+(.+)$/i.exec(request.headers.authorization ?? "")
|
||||||
if (match) return decodeCredential(match[1])
|
if (match) return decodeCredential(match[1])
|
||||||
|
|
@ -44,13 +44,17 @@ export const authorizationLayer = Layer.effect(
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const config = yield* ServerAuth.Config
|
const config = yield* ServerAuth.Config
|
||||||
if (!ServerAuth.required(config)) return Authorization.of((effect) => effect)
|
if (!ServerAuth.required(config)) return Authorization.of((effect) => effect)
|
||||||
return Authorization.of((effect) =>
|
return Authorization.of((effect, options) =>
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const request = yield* HttpServerRequest.HttpServerRequest
|
const request = yield* HttpServerRequest.HttpServerRequest
|
||||||
// Browsers cannot set headers on WebSocket upgrades, so a ticketed PTY connect skips
|
// Browsers cannot set headers on WebSocket upgrades, so a ticketed PTY connect skips
|
||||||
// credential checks here; the connect handler consumes and validates the ticket.
|
// credential checks here; the connect handler consumes and validates the ticket.
|
||||||
if (hasPtyConnectTicketURL(new URL(request.url, "http://localhost"))) return yield* effect
|
if (hasPtyConnectTicketURL(new URL(request.url, "http://opencode.invalid"))) return yield* effect
|
||||||
if (yield* authorizedRequest(request, config)) 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) =>
|
yield* HttpEffect.appendPreResponseHandler((_request, response) =>
|
||||||
Effect.succeed(HttpServerResponse.setHeader(response, "www-authenticate", WWW_AUTHENTICATE)),
|
Effect.succeed(HttpServerResponse.setHeader(response, "www-authenticate", WWW_AUTHENTICATE)),
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -130,9 +130,39 @@ function bind(hostname: string, port: number) {
|
||||||
const parentScope = yield* Scope.Scope
|
const parentScope = yield* Scope.Scope
|
||||||
const serverScope = yield* Scope.fork(parentScope)
|
const serverScope = yield* Scope.fork(parentScope)
|
||||||
const server = createServer()
|
const server = createServer()
|
||||||
|
const sockets = new Set<import("node:net").Socket>()
|
||||||
|
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* () {
|
return yield* Effect.gen(function* () {
|
||||||
const http = yield* NodeHttpServer.make(() => server, { port, host: hostname })
|
const http = yield* NodeHttpServer.make(() => server, { port, host: hostname })
|
||||||
yield* Effect.addFinalizer(() => Effect.sync(() => server.closeAllConnections()))
|
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 }
|
return { http, server, scope: serverScope }
|
||||||
}).pipe(
|
}).pipe(
|
||||||
Effect.provideService(Scope.Scope, serverScope),
|
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
|
* 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
|
* 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) {
|
const installRestartContinuity = Effect.fnUntraced(function* (restart: SessionRestart.Interface) {
|
||||||
yield* Effect.forkScoped(restart.resumeSuspendedSessions)
|
yield* Effect.forkScoped(restart.resumeSuspendedSessions)
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ import { LayerNode } from "@opencode-ai/util/effect/layer-node"
|
||||||
import { httpClient } from "@opencode-ai/util/effect/app-node-platform"
|
import { httpClient } from "@opencode-ai/util/effect/app-node-platform"
|
||||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||||
import { Bus } from "@opencode-ai/core/bus"
|
import { Bus } from "@opencode-ai/core/bus"
|
||||||
|
import { BrowserHost } from "@opencode-ai/core/browser-host"
|
||||||
import { EventLogger } from "@opencode-ai/core/event-logger"
|
import { EventLogger } from "@opencode-ai/core/event-logger"
|
||||||
import { FileSystemSearch } from "@opencode-ai/core/filesystem/search"
|
import { FileSystemSearch } from "@opencode-ai/core/filesystem/search"
|
||||||
import { Observability } from "@opencode-ai/util/observability"
|
import { Observability } from "@opencode-ai/util/observability"
|
||||||
|
|
@ -40,11 +41,13 @@ import { layer } from "./location"
|
||||||
import { formLocationLayer } from "./middleware/form-location"
|
import { formLocationLayer } from "./middleware/form-location"
|
||||||
import { sessionLocationLayer } from "./middleware/session-location"
|
import { sessionLocationLayer } from "./middleware/session-location"
|
||||||
import { ServerInfo } from "./server-info"
|
import { ServerInfo } from "./server-info"
|
||||||
|
import { BrowserTunnelServer } from "./browser-tunnel"
|
||||||
import type { ServerOptions } from "./options"
|
import type { ServerOptions } from "./options"
|
||||||
|
|
||||||
const applicationServices = LayerNode.group([
|
const applicationServices = LayerNode.group([
|
||||||
Database.node,
|
Database.node,
|
||||||
Bus.node,
|
Bus.node,
|
||||||
|
BrowserHost.node,
|
||||||
EventLogger.node,
|
EventLogger.node,
|
||||||
httpClient,
|
httpClient,
|
||||||
Job.node,
|
Job.node,
|
||||||
|
|
@ -131,8 +134,11 @@ function makeRoutes<AuthError, AuthServices>(
|
||||||
return serviceLayer.pipe(
|
return serviceLayer.pipe(
|
||||||
Layer.flatMap((context) => {
|
Layer.flatMap((context) => {
|
||||||
const services = Layer.succeedContext(context)
|
const services = Layer.succeedContext(context)
|
||||||
|
const browserTunnel = BrowserTunnelServer.layer.pipe(Layer.provide(services))
|
||||||
const requestServices = Layer.merge(
|
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),
|
ServerInfo.layer(serviceURLs, options.app),
|
||||||
)
|
)
|
||||||
return HttpApiBuilder.layer(Api, { openapiPath: "/openapi.json" }).pipe(
|
return HttpApiBuilder.layer(Api, { openapiPath: "/openapi.json" }).pipe(
|
||||||
|
|
@ -144,6 +150,7 @@ function makeRoutes<AuthError, AuthServices>(
|
||||||
Layer.provide(schemaErrorLayer),
|
Layer.provide(schemaErrorLayer),
|
||||||
Layer.provide(auth),
|
Layer.provide(auth),
|
||||||
HttpRouter.provideRequest(requestServices),
|
HttpRouter.provideRequest(requestServices),
|
||||||
|
Layer.provideMerge(browserTunnel),
|
||||||
Layer.provideMerge(services),
|
Layer.provideMerge(services),
|
||||||
Layer.provideMerge(HttpRouter.layer),
|
Layer.provideMerge(HttpRouter.layer),
|
||||||
)
|
)
|
||||||
|
|
|
||||||
120
packages/server/test/browser.test.ts
Normal file
120
packages/server/test/browser.test.ts
Normal file
|
|
@ -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<string | Uint8Array | typeof end>()
|
||||||
|
const outbound = yield* Queue.unbounded<string | Uint8Array | Socket.CloseEvent>()
|
||||||
|
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<typeof makeSocket>) {
|
||||||
|
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<ReturnType<typeof createServer>, 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()),
|
||||||
|
)
|
||||||
Loading…
Add table
Add a link
Reference in a new issue