diff --git a/apps/server/src/git/GitManager.ts b/apps/server/src/git/GitManager.ts index 5bcc110ba42b..3ffec1cd56c9 100644 --- a/apps/server/src/git/GitManager.ts +++ b/apps/server/src/git/GitManager.ts @@ -1,4 +1,3 @@ -import * as LookupResultCache from "./LookupResultCache.ts"; import * as Arr from "effect/Array"; import * as Cache from "effect/Cache"; import * as Context from "effect/Context"; @@ -66,6 +65,7 @@ import { import * as ProjectSetupScriptRunner from "../project/ProjectSetupScriptRunner.ts"; import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; import { extractBranchNameFromRemoteRef } from "./remoteRefs.ts"; +import { detachStackFrame } from "./detachStackFrame.ts"; import * as ServerSettings from "../serverSettings.ts"; import type { GitManagerServiceError } from "@t3tools/contracts"; import * as GitVcsDriver from "../vcs/GitVcsDriver.ts"; @@ -1093,7 +1093,7 @@ export const make = Effect.gen(function* () { prLookupFailureStreakByKey.set(key, streak); return prLookupFailureTtl(streak); }; - const prLookupCache = yield* LookupResultCache.make( + const prLookupCache = yield* Cache.makeWith( (key: string) => { const [ cwd = "", @@ -1141,6 +1141,7 @@ export const make = Effect.gen(function* () { }, }, ); + const getPrLookup = (key: string) => detachStackFrame(Cache.get(prLookupCache, key)); // A transient lookup failure (rate limit, network blip) must not clear an // already-known PR badge, so the last successful answer per branch sticks // around as the fallback. Keep the resolved head context with it so a @@ -1214,14 +1215,14 @@ export const make = Effect.gen(function* () { const branchKey = `${cwd}\u0000${details.branch}`; const cacheKey = prLookupCacheKey(cwd, details); if (refreshMissingPullRequest) { - const cached = yield* prLookupCache - .getOption(cacheKey) - .pipe(Effect.orElseSucceed(() => Option.none())); + const cached = yield* Cache.getOption(prLookupCache, cacheKey).pipe( + Effect.orElseSucceed(() => Option.none()), + ); if (Option.isSome(cached) && cached.value.latest === null) { - yield* prLookupCache.invalidate(cacheKey); + yield* Cache.invalidate(prLookupCache, cacheKey); } } - return yield* prLookupCache.get(cacheKey).pipe( + return yield* getPrLookup(cacheKey).pipe( Effect.map(({ latest, headContext }) => { if (!latest) return { pr: null, headContext }; // On the default branch, only surface open PRs. @@ -2226,12 +2227,12 @@ export const make = Effect.gen(function* () { if (options?.refresh) { // A completed turn can create a PR or reuse a merged PR's branch. // Refresh successful answers, but keep failed lookups' retry backoff. - const cached = yield* prLookupCache - .getOption(cacheKey) - .pipe(Effect.orElseSucceed(() => Option.none())); - if (Option.isSome(cached)) yield* prLookupCache.invalidate(cacheKey); + const cached = yield* Cache.getOption(prLookupCache, cacheKey).pipe( + Effect.orElseSucceed(() => Option.none()), + ); + if (Option.isSome(cached)) yield* Cache.invalidate(prLookupCache, cacheKey); } - let cached = yield* prLookupCache.get(cacheKey); + let cached = yield* getPrLookup(cacheKey); // The cached head context may have resolved on a different remote than // the saved upstream: a branch tracking origin/main but pushed to a fork // is looked up on the fork. Verify against the remote the lookup used. @@ -2258,8 +2259,8 @@ export const make = Effect.gen(function* () { }); } if (!hasSameIdentity(cached.headContext, currentIdentity)) { - yield* prLookupCache.invalidate(cacheKey); - cached = yield* prLookupCache.get(cacheKey); + yield* Cache.invalidate(prLookupCache, cacheKey); + cached = yield* getPrLookup(cacheKey); const refreshedIdentity = yield* resolvePrLookupRepositoryIdentity( cacheCwd, branch, diff --git a/apps/server/src/git/LookupResultCache.test.ts b/apps/server/src/git/LookupResultCache.test.ts deleted file mode 100644 index d9b1c31f18b7..000000000000 --- a/apps/server/src/git/LookupResultCache.test.ts +++ /dev/null @@ -1,112 +0,0 @@ -import { expect, 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 Option from "effect/Option"; -import { TestClock } from "effect/testing"; -import * as LookupResultCache from "./LookupResultCache.ts"; - -it.effect("shares pending lookups and expires results from completion", () => - Effect.gen(function* () { - const started = yield* Deferred.make(); - const finish = yield* Deferred.make(); - let calls = 0; - const cache = yield* LookupResultCache.make( - (_key: string) => - Effect.gen(function* () { - calls++; - yield* Deferred.succeed(started, undefined); - yield* Deferred.await(finish); - return calls; - }), - { capacity: 2, timeToLive: () => "1 minute" }, - ); - const first = yield* Effect.forkChild(cache.get("a")); - yield* Deferred.await(started); - const second = yield* Effect.forkChild(cache.get("a")); - yield* TestClock.adjust("2 minutes"); - yield* Deferred.succeed(finish, undefined); - expect(yield* Fiber.join(first)).toBe(1); - expect(yield* Fiber.join(second)).toBe(1); - expect(yield* cache.get("a")).toBe(1); - yield* TestClock.adjust("1 minute"); - expect(yield* cache.getOption("a")).toEqual(Option.none()); - expect(yield* cache.get("a")).toBe(2); - }).pipe(Effect.scoped), -); - -it.effect("preserves failure TTL and invalidates failed entries", () => - Effect.gen(function* () { - let calls = 0; - const cache = yield* LookupResultCache.make( - (_key: string) => - Effect.suspend(() => { - calls++; - return Effect.fail("unavailable"); - }), - { capacity: 2, timeToLive: (exit) => (Exit.isFailure(exit) ? "10 seconds" : "1 minute") }, - ); - expect(yield* Effect.flip(cache.get("a"))).toBe("unavailable"); - expect(yield* Effect.flip(cache.getOption("a"))).toBe("unavailable"); - expect(calls).toBe(1); - yield* TestClock.adjust("10 seconds"); - yield* Effect.exit(cache.get("a")); - expect(calls).toBe(2); - yield* cache.invalidate("a"); - yield* Effect.exit(cache.get("a")); - expect(calls).toBe(3); - }).pipe(Effect.scoped), -); - -it.effect("evicts the least recently read result and cannot repopulate an invalidated lookup", () => - Effect.gen(function* () { - const started = yield* Deferred.make(); - const finish = yield* Deferred.make(); - let calls = 0; - const cache = yield* LookupResultCache.make( - (key: string) => - Effect.gen(function* () { - const result = ++calls; - if (key === "pending") { - yield* Deferred.succeed(started, undefined); - yield* Deferred.await(finish); - } - return result; - }), - { capacity: 2, timeToLive: () => "1 minute" }, - ); - yield* cache.get("a"); - yield* cache.get("b"); - yield* cache.get("a"); - yield* cache.get("c"); - expect(yield* cache.getOption("b")).toEqual(Option.none()); - const pending = yield* Effect.forkChild(cache.get("pending")); - yield* Deferred.await(started); - yield* cache.invalidate("pending"); - yield* Deferred.succeed(finish, undefined); - yield* Fiber.join(pending); - expect(yield* cache.getOption("pending")).toEqual(Option.none()); - }).pipe(Effect.scoped), -); - -it.effect("a cancelled caller does not cancel another waiter", () => - Effect.gen(function* () { - const started = yield* Deferred.make(); - const finish = yield* Deferred.make(); - const cache = yield* LookupResultCache.make( - (_key: string) => - Effect.gen(function* () { - yield* Deferred.succeed(started, undefined); - return yield* Deferred.await(finish); - }), - { capacity: 2, timeToLive: () => "1 minute" }, - ); - const first = yield* Effect.forkChild(cache.get("a")); - yield* Deferred.await(started); - yield* Fiber.interrupt(first); - const second = yield* Effect.forkChild(cache.get("a")); - yield* Deferred.succeed(finish, 42); - expect(yield* Fiber.join(second)).toBe(42); - }).pipe(Effect.scoped), -); diff --git a/apps/server/src/git/LookupResultCache.ts b/apps/server/src/git/LookupResultCache.ts deleted file mode 100644 index 9ce31e760e83..000000000000 --- a/apps/server/src/git/LookupResultCache.ts +++ /dev/null @@ -1,90 +0,0 @@ -import * as Clock from "effect/Clock"; -import * as Deferred from "effect/Deferred"; -import * as Duration from "effect/Duration"; -import * as Effect from "effect/Effect"; -import * as Exit from "effect/Exit"; -import * as Option from "effect/Option"; -import * as References from "effect/References"; - -const snapshotStack = ( - frame: References.StackFrame | undefined, -): References.StackFrame | undefined => { - if (frame === undefined) return undefined; - const stack = frame.stack(); - return { name: frame.name, stack: () => stack, parent: snapshotStack(frame.parent) }; -}; - -// Completed Effect cache fibers can retain a caller's lazy trace stack and its -// request snapshot. Store only the result; running lookups belong to the service scope. -export const make = Effect.fnUntraced(function* ( - lookup: (key: Key) => Effect.Effect, - options: { - readonly capacity: number; - readonly timeToLive: (exit: Exit.Exit, key: Key) => Duration.Input; - }, -) { - const scope = yield* Effect.scope; - const clock = yield* Clock.Clock; - const entries = new Map< - Key, - { - readonly result: Deferred.Deferred; - expiresAt: number; - } - >(); - const existing = (key: Key) => { - const entry = entries.get(key); - if (entry === undefined) return undefined; - if (entry.expiresAt <= clock.currentTimeMillisUnsafe()) { - entries.delete(key); - return undefined; - } - entries.delete(key); - entries.set(key, entry); - return entry; - }; - const get = Effect.fnUntraced(function* (key: Key) { - return yield* Effect.uninterruptibleMask((restore) => - Effect.gen(function* () { - const cached = existing(key); - if (cached !== undefined) return yield* restore(Deferred.await(cached.result)); - const entry = { result: yield* Deferred.make(), expiresAt: Infinity }; - entries.set(key, entry); - if (entries.size > options.capacity) { - const oldest = entries.keys().next(); - if (!oldest.done) entries.delete(oldest.value); - } - const stack = snapshotStack(yield* References.CurrentStackFrame); - yield* Effect.forkIn( - Effect.exit(Effect.interruptible(Effect.suspend(() => lookup(key)))).pipe( - Effect.provideService(References.CurrentStackFrame, stack), - Effect.flatMap((exit) => - Effect.gen(function* () { - entry.expiresAt = - clock.currentTimeMillisUnsafe() + - Duration.toMillis(Duration.fromInputUnsafe(options.timeToLive(exit, key))); - yield* Deferred.done(entry.result, exit); - }), - ), - ), - scope, - ); - return yield* restore(Deferred.await(entry.result)); - }), - ); - }); - return { - get, - getOption: (key: Key) => - Effect.suspend(() => { - const entry = existing(key); - return entry === undefined - ? Effect.succeed(Option.none()) - : Effect.map(Deferred.await(entry.result), Option.some); - }), - invalidate: (key: Key) => - Effect.sync(() => { - entries.delete(key); - }), - }; -}); diff --git a/apps/server/src/git/LookupResultCache.memory.test.ts b/apps/server/src/git/detachStackFrame.memory.test.ts similarity index 82% rename from apps/server/src/git/LookupResultCache.memory.test.ts rename to apps/server/src/git/detachStackFrame.memory.test.ts index 7effc2c6a29b..2edc1bd483da 100644 --- a/apps/server/src/git/LookupResultCache.memory.test.ts +++ b/apps/server/src/git/detachStackFrame.memory.test.ts @@ -5,10 +5,10 @@ import * as NodeUtil from "node:util"; import { expect, it } from "vite-plus/test"; const execFile = NodeUtil.promisify(NodeChildProcess.execFile); it.each([false, true])( - "releases the traced caller snapshot while caching a result (failure=%s)", + "a detached cache lookup releases the traced caller snapshot (failure=%s)", async (failure) => { const fixture = NodeURL.fileURLToPath( - new URL("./testing/LookupRetention.fixture.mjs", import.meta.url), + new URL("./testing/StackRetention.fixture.mjs", import.meta.url), ); const { stdout } = await execFile(process.execPath, [ "--expose-gc", diff --git a/apps/server/src/git/detachStackFrame.ts b/apps/server/src/git/detachStackFrame.ts new file mode 100644 index 000000000000..a653eeedafd1 --- /dev/null +++ b/apps/server/src/git/detachStackFrame.ts @@ -0,0 +1,20 @@ +import * as Effect from "effect/Effect"; +import * as References from "effect/References"; + +const snapshotStack = ( + frame: References.StackFrame | undefined, +): References.StackFrame | undefined => { + if (frame === undefined) return undefined; + const stack = frame.stack(); + return { name: frame.name, stack: () => stack, parent: snapshotStack(frame.parent) }; +}; + +/** + * Wrap a `Cache.get` so its lookup fiber inherits a materialized stack. The lazy + * caller frame would otherwise retain the caller's request snapshot for the TTL. + */ +export const detachStackFrame = (effect: Effect.Effect) => + Effect.gen(function* () { + const frame = snapshotStack(yield* References.CurrentStackFrame); + return yield* Effect.provideService(effect, References.CurrentStackFrame, frame); + }); diff --git a/apps/server/src/git/testing/LookupRetention.fixture.mjs b/apps/server/src/git/testing/LookupRetention.fixture.mjs deleted file mode 100644 index 9ad6092f9ce9..000000000000 --- a/apps/server/src/git/testing/LookupRetention.fixture.mjs +++ /dev/null @@ -1,41 +0,0 @@ -import { Cache, Effect, Exit, Scope } from "effect"; -import * as LookupResultCache from "../LookupResultCache.ts"; -const scope = Effect.runSync(Scope.make()); -const original = process.argv[2] === "original"; -const failure = process.argv.includes("failure"); -const lookup = () => (failure ? Effect.fail("unavailable") : Effect.succeed(1)); -const cache = await Effect.runPromise( - original - ? Cache.makeWith(lookup, { capacity: 8, timeToLive: () => "1 hour" }) - : LookupResultCache.make(lookup, { - capacity: 8, - timeToLive: () => "1 hour", - }).pipe(Effect.provideService(Scope.Scope, scope)), -); -let reference; -async function populate() { - await Effect.runPromise( - Effect.gen(function* () { - const snapshot = { values: Array.from({ length: 250000 }, (_, i) => i) }; - reference = new WeakRef(snapshot); - const request = Effect.fn("LookupRetention.request")(function* () { - yield* Effect.exit(original ? Cache.get(cache, "key") : cache.get("key")); - return snapshot.values.length; - }); - yield* request(); - }), - ); -} -await populate(); -for (let i = 0; i < 12; i++) { - await new Promise((resolve) => setImmediate(resolve)); - global.gc(); -} -const retained = reference.deref() !== undefined; -const cachedResult = await Effect.runPromise( - (original ? Cache.get(cache, "key") : cache.get("key")).pipe( - Effect.catch(() => Effect.succeed("unavailable")), - ), -); -await Effect.runPromise(Scope.close(scope, Exit.void)); -process.stdout.write(JSON.stringify({ retained, cachedResult })); diff --git a/apps/server/src/git/testing/StackRetention.fixture.mjs b/apps/server/src/git/testing/StackRetention.fixture.mjs new file mode 100644 index 000000000000..1f04994a197e --- /dev/null +++ b/apps/server/src/git/testing/StackRetention.fixture.mjs @@ -0,0 +1,31 @@ +import { Cache, Effect } from "effect"; +import { detachStackFrame } from "../detachStackFrame.ts"; +// `original` runs the bare Effect Cache to show the leak this helper prevents. +const original = process.argv.includes("original"); +const failure = process.argv.includes("failure"); +const lookup = () => (failure ? Effect.fail("unavailable") : Effect.succeed(1)); +const cache = await Effect.runPromise( + Cache.makeWith(lookup, { capacity: 8, timeToLive: () => "1 hour" }), +); +const get = () => (original ? Cache.get(cache, "key") : detachStackFrame(Cache.get(cache, "key"))); +let reference; +await Effect.runPromise( + Effect.gen(function* () { + const snapshot = { values: Array.from({ length: 250000 }, (_, i) => i) }; + reference = new WeakRef(snapshot); + const request = Effect.fn("StackRetention.request")(function* () { + yield* Effect.exit(get()); + return snapshot.values.length; + }); + yield* request(); + }), +); +for (let i = 0; i < 12; i++) { + await new Promise((resolve) => setImmediate(resolve)); + global.gc(); +} +const retained = reference.deref() !== undefined; +const cachedResult = await Effect.runPromise( + get().pipe(Effect.catch(() => Effect.succeed("unavailable"))), +); +process.stdout.write(JSON.stringify({ retained, cachedResult }));