mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-07 00:03:28 +00:00
serve: add socket listener mode
This commit is contained in:
parent
042e6a5c86
commit
2719140b56
3 changed files with 104 additions and 7 deletions
|
|
@ -1,12 +1,17 @@
|
|||
import { Effect } from "effect"
|
||||
import { Server } from "../../server/server"
|
||||
import { effectCmd } from "../effect-cmd"
|
||||
import { effectCmd, fail } from "../effect-cmd"
|
||||
import { withNetworkOptions, resolveNetworkOptions } from "../network"
|
||||
import { Flag } from "@opencode-ai/core/flag/flag"
|
||||
|
||||
export const ServeCommand = effectCmd({
|
||||
command: "serve",
|
||||
builder: (yargs) => withNetworkOptions(yargs),
|
||||
builder: (yargs) =>
|
||||
withNetworkOptions(yargs)
|
||||
.option("socket", {
|
||||
type: "string",
|
||||
describe: "Unix socket path or Windows named pipe name/path to listen on",
|
||||
}),
|
||||
describe: "starts a headless opencode server",
|
||||
// Server loads instances per-request via x-opencode-directory header — no
|
||||
// need for an ambient project InstanceContext at startup.
|
||||
|
|
@ -15,10 +20,37 @@ export const ServeCommand = effectCmd({
|
|||
if (!Flag.OPENCODE_SERVER_PASSWORD) {
|
||||
console.log("Warning: OPENCODE_SERVER_PASSWORD is not set; server is unsecured.")
|
||||
}
|
||||
if (args.socket) {
|
||||
const conflicts = explicitNetworkConflicts()
|
||||
if (conflicts.length) yield* fail(`--socket cannot be used with ${conflicts.join(", ")}`)
|
||||
}
|
||||
const opts = yield* resolveNetworkOptions(args)
|
||||
const server = yield* Effect.promise(() => Server.listen(opts))
|
||||
const server = yield* Effect.promise(() =>
|
||||
Server.listen(args.socket ? { ...opts, socket: resolveSocketPath(args.socket) } : opts),
|
||||
)
|
||||
if (server.socket) {
|
||||
console.log(`opencode server listening on socket ${server.socket}`)
|
||||
yield* Effect.never
|
||||
}
|
||||
console.log(`opencode server listening on http://${server.hostname}:${server.port}`)
|
||||
|
||||
yield* Effect.never
|
||||
}),
|
||||
})
|
||||
|
||||
function resolveSocketPath(input: string) {
|
||||
if (process.platform !== "win32") return input
|
||||
const lower = input.toLowerCase()
|
||||
if (lower.startsWith("\\\\.\\pipe\\") || lower.startsWith("\\\\?\\pipe\\")) return input
|
||||
const name = input
|
||||
.replace(/^[a-zA-Z]:/, (drive) => drive.slice(0, 1))
|
||||
.replace(/[\\/:]+/g, "-")
|
||||
.replace(/^-+|-+$/g, "")
|
||||
return `\\\\.\\pipe\\${name || "opencode"}`
|
||||
}
|
||||
|
||||
function explicitNetworkConflicts() {
|
||||
return ["--port", "--hostname", "--mdns", "--mdns-domain"].filter((flag) =>
|
||||
process.argv.some((arg) => arg === flag || arg.startsWith(`${flag}=`)),
|
||||
)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@ export type Listener = {
|
|||
hostname: string
|
||||
port: number
|
||||
url: URL
|
||||
socket?: string
|
||||
stop: (close?: boolean) => Promise<void>
|
||||
}
|
||||
|
||||
|
|
@ -34,6 +35,7 @@ type ServerApp = {
|
|||
type ListenOptions = CorsOptions & {
|
||||
port: number
|
||||
hostname: string
|
||||
socket?: string
|
||||
mdns?: boolean
|
||||
mdnsDomain?: string
|
||||
}
|
||||
|
|
@ -78,6 +80,7 @@ export async function listen(opts: ListenOptions): Promise<Listener> {
|
|||
hostname: listener.hostname,
|
||||
port: listener.port,
|
||||
url: listener.url,
|
||||
socket: listener.socket,
|
||||
stop: (close?: boolean) => Effect.runPromiseExit(listener.stop(close)).then(() => undefined),
|
||||
}
|
||||
}
|
||||
|
|
@ -85,6 +88,20 @@ export async function listen(opts: ListenOptions): Promise<Listener> {
|
|||
const listenEffect: (opts: ListenOptions) => Effect.Effect<EffectListener, unknown> = Effect.fn("Server.listen")(
|
||||
function* (opts: ListenOptions) {
|
||||
const state = yield* startWithPortFallback(opts)
|
||||
if (opts.socket) {
|
||||
const address = yield* unixAddress(state)
|
||||
const listenerUrl = makeURL("localhost", 0)
|
||||
url = listenerUrl
|
||||
|
||||
return {
|
||||
hostname: "localhost",
|
||||
port: 0,
|
||||
socket: address.path,
|
||||
url: listenerUrl,
|
||||
stop: yield* makeStop(state, Effect.void),
|
||||
}
|
||||
}
|
||||
|
||||
const address = yield* tcpAddress(state)
|
||||
const listenerUrl = makeURL(opts.hostname, address.port)
|
||||
url = listenerUrl
|
||||
|
|
@ -107,7 +124,7 @@ function listenerLayer(opts: ListenOptions, port: number) {
|
|||
disableListenLog: true,
|
||||
}).pipe(
|
||||
Layer.provideMerge(WebSocketTracker.layer),
|
||||
Layer.provideMerge(serverLayer({ port, hostname: opts.hostname })),
|
||||
Layer.provideMerge(serverLayer(opts.socket ? { socket: opts.socket } : { 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
|
||||
|
|
@ -118,6 +135,7 @@ function listenerLayer(opts: ListenOptions, port: number) {
|
|||
}
|
||||
|
||||
function startWithPortFallback(opts: ListenOptions) {
|
||||
if (opts.socket) return startListener(opts, opts.port)
|
||||
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.
|
||||
|
|
@ -148,10 +166,18 @@ function tcpAddress(state: ListenerState) {
|
|||
})
|
||||
}
|
||||
|
||||
function unixAddress(state: ListenerState) {
|
||||
return Effect.gen(function* () {
|
||||
if (state.server.address._tag === "UnixAddress") 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)
|
||||
if (port) result.port = String(port)
|
||||
return result
|
||||
}
|
||||
|
||||
|
|
@ -188,7 +214,7 @@ function forceClose(state: ListenerState) {
|
|||
return Effect.all([state.http.closeAll, state.websockets.closeAll], { concurrency: "unbounded", discard: true })
|
||||
}
|
||||
|
||||
function serverLayer(opts: { port: number; hostname: string }) {
|
||||
function serverLayer(opts: { port: number; hostname: string } | { socket: string }) {
|
||||
const server = createServer()
|
||||
const serverRef = { closeStarted: false, forceStop: false }
|
||||
const close = server.close.bind(server)
|
||||
|
|
@ -203,7 +229,10 @@ function serverLayer(opts: { port: number; hostname: string }) {
|
|||
}) as typeof server.close
|
||||
|
||||
return Layer.mergeAll(
|
||||
NodeHttpServer.layer(() => server, { port: opts.port, host: opts.hostname, gracefulShutdownTimeout: "1 second" }),
|
||||
NodeHttpServer.layer(() => server, {
|
||||
...("socket" in opts ? { path: opts.socket } : { port: opts.port, host: opts.hostname }),
|
||||
gracefulShutdownTimeout: "1 second",
|
||||
}),
|
||||
Layer.succeed(ListenerServerService)(
|
||||
ListenerServerService.of({
|
||||
closeAll: Effect.sync(() => {
|
||||
|
|
|
|||
|
|
@ -1,5 +1,7 @@
|
|||
import { afterEach, describe, expect, test } from "bun:test"
|
||||
import http from "node:http"
|
||||
import net from "node:net"
|
||||
import path from "node:path"
|
||||
import { Flag } from "@opencode-ai/core/flag/flag"
|
||||
import * as Log from "@opencode-ai/core/util/log"
|
||||
import { Server } from "../../src/server/server"
|
||||
|
|
@ -285,6 +287,25 @@ describe("HttpApi Server.listen", () => {
|
|||
).rejects.toThrow()
|
||||
})
|
||||
|
||||
test("listens on a socket path", async () => {
|
||||
await using tmp = await tmpdir({ git: true, config: { formatter: false, lsp: false } })
|
||||
const socket =
|
||||
process.platform === "win32"
|
||||
? `\\\\.\\pipe\\opencode-test-${process.pid}-${Date.now()}`
|
||||
: path.join(tmp.path, "opencode.sock")
|
||||
const listener = await Server.listen({ hostname: "127.0.0.1", port: 0, socket })
|
||||
try {
|
||||
expect(listener.port).toBe(0)
|
||||
expect(listener.socket).toBe(socket)
|
||||
|
||||
const response = await requestSocketRoot(socket)
|
||||
expect(response.statusCode).toBe(200)
|
||||
expect(response.body).toContain("OpenCode")
|
||||
} finally {
|
||||
await stop(listener, "timed out cleaning up socket listener")
|
||||
}
|
||||
})
|
||||
|
||||
test("default in-process handler does not emit Effect HTTP response logs", async () => {
|
||||
let output = ""
|
||||
// oxlint-disable-next-line typescript-eslint/unbound-method -- restored in finally after temporarily capturing stderr.
|
||||
|
|
@ -433,3 +454,18 @@ function occupyPort(port: number) {
|
|||
server.listen(port, "127.0.0.1", () => resolve(server))
|
||||
})
|
||||
}
|
||||
|
||||
function requestSocketRoot(socket: string) {
|
||||
return new Promise<{ statusCode?: number; body: string }>((resolve, reject) => {
|
||||
const request = http.request(
|
||||
{ socketPath: socket, path: "/", method: "GET" },
|
||||
(response) => {
|
||||
const chunks: Buffer[] = []
|
||||
response.on("data", (chunk: Buffer) => chunks.push(chunk))
|
||||
response.on("end", () => resolve({ statusCode: response.statusCode, body: Buffer.concat(chunks).toString() }))
|
||||
},
|
||||
)
|
||||
request.on("error", reject)
|
||||
request.end()
|
||||
})
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue