diff --git a/packages/core/src/location-services.ts b/packages/core/src/location-services.ts index a2d0f67113d0..f44685cc7a6e 100644 --- a/packages/core/src/location-services.ts +++ b/packages/core/src/location-services.ts @@ -1,5 +1,4 @@ -import { Duration, Effect, Layer, LayerMap } from "effect" -import { existsSync } from "fs" +import { Duration, Effect, Exit, Layer, LayerMap, MutableHashMap, Option } from "effect" import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { Instance } from "./instance.js" import { Location } from "./location.js" @@ -16,20 +15,48 @@ export function buildLocationServiceMap( return Layer.effect( LocationServiceMap.Service, Effect.gen(function* () { - const inner = yield* LayerMap.make((ref: Location.Ref) => Instance.layer(ref, { replacements: bindings }), { - // Workspace-placed directories exist only inside the workspace, so a - // local stat consults the wrong filesystem. Workspace liveness is - // owned by placement; do not probe the sandbox here, which would - // provision lazily-idle workspaces. - idleTimeToLive: (ref) => - ref.workspaceID !== undefined || existsSync(ref.directory) ? Duration.infinity : Duration.zero, - }) + const owner = yield* Effect.scope + const booting = MutableHashMap.empty() + const inner: LayerMap.LayerMap = yield* LayerMap.make( + (ref: Location.Ref) => { + const build = {} + MutableHashMap.set(booting, ref, build) + return Layer.fromBuild((memoMap, scope) => + Effect.suspend(() => + Layer.buildWithMemoMap(Instance.layer(ref, { replacements: bindings }), memoMap, scope), + ).pipe( + Effect.onExit((exit) => { + const finish = Effect.suspend(() => { + // An explicitly invalidated build must not evict its replacement. + if (Option.getOrUndefined(MutableHashMap.get(booting, ref)) !== build) return Effect.void + MutableHashMap.remove(booting, ref) + // Evict once per failed build, before its result reaches borrowers. + return Exit.isFailure(exit) ? inner.invalidate(ref) : Effect.void + }) + // With no borrowers, invalidation closes the entry's scope and + // joins this lookup fiber. Let the owner finish that cleanup. + return Exit.isFailure(exit) + ? finish.pipe(Effect.forkIn(owner, { startImmediately: true }), Effect.asVoid) + : finish + }), + ), + ) + }, + // Retain healthy graphs. Boot failures, not local filesystem probes, + // decide whether a location (including workspace placement) can retry. + { idleTimeToLive: Duration.infinity }, + ) const map = { ...inner, get: (ref: Location.Ref) => inner.get(LocationServiceMap.canonical(ref)), contextEffect: (ref: Location.Ref) => inner.contextEffect(LocationServiceMap.canonical(ref)), contextEffectOption: (ref: Location.Ref) => inner.contextEffectOption(LocationServiceMap.canonical(ref)), - invalidate: (ref: Location.Ref) => inner.invalidate(LocationServiceMap.canonical(ref)), + invalidate: (ref: Location.Ref) => + Effect.suspend(() => { + const key = LocationServiceMap.canonical(ref) + MutableHashMap.remove(booting, key) + return inner.invalidate(key) + }), } // Cached instances borrow their owner instead of retaining its Layer scope. const bindings: LayerNode.Replacements = [ diff --git a/packages/core/test/location-layer.test.ts b/packages/core/test/location-layer.test.ts index 355351bb77b6..41216cba752b 100644 --- a/packages/core/test/location-layer.test.ts +++ b/packages/core/test/location-layer.test.ts @@ -3,7 +3,24 @@ import path from "path" import { describe, expect } from "bun:test" import { Config } from "@opencode-ai/schema/config" import { Money } from "@opencode-ai/schema/money" -import { DateTime, Duration, Effect, Equal, Hash, Layer, LayerMap, Option, RcMap, Schema, Stream } from "effect" +import { + Cause, + DateTime, + Deferred, + Duration, + Effect, + Equal, + Exit, + Fiber, + Hash, + Layer, + LayerMap, + Option, + RcMap, + Schema, + Scope, + Stream, +} from "effect" import { TestClock } from "effect/testing" import { Agent } from "@opencode-ai/core/agent" import { Catalog } from "@opencode-ai/core/catalog" @@ -13,6 +30,7 @@ import { Global } from "@opencode-ai/util/global" import { LocationServiceMap, type LocationServices } from "@opencode-ai/core/location-services" import { LocationActivity } from "@opencode-ai/core/location-activity" import { Location } from "@opencode-ai/core/location" +import { LocationWatcher } from "@opencode-ai/core/filesystem/location-watcher" import { Plugin } from "@opencode-ai/core/plugin" import { Model } from "@opencode-ai/core/model" import { Project } from "@opencode-ai/core/project" @@ -22,7 +40,7 @@ import { Session } from "@opencode-ai/core/session" import { Workspace } from "@opencode-ai/core/workspace" import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionRunnerModel } from "@opencode-ai/core/session/runner/model" -import { tmpdir } from "./fixture/tmpdir" +import { tmpdir, tmpdirScoped } from "./fixture/tmpdir" import { tempGlobalLayer } from "./fixture/global" import { offlineModels } from "./fixture/models" import { testEffect } from "./lib/effect" @@ -61,6 +79,177 @@ const itWithActivity = testEffect( ) describe("LocationServiceMap", () => { + for (const failure of ["file", "permissions", "config reference"] as const) { + for (const invalidate of [false, true]) { + // The file-path fixture boots on Windows rather than failing during + // discovery. The config-reference case covers repair on every OS. + // Windows does not enforce POSIX directory modes, and root bypasses them. + const test = + (failure === "file" && process.platform === "win32") || + (failure === "permissions" && (process.platform === "win32" || process.getuid?.() === 0)) + ? it.live.skip + : it.live + test(`retries after repairing ${failure}${invalidate ? " with explicit invalidation" : ""}`, () => + Effect.gen(function* () { + const dir = yield* tmpdirScoped() + const directory = path.join(dir.path, "repaired") + const ref = Location.Ref.make({ directory: AbsolutePath.make(directory) }) + const locations = yield* LocationServiceMap.Service + const load = Location.Service.pipe(Effect.provide(locations.get(ref)), Effect.scoped) + + if (failure === "file") yield* Effect.promise(() => fs.writeFile(directory, "file")) + if (failure === "permissions") { + yield* Effect.promise(() => fs.mkdir(directory, { mode: 0o000 })) + yield* Effect.addFinalizer(() => Effect.promise(() => fs.chmod(directory, 0o755))) + } + if (failure === "config reference") { + yield* Effect.promise(() => fs.mkdir(directory)) + yield* Effect.promise(() => + fs.writeFile(path.join(directory, "opencode.json"), JSON.stringify({ username: "{file:username.txt}" })), + ) + } + const first = yield* Effect.exit(load) + expect(Exit.isFailure(first)).toBe(true) + if (failure === "config reference" && Exit.isFailure(first)) { + expect(Cause.squash(first.cause)).toMatchObject({ + name: "ConfigInvalidError", + data: { message: expect.stringContaining('bad file reference: "{file:username.txt}"') }, + }) + } + if (!invalidate) expect(yield* locations.contextEffectOption(ref).pipe(Effect.scoped)).toEqual(Option.none()) + + if (failure === "file") { + yield* Effect.promise(() => fs.rm(directory)) + yield* Effect.promise(() => fs.mkdir(directory)) + } + if (failure === "permissions") yield* Effect.promise(() => fs.chmod(directory, 0o755)) + if (failure === "config reference") { + yield* Effect.promise(() => fs.writeFile(path.join(directory, "username.txt"), "test-user")) + } + expect((yield* Effect.promise(() => fs.stat(directory))).isDirectory()).toBe(true) + if (invalidate) yield* locations.invalidate(ref) + const repaired = yield* Effect.exit(load) + expect(Exit.isSuccess(repaired)).toBe(true) + // A successful graph remains cached after its last borrower releases. + expect(yield* load).toBe(yield* repaired) + })) + } + } + + for (const failure of ["file", "missing", "config reference"] as const) { + // A file-path Location boots on Windows; use the missing config reference there. + const test = failure === "file" && process.platform === "win32" ? it.live.skip : it.live + test(`keeps the repaired graph after concurrent ${failure} failures release`, () => + Effect.gen(function* () { + const dir = yield* tmpdirScoped() + const directory = path.join(dir.path, "concurrent") + const ref = Location.Ref.make({ directory: AbsolutePath.make(directory) }) + const locations = yield* LocationServiceMap.Service + if (failure === "file") yield* Effect.promise(() => fs.writeFile(directory, "file")) + if (failure === "config reference") { + yield* Effect.promise(() => fs.mkdir(directory)) + yield* Effect.promise(() => + fs.writeFile(path.join(directory, "opencode.json"), JSON.stringify({ username: "{file:username.txt}" })), + ) + } + + const scopes = yield* Effect.forEach(Array.from({ length: 8 }), () => + Effect.acquireRelease(Scope.make(), (scope) => Scope.close(scope, Exit.void)), + ) + const failures = yield* Effect.forEach( + scopes, + (scope) => locations.contextEffect(ref).pipe(Scope.provide(scope), Effect.exit), + { concurrency: "unbounded" }, + ) + expect(failures.every(Exit.isFailure)).toBe(true) + if (failure === "file") yield* Effect.promise(() => fs.rm(directory)) + if (failure !== "config reference") yield* Effect.promise(() => fs.mkdir(directory)) + if (failure === "config reference") { + yield* Effect.promise(() => fs.writeFile(path.join(directory, "username.txt"), "test-user")) + } + const repaired = yield* locations.contextEffect(ref) + + yield* Effect.forEach(scopes, (scope) => Scope.close(scope, Exit.void)) + expect(yield* locations.contextEffect(ref)).toBe(repaired) + expect(Option.getOrThrow(yield* locations.contextEffectOption(ref))).toBe(repaired) + })) + } + + for (const disposition of ["retry", "invalidate", "interrupt"] as const) { + testEffect(Layer.empty).live( + disposition === "invalidate" + ? "does not let an invalidated boot failure evict its replacement" + : disposition === "interrupt" + ? "finishes a failed boot after its acquisition scopes close" + : "shares a failed boot across acquisition APIs and retries", + () => + Effect.gen(function* () { + const entered = yield* Deferred.make() + const release = yield* Deferred.make() + const builds = { started: 0 } + const finalized: number[] = [] + const layer = AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, LocationServiceMap.node]), [ + Global.node.replace(tempGlobalLayer), + offlineModels, + LocationWatcher.node.replace( + LocationWatcher.node.mapLayer((layer) => + layer.pipe( + Layer.tap(() => + Effect.gen(function* () { + const build = ++builds.started + yield* Effect.addFinalizer(() => Effect.sync(() => finalized.push(build))) + if (build !== 1) return + yield* Deferred.succeed(entered, undefined) + yield* Deferred.await(release) + yield* Effect.die("first boot failed") + }), + ), + ), + ), + ), + ]) + yield* Effect.gen(function* () { + const dir = yield* tmpdirScoped() + const ref = Location.Ref.make({ directory: AbsolutePath.make(dir.path) }) + const locations = yield* LocationServiceMap.Service + const first = yield* locations.contextEffect(ref).pipe(Effect.scoped, Effect.exit, Effect.forkScoped) + yield* Effect.addFinalizer(() => Deferred.succeed(release, undefined)) + yield* Deferred.await(entered) + const scope = yield* Effect.acquireRelease(Scope.make(), (scope) => Scope.close(scope, Exit.void)) + const second = yield* Location.Service.pipe( + Effect.provide(locations.get(ref)), + Effect.exit, + Effect.forkScoped({ startImmediately: true }), + ) + const third = yield* locations + .contextEffectOption(ref) + .pipe(Scope.provide(scope), Effect.exit, Effect.forkScoped({ startImmediately: true })) + if (disposition === "invalidate") yield* locations.invalidate(ref) + if (disposition === "interrupt") { + yield* Fiber.interrupt(first) + yield* Fiber.interrupt(second) + yield* Scope.close(scope, Exit.void) + } + const replacement = disposition === "invalidate" ? yield* locations.contextEffect(ref) : undefined + yield* Deferred.succeed(release, undefined) + if (disposition !== "interrupt") { + expect(Exit.isFailure(yield* Fiber.join(first))).toBe(true) + expect(Exit.isFailure(yield* Fiber.join(second))).toBe(true) + } + expect(Exit.isFailure(yield* Fiber.join(third).pipe(Effect.timeout("2 seconds")))).toBe(true) + expect(finalized).toEqual([1]) + expect(builds.started).toBe(disposition === "invalidate" ? 2 : 1) + const recovered = yield* locations.contextEffect(ref) + if (replacement) expect(recovered).toBe(replacement) + yield* Scope.close(scope, Exit.void) + expect(yield* locations.contextEffect(ref)).toBe(recovered) + expect(builds.started).toBe(2) + }).pipe(Effect.provide(layer)) + expect(finalized).toEqual([1, 2]) + }), + ) + } + itWithActivity.effect("does not refresh lifetime from inferred Session routing", () => Effect.gen(function* () { const locations = yield* LocationServiceMap.Service @@ -150,8 +339,8 @@ describe("LocationServiceMap", () => { expect(location.directory).toBe(directory) expect(Array.from(yield* RcMap.keys(locations.rcMap))).toEqual([workspaceRef]) - // A local ref with the same missing directory keeps the existing - // behavior: dropped as soon as it goes idle so a retry can rebuild it. + // A local ref with the same missing directory is dropped after its + // boot failure so a retry can rebuild it. yield* Location.Service.pipe(Effect.provide(locations.get(localRef)), Effect.scoped, Effect.exit) expect(Array.from(yield* RcMap.keys(locations.rcMap))).toEqual([workspaceRef]) }),