From 075fcc4b05118f6d71a5aeb8b3f192051e093afe Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 02:18:11 -0700 Subject: [PATCH 1/4] refactor(server): share one keyed lock that releases idle keys Several services kept a hand-rolled Map for per-key mutual exclusion and never removed entries, so the maps grew with every thread, cwd, ssh target, session, and update lock key ever seen. Add KeyedLock in packages/shared, an RcMap of Semaphore(1) whose entry lives only while someone holds or waits on its key. KeyedSerialExecutor now uses it, and the terminal Manager, VcsStatusBroadcaster, ssh tunnel manager, provider maintenance coordinator, CheckpointService, and OpenCode2 session gates move onto it, keeping one FIFO permit per key. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Adapters/OpenCode2AdapterV2.ts | 10 +- .../src/orchestration-v2/CheckpointService.ts | 27 +--- .../KeyedSerialExecutor.test.ts | 67 ---------- .../orchestration-v2/KeyedSerialExecutor.ts | 54 ++------ .../providerMaintenanceCommandCoordinator.ts | 33 +---- .../providerMaintenanceRunner.test.ts | 34 ++--- apps/server/src/terminal/Manager.ts | 25 +--- apps/server/src/vcs/VcsStatusBroadcaster.ts | 14 +-- packages/shared/package.json | 4 + packages/shared/src/KeyedLock.test.ts | 117 ++++++++++++++++++ packages/shared/src/KeyedLock.ts | 38 ++++++ packages/ssh/src/tunnel.ts | 11 +- 12 files changed, 208 insertions(+), 226 deletions(-) delete mode 100644 apps/server/src/orchestration-v2/KeyedSerialExecutor.test.ts create mode 100644 packages/shared/src/KeyedLock.test.ts create mode 100644 packages/shared/src/KeyedLock.ts diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCode2AdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/OpenCode2AdapterV2.ts index d1e81eb725c4..bcc9bcc646bb 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCode2AdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCode2AdapterV2.ts @@ -78,6 +78,7 @@ import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import { buildRuntimeInstructions } from "../../provider/RuntimeInstructions.ts"; import { t3OrchestrationSystemPrompt } from "../../provider/T3OrchestrationInstructions.ts"; import { SKILL_MENTION_PATTERN } from "@t3tools/shared/composerInlineTokens"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import { getModelSelectionStringOptionValue, modelSelectionsEqual } from "@t3tools/shared/model"; import { causeErrorTag } from "@t3tools/shared/observability"; @@ -887,18 +888,13 @@ export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: Pr const stagedReverts = new Set(); // Starting a turn and cutting the history take turns on a session: each // checks that the other is not running before its own requests yield. - const sessionGates = new Map(); + const sessionGates = yield* KeyedLock.make(); const exclusive = (providerThread: OrchestrationV2ProviderThread) => (effect: Effect.Effect) => { const sessionId = providerThread.nativeThreadRef?.nativeId; if (sessionId == null) return effect; - let gate = sessionGates.get(sessionId); - if (gate === undefined) { - gate = Semaphore.makeUnsafe(1); - sessionGates.set(sessionId, gate); - } - return gate.withPermit(effect); + return sessionGates.withLock(sessionId, effect); }; const emit = (event: ProviderAdapter.ProviderAdapterV2Event) => Queue.offer(events, event).pipe(Effect.asVoid); diff --git a/apps/server/src/orchestration-v2/CheckpointService.ts b/apps/server/src/orchestration-v2/CheckpointService.ts index 1471de750357..1338f68328d1 100644 --- a/apps/server/src/orchestration-v2/CheckpointService.ts +++ b/apps/server/src/orchestration-v2/CheckpointService.ts @@ -14,10 +14,9 @@ import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Encoding from "effect/Encoding"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import * as Layer from "effect/Layer"; -import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; -import * as Semaphore from "effect/Semaphore"; import { parseTurnDiffFilesFromNumstat } from "../checkpointing/Diffs.ts"; import * as CheckpointStore from "../checkpointing/CheckpointStore.ts"; @@ -249,29 +248,9 @@ export const layer: Layer.Layer< Effect.gen(function* () { const checkpointStore = yield* CheckpointStore.CheckpointStore; const idAllocator = yield* IdAllocator.IdAllocatorV2; - const workspaceSemaphores = yield* Ref.make(new Map()); - - const getWorkspaceSemaphore = (cwd: string) => - Effect.gen(function* () { - const existing = (yield* Ref.get(workspaceSemaphores)).get(cwd); - if (existing !== undefined) { - return existing; - } - - const created = yield* Semaphore.make(1); - return yield* Ref.modify(workspaceSemaphores, (current) => { - const concurrent = current.get(cwd); - if (concurrent !== undefined) { - return [concurrent, current]; - } - const updated = new Map(current); - updated.set(cwd, created); - return [created, updated]; - }); - }); - + const workspaceLocks = yield* KeyedLock.make(); const withWorkspaceLock = (cwd: string, effect: Effect.Effect) => - Effect.flatMap(getWorkspaceSemaphore(cwd), (semaphore) => semaphore.withPermits(1)(effect)); + workspaceLocks.withLock(cwd, effect); const isGitCheckpointable = (cwd: string) => checkpointStore.isGitRepository(cwd).pipe(Effect.orElseSucceed(() => false)); diff --git a/apps/server/src/orchestration-v2/KeyedSerialExecutor.test.ts b/apps/server/src/orchestration-v2/KeyedSerialExecutor.test.ts deleted file mode 100644 index 6f376b9f330e..000000000000 --- a/apps/server/src/orchestration-v2/KeyedSerialExecutor.test.ts +++ /dev/null @@ -1,67 +0,0 @@ -import { describe, expect, it } from "@effect/vitest"; - -import * as Deferred from "effect/Deferred"; -import * as Effect from "effect/Effect"; -import * as Fiber from "effect/Fiber"; -import * as Queue from "effect/Queue"; -import * as Ref from "effect/Ref"; - -import { makeKeyedSerialExecutor } from "./KeyedSerialExecutor.ts"; - -describe("makeKeyedSerialExecutor", () => { - it.effect("allows unrelated keys to run concurrently", () => - Effect.gen(function* () { - const executor = yield* makeKeyedSerialExecutor(); - const arrivals = yield* Queue.unbounded(); - const release = yield* Deferred.make(); - const run = (key: string) => - executor.withLock( - key, - Queue.offer(arrivals, key).pipe(Effect.andThen(Deferred.await(release))), - ); - - const fibers = yield* Effect.forEach(["a", "b"], run, { - concurrency: "unbounded", - discard: false, - }).pipe(Effect.forkChild); - const observed = [yield* Queue.take(arrivals), yield* Queue.take(arrivals)]; - yield* Deferred.succeed(release, undefined); - yield* Fiber.join(fibers); - expect(new Set(observed)).toEqual(new Set(["a", "b"])); - }).pipe(Effect.timeout("1 second")), - ); - - it.effect("serializes work for the same key", () => - Effect.gen(function* () { - const executor = yield* makeKeyedSerialExecutor(); - const firstEntered = yield* Deferred.make(); - const releaseFirst = yield* Deferred.make(); - const events = yield* Ref.make>([]); - - const first = yield* executor - .withLock( - "thread", - Ref.update(events, (current) => [...current, "first:start"]).pipe( - Effect.andThen(Deferred.succeed(firstEntered, undefined)), - Effect.andThen(Deferred.await(releaseFirst)), - Effect.andThen(Ref.update(events, (current) => [...current, "first:end"])), - ), - ) - .pipe(Effect.forkChild); - yield* Deferred.await(firstEntered); - const second = yield* executor - .withLock( - "thread", - Ref.update(events, (current) => [...current, "second:start"]), - ) - .pipe(Effect.forkChild); - - yield* Effect.yieldNow; - expect(yield* Ref.get(events)).toEqual(["first:start"]); - yield* Deferred.succeed(releaseFirst, undefined); - yield* Fiber.join(first); - yield* Fiber.join(second); - expect(yield* Ref.get(events)).toEqual(["first:start", "first:end", "second:start"]); - }).pipe(Effect.timeout("1 second")), - ); -}); diff --git a/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts b/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts index 93dbe04bbbfd..2c8d4ee68a23 100644 --- a/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts +++ b/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts @@ -1,11 +1,6 @@ -import * as Effect from "effect/Effect"; -import * as Ref from "effect/Ref"; -import * as Semaphore from "effect/Semaphore"; - -interface LockEntry { - readonly semaphore: Semaphore.Semaphore; - readonly users: number; -} +import * as KeyedLock from "@t3tools/shared/KeyedLock"; +import type * as Effect from "effect/Effect"; +import type * as Scope from "effect/Scope"; export interface KeyedSerialExecutor { readonly withLock: (key: Key, effect: Effect.Effect) => Effect.Effect; @@ -15,41 +10,8 @@ export interface KeyedSerialExecutor { * Serializes work that targets the same domain identity without coupling * unrelated identities to a process-wide mutex. */ -export const makeKeyedSerialExecutor = (): Effect.Effect> => - Effect.gen(function* () { - const locks = yield* Ref.make(new Map()); - - const acquireLock = (key: Key) => - Effect.gen(function* () { - const candidate = yield* Semaphore.make(1); - return yield* Ref.modify(locks, (current) => { - const existing = current.get(key); - const semaphore = existing?.semaphore ?? candidate; - const next = new Map(current); - next.set(key, { semaphore, users: (existing?.users ?? 0) + 1 }); - return [semaphore, next] as const; - }); - }); - - const releaseLock = (key: Key) => - Ref.update(locks, (current) => { - const existing = current.get(key); - if (existing === undefined) return current; - const next = new Map(current); - if (existing.users === 1) { - next.delete(key); - } else { - next.set(key, { ...existing, users: existing.users - 1 }); - } - return next; - }); - - return { - withLock: (key, effect) => - Effect.acquireUseRelease( - acquireLock(key), - (semaphore) => semaphore.withPermit(effect), - () => releaseLock(key), - ), - } satisfies KeyedSerialExecutor; - }); +export const makeKeyedSerialExecutor = (): Effect.Effect< + KeyedSerialExecutor, + never, + Scope.Scope +> => KeyedLock.make(); diff --git a/apps/server/src/provider/providerMaintenanceCommandCoordinator.ts b/apps/server/src/provider/providerMaintenanceCommandCoordinator.ts index 7c456c3c484d..f2fc03659be0 100644 --- a/apps/server/src/provider/providerMaintenanceCommandCoordinator.ts +++ b/apps/server/src/provider/providerMaintenanceCommandCoordinator.ts @@ -1,6 +1,6 @@ +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import * as Effect from "effect/Effect"; import * as Ref from "effect/Ref"; -import * as Semaphore from "effect/Semaphore"; export interface ProviderMaintenanceCommandCoordinatorShape { readonly withCommandLock: (input: { @@ -15,7 +15,7 @@ export const makeProviderMaintenanceCommandCoordinator = Effect.fn( "makeProviderMaintenanceCommandCoordinator", )(function* (input: { readonly makeAlreadyRunningError: (targetKey: string) => E }) { const runningTargetsRef = yield* Ref.make>(new Set()); - const locksRef = yield* Ref.make>(new Map()); + const locks = yield* KeyedLock.make(); const acquireTarget = Effect.fn("acquireTarget")(function* (targetKey: string) { return yield* Ref.modify(runningTargetsRef, (runningTargets) => { @@ -35,24 +35,6 @@ export const makeProviderMaintenanceCommandCoordinator = Effect.fn( return next; }); - const getLock = Effect.fn("getProviderMaintenanceCommandLock")(function* (lockKey: string) { - const existing = (yield* Ref.get(locksRef)).get(lockKey); - if (existing) { - return existing; - } - - const lock = yield* Semaphore.make(1); - return yield* Ref.modify(locksRef, (locks) => { - const current = locks.get(lockKey); - if (current) { - return [current, locks] as const; - } - const next = new Map(locks); - next.set(lockKey, lock); - return [lock, next] as const; - }); - }); - const withCommandLock: ProviderMaintenanceCommandCoordinatorShape["withCommandLock"] = ({ targetKey, lockKey, @@ -65,13 +47,10 @@ export const makeProviderMaintenanceCommandCoordinator = Effect.fn( return yield* Effect.fail(input.makeAlreadyRunningError(targetKey)); } - return yield* Effect.gen(function* () { - const lock = yield* getLock(lockKey); - if (onQueued) { - yield* onQueued; - } - return yield* lock.withPermits(1)(run); - }).pipe(Effect.ensuring(releaseTarget(targetKey))); + return yield* (onQueued ?? Effect.void).pipe( + Effect.andThen(locks.withLock(lockKey, run)), + Effect.ensuring(releaseTarget(targetKey)), + ); }); return { diff --git a/apps/server/src/provider/providerMaintenanceRunner.test.ts b/apps/server/src/provider/providerMaintenanceRunner.test.ts index 64554ce263d3..4c88e8bdf78c 100644 --- a/apps/server/src/provider/providerMaintenanceRunner.test.ts +++ b/apps/server/src/provider/providerMaintenanceRunner.test.ts @@ -7,6 +7,7 @@ import { } from "@t3tools/contracts"; import { ServerProviderUpdateError } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; +import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; @@ -223,24 +224,27 @@ const makeTestRunner = ( })), }, ) => - Effect.service(ProviderMaintenanceRunner.ProviderMaintenanceRunner).pipe( - Effect.provide( - ProviderMaintenanceRunner.layer.pipe( - Layer.provide( - Layer.mergeAll( - Layer.succeed(ProviderRegistry.ProviderRegistry, registry), - Layer.succeed(ModelManifest.ModelManifest, { - current: Effect.succeed(manifest), - refresh: Effect.succeed(manifest), - forceRefresh: Effect.succeed(manifest), - refreshInBackground: Effect.void, - }), - // Fresh per runner so a version cached by one test cannot leak into another. - Layer.sync(ProviderVersionCache, () => new Map()), - ), + // Built in the test's scope: the runner must stay open while the test uses it. + Layer.build( + ProviderMaintenanceRunner.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.succeed(ProviderRegistry.ProviderRegistry, registry), + Layer.succeed(ModelManifest.ModelManifest, { + current: Effect.succeed(manifest), + refresh: Effect.succeed(manifest), + forceRefresh: Effect.succeed(manifest), + refreshInBackground: Effect.void, + }), + // Fresh per runner so a version cached by one test cannot leak into another. + Layer.sync(ProviderVersionCache, () => new Map()), ), ), ), + ).pipe( + Effect.map((context) => + Context.get(context, ProviderMaintenanceRunner.ProviderMaintenanceRunner), + ), ); describe("providerMaintenanceRunner", () => { diff --git a/apps/server/src/terminal/Manager.ts b/apps/server/src/terminal/Manager.ts index 6122c8606f0b..8766255b7226 100644 --- a/apps/server/src/terminal/Manager.ts +++ b/apps/server/src/terminal/Manager.ts @@ -40,6 +40,7 @@ import { ProviderInstanceId, } from "@t3tools/contracts"; import { makeKeyedCoalescingWorker } from "@t3tools/shared/KeyedCoalescingWorker"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess"; import { mergePathEntries } from "@t3tools/shared/shell"; @@ -58,7 +59,6 @@ import * as Option from "effect/Option"; import * as Path from "effect/Path"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; -import * as Semaphore from "effect/Semaphore"; import * as SynchronizedRef from "effect/SynchronizedRef"; import * as ServerConfig from "../config.ts"; @@ -1569,7 +1569,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func sessions: new Map(), killFibers: new Map(), }); - const threadLocksRef = yield* SynchronizedRef.make(new Map()); + const threadLocks = yield* KeyedLock.make(); const terminalEventListeners = new Set<(event: TerminalEvent) => Effect.Effect>(); const workerScope = yield* Scope.make("sequential"); yield* Effect.addFinalizer(() => Scope.close(workerScope, Exit.void)); @@ -1598,29 +1598,10 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func f: (state: TerminalManagerState) => readonly [A, TerminalManagerState], ) => SynchronizedRef.modify(managerStateRef, f); - const getThreadSemaphore = (threadId: string) => - SynchronizedRef.modifyEffect(threadLocksRef, (current) => { - const existing: Option.Option = Option.fromNullishOr( - current.get(threadId), - ); - return Option.match(existing, { - onNone: () => - Semaphore.make(1).pipe( - Effect.map((semaphore) => { - const next = new Map(current); - next.set(threadId, semaphore); - return [semaphore, next] as const; - }), - ), - onSome: (semaphore) => Effect.succeed([semaphore, current] as const), - }); - }); - const withThreadLock = ( threadId: string, effect: Effect.Effect, - ): Effect.Effect => - Effect.flatMap(getThreadSemaphore(threadId), (semaphore) => semaphore.withPermit(effect)); + ): Effect.Effect => threadLocks.withLock(threadId, effect); const clearKillFiber = Effect.fn("terminal.clearKillFiber")(function* ( process: PtyAdapter.PtyProcess | null, diff --git a/apps/server/src/vcs/VcsStatusBroadcaster.ts b/apps/server/src/vcs/VcsStatusBroadcaster.ts index 118e2bd29807..a009bbdb63ae 100644 --- a/apps/server/src/vcs/VcsStatusBroadcaster.ts +++ b/apps/server/src/vcs/VcsStatusBroadcaster.ts @@ -10,7 +10,6 @@ import * as PubSub from "effect/PubSub"; import * as Ref from "effect/Ref"; import * as Schedule from "effect/Schedule"; import * as Scope from "effect/Scope"; -import * as Semaphore from "effect/Semaphore"; import * as Stream from "effect/Stream"; import * as SynchronizedRef from "effect/SynchronizedRef"; import type { @@ -22,6 +21,7 @@ import type { VcsStatusStreamEvent, } from "@t3tools/contracts"; import { mergeGitStatusParts } from "@t3tools/shared/git"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import * as BackgroundPolicy from "../background/BackgroundPolicy.ts"; @@ -234,15 +234,9 @@ export const make = Effect.gen(function* () { // One permit per cwd for remote reads that write the cache. Without it a // periodic poll that started before `gh pr create` can finish after the // turn-end refresh and overwrite the fresh PR with its stale `pr: null`. - const remoteWriteLocks = new Map(); - const withRemoteWriteLock = (cwd: string, effect: Effect.Effect) => { - let lock = remoteWriteLocks.get(cwd); - if (lock === undefined) { - lock = Semaphore.makeUnsafe(1); - remoteWriteLocks.set(cwd, lock); - } - return lock.withPermits(1)(effect); - }; + const remoteWriteLocks = yield* KeyedLock.make(); + const withRemoteWriteLock = (cwd: string, effect: Effect.Effect) => + remoteWriteLocks.withLock(cwd, effect); const pollersRef = yield* SynchronizedRef.make(new Map()); const getCachedStatus = Effect.fn("VcsStatusBroadcaster.getCachedStatus")(function* ( diff --git a/packages/shared/package.json b/packages/shared/package.json index b67a205e5bb3..0daaf6dd1993 100644 --- a/packages/shared/package.json +++ b/packages/shared/package.json @@ -95,6 +95,10 @@ "types": "./src/KeyedCoalescingWorker.ts", "import": "./src/KeyedCoalescingWorker.ts" }, + "./KeyedLock": { + "types": "./src/KeyedLock.ts", + "import": "./src/KeyedLock.ts" + }, "./schemaJson": { "types": "./src/schemaJson.ts", "import": "./src/schemaJson.ts" diff --git a/packages/shared/src/KeyedLock.test.ts b/packages/shared/src/KeyedLock.test.ts new file mode 100644 index 000000000000..ff28d9a4ae20 --- /dev/null +++ b/packages/shared/src/KeyedLock.test.ts @@ -0,0 +1,117 @@ +import { assert, describe, it } from "@effect/vitest"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import * as Queue from "effect/Queue"; +import * as Ref from "effect/Ref"; +import * as Scope from "effect/Scope"; + +import * as KeyedLock from "./KeyedLock.ts"; + +describe("KeyedLock", () => { + it.effect("runs one holder of a key at a time, in arrival order", () => + Effect.gen(function* () { + const lock = yield* KeyedLock.make(); + const events = yield* Ref.make>([]); + const record = (event: string) => Ref.update(events, (current) => [...current, event]); + const firstEntered = yield* Deferred.make(); + const releaseFirst = yield* Deferred.make(); + + const first = yield* lock + .withLock( + "key", + record("first:start").pipe( + Effect.andThen(Deferred.succeed(firstEntered, undefined)), + Effect.andThen(Deferred.await(releaseFirst)), + Effect.andThen(record("first:end")), + ), + ) + .pipe(Effect.forkChild); + yield* Deferred.await(firstEntered); + const second = yield* lock.withLock("key", record("second")).pipe(Effect.forkChild); + const third = yield* lock.withLock("key", record("third")).pipe(Effect.forkChild); + yield* Effect.yieldNow; + assert.deepStrictEqual(yield* Ref.get(events), ["first:start"]); + + yield* Deferred.succeed(releaseFirst, undefined); + yield* Fiber.joinAll([first, second, third]); + assert.deepStrictEqual(yield* Ref.get(events), [ + "first:start", + "first:end", + "second", + "third", + ]); + }), + ); + + it.effect("runs holders of different keys concurrently", () => + Effect.gen(function* () { + const lock = yield* KeyedLock.make(); + const entered = yield* Queue.unbounded(); + const release = yield* Deferred.make(); + const hold = (key: string) => + lock.withLock(key, Queue.offer(entered, key).pipe(Effect.andThen(Deferred.await(release)))); + + const holders = yield* Effect.forEach(["a", "b"], hold, { concurrency: "unbounded" }).pipe( + Effect.forkChild, + ); + // Both enter while the other still holds its key. + const both = new Set([yield* Queue.take(entered), yield* Queue.take(entered)]); + assert.deepStrictEqual(both, new Set(["a", "b"])); + yield* Deferred.succeed(release, undefined); + yield* Fiber.join(holders); + }), + ); + + it.effect("keeps no lock for a key once its holders and waiters are gone", () => + Effect.gen(function* () { + const lock = yield* KeyedLock.make(); + const entered = yield* Deferred.make(); + const release = yield* Deferred.make(); + + const holder = yield* lock + .withLock( + "key", + Deferred.succeed(entered, undefined).pipe(Effect.andThen(Deferred.await(release))), + ) + .pipe(Effect.forkChild); + yield* Deferred.await(entered); + // A waiter that gives up and one that runs. + const abandoned = yield* lock.withLock("key", Effect.void).pipe(Effect.forkChild); + const waiting = yield* lock.withLock("key", Effect.void).pipe(Effect.forkChild); + yield* Effect.yieldNow; + assert.deepStrictEqual(yield* lock.activeKeys, ["key"]); + + yield* Fiber.interrupt(abandoned); + yield* Deferred.succeed(release, undefined); + yield* Fiber.joinAll([holder, waiting]); + // A holder that fails releases its key too. + assert.strictEqual( + yield* lock.withLock("other", Effect.fail("boom")).pipe(Effect.flip), + "boom", + ); + + assert.deepStrictEqual(yield* lock.activeKeys, []); + }), + ); + + it.effect("runs the effect in the caller's scope", () => + Effect.gen(function* () { + const lock = yield* KeyedLock.make(); + const finalized = yield* Ref.make(false); + const callerScope = yield* Scope.make(); + + yield* lock + .withLock( + "key", + Effect.addFinalizer(() => Ref.set(finalized, true)), + ) + .pipe(Scope.provide(callerScope)); + // The finalizer belongs to the caller, so releasing the lock leaves it. + assert.isFalse(yield* Ref.get(finalized)); + yield* Scope.close(callerScope, Exit.void); + assert.isTrue(yield* Ref.get(finalized)); + }), + ); +}); diff --git a/packages/shared/src/KeyedLock.ts b/packages/shared/src/KeyedLock.ts new file mode 100644 index 000000000000..9278492304dd --- /dev/null +++ b/packages/shared/src/KeyedLock.ts @@ -0,0 +1,38 @@ +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as RcMap from "effect/RcMap"; +import * as Scope from "effect/Scope"; +import * as Semaphore from "effect/Semaphore"; + +export interface KeyedLock { + /** Runs `effect` once no other holder of `key` is running, in arrival order. */ + readonly withLock: (key: Key, effect: Effect.Effect) => Effect.Effect; + /** The keys someone holds or waits on right now. */ + readonly activeKeys: Effect.Effect>; +} + +/** + * Mutual exclusion per key. A key's lock exists only while someone holds or + * waits on it, so locks for keys that come and go never accumulate. Not + * reentrant: taking a key while already holding it deadlocks. + * + * Keys compare by `Equal` (by value for strings and numbers). The lock lives as + * long as the scope it is made in; after that scope closes, `withLock` + * interrupts. + */ +export const make = (): Effect.Effect, never, Scope.Scope> => + Effect.map(RcMap.make({ lookup: (_key: Key) => Semaphore.make(1) }), (locks) => ({ + // The handle on the key's lock lives in its own scope, so `effect` still + // runs in (and adds finalizers to) the caller's scope, not the lock's. + withLock: (key, effect) => + Effect.acquireUseRelease( + Scope.make(), + (handle) => + RcMap.get(locks, key).pipe( + Scope.provide(handle), + Effect.flatMap((semaphore) => semaphore.withPermit(effect)), + ), + (handle) => Scope.close(handle, Exit.void), + ), + activeKeys: RcMap.keys(locks).pipe(Effect.map((keys) => Array.from(keys))), + })); diff --git a/packages/ssh/src/tunnel.ts b/packages/ssh/src/tunnel.ts index 6c58cf7174eb..37d539c806b4 100644 --- a/packages/ssh/src/tunnel.ts +++ b/packages/ssh/src/tunnel.ts @@ -7,6 +7,7 @@ import { waitForHttpReady as waitForHttpReadyShared, } from "@t3tools/shared/httpReadiness"; import { cliReleaseDownloadBaseUrl } from "@t3tools/shared/cliRelease"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import * as NetService from "@t3tools/shared/Net"; import { extractJsonObject, fromLenientJson } from "@t3tools/shared/schemaJson"; import { satisfiesSemverRange } from "@t3tools/shared/semver"; @@ -18,7 +19,6 @@ import * as Layer from "effect/Layer"; import * as Path from "effect/Path"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; -import * as Semaphore from "effect/Semaphore"; import * as Stream from "effect/Stream"; import { HttpClient } from "effect/unstable/http"; import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; @@ -1328,7 +1328,7 @@ const makeSshEnvironmentManager = Effect.fn("ssh/tunnel.SshEnvironmentManager.ma ): Effect.fn.Return { const managerScope = yield* Scope.Scope; const tunnels = new Map(); - const targetLocks = new Map(); + const targetLocks = yield* KeyedLock.make(); const authSecrets = new Map(); // Keep one lock per target so reconnect cannot reuse a server while stop is pending. @@ -1336,12 +1336,7 @@ const makeSshEnvironmentManager = Effect.fn("ssh/tunnel.SshEnvironmentManager.ma key: string, effect: Effect.Effect, ): Effect.fn.Return { - let lock = targetLocks.get(key); - if (lock === undefined) { - lock = Semaphore.makeUnsafe(1); - targetLocks.set(key, lock); - } - return yield* lock.withPermits(1)(effect); + return yield* targetLocks.withLock(key, effect); }); const closeTunnelEntry = Effect.fn("ssh/tunnel.closeTunnelEntry")(function* ( From fe600b50906657f5676a60ce7229af430b1e8a20 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 20:30:49 -0700 Subject: [PATCH 2/4] fix(shared): keyed locks keep working after the scope that made them ProviderMaintenanceRunner is built per WebSocket connection inside an Effect.provide whose scope closes before any request is served. The RcMap-backed lock closed with it, so provider updates were interrupted after being marked queued and stayed queued. KeyedLock is now a refcounted Map of semaphores, as KeyedSerialExecutor was before, so it needs no scope. The maintenance runner tests go back to main's harness, which builds the runner the way ws.ts does. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../orchestration-v2/KeyedSerialExecutor.ts | 8 +-- .../providerMaintenanceRunner.test.ts | 34 +++++------ packages/shared/src/KeyedLock.test.ts | 8 +++ packages/shared/src/KeyedLock.ts | 60 ++++++++++++------- 4 files changed, 63 insertions(+), 47 deletions(-) diff --git a/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts b/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts index 2c8d4ee68a23..4096a9ef5fcb 100644 --- a/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts +++ b/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts @@ -1,6 +1,5 @@ import * as KeyedLock from "@t3tools/shared/KeyedLock"; import type * as Effect from "effect/Effect"; -import type * as Scope from "effect/Scope"; export interface KeyedSerialExecutor { readonly withLock: (key: Key, effect: Effect.Effect) => Effect.Effect; @@ -10,8 +9,5 @@ export interface KeyedSerialExecutor { * Serializes work that targets the same domain identity without coupling * unrelated identities to a process-wide mutex. */ -export const makeKeyedSerialExecutor = (): Effect.Effect< - KeyedSerialExecutor, - never, - Scope.Scope -> => KeyedLock.make(); +export const makeKeyedSerialExecutor = (): Effect.Effect> => + KeyedLock.make(); diff --git a/apps/server/src/provider/providerMaintenanceRunner.test.ts b/apps/server/src/provider/providerMaintenanceRunner.test.ts index 4c88e8bdf78c..64554ce263d3 100644 --- a/apps/server/src/provider/providerMaintenanceRunner.test.ts +++ b/apps/server/src/provider/providerMaintenanceRunner.test.ts @@ -7,7 +7,6 @@ import { } from "@t3tools/contracts"; import { ServerProviderUpdateError } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; -import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; @@ -224,27 +223,24 @@ const makeTestRunner = ( })), }, ) => - // Built in the test's scope: the runner must stay open while the test uses it. - Layer.build( - ProviderMaintenanceRunner.layer.pipe( - Layer.provide( - Layer.mergeAll( - Layer.succeed(ProviderRegistry.ProviderRegistry, registry), - Layer.succeed(ModelManifest.ModelManifest, { - current: Effect.succeed(manifest), - refresh: Effect.succeed(manifest), - forceRefresh: Effect.succeed(manifest), - refreshInBackground: Effect.void, - }), - // Fresh per runner so a version cached by one test cannot leak into another. - Layer.sync(ProviderVersionCache, () => new Map()), + Effect.service(ProviderMaintenanceRunner.ProviderMaintenanceRunner).pipe( + Effect.provide( + ProviderMaintenanceRunner.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.succeed(ProviderRegistry.ProviderRegistry, registry), + Layer.succeed(ModelManifest.ModelManifest, { + current: Effect.succeed(manifest), + refresh: Effect.succeed(manifest), + forceRefresh: Effect.succeed(manifest), + refreshInBackground: Effect.void, + }), + // Fresh per runner so a version cached by one test cannot leak into another. + Layer.sync(ProviderVersionCache, () => new Map()), + ), ), ), ), - ).pipe( - Effect.map((context) => - Context.get(context, ProviderMaintenanceRunner.ProviderMaintenanceRunner), - ), ); describe("providerMaintenanceRunner", () => { diff --git a/packages/shared/src/KeyedLock.test.ts b/packages/shared/src/KeyedLock.test.ts index ff28d9a4ae20..533c8a0be45b 100644 --- a/packages/shared/src/KeyedLock.test.ts +++ b/packages/shared/src/KeyedLock.test.ts @@ -96,6 +96,14 @@ describe("KeyedLock", () => { }), ); + it.effect("keeps working after the scope that made it closes", () => + Effect.gen(function* () { + // Services built per request hand their lock to work that outlives the build. + const lock = yield* Effect.scoped(KeyedLock.make()); + assert.strictEqual(yield* lock.withLock("key", Effect.succeed("ran")), "ran"); + }), + ); + it.effect("runs the effect in the caller's scope", () => Effect.gen(function* () { const lock = yield* KeyedLock.make(); diff --git a/packages/shared/src/KeyedLock.ts b/packages/shared/src/KeyedLock.ts index 9278492304dd..8cce5be30150 100644 --- a/packages/shared/src/KeyedLock.ts +++ b/packages/shared/src/KeyedLock.ts @@ -1,7 +1,4 @@ import * as Effect from "effect/Effect"; -import * as Exit from "effect/Exit"; -import * as RcMap from "effect/RcMap"; -import * as Scope from "effect/Scope"; import * as Semaphore from "effect/Semaphore"; export interface KeyedLock { @@ -11,28 +8,47 @@ export interface KeyedLock { readonly activeKeys: Effect.Effect>; } +interface LockEntry { + readonly semaphore: Semaphore.Semaphore; + users: number; +} + /** * Mutual exclusion per key. A key's lock exists only while someone holds or * waits on it, so locks for keys that come and go never accumulate. Not * reentrant: taking a key while already holding it deadlocks. * - * Keys compare by `Equal` (by value for strings and numbers). The lock lives as - * long as the scope it is made in; after that scope closes, `withLock` - * interrupts. + * Keys compare like `Map` keys (by value for strings and numbers). The lock is + * not tied to a scope, so it keeps working for as long as anyone references it. */ -export const make = (): Effect.Effect, never, Scope.Scope> => - Effect.map(RcMap.make({ lookup: (_key: Key) => Semaphore.make(1) }), (locks) => ({ - // The handle on the key's lock lives in its own scope, so `effect` still - // runs in (and adds finalizers to) the caller's scope, not the lock's. - withLock: (key, effect) => - Effect.acquireUseRelease( - Scope.make(), - (handle) => - RcMap.get(locks, key).pipe( - Scope.provide(handle), - Effect.flatMap((semaphore) => semaphore.withPermit(effect)), - ), - (handle) => Scope.close(handle, Exit.void), - ), - activeKeys: RcMap.keys(locks).pipe(Effect.map((keys) => Array.from(keys))), - })); +export const make = (): Effect.Effect> => + Effect.sync(() => { + const locks = new Map(); + + const acquire = (key: Key) => + Effect.sync(() => { + let entry = locks.get(key); + if (entry === undefined) { + entry = { semaphore: Semaphore.makeUnsafe(1), users: 0 }; + locks.set(key, entry); + } + entry.users += 1; + return entry; + }); + + const release = (key: Key, entry: LockEntry) => + Effect.sync(() => { + entry.users -= 1; + if (entry.users === 0) locks.delete(key); + }); + + return { + withLock: (key, effect) => + Effect.acquireUseRelease( + acquire(key), + (entry) => entry.semaphore.withPermit(effect), + (entry) => release(key, entry), + ), + activeKeys: Effect.sync(() => Array.from(locks.keys())), + }; + }); From 617478cf40697c5e0f695049d546512b20752257 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 21:16:10 -0700 Subject: [PATCH 3/4] test(shared): keyed lock keeps a key while others still need it Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/KeyedLock.test.ts | 4 +++- packages/shared/src/KeyedLock.ts | 2 +- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/packages/shared/src/KeyedLock.test.ts b/packages/shared/src/KeyedLock.test.ts index 533c8a0be45b..046729af5568 100644 --- a/packages/shared/src/KeyedLock.test.ts +++ b/packages/shared/src/KeyedLock.test.ts @@ -10,7 +10,7 @@ import * as Scope from "effect/Scope"; import * as KeyedLock from "./KeyedLock.ts"; describe("KeyedLock", () => { - it.effect("runs one holder of a key at a time, in arrival order", () => + it.effect("runs one holder of a key at a time, waiters in the order they queued", () => Effect.gen(function* () { const lock = yield* KeyedLock.make(); const events = yield* Ref.make>([]); @@ -84,6 +84,8 @@ describe("KeyedLock", () => { assert.deepStrictEqual(yield* lock.activeKeys, ["key"]); yield* Fiber.interrupt(abandoned); + // The holder and the other waiter still need the lock. + assert.deepStrictEqual(yield* lock.activeKeys, ["key"]); yield* Deferred.succeed(release, undefined); yield* Fiber.joinAll([holder, waiting]); // A holder that fails releases its key too. diff --git a/packages/shared/src/KeyedLock.ts b/packages/shared/src/KeyedLock.ts index 8cce5be30150..f9bf18676358 100644 --- a/packages/shared/src/KeyedLock.ts +++ b/packages/shared/src/KeyedLock.ts @@ -2,7 +2,7 @@ import * as Effect from "effect/Effect"; import * as Semaphore from "effect/Semaphore"; export interface KeyedLock { - /** Runs `effect` once no other holder of `key` is running, in arrival order. */ + /** Runs `effect` once no other holder of `key` is running. Queued waiters go first-in, first-out. */ readonly withLock: (key: Key, effect: Effect.Effect) => Effect.Effect; /** The keys someone holds or waits on right now. */ readonly activeKeys: Effect.Effect>; From 3cf578853f656ccf6378ab087b2c300418ffecf3 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 21:24:04 -0700 Subject: [PATCH 4/4] refactor(server): use KeyedLock directly and keep its state in a Ref Remove the KeyedSerialExecutor alias module; its consumers import @t3tools/shared/KeyedLock. KeyedLock keeps its entries in a Ref and makes semaphores with Semaphore.make instead of mutating a Map inside Effect.sync. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../orchestration-v2/KeyedSerialExecutor.ts | 13 ----- .../ProviderSessionManager.ts | 8 ++-- .../orchestration-v2/ThreadCommandExecutor.ts | 7 ++- .../legacy/LegacyV1ThreadImporter.ts | 4 +- apps/server/src/project/ProjectService.ts | 6 +-- packages/shared/src/KeyedLock.ts | 48 +++++++++++-------- 6 files changed, 40 insertions(+), 46 deletions(-) delete mode 100644 apps/server/src/orchestration-v2/KeyedSerialExecutor.ts diff --git a/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts b/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts deleted file mode 100644 index 4096a9ef5fcb..000000000000 --- a/apps/server/src/orchestration-v2/KeyedSerialExecutor.ts +++ /dev/null @@ -1,13 +0,0 @@ -import * as KeyedLock from "@t3tools/shared/KeyedLock"; -import type * as Effect from "effect/Effect"; - -export interface KeyedSerialExecutor { - readonly withLock: (key: Key, effect: Effect.Effect) => Effect.Effect; -} - -/** - * Serializes work that targets the same domain identity without coupling - * unrelated identities to a process-wide mutex. - */ -export const makeKeyedSerialExecutor = (): Effect.Effect> => - KeyedLock.make(); diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 20e240e79c58..55278b810e39 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -1,3 +1,4 @@ +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import { ModelSelection, @@ -36,7 +37,6 @@ import * as ServerSettings from "../serverSettings.ts"; import * as McpSessionRegistry from "../mcp/McpSessionRegistry.ts"; import * as EventSink from "./EventSink.ts"; import * as IdAllocator from "./IdAllocator.ts"; -import { makeKeyedSerialExecutor } from "./KeyedSerialExecutor.ts"; import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; import { ProviderAdapterEventStreamError, @@ -386,9 +386,9 @@ export const layerWithOptions = ( // cannot drop cleanup for threads only the earlier session served. const releaseRecordRetries = yield* FiberSet.make(); const nextSubscriberId = yield* Ref.make(0); - const sessionOpen = yield* makeKeyedSerialExecutor(); + const sessionOpen = yield* KeyedLock.make(); // Orders a thread's attach against a detach unloading it on the same session. - const threadAttachment = yield* makeKeyedSerialExecutor(); + const threadAttachment = yield* KeyedLock.make(); const threadAttachmentKey = (input: { readonly providerSessionId: ProviderSessionId; readonly threadId: ThreadId; @@ -427,7 +427,7 @@ export const layerWithOptions = ( }; const isMcpCredentialReserved = (threadId: ThreadId, mcpCredentialId: string) => (mcpCredentialReservations.get(mcpReservationKey(threadId, mcpCredentialId)) ?? 0) > 0; - const mcpPrepareLock = yield* makeKeyedSerialExecutor(); + const mcpPrepareLock = yield* KeyedLock.make(); /** * Resolves (or mints) the thread's MCP credential and returns it with a * reservation held; the caller must drop the reservation exactly once. diff --git a/apps/server/src/orchestration-v2/ThreadCommandExecutor.ts b/apps/server/src/orchestration-v2/ThreadCommandExecutor.ts index c9863776d44a..cec1feb122fc 100644 --- a/apps/server/src/orchestration-v2/ThreadCommandExecutor.ts +++ b/apps/server/src/orchestration-v2/ThreadCommandExecutor.ts @@ -1,13 +1,12 @@ import type { ThreadId } from "@t3tools/contracts"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import * as Context from "effect/Context"; import * as Layer from "effect/Layer"; -import { makeKeyedSerialExecutor, type KeyedSerialExecutor } from "./KeyedSerialExecutor.ts"; - /** Shared by thread commands and project deletion so both plan against current thread state. */ export class ThreadCommandExecutor extends Context.Service< ThreadCommandExecutor, - KeyedSerialExecutor + KeyedLock.KeyedLock >()("t3/orchestration-v2/ThreadCommandExecutor") {} -export const layer = Layer.effect(ThreadCommandExecutor, makeKeyedSerialExecutor()); +export const layer = Layer.effect(ThreadCommandExecutor, KeyedLock.make()); diff --git a/apps/server/src/orchestration-v2/legacy/LegacyV1ThreadImporter.ts b/apps/server/src/orchestration-v2/legacy/LegacyV1ThreadImporter.ts index 20988ac44f00..9a87bec48bbb 100644 --- a/apps/server/src/orchestration-v2/legacy/LegacyV1ThreadImporter.ts +++ b/apps/server/src/orchestration-v2/legacy/LegacyV1ThreadImporter.ts @@ -21,6 +21,7 @@ import { ThreadPullRequestLink, TurnItemId, } from "@t3tools/contracts"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -30,7 +31,6 @@ import * as Schema from "effect/Schema"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as EventSink from "../EventSink.ts"; -import { makeKeyedSerialExecutor } from "../KeyedSerialExecutor.ts"; import { randomUuidV4 } from "../RandomUuid.ts"; const IMPORT_EVENT_PREFIX = "migration:v1"; @@ -347,7 +347,7 @@ function chunks(items: ReadonlyArray, size: number): Array(); + const transcriptImports = yield* KeyedLock.make(); const listMessages = (threadId: ThreadId) => sql` diff --git a/apps/server/src/project/ProjectService.ts b/apps/server/src/project/ProjectService.ts index 5d0e6b68e626..247ff0aaa925 100644 --- a/apps/server/src/project/ProjectService.ts +++ b/apps/server/src/project/ProjectService.ts @@ -8,6 +8,7 @@ import { type ProjectSnapshot, type ThreadId, } from "@t3tools/contracts"; +import * as KeyedLock from "@t3tools/shared/KeyedLock"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -18,7 +19,6 @@ import * as Schema from "effect/Schema"; import * as EventSink from "../orchestration-v2/EventSink.ts"; import * as IdAllocator from "../orchestration-v2/IdAllocator.ts"; -import { makeKeyedSerialExecutor } from "../orchestration-v2/KeyedSerialExecutor.ts"; import * as LegacyV1ThreadImporter from "../orchestration-v2/legacy/LegacyV1ThreadImporter.ts"; import * as ProjectionStore from "../orchestration-v2/ProjectionStore.ts"; import { @@ -155,8 +155,8 @@ export const make = Effect.gen(function* () { const threadCommands = yield* ThreadCommandExecutor.ThreadCommandExecutor; // Commands for one project run in order. Commands that claim a workspace root // also hold that root, so two projects cannot both claim it. - const projectLocks = yield* makeKeyedSerialExecutor(); - const workspaceLocks = yield* makeKeyedSerialExecutor(); + const projectLocks = yield* KeyedLock.make(); + const workspaceLocks = yield* KeyedLock.make(); const toProject = ( row: ProjectStore.ProjectRow, diff --git a/packages/shared/src/KeyedLock.ts b/packages/shared/src/KeyedLock.ts index f9bf18676358..780ee3a9bbca 100644 --- a/packages/shared/src/KeyedLock.ts +++ b/packages/shared/src/KeyedLock.ts @@ -1,4 +1,5 @@ import * as Effect from "effect/Effect"; +import * as Ref from "effect/Ref"; import * as Semaphore from "effect/Semaphore"; export interface KeyedLock { @@ -10,7 +11,7 @@ export interface KeyedLock { interface LockEntry { readonly semaphore: Semaphore.Semaphore; - users: number; + readonly users: number; } /** @@ -22,33 +23,40 @@ interface LockEntry { * not tied to a scope, so it keeps working for as long as anyone references it. */ export const make = (): Effect.Effect> => - Effect.sync(() => { - const locks = new Map(); + Effect.gen(function* () { + const locks = yield* Ref.make>(new Map()); const acquire = (key: Key) => - Effect.sync(() => { - let entry = locks.get(key); - if (entry === undefined) { - entry = { semaphore: Semaphore.makeUnsafe(1), users: 0 }; - locks.set(key, entry); - } - entry.users += 1; - return entry; - }); + Effect.flatMap(Semaphore.make(1), (candidate) => + Ref.modify(locks, (current) => { + const existing = current.get(key); + const semaphore = existing?.semaphore ?? candidate; + const next = new Map(current); + next.set(key, { semaphore, users: (existing?.users ?? 0) + 1 }); + return [semaphore, next] as const; + }), + ); - const release = (key: Key, entry: LockEntry) => - Effect.sync(() => { - entry.users -= 1; - if (entry.users === 0) locks.delete(key); + const release = (key: Key) => + Ref.update(locks, (current) => { + const existing = current.get(key); + if (existing === undefined) return current; + const next = new Map(current); + if (existing.users === 1) { + next.delete(key); + } else { + next.set(key, { ...existing, users: existing.users - 1 }); + } + return next; }); return { withLock: (key, effect) => Effect.acquireUseRelease( acquire(key), - (entry) => entry.semaphore.withPermit(effect), - (entry) => release(key, entry), + (semaphore) => semaphore.withPermit(effect), + () => release(key), ), - activeKeys: Effect.sync(() => Array.from(locks.keys())), - }; + activeKeys: Effect.map(Ref.get(locks), (current) => Array.from(current.keys())), + } satisfies KeyedLock; });