refactor: simplify instance store concurrency
This commit is contained in:
parent
c565bd54e2
commit
8a63cbe79c
2 changed files with 264 additions and 76 deletions
|
|
@ -3,9 +3,7 @@ import { WorkspaceContext } from "@/control-plane/workspace-context"
|
||||||
import { disposeInstance } from "@/effect/instance-registry"
|
import { disposeInstance } from "@/effect/instance-registry"
|
||||||
import { makeRuntime } from "@/effect/run-service"
|
import { makeRuntime } from "@/effect/run-service"
|
||||||
import { AppFileSystem } from "@opencode-ai/core/filesystem"
|
import { AppFileSystem } from "@opencode-ai/core/filesystem"
|
||||||
import * as Log from "@opencode-ai/core/util/log"
|
import { Context, Deferred, Effect, Exit, Layer, Scope } from "effect"
|
||||||
import { Context, Effect, Layer } from "effect"
|
|
||||||
import { iife } from "@/util/iife"
|
|
||||||
import { context, type InstanceContext } from "./instance-context"
|
import { context, type InstanceContext } from "./instance-context"
|
||||||
import * as Project from "./project"
|
import * as Project from "./project"
|
||||||
|
|
||||||
|
|
@ -25,13 +23,18 @@ export interface Interface {
|
||||||
|
|
||||||
export class Service extends Context.Service<Service, Interface>()("@opencode/InstanceStore") {}
|
export class Service extends Context.Service<Service, Interface>()("@opencode/InstanceStore") {}
|
||||||
|
|
||||||
|
interface Entry {
|
||||||
|
readonly deferred: Deferred.Deferred<InstanceContext>
|
||||||
|
}
|
||||||
|
|
||||||
export const layer: Layer.Layer<Service, never, Project.Service> = Layer.effect(
|
export const layer: Layer.Layer<Service, never, Project.Service> = Layer.effect(
|
||||||
Service,
|
Service,
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const project = yield* Project.Service
|
const project = yield* Project.Service
|
||||||
const cache = new Map<string, Promise<InstanceContext>>()
|
const scope = yield* Scope.Scope
|
||||||
|
const cache = new Map<string, Entry>()
|
||||||
const disposal = {
|
const disposal = {
|
||||||
all: undefined as Promise<void> | undefined,
|
all: undefined as Deferred.Deferred<void> | undefined,
|
||||||
}
|
}
|
||||||
|
|
||||||
const boot = Effect.fn("InstanceStore.boot")(function* (input: LoadInput & { directory: string }) {
|
const boot = Effect.fn("InstanceStore.boot")(function* (input: LoadInput & { directory: string }) {
|
||||||
|
|
@ -54,91 +57,128 @@ export const layer: Layer.Layer<Service, never, Project.Service> = Layer.effect(
|
||||||
return ctx
|
return ctx
|
||||||
})
|
})
|
||||||
|
|
||||||
function track(directory: string, next: Promise<InstanceContext>) {
|
const removeEntry = (directory: string, entry: Entry) =>
|
||||||
const task = next.catch((error) => {
|
Effect.sync(() => {
|
||||||
if (cache.get(directory) === task) cache.delete(directory)
|
if (cache.get(directory) !== entry) return false
|
||||||
throw error
|
cache.delete(directory)
|
||||||
|
return true
|
||||||
})
|
})
|
||||||
cache.set(directory, task)
|
|
||||||
return task
|
const completeLoad = Effect.fnUntraced(function* (directory: string, input: LoadInput, entry: Entry) {
|
||||||
}
|
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 emitDisposed = (input: { directory: string; project?: string }) =>
|
||||||
|
Effect.sync(() =>
|
||||||
|
GlobalBus.emit("event", {
|
||||||
|
directory: input.directory,
|
||||||
|
project: input.project,
|
||||||
|
workspace: WorkspaceContext.workspaceID,
|
||||||
|
payload: {
|
||||||
|
type: "server.instance.disposed",
|
||||||
|
properties: {
|
||||||
|
directory: input.directory,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
const disposeContext = Effect.fn("InstanceStore.disposeContext")(function* (ctx: InstanceContext) {
|
||||||
|
yield* Effect.logInfo("disposing instance", { directory: ctx.directory })
|
||||||
|
yield* Effect.promise(() => disposeInstance(ctx.directory))
|
||||||
|
yield* emitDisposed({ directory: ctx.directory, project: ctx.project.id })
|
||||||
|
})
|
||||||
|
|
||||||
|
const disposeEntry = Effect.fnUntraced(function* (directory: string, entry: Entry, ctx: InstanceContext) {
|
||||||
|
if (cache.get(directory) !== entry) return false
|
||||||
|
yield* disposeContext(ctx)
|
||||||
|
if (cache.get(directory) !== entry) return false
|
||||||
|
cache.delete(directory)
|
||||||
|
return true
|
||||||
|
})
|
||||||
|
|
||||||
const load = Effect.fn("InstanceStore.load")(function* (input: LoadInput) {
|
const load = Effect.fn("InstanceStore.load")(function* (input: LoadInput) {
|
||||||
const directory = AppFileSystem.resolve(input.directory)
|
const directory = AppFileSystem.resolve(input.directory)
|
||||||
const existing = cache.get(directory)
|
return yield* Effect.uninterruptibleMask((restore) =>
|
||||||
if (existing) return yield* Effect.promise(() => existing)
|
Effect.gen(function* () {
|
||||||
|
const existing = cache.get(directory)
|
||||||
|
if (existing) return yield* restore(Deferred.await(existing.deferred))
|
||||||
|
|
||||||
Log.Default.info("creating instance", { directory })
|
const entry: Entry = { deferred: Deferred.makeUnsafe<InstanceContext>() }
|
||||||
return yield* Effect.promise(() => track(directory, Effect.runPromise(boot({ ...input, directory }))))
|
cache.set(directory, entry)
|
||||||
|
yield* Effect.gen(function* () {
|
||||||
|
yield* Effect.logInfo("creating instance", { directory })
|
||||||
|
yield* completeLoad(directory, input, entry)
|
||||||
|
}).pipe(Effect.forkIn(scope, { startImmediately: true }))
|
||||||
|
return yield* restore(Deferred.await(entry.deferred))
|
||||||
|
}),
|
||||||
|
)
|
||||||
})
|
})
|
||||||
|
|
||||||
const reload = Effect.fn("InstanceStore.reload")(function* (input: LoadInput) {
|
const reload = Effect.fn("InstanceStore.reload")(function* (input: LoadInput) {
|
||||||
const directory = AppFileSystem.resolve(input.directory)
|
const directory = AppFileSystem.resolve(input.directory)
|
||||||
Log.Default.info("reloading instance", { directory })
|
return yield* Effect.uninterruptibleMask((restore) =>
|
||||||
yield* Effect.promise(() => disposeInstance(directory))
|
Effect.gen(function* () {
|
||||||
cache.delete(directory)
|
const previous = cache.get(directory)
|
||||||
const next = track(directory, Effect.runPromise(boot({ ...input, directory })))
|
const entry: Entry = { deferred: Deferred.makeUnsafe<InstanceContext>() }
|
||||||
|
cache.set(directory, entry)
|
||||||
GlobalBus.emit("event", {
|
yield* Effect.gen(function* () {
|
||||||
directory,
|
yield* Effect.logInfo("reloading instance", { directory })
|
||||||
project: input.project?.id,
|
if (previous) yield* Deferred.await(previous.deferred).pipe(Effect.exit, Effect.asVoid)
|
||||||
workspace: WorkspaceContext.workspaceID,
|
yield* Effect.promise(() => disposeInstance(directory))
|
||||||
payload: {
|
yield* emitDisposed({ directory, project: input.project?.id })
|
||||||
type: "server.instance.disposed",
|
yield* completeLoad(directory, input, entry)
|
||||||
properties: {
|
}).pipe(Effect.forkIn(scope, { startImmediately: true }))
|
||||||
directory,
|
return yield* restore(Deferred.await(entry.deferred))
|
||||||
},
|
}),
|
||||||
},
|
)
|
||||||
})
|
|
||||||
|
|
||||||
return yield* Effect.promise(() => next)
|
|
||||||
})
|
})
|
||||||
|
|
||||||
const dispose = Effect.fn("InstanceStore.dispose")(function* (ctx: InstanceContext) {
|
const dispose = Effect.fn("InstanceStore.dispose")(function* (ctx: InstanceContext) {
|
||||||
Log.Default.info("disposing instance", { directory: ctx.directory })
|
const entry = cache.get(ctx.directory)
|
||||||
yield* Effect.promise(() => disposeInstance(ctx.directory))
|
if (!entry) return yield* disposeContext(ctx)
|
||||||
cache.delete(ctx.directory)
|
|
||||||
|
|
||||||
GlobalBus.emit("event", {
|
const exit = yield* Deferred.await(entry.deferred).pipe(Effect.exit)
|
||||||
directory: ctx.directory,
|
if (Exit.isFailure(exit)) return yield* removeEntry(ctx.directory, entry).pipe(Effect.asVoid)
|
||||||
project: ctx.project.id,
|
if (exit.value !== ctx) return
|
||||||
workspace: WorkspaceContext.workspaceID,
|
yield* disposeEntry(ctx.directory, entry, ctx).pipe(Effect.asVoid)
|
||||||
payload: {
|
|
||||||
type: "server.instance.disposed",
|
|
||||||
properties: {
|
|
||||||
directory: ctx.directory,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
})
|
|
||||||
})
|
})
|
||||||
|
|
||||||
const disposeAll = Effect.fn("InstanceStore.disposeAll")(function* () {
|
const disposeAll = Effect.fn("InstanceStore.disposeAll")(function* () {
|
||||||
if (disposal.all) return yield* Effect.promise(() => disposal.all!)
|
return yield* Effect.uninterruptibleMask((restore) =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const existing = disposal.all
|
||||||
|
if (existing) return yield* restore(Deferred.await(existing))
|
||||||
|
|
||||||
disposal.all = iife(async () => {
|
const done = Deferred.makeUnsafe<void>()
|
||||||
Log.Default.info("disposing all instances")
|
const entries = [...cache.entries()]
|
||||||
const entries = [...cache.entries()]
|
disposal.all = done
|
||||||
for (const [key, value] of entries) {
|
const exit = yield* Effect.gen(function* () {
|
||||||
if (cache.get(key) !== value) continue
|
yield* Effect.logInfo("disposing all instances")
|
||||||
|
yield* Effect.forEach(
|
||||||
const ctx = await value.catch((error) => {
|
entries,
|
||||||
Log.Default.warn("instance dispose failed", { key, error })
|
(item) =>
|
||||||
return undefined
|
Effect.gen(function* () {
|
||||||
})
|
const exit = yield* Deferred.await(item[1].deferred).pipe(Effect.exit)
|
||||||
|
if (Exit.isFailure(exit)) {
|
||||||
if (!ctx) {
|
yield* Effect.logWarning("instance dispose failed", { key: item[0], cause: exit.cause })
|
||||||
if (cache.get(key) === value) cache.delete(key)
|
yield* removeEntry(item[0], item[1])
|
||||||
continue
|
return
|
||||||
|
}
|
||||||
|
yield* disposeEntry(item[0], item[1], exit.value)
|
||||||
|
}),
|
||||||
|
{ discard: true },
|
||||||
|
)
|
||||||
|
}).pipe(Effect.exit)
|
||||||
|
yield* Deferred.done(done, exit).pipe(Effect.asVoid)
|
||||||
|
if (disposal.all === done) {
|
||||||
|
disposal.all = undefined
|
||||||
}
|
}
|
||||||
|
return yield* restore(Deferred.await(done))
|
||||||
if (cache.get(key) !== value) continue
|
}),
|
||||||
await Effect.runPromise(dispose(ctx))
|
)
|
||||||
}
|
|
||||||
}).finally(() => {
|
|
||||||
disposal.all = undefined
|
|
||||||
})
|
|
||||||
|
|
||||||
return yield* Effect.promise(() => disposal.all!)
|
|
||||||
})
|
})
|
||||||
|
|
||||||
yield* Effect.addFinalizer(() => disposeAll().pipe(Effect.ignore))
|
yield* Effect.addFinalizer(() => disposeAll().pipe(Effect.ignore))
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
import { afterEach, describe, expect } from "bun:test"
|
import { afterEach, describe, expect } from "bun:test"
|
||||||
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
|
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
|
||||||
import { Effect, Layer } from "effect"
|
import { Effect, Fiber, Layer } from "effect"
|
||||||
|
import { registerDisposer } from "../../src/effect/instance-registry"
|
||||||
import { Instance } from "../../src/project/instance"
|
import { Instance } from "../../src/project/instance"
|
||||||
import { InstanceStore } from "../../src/project/instance-store"
|
import { InstanceStore } from "../../src/project/instance-store"
|
||||||
import { tmpdirScoped } from "../fixture/fixture"
|
import { tmpdirScoped } from "../fixture/fixture"
|
||||||
|
|
@ -67,6 +68,153 @@ describe("InstanceStore", () => {
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
it.live("dedupes concurrent loads while init is in flight", () =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const dir = yield* tmpdirScoped({ git: true })
|
||||||
|
const store = yield* InstanceStore.Service
|
||||||
|
const started = Promise.withResolvers<void>()
|
||||||
|
const release = Promise.withResolvers<void>()
|
||||||
|
let initialized = 0
|
||||||
|
|
||||||
|
const first = yield* store
|
||||||
|
.load({
|
||||||
|
directory: dir,
|
||||||
|
init: async () => {
|
||||||
|
initialized++
|
||||||
|
started.resolve()
|
||||||
|
await release.promise
|
||||||
|
},
|
||||||
|
})
|
||||||
|
.pipe(Effect.forkScoped)
|
||||||
|
|
||||||
|
yield* Effect.promise(() => started.promise)
|
||||||
|
|
||||||
|
const second = yield* store
|
||||||
|
.load({
|
||||||
|
directory: dir,
|
||||||
|
init: async () => {
|
||||||
|
initialized++
|
||||||
|
},
|
||||||
|
})
|
||||||
|
.pipe(Effect.forkScoped)
|
||||||
|
|
||||||
|
expect(initialized).toBe(1)
|
||||||
|
release.resolve()
|
||||||
|
|
||||||
|
const [firstCtx, secondCtx] = yield* Effect.all([Fiber.join(first), Fiber.join(second)])
|
||||||
|
expect(secondCtx).toBe(firstCtx)
|
||||||
|
expect(initialized).toBe(1)
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
it.live("removes failed loads from the cache", () =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const dir = yield* tmpdirScoped({ git: true })
|
||||||
|
const store = yield* InstanceStore.Service
|
||||||
|
let attempts = 0
|
||||||
|
|
||||||
|
const failed = yield* store
|
||||||
|
.load({
|
||||||
|
directory: dir,
|
||||||
|
init: async () => {
|
||||||
|
attempts++
|
||||||
|
throw new Error("init failed")
|
||||||
|
},
|
||||||
|
})
|
||||||
|
.pipe(
|
||||||
|
Effect.as(false),
|
||||||
|
Effect.catchCause(() => Effect.succeed(true)),
|
||||||
|
)
|
||||||
|
|
||||||
|
expect(failed).toBe(true)
|
||||||
|
|
||||||
|
const ctx = yield* store.load({
|
||||||
|
directory: dir,
|
||||||
|
init: async () => {
|
||||||
|
attempts++
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
|
expect(ctx.directory).toBe(dir)
|
||||||
|
expect(attempts).toBe(2)
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
it.live("reload replaces the cached context", () =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const dir = yield* tmpdirScoped({ git: true })
|
||||||
|
const store = yield* InstanceStore.Service
|
||||||
|
|
||||||
|
const first = yield* store.load({ directory: dir })
|
||||||
|
const second = yield* store.reload({ directory: dir })
|
||||||
|
const cached = yield* store.load({ directory: dir })
|
||||||
|
|
||||||
|
expect(second).not.toBe(first)
|
||||||
|
expect(cached).toBe(second)
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
it.live("stale dispose does not delete an in-flight reload", () =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const dir = yield* tmpdirScoped({ git: true })
|
||||||
|
const store = yield* InstanceStore.Service
|
||||||
|
const reloading = Promise.withResolvers<void>()
|
||||||
|
const releaseReload = Promise.withResolvers<void>()
|
||||||
|
const disposed: Array<string> = []
|
||||||
|
const off = registerDisposer(async (directory) => {
|
||||||
|
disposed.push(directory)
|
||||||
|
})
|
||||||
|
yield* Effect.addFinalizer(() => Effect.sync(off))
|
||||||
|
|
||||||
|
const first = yield* store.load({ directory: dir })
|
||||||
|
const reload = yield* store
|
||||||
|
.reload({
|
||||||
|
directory: dir,
|
||||||
|
init: async () => {
|
||||||
|
reloading.resolve()
|
||||||
|
await releaseReload.promise
|
||||||
|
},
|
||||||
|
})
|
||||||
|
.pipe(Effect.forkScoped)
|
||||||
|
|
||||||
|
yield* Effect.promise(() => reloading.promise)
|
||||||
|
const staleDispose = yield* store.dispose(first).pipe(Effect.forkScoped)
|
||||||
|
releaseReload.resolve()
|
||||||
|
|
||||||
|
const second = yield* Fiber.join(reload)
|
||||||
|
yield* Fiber.join(staleDispose)
|
||||||
|
|
||||||
|
expect(disposed).toEqual([dir])
|
||||||
|
expect(yield* store.load({ directory: dir })).toBe(second)
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
it.live("dedupes concurrent disposeAll calls", () =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const dir = yield* tmpdirScoped({ git: true })
|
||||||
|
const store = yield* InstanceStore.Service
|
||||||
|
const disposing = Promise.withResolvers<void>()
|
||||||
|
const releaseDispose = Promise.withResolvers<void>()
|
||||||
|
const disposed: Array<string> = []
|
||||||
|
const off = registerDisposer(async (directory) => {
|
||||||
|
disposed.push(directory)
|
||||||
|
disposing.resolve()
|
||||||
|
await releaseDispose.promise
|
||||||
|
})
|
||||||
|
yield* Effect.addFinalizer(() => Effect.sync(off))
|
||||||
|
|
||||||
|
yield* store.load({ directory: dir })
|
||||||
|
const first = yield* store.disposeAll().pipe(Effect.forkScoped)
|
||||||
|
yield* Effect.promise(() => disposing.promise)
|
||||||
|
const second = yield* store.disposeAll().pipe(Effect.forkScoped)
|
||||||
|
|
||||||
|
expect(disposed).toEqual([dir])
|
||||||
|
releaseDispose.resolve()
|
||||||
|
yield* Effect.all([Fiber.join(first), Fiber.join(second)])
|
||||||
|
expect(disposed).toEqual([dir])
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
it.live("keeps Instance.provide as the legacy ALS wrapper", () =>
|
it.live("keeps Instance.provide as the legacy ALS wrapper", () =>
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const dir = yield* tmpdirScoped({ git: true })
|
const dir = yield* tmpdirScoped({ git: true })
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue