cloudaxe-opencode/packages/opencode/src/server/server.ts

268 lines
8.9 KiB
TypeScript
Raw Normal View History

2026-05-15 02:05:12 +00:00
import "./init-projectors"
import { NodeHttpServer } from "@effect/platform-node"
import * as Log from "@opencode-ai/core/util/log"
import { ConfigProvider, Context, Effect, Exit, Layer, Scope } from "effect"
import { HttpRouter, HttpServer } from "effect/unstable/http"
import { OpenApi } from "effect/unstable/httpapi"
import { createServer } from "node:http"
import { MDNS } from "./mdns"
import { HttpApiApp } from "./routes/instance/httpapi/server"
import { disposeMiddleware } from "./routes/instance/httpapi/lifecycle"
import { WebSocketTracker } from "./routes/instance/httpapi/websocket-tracker"
import { PublicApi } from "./routes/instance/httpapi/public"
import type { CorsOptions } from "./cors"
import { lazy } from "@/util/lazy"
2025-06-09 18:01:11 +00:00
// @ts-ignore This global is needed to prevent ai-sdk from logging warnings to stdout https://github.com/vercel/ai/blob/2dc67e0ef538307f21368db32d5a12345d98831b/packages/ai/src/logger/log-warnings.ts#L85
globalThis.AI_SDK_LOG_WARNINGS = false
const log = Log.create({ service: "server" })
export type TcpListener = {
type: "tcp"
hostname: string
port: number
url: URL
stop: (close?: boolean) => Promise<void>
}
export type SocketListener = {
type: "socket"
socket: string
url: URL
stop: (close?: boolean) => Promise<void>
}
export type Listener = TcpListener | SocketListener
type ServerApp = {
fetch(request: Request): Response | Promise<Response>
request(input: string | URL | Request, init?: RequestInit): Response | Promise<Response>
}
export type TcpListenOptions = CorsOptions & {
type: "tcp"
port: number
hostname: string
2026-05-16 04:51:18 +00:00
mdns?: true | { domain: string }
}
export type SocketListenOptions = CorsOptions & {
type: "socket"
socket: string
}
export type ListenOptions = TcpListenOptions | SocketListenOptions
type ListenerState = {
scope: Scope.Scope
server: Context.Service.Shape<typeof HttpServer.HttpServer>
http: ListenerServer
websockets: WebSocketTracker.Interface
}
type EffectTcpListener = Omit<TcpListener, "stop"> & {
stop: (close?: boolean) => Effect.Effect<void>
}
type EffectSocketListener = Omit<SocketListener, "stop"> & {
stop: (close?: boolean) => Effect.Effect<void>
}
type EffectListener = EffectTcpListener | EffectSocketListener
interface ListenerServer {
readonly closeAll: Effect.Effect<void>
}
class ListenerServerService extends Context.Service<ListenerServerService, ListenerServer>()(
"@opencode/ListenerServer",
) {}
export const Default = lazy(() => {
const handler = HttpApiApp.webHandler().handler
const app: ServerApp = {
fetch: (request: Request) => handler(request, HttpApiApp.context),
request(input, init) {
return app.fetch(input instanceof Request ? input : new Request(new URL(input, "http://localhost"), init))
},
}
return { app }
})
export async function openapi() {
return OpenApi.fromApi(PublicApi)
}
export let url: URL
export function listen(opts: TcpListenOptions): Promise<TcpListener>
export function listen(opts: SocketListenOptions): Promise<SocketListener>
export async function listen(opts: ListenOptions): Promise<Listener> {
const listener = await Effect.runPromise(listenEffect(opts))
const stop = (close?: boolean) => Effect.runPromiseExit(listener.stop(close)).then(() => undefined)
if (listener.type === "socket") {
return {
type: "socket" as const,
socket: listener.socket,
url: listener.url,
stop,
}
}
return {
type: "tcp" as const,
hostname: listener.hostname,
port: listener.port,
url: listener.url,
stop,
}
}
const listenEffect: (opts: ListenOptions) => Effect.Effect<EffectListener, unknown> = Effect.fn("Server.listen")(
function* (opts: ListenOptions) {
const state = yield* startWithPortFallback(opts)
if (opts.type === "socket") {
const listenerUrl = makeURL("localhost", 0)
url = listenerUrl
return {
type: "socket" as const,
socket: opts.socket,
url: listenerUrl,
stop: yield* makeStop(state, Effect.void),
}
}
2026-05-16 04:23:09 +00:00
const address = yield* tcpAddress(state)
const listenerUrl = makeURL(opts.hostname, address.port)
url = listenerUrl
const unpublishMdns = yield* setupMdns(opts, address.port, state.scope)
return {
type: "tcp" as const,
hostname: opts.hostname,
port: address.port,
url: listenerUrl,
stop: yield* makeStop(state, unpublishMdns),
}
},
)
function listenerLayer(opts: ListenOptions) {
return HttpRouter.serve(HttpApiApp.createRoutes(opts), {
middleware: disposeMiddleware,
disableLogger: true,
disableListenLog: true,
}).pipe(
Layer.provideMerge(WebSocketTracker.layer),
Layer.provideMerge(serverLayer(opts)),
// Install a fresh `ConfigProvider` per listener so `Config.string(...)`
// reads reflect the current `process.env`. Effect's default
// `ConfigProvider` snapshots `process.env` on first read and caches the
// result on a module-singleton Reference; without overriding it here,
// every later `Server.listen()` keeps observing that initial snapshot.
Layer.provide(ConfigProvider.layer(ConfigProvider.fromEnv())),
)
}
function startWithPortFallback(opts: ListenOptions) {
if (opts.type === "socket" || opts.port !== 0) return startListener(opts)
// Match the legacy listener port-resolution behavior: explicit `0` prefers
// 4096 first, then any free port.
return startListener({ ...opts, port: 4096 }).pipe(Effect.catch(() => startListener(opts)))
}
function startListener(opts: ListenOptions) {
const scope = Scope.makeUnsafe()
return Layer.buildWithMemoMap(listenerLayer(opts), Layer.makeMemoMapUnsafe(), scope).pipe(
Effect.provide(HttpApiApp.context),
Effect.onError(() => Scope.close(scope, Exit.void).pipe(Effect.ignore)),
Effect.map(
(ctx): ListenerState => ({
scope,
server: Context.get(ctx, HttpServer.HttpServer),
http: Context.get(ctx, ListenerServerService),
websockets: Context.get(ctx, WebSocketTracker.Service),
}),
),
)
}
function tcpAddress(state: ListenerState) {
2026-05-16 04:23:09 +00:00
return Effect.gen(function* () {
if (state.server.address._tag === "TcpAddress") return state.server.address
2026-05-16 04:23:09 +00:00
yield* Scope.close(state.scope, Exit.void).pipe(Effect.ignore)
return yield* Effect.die(new Error(`Unexpected HttpServer address tag: ${state.server.address._tag}`))
})
}
function makeURL(hostname: string, port: number) {
const result = new URL("http://localhost")
result.hostname = hostname
result.port = String(port)
return result
}
function setupMdns(opts: TcpListenOptions, port: number, scope: Scope.Scope) {
return Effect.gen(function* () {
const publish =
opts.mdns && port && opts.hostname !== "127.0.0.1" && opts.hostname !== "localhost" && opts.hostname !== "::1"
if (publish) {
const unpublish = yield* Effect.cached(Effect.sync(() => MDNS.unpublish()))
2026-05-16 04:51:18 +00:00
yield* Effect.sync(() => MDNS.publish(port, opts.mdns === true ? undefined : opts.mdns.domain))
yield* Scope.addFinalizer(scope, unpublish)
return unpublish
}
if (opts.mdns) log.warn("mDNS enabled but hostname is loopback; skipping mDNS publish")
return Effect.void
})
}
function makeStop(state: ListenerState, unpublishMdns: Effect.Effect<void>) {
return Effect.gen(function* () {
const forceCloseOnce = yield* Effect.cached(forceClose(state).pipe(Effect.ignore))
const closeScopeOnce = yield* Effect.cached(Scope.close(state.scope, Exit.void).pipe(Effect.ignore))
return (close?: boolean) =>
Effect.gen(function* () {
yield* unpublishMdns
if (close) yield* forceCloseOnce
yield* closeScopeOnce
})
})
}
function forceClose(state: ListenerState) {
return Effect.all([state.http.closeAll, state.websockets.closeAll], { concurrency: "unbounded", discard: true })
}
function serverLayer(opts: ListenOptions) {
const server = createServer()
const serverRef = { closeStarted: false, forceStop: false }
const close = server.close.bind(server)
// Keep shutdown owned by NodeHttpServer, but honor listener.stop(true) by
// force-closing active HTTP sockets when its finalizer calls server.close().
// oxlint-disable-next-line typescript-eslint/no-unsafe-type-assertion -- Node's overloads don't preserve a monkey-patched method assignment.
server.close = ((callback?: Parameters<typeof server.close>[0]) => {
serverRef.closeStarted = true
const result = close(callback)
if (serverRef.forceStop) server.closeAllConnections()
return result
}) as typeof server.close
return Layer.mergeAll(
2026-05-16 04:49:27 +00:00
NodeHttpServer.layer(() => server, {
...(opts.type === "socket" ? { path: opts.socket } : { port: opts.port, host: opts.hostname }),
gracefulShutdownTimeout: "1 second",
}),
Layer.succeed(ListenerServerService)(
ListenerServerService.of({
closeAll: Effect.sync(() => {
serverRef.forceStop = true
if (serverRef.closeStarted) server.closeAllConnections()
}),
}),
),
)
2025-05-18 01:31:42 +00:00
}
export * as Server from "./server"