|
|
@@ -1,15 +1,17 @@
|
|
|
+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 * as HttpApiServer from "#httpapi-server"
|
|
|
+import { createServer } from "node:http"
|
|
|
import { MDNS } from "./mdns"
|
|
|
import { initProjectors } from "./projectors"
|
|
|
-import { ExperimentalHttpApiServer } from "./routes/instance/httpapi/server"
|
|
|
+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"
|
|
|
|
|
|
// @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
|
|
|
@@ -36,19 +38,34 @@ type ListenOptions = CorsOptions & {
|
|
|
mdns?: boolean
|
|
|
mdnsDomain?: string
|
|
|
}
|
|
|
+type ListenerState = {
|
|
|
+ scope: Scope.Scope
|
|
|
+ server: Context.Service.Shape<typeof HttpServer.HttpServer>
|
|
|
+ http: ListenerServer
|
|
|
+ websockets: WebSocketTracker.Interface
|
|
|
+}
|
|
|
+type EffectListener = Omit<Listener, "stop"> & {
|
|
|
+ stop: (close?: boolean) => Effect.Effect<void>
|
|
|
+}
|
|
|
+
|
|
|
+interface ListenerServer {
|
|
|
+ readonly closeAll: Effect.Effect<void>
|
|
|
+}
|
|
|
+
|
|
|
+class ListenerServerService extends Context.Service<ListenerServerService, ListenerServer>()(
|
|
|
+ "@opencode/ListenerServer",
|
|
|
+) {}
|
|
|
|
|
|
-const defaultHttpApi = (() => {
|
|
|
- const handler = ExperimentalHttpApiServer.webHandler().handler
|
|
|
+export const Default = lazy(() => {
|
|
|
+ const handler = HttpApiApp.webHandler().handler
|
|
|
const app: ServerApp = {
|
|
|
- fetch: (request: Request) => handler(request, ExperimentalHttpApiServer.context),
|
|
|
+ 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 const Default = () => defaultHttpApi
|
|
|
+})
|
|
|
|
|
|
export async function openapi() {
|
|
|
return OpenApi.fromApi(PublicApi)
|
|
|
@@ -57,102 +74,146 @@ export async function openapi() {
|
|
|
export let url: URL
|
|
|
|
|
|
export async function listen(opts: ListenOptions): Promise<Listener> {
|
|
|
- log.info("server backend", { "opencode.server.runtime": HttpApiServer.name })
|
|
|
-
|
|
|
- const buildLayer = (port: number) =>
|
|
|
- HttpRouter.serve(ExperimentalHttpApiServer.createRoutes(opts), {
|
|
|
- middleware: disposeMiddleware,
|
|
|
- disableLogger: true,
|
|
|
- disableListenLog: true,
|
|
|
- }).pipe(
|
|
|
- Layer.provideMerge(WebSocketTracker.layer),
|
|
|
- Layer.provideMerge(HttpApiServer.layer({ port, hostname: opts.hostname })),
|
|
|
- // 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())),
|
|
|
- )
|
|
|
-
|
|
|
- const start = async (port: number) => {
|
|
|
- const scope = Scope.makeUnsafe()
|
|
|
- try {
|
|
|
- const layer = buildLayer(port) as Layer.Layer<
|
|
|
- HttpServer.HttpServer | WebSocketTracker.Service | HttpApiServer.Service,
|
|
|
- unknown,
|
|
|
- never
|
|
|
- >
|
|
|
- const ctx = await Effect.runPromise(Layer.buildWithMemoMap(layer, Layer.makeMemoMapUnsafe(), scope))
|
|
|
- return { scope, ctx }
|
|
|
- } catch (err) {
|
|
|
- await Effect.runPromise(Scope.close(scope, Exit.void)).catch(() => undefined)
|
|
|
- throw err
|
|
|
- }
|
|
|
+ const listener = await Effect.runPromise(listenEffect(opts))
|
|
|
+ return {
|
|
|
+ hostname: listener.hostname,
|
|
|
+ port: listener.port,
|
|
|
+ url: listener.url,
|
|
|
+ stop: (close?: boolean) => Effect.runPromiseExit(listener.stop(close)).then(() => undefined),
|
|
|
}
|
|
|
+}
|
|
|
+
|
|
|
+const listenEffect: (opts: ListenOptions) => Effect.Effect<EffectListener, unknown> = Effect.fn("Server.listen")(
|
|
|
+ function* (opts: ListenOptions) {
|
|
|
+ const state = yield* startWithPortFallback(opts)
|
|
|
+ const address = yield* tcpAddress(state)
|
|
|
+ const listenerUrl = makeURL(opts.hostname, address.port)
|
|
|
+ url = listenerUrl
|
|
|
|
|
|
- // Match the legacy adapter port-resolution behavior: explicit `0` prefers
|
|
|
+ const unpublishMdns = yield* setupMdns(opts, address.port, state.scope)
|
|
|
+
|
|
|
+ return {
|
|
|
+ hostname: opts.hostname,
|
|
|
+ port: address.port,
|
|
|
+ url: listenerUrl,
|
|
|
+ stop: yield* makeStop(state, unpublishMdns),
|
|
|
+ }
|
|
|
+ },
|
|
|
+)
|
|
|
+
|
|
|
+function listenerLayer(opts: ListenOptions, port: number) {
|
|
|
+ return HttpRouter.serve(HttpApiApp.createRoutes(opts), {
|
|
|
+ middleware: disposeMiddleware,
|
|
|
+ disableLogger: true,
|
|
|
+ disableListenLog: true,
|
|
|
+ }).pipe(
|
|
|
+ Layer.provideMerge(WebSocketTracker.layer),
|
|
|
+ Layer.provideMerge(serverLayer({ port, hostname: opts.hostname })),
|
|
|
+ // 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.port !== 0) return startListener(opts, opts.port)
|
|
|
+ // Match the legacy listener port-resolution behavior: explicit `0` prefers
|
|
|
// 4096 first, then any free port.
|
|
|
- let resolved: Awaited<ReturnType<typeof start>> | undefined
|
|
|
- if (opts.port === 0) {
|
|
|
- resolved = await start(4096).catch(() => undefined)
|
|
|
- if (!resolved) resolved = await start(0)
|
|
|
- } else {
|
|
|
- resolved = await start(opts.port)
|
|
|
- }
|
|
|
- if (!resolved) throw new Error(`Failed to start server on port ${opts.port}`)
|
|
|
+ return startListener(opts, 4096).pipe(Effect.catch(() => startListener(opts, 0)))
|
|
|
+}
|
|
|
|
|
|
- const server = Context.get(resolved.ctx, HttpServer.HttpServer)
|
|
|
- if (server.address._tag !== "TcpAddress") {
|
|
|
- await Effect.runPromise(Scope.close(resolved.scope, Exit.void))
|
|
|
- throw new Error(`Unexpected HttpServer address tag: ${server.address._tag}`)
|
|
|
- }
|
|
|
- const port = server.address.port
|
|
|
-
|
|
|
- const innerUrl = new URL("http://localhost")
|
|
|
- innerUrl.hostname = opts.hostname
|
|
|
- innerUrl.port = String(port)
|
|
|
- url = innerUrl
|
|
|
-
|
|
|
- const mdns =
|
|
|
- opts.mdns && port && opts.hostname !== "127.0.0.1" && opts.hostname !== "localhost" && opts.hostname !== "::1"
|
|
|
- if (mdns) {
|
|
|
- MDNS.publish(port, opts.mdnsDomain)
|
|
|
- } else if (opts.mdns) {
|
|
|
- log.warn("mDNS enabled but hostname is loopback; skipping mDNS publish")
|
|
|
- }
|
|
|
+function startListener(opts: ListenOptions, port: number) {
|
|
|
+ const scope = Scope.makeUnsafe()
|
|
|
+ return Layer.buildWithMemoMap(listenerLayer(opts, port), 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),
|
|
|
+ }),
|
|
|
+ ),
|
|
|
+ )
|
|
|
+}
|
|
|
|
|
|
- let forceStopPromise: Promise<void> | undefined
|
|
|
- let stopPromise: Promise<void> | undefined
|
|
|
- let mdnsUnpublished = false
|
|
|
- const unpublish = () => {
|
|
|
- if (!mdns || mdnsUnpublished) return
|
|
|
- mdnsUnpublished = true
|
|
|
- MDNS.unpublish()
|
|
|
- }
|
|
|
- const forceStop = () => {
|
|
|
- forceStopPromise ??= Effect.runPromiseExit(
|
|
|
+function tcpAddress(state: ListenerState) {
|
|
|
+ return Effect.gen(function* () {
|
|
|
+ if (state.server.address._tag === "TcpAddress") return state.server.address
|
|
|
+ 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: ListenOptions, 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()))
|
|
|
+ yield* Effect.sync(() => MDNS.publish(port, opts.mdnsDomain))
|
|
|
+ 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* Context.get(resolved!.ctx, HttpApiServer.Service).closeAll
|
|
|
- yield* Context.get(resolved!.ctx, WebSocketTracker.Service).closeAll
|
|
|
- }),
|
|
|
- ).then(() => undefined)
|
|
|
- return forceStopPromise
|
|
|
- }
|
|
|
+ yield* unpublishMdns
|
|
|
+ if (close) yield* forceCloseOnce
|
|
|
+ yield* closeScopeOnce
|
|
|
+ })
|
|
|
+ })
|
|
|
+}
|
|
|
|
|
|
- return {
|
|
|
- hostname: opts.hostname,
|
|
|
- port,
|
|
|
- url: innerUrl,
|
|
|
- stop: (close?: boolean) => {
|
|
|
- unpublish()
|
|
|
- const requested = close ? forceStop() : Promise.resolve()
|
|
|
- stopPromise ??= requested
|
|
|
- .then(() => Effect.runPromiseExit(Scope.close(resolved!.scope, Exit.void)))
|
|
|
- .then(() => undefined)
|
|
|
- return requested.then(() => stopPromise!)
|
|
|
- },
|
|
|
- }
|
|
|
+function forceClose(state: ListenerState) {
|
|
|
+ return Effect.all([state.http.closeAll, state.websockets.closeAll], { concurrency: "unbounded", discard: true })
|
|
|
+}
|
|
|
+
|
|
|
+function serverLayer(opts: { port: number; hostname: string }) {
|
|
|
+ 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(
|
|
|
+ NodeHttpServer.layer(() => server, { 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()
|
|
|
+ }),
|
|
|
+ }),
|
|
|
+ ),
|
|
|
+ )
|
|
|
}
|
|
|
|
|
|
export * as Server from "./server"
|