From 566fa7d254db4be17e58588e197295ed481881bd Mon Sep 17 00:00:00 2001 From: dhaern Date: Sat, 26 Sep 2026 08:04:23 +0000 Subject: [PATCH 1/2] fix(opencode): share app graph and correct listener freshness boundary The original fresh boundary omitted router-registering middleware. Keep corsVaryFix fresh per listener while retaining shared app services, and build Default without persistent test auth. --- .../server/routes/instance/httpapi/server.ts | 14 +- packages/opencode/src/server/server.ts | 9 +- .../test/server/httpapi-listen.test.ts | 142 ++++++++++++++++++ 3 files changed, 153 insertions(+), 12 deletions(-) diff --git a/packages/opencode/src/server/routes/instance/httpapi/server.ts b/packages/opencode/src/server/routes/instance/httpapi/server.ts index fb9d2db65621..fc5addd4adf7 100644 --- a/packages/opencode/src/server/routes/instance/httpapi/server.ts +++ b/packages/opencode/src/server/routes/instance/httpapi/server.ts @@ -273,19 +273,15 @@ export function createRoutes( ): Layer.Layer { const locationServiceMapV2 = buildLocationServiceMap() - return Layer.mergeAll( - rootApiRoutes, - eventApiRoutes, - ptyConnectApiRoutes, - instanceRoutes, - serverRoutes, - docRoute, - uiRoute, + // Every router-registering layer must be fresh per listener. + // cors(corsOptions) already is because each call creates a new layer. + return Layer.fresh( + Layer.mergeAll(rootApiRoutes, eventApiRoutes, ptyConnectApiRoutes, instanceRoutes, serverRoutes, docRoute, uiRoute), ).pipe( Layer.provide([ errorLayer, compressionLayer, - corsVaryFix, + Layer.fresh(corsVaryFix), fenceLayer, cors(corsOptions), AppNodeBuilderV1.build(MoveSession.node, [[LocationServiceMap.node, locationServiceMapV2]]), diff --git a/packages/opencode/src/server/server.ts b/packages/opencode/src/server/server.ts index 440b992c1557..f7e3be4ed611 100644 --- a/packages/opencode/src/server/server.ts +++ b/packages/opencode/src/server/server.ts @@ -2,6 +2,7 @@ import "./init-projectors" import { NodeHttpServer } from "@effect/platform-node" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" +import { memoMap } from "@opencode-ai/core/effect/memo-map" import { ConfigProvider, Context, Effect, Exit, Layer, Scope } from "effect" import { HttpRouter, HttpServer } from "effect/unstable/http" import { OpenApi } from "effect/unstable/httpapi" @@ -98,12 +99,14 @@ const listenEffect: (opts: ListenOptions) => Effect.Effect Scope.close(scope, Exit.void).pipe(Effect.ignore)), Effect.map( diff --git a/packages/opencode/test/server/httpapi-listen.test.ts b/packages/opencode/test/server/httpapi-listen.test.ts index 585c59cb4317..ee9a95693d7d 100644 --- a/packages/opencode/test/server/httpapi-listen.test.ts +++ b/packages/opencode/test/server/httpapi-listen.test.ts @@ -3,6 +3,7 @@ import net from "node:net" import path from "node:path" import { pathToFileURL } from "node:url" import { Flag } from "@opencode-ai/core/flag/flag" +import { InstanceRuntime } from "../../src/project/instance-runtime" import { Server } from "../../src/server/server" import { PtyPaths } from "../../src/server/routes/instance/httpapi/groups/pty" import { withTimeout } from "../../src/util/timeout" @@ -166,7 +167,148 @@ async function openPtySocket(listener: Awaited> } } +async function markerInstance() { + return tmpdir({ + init: async (directory) => { + const plugin = path.join(directory, "plugin.ts") + const initialized = path.join(directory, "initialized.txt") + await Bun.write( + plugin, + [ + 'import { appendFileSync } from "node:fs"', + "export default async function plugin() {", + ` appendFileSync(${JSON.stringify(initialized)}, "initialized\\n")`, + " return {}", + "}", + ].join("\n"), + ) + await Bun.write( + path.join(directory, "opencode.json"), + JSON.stringify({ formatter: false, lsp: false, plugin: [pathToFileURL(plugin).href] }), + ) + return initialized + }, + }) +} + +async function requestConfig(listener: Awaited>, directory: string) { + const response = await fetch(new URL("/config", listener.url), { + headers: { authorization: authorization(), "x-opencode-directory": directory }, + }) + expect(response.status).toBe(200) +} + describe("HttpApi Server.listen", () => { + for (const order of ["runtime-first", "listener-first", "concurrent"] as const) { + test(`shares one bootstrapped instance with AppRuntime (${order})`, async () => { + await using tmp = await markerInstance() + const previous = process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS + process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS = "1" + let listener: Awaited> | undefined + try { + const active = await startListener() + listener = active + const load = () => InstanceRuntime.load({ directory: tmp.path }) + const request = () => requestConfig(active, tmp.path) + if (order === "runtime-first") { + await load() + await request() + } else if (order === "listener-first") { + await request() + await load() + } else { + await Promise.all([load(), request()]) + } + expect(await Bun.file(tmp.extra).text()).toBe("initialized\n") + } finally { + if (listener) await stop(listener, "timed out cleaning up shared-graph listener") + if (previous === undefined) delete process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS + else process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS = previous + } + }) + } + + test("keeps the AppRuntime instance alive across listener restarts", async () => { + await using tmp = await markerInstance() + const previous = process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS + process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS = "1" + let listener: Awaited> | undefined + try { + const first = await InstanceRuntime.load({ directory: tmp.path }) + listener = await startListener() + await requestConfig(listener, tmp.path) + await stop(listener, "timed out stopping first shared-graph listener") + listener = await startListener() + await requestConfig(listener, tmp.path) + expect(await InstanceRuntime.load({ directory: tmp.path })).toBe(first) + expect(await Bun.file(tmp.extra).text()).toBe("initialized\n") + } finally { + if (listener) await stop(listener, "timed out cleaning up restarted shared-graph listener") + if (previous === undefined) delete process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS + else process.env.OPENCODE_DISABLE_DEFAULT_PLUGINS = previous + } + }) + + test("keeps Vary: Origin on preflight after Default and across listeners", async () => { + const headers = { + origin: "http://localhost:3000", + "access-control-request-method": "POST", + "access-control-request-headers": "content-type, x-opencode-directory", + } + const preflight = (response: Response) => { + expect([200, 204]).toContain(response.status) + expect(response.headers.get("access-control-allow-origin")).toBe(headers.origin) + expect((response.headers.get("vary") ?? "").toLowerCase()).toContain("origin") + } + preflight(await Server.Default().app.request("/global/config", { method: "OPTIONS", headers })) + const first = await startNoAuthListener() + let second: Awaited> | undefined + try { + second = await startNoAuthListener() + for (const listener of [first, second]) { + preflight(await fetch(new URL("/global/config", listener.url), { method: "OPTIONS", headers })) + } + } finally { + if (second) await stop(second, "timed out cleaning up second CORS listener") + await stop(first, "timed out cleaning up first CORS listener") + } + }) + + test("uses fresh password auth after Default() builds the shared graph", async () => { + const initial = await Server.Default().app.request("/status") + expect(initial.status).toBe(200) + const listener = await startListener() + try { + expect((await fetch(new URL("/status", listener.url))).status).toBe(401) + expect( + (await fetch(new URL("/status", listener.url), { headers: { authorization: authorization() } })).status, + ).toBe(200) + } finally { + await stop(listener, "timed out cleaning up fresh-auth listener") + } + }) + + testPty("stop(true) closes only its own listener's websockets", async () => { + await using tmp = await tmpdir({ git: true, config: { formatter: false, lsp: false } }) + const first = await startListener() + let second: Awaited> | undefined + try { + second = await startListener() + const one = await openPtySocket(first, tmp.path) + const two = await openPtySocket(second, tmp.path) + await stop(first, "timed out stopping first listener with two websockets") + await withTimeout(one.closed, 5_000, "first listener websocket stayed open") + expect(two.ws.readyState).toBe(WebSocket.OPEN) + const message = waitForMessage(two.ws, (data) => data.includes("still-open")) + two.ws.send("still-open\n") + expect(await message).toContain("still-open") + two.ws.close(1000) + } finally { + if (second) await stop(second, "timed out cleaning up second websocket listener") + await stop(first, "timed out cleaning up first websocket listener") + } + }) + testPty("serves HTTP routes and upgrades PTY websocket through Server.listen", async () => { await using tmp = await tmpdir({ config: { formatter: false, lsp: false } }) const listener = await startListener() From d3c74c61e74237e11b9baf099c819fd113f1bdf2 Mon Sep 17 00:00:00 2001 From: dhaern Date: Sat, 26 Sep 2026 08:04:23 +0000 Subject: [PATCH 2/2] fix(opencode): complete interrupted instance loads --- .../opencode/src/project/instance-store.ts | 83 +++++++++----- .../opencode/test/project/instance.test.ts | 104 +++++++++++++++++- 2 files changed, 155 insertions(+), 32 deletions(-) diff --git a/packages/opencode/src/project/instance-store.ts b/packages/opencode/src/project/instance-store.ts index 720549ddaff7..94d84c89cf6a 100644 --- a/packages/opencode/src/project/instance-store.ts +++ b/packages/opencode/src/project/instance-store.ts @@ -31,7 +31,9 @@ export class Service extends Context.Service()("@opencode/In export const use = serviceUse(Service) interface Entry { - readonly deferred: Deferred.Deferred + // Carry the Exit as a value so one interrupted waiter cannot prevent other + // Deferred subscribers from being notified of the same interruption. + readonly deferred: Deferred.Deferred> } const layer: Layer.Layer = Layer.effect( @@ -69,12 +71,17 @@ const layer: Layer.Layer - Effect.gen(function* () { - const exit = yield* Effect.exit(boot({ ...input, directory })) - if (Exit.isFailure(exit)) yield* removeEntry(directory, entry) - yield* Deferred.done(entry.deferred, exit).pipe(Effect.asVoid) - }) + const completeLoad = (directory: string, input: LoadInput, entry: Entry, prepare: Effect.Effect) => + prepare.pipe( + Effect.andThen(boot({ ...input, directory })), + Effect.onExit((exit) => + Effect.gen(function* () { + if (Exit.isFailure(exit)) yield* removeEntry(directory, entry) + yield* Deferred.succeed(entry.deferred, exit) + }), + ), + Effect.asVoid, + ) const emitDisposed = (input: { directory: string; project?: string }) => Effect.sync(() => @@ -110,15 +117,18 @@ const layer: Layer.Layer Effect.gen(function* () { const existing = cache.get(directory) - if (existing) return yield* restore(Deferred.await(existing.deferred)) + if (existing) { + const exit = yield* restore(Deferred.await(existing.deferred)) + return yield* exit + } - const entry: Entry = { deferred: Deferred.makeUnsafe() } + const entry: Entry = { deferred: Deferred.makeUnsafe>() } cache.set(directory, entry) - yield* Effect.gen(function* () { - yield* Effect.logInfo("creating instance", { directory: directory }) - yield* completeLoad(directory, input, entry) - }).pipe(Effect.forkIn(scope, { startImmediately: true })) - return yield* restore(Deferred.await(entry.deferred)) + yield* completeLoad(directory, input, entry, Effect.logInfo("creating instance", { directory })).pipe( + Effect.forkIn(scope, { startImmediately: true }), + ) + const exit = yield* restore(Deferred.await(entry.deferred)) + return yield* exit }), ).pipe(Effect.withSpan("InstanceStore.load")) } @@ -128,28 +138,38 @@ const layer: Layer.Layer Effect.gen(function* () { const previous = cache.get(directory) - const entry: Entry = { deferred: Deferred.makeUnsafe() } + const entry: Entry = { deferred: Deferred.makeUnsafe>() } cache.set(directory, entry) - yield* Effect.gen(function* () { - yield* Effect.logInfo("reloading instance", { directory: directory }) - if (previous) { - yield* Deferred.await(previous.deferred).pipe(Effect.ignore) + yield* completeLoad( + directory, + input, + entry, + Effect.gen(function* () { + yield* Effect.logInfo("reloading instance", { directory }) + if (!previous) return + yield* (yield* Deferred.await(previous.deferred)).pipe(Effect.ignore) yield* Effect.promise(() => runDisposers(directory)) yield* emitDisposed({ directory, project: input.project?.id }) - } - yield* completeLoad(directory, input, entry) - }).pipe(Effect.forkIn(scope, { startImmediately: true })) - return yield* restore(Deferred.await(entry.deferred)) + }), + ).pipe(Effect.forkIn(scope, { startImmediately: true })) + const exit = yield* restore(Deferred.await(entry.deferred)) + return yield* exit }), ).pipe(Effect.withSpan("InstanceStore.reload")) } const dispose = Effect.fn("InstanceStore.dispose")(function* (ctx: InstanceContext) { const entry = cache.get(ctx.directory) - if (!entry) return yield* disposeContext(ctx) - - const exit = yield* Deferred.await(entry.deferred).pipe(Effect.exit) - if (Exit.isFailure(exit)) return yield* removeEntry(ctx.directory, entry).pipe(Effect.asVoid) + if (!entry) { + yield* disposeContext(ctx) + return + } + + const exit = yield* Deferred.await(entry.deferred) + if (Exit.isFailure(exit)) { + yield* removeEntry(ctx.directory, entry).pipe(Effect.asVoid) + return + } if (exit.value !== ctx) return yield* disposeEntry(ctx.directory, entry, ctx).pipe(Effect.asVoid) }) @@ -158,8 +178,11 @@ const layer: Layer.Layer Effect.gen(function* () { - const exit = yield* Deferred.await(item[1].deferred).pipe(Effect.exit) + const exit = yield* Deferred.await(item[1].deferred) if (Exit.isFailure(exit)) { yield* Effect.logWarning("instance dispose failed", { key: item[0], cause: exit.cause }) yield* removeEntry(item[0], item[1]) diff --git a/packages/opencode/test/project/instance.test.ts b/packages/opencode/test/project/instance.test.ts index f78b99ef7d9b..4f6af7bfc3e3 100644 --- a/packages/opencode/test/project/instance.test.ts +++ b/packages/opencode/test/project/instance.test.ts @@ -1,13 +1,13 @@ import { describe, expect } from "bun:test" import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" -import { Deferred, Effect, Fiber, Layer } from "effect" +import { Cause, Context, Deferred, Effect, Exit, Fiber, Layer, Scope } from "effect" import { InstanceRef } from "../../src/effect/instance-ref" import { registerDisposer } from "../../src/effect/instance-registry" import { InstanceBootstrap } from "../../src/project/bootstrap" import { InstanceStore } from "../../src/project/instance-store" import { tmpdirScoped } from "../fixture/fixture" -import { testEffect } from "../lib/effect" +import { awaitWithTimeout, testEffect } from "../lib/effect" let bootstrapRun: Effect.Effect = Effect.void const noopBootstrap = Layer.succeed( @@ -39,6 +39,106 @@ const registerDisposerScoped = (disposer: (directory: string) => Promise) ) describe("InstanceStore", () => { + it.live("closing the store scope interrupts an in-flight load and its waiters", () => + Effect.gen(function* () { + const dir = yield* tmpdirScoped({ git: true }) + const started = yield* Deferred.make() + const release = yield* Deferred.make() + const storeScope = yield* Scope.make() + const testLayer = LayerNode.compile(LayerNode.group([InstanceStore.node, CrossSpawnSpawner.node]), [ + [ + InstanceStore.bootstrapNode, + Layer.succeed( + InstanceBootstrap.Service, + InstanceBootstrap.Service.of({ + run: Effect.gen(function* () { + yield* Deferred.succeed(started, undefined) + yield* Deferred.await(release) + }), + }), + ), + ], + ]) + const services = yield* Layer.buildWithMemoMap(testLayer, Layer.makeMemoMapUnsafe(), storeScope) + const store = Context.get(services, InstanceStore.Service) + const load = yield* store.load({ directory: dir }).pipe(Effect.exit, Effect.forkScoped) + yield* Deferred.await(started) + const waiter = yield* store + .load({ directory: dir }) + .pipe(Effect.exit, Effect.forkScoped({ startImmediately: true })) + const closing = yield* Scope.close(storeScope, Exit.void).pipe(Effect.forkScoped) + const exit = yield* awaitWithTimeout(Fiber.join(waiter), "load waiter never saw scope interruption").pipe( + Effect.ensuring(Deferred.succeed(release, undefined)), + ) + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true) + yield* awaitWithTimeout(Fiber.join(closing), "store scope close remained blocked") + const initial = yield* awaitWithTimeout(Fiber.join(load), "initial load remained blocked") + expect(Exit.isFailure(initial)).toBe(true) + if (Exit.isFailure(initial)) expect(Cause.hasInterruptsOnly(initial.cause)).toBe(true) + }), + ) + + it.live("evicts interrupted bootstraps so the same directory can retry", () => + Effect.gen(function* () { + const dir = yield* tmpdirScoped({ git: true }) + const store = yield* InstanceStore.Service + let attempts = 0 + yield* setBootstrap( + Effect.gen(function* () { + attempts++ + if (attempts === 1) yield* Effect.interrupt + }), + ) + const first = yield* Effect.exit(store.load({ directory: dir })) + expect(Exit.isFailure(first)).toBe(true) + if (Exit.isFailure(first)) expect(Cause.hasInterruptsOnly(first.cause)).toBe(true) + expect((yield* store.load({ directory: dir })).directory).toBe(dir) + expect(attempts).toBe(2) + }), + ) + + it.live("disposes an uncached context after its cache entry is removed", () => + Effect.gen(function* () { + const dir = yield* tmpdirScoped({ git: true }) + const store = yield* InstanceStore.Service + const ctx = yield* store.load({ directory: dir }) + yield* store.dispose(ctx) + yield* store.dispose(ctx) + expect(yield* store.load({ directory: dir })).not.toBe(ctx) + }), + ) + + it.live("unblocks all disposal paths when an in-flight reload is interrupted", () => + Effect.gen(function* () { + const dir = yield* tmpdirScoped({ git: true }) + const store = yield* InstanceStore.Service + const first = yield* store.load({ directory: dir }) + const started = yield* Deferred.make() + const release = yield* Deferred.make() + yield* setBootstrap( + Effect.gen(function* () { + yield* Deferred.succeed(started, undefined) + yield* Deferred.await(release) + yield* Effect.interrupt + }), + ) + const reloading = yield* store.reload({ directory: dir }).pipe(Effect.exit, Effect.forkScoped) + yield* Deferred.await(started) + const disposing = yield* store.dispose(first).pipe(Effect.forkScoped({ startImmediately: true })) + const directory = yield* store.disposeDirectory(dir).pipe(Effect.forkScoped({ startImmediately: true })) + const all = yield* store.disposeAll().pipe(Effect.forkScoped({ startImmediately: true })) + yield* Deferred.succeed(release, undefined) + const exit = yield* awaitWithTimeout(Fiber.join(reloading), "reload remained blocked after interruption") + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true) + yield* awaitWithTimeout( + Effect.all([Fiber.join(disposing), Fiber.join(directory), Fiber.join(all)]), + "dispose remained blocked", + ) + }), + ) + it.live("loads instance context", () => Effect.gen(function* () { const dir = yield* tmpdirScoped({ git: true })