Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
83 changes: 53 additions & 30 deletions packages/opencode/src/project/instance-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,9 @@ export class Service extends Context.Service<Service, Interface>()("@opencode/In
export const use = serviceUse(Service)

interface Entry {
readonly deferred: Deferred.Deferred<InstanceContext>
// 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<Exit.Exit<InstanceContext>>
}

const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Service> = Layer.effect(
Expand Down Expand Up @@ -69,12 +71,17 @@ const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Ser
return true
})

const completeLoad = (directory: string, input: LoadInput, entry: Entry) =>
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<void>) =>
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(() =>
Expand Down Expand Up @@ -110,15 +117,18 @@ const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Ser
return Effect.uninterruptibleMask((restore) =>
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<InstanceContext>() }
const entry: Entry = { deferred: Deferred.makeUnsafe<Exit.Exit<InstanceContext>>() }
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"))
}
Expand All @@ -128,28 +138,38 @@ const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Ser
return Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const previous = cache.get(directory)
const entry: Entry = { deferred: Deferred.makeUnsafe<InstanceContext>() }
const entry: Entry = { deferred: Deferred.makeUnsafe<Exit.Exit<InstanceContext>>() }
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)
})
Expand All @@ -158,8 +178,11 @@ const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Ser
const directory = FSUtil.resolve(input)
const entry = cache.get(directory)
if (!entry) return
const exit = yield* Deferred.await(entry.deferred).pipe(Effect.exit)
if (Exit.isFailure(exit)) return yield* removeEntry(directory, entry).pipe(Effect.asVoid)
const exit = yield* Deferred.await(entry.deferred)
if (Exit.isFailure(exit)) {
yield* removeEntry(directory, entry).pipe(Effect.asVoid)
return
}
yield* disposeEntry(directory, entry, exit.value).pipe(Effect.asVoid)
})

Expand All @@ -169,7 +192,7 @@ const layer: Layer.Layer<Service, never, Project.Service | InstanceBootstrap.Ser
[...cache.entries()],
(item) =>
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])
Expand Down
14 changes: 5 additions & 9 deletions packages/opencode/src/server/routes/instance/httpapi/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -273,19 +273,15 @@ export function createRoutes(
): Layer.Layer<never, EffectConfig.ConfigError, RouteRequirements> {
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]]),
Expand Down
9 changes: 6 additions & 3 deletions packages/opencode/src/server/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -98,12 +99,14 @@ const listenEffect: (opts: ListenOptions) => Effect.Effect<EffectListener, unkno
)

function listenerLayer(opts: ListenOptions, port: number) {
return HttpRouter.serve(HttpApiApp.createRoutes(opts), {
// Route registrations and the router are listener-owned; the app services
// provided inside createRoutes must remain outside its fresh boundary.
return HttpRouter.serve(HttpApiApp.createRoutes(opts).pipe(Layer.provideMerge(Layer.fresh(HttpRouter.layer))), {
middleware: disposeMiddleware,
disableLogger: true,
disableListenLog: true,
}).pipe(
Layer.provideMerge(AppNodeBuilder.build(WebSocketTracker.node)),
Layer.provideMerge(Layer.fresh(AppNodeBuilder.build(WebSocketTracker.node))),
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
Expand All @@ -123,7 +126,7 @@ function startWithPortFallback(opts: ListenOptions) {

function startListener(opts: ListenOptions, port: number) {
const scope = Scope.makeUnsafe()
return Layer.buildWithMemoMap(listenerLayer(opts, port), Layer.makeMemoMapUnsafe(), scope).pipe(
return Layer.buildWithMemoMap(listenerLayer(opts, port), memoMap, scope).pipe(
Effect.provide(HttpApiApp.context),
Effect.onError(() => Scope.close(scope, Exit.void).pipe(Effect.ignore)),
Effect.map(
Expand Down
104 changes: 102 additions & 2 deletions packages/opencode/test/project/instance.test.ts
Original file line number Diff line number Diff line change
@@ -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<void> = Effect.void
const noopBootstrap = Layer.succeed(
Expand Down Expand Up @@ -39,6 +39,106 @@ const registerDisposerScoped = (disposer: (directory: string) => Promise<void>)
)

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<void>()
const release = yield* Deferred.make<void>()
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<void>()
const release = yield* Deferred.make<void>()
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 })
Expand Down
Loading
Loading