From e47951844ad27153a093ddf818ee89285e2551bf Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 23 Sep 2026 13:24:07 -0700 Subject: [PATCH 1/3] fix(server): release caller snapshots from cached PR lookups --- apps/server/src/git/GitManager.ts | 27 +++-- .../src/git/LookupResultCache.memory.test.ts | 23 ++++ apps/server/src/git/LookupResultCache.test.ts | 112 ++++++++++++++++++ apps/server/src/git/LookupResultCache.ts | 90 ++++++++++++++ .../git/testing/LookupRetention.fixture.mjs | 41 +++++++ 5 files changed, 280 insertions(+), 13 deletions(-) create mode 100644 apps/server/src/git/LookupResultCache.memory.test.ts create mode 100644 apps/server/src/git/LookupResultCache.test.ts create mode 100644 apps/server/src/git/LookupResultCache.ts create mode 100644 apps/server/src/git/testing/LookupRetention.fixture.mjs diff --git a/apps/server/src/git/GitManager.ts b/apps/server/src/git/GitManager.ts index d1513839ea6b..5bcc110ba42b 100644 --- a/apps/server/src/git/GitManager.ts +++ b/apps/server/src/git/GitManager.ts @@ -1,3 +1,4 @@ +import * as LookupResultCache from "./LookupResultCache.ts"; import * as Arr from "effect/Array"; import * as Cache from "effect/Cache"; import * as Context from "effect/Context"; @@ -1092,7 +1093,7 @@ export const make = Effect.gen(function* () { prLookupFailureStreakByKey.set(key, streak); return prLookupFailureTtl(streak); }; - const prLookupCache = yield* Cache.makeWith( + const prLookupCache = yield* LookupResultCache.make( (key: string) => { const [ cwd = "", @@ -1213,14 +1214,14 @@ export const make = Effect.gen(function* () { const branchKey = `${cwd}\u0000${details.branch}`; const cacheKey = prLookupCacheKey(cwd, details); if (refreshMissingPullRequest) { - const cached = yield* Cache.getOption(prLookupCache, cacheKey).pipe( - Effect.orElseSucceed(() => Option.none()), - ); + const cached = yield* prLookupCache + .getOption(cacheKey) + .pipe(Effect.orElseSucceed(() => Option.none())); if (Option.isSome(cached) && cached.value.latest === null) { - yield* Cache.invalidate(prLookupCache, cacheKey); + yield* prLookupCache.invalidate(cacheKey); } } - return yield* Cache.get(prLookupCache, cacheKey).pipe( + return yield* prLookupCache.get(cacheKey).pipe( Effect.map(({ latest, headContext }) => { if (!latest) return { pr: null, headContext }; // On the default branch, only surface open PRs. @@ -2225,12 +2226,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* Cache.getOption(prLookupCache, cacheKey).pipe( - Effect.orElseSucceed(() => Option.none()), - ); - if (Option.isSome(cached)) yield* Cache.invalidate(prLookupCache, cacheKey); + const cached = yield* prLookupCache + .getOption(cacheKey) + .pipe(Effect.orElseSucceed(() => Option.none())); + if (Option.isSome(cached)) yield* prLookupCache.invalidate(cacheKey); } - let cached = yield* Cache.get(prLookupCache, cacheKey); + let cached = yield* prLookupCache.get(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. @@ -2257,8 +2258,8 @@ export const make = Effect.gen(function* () { }); } if (!hasSameIdentity(cached.headContext, currentIdentity)) { - yield* Cache.invalidate(prLookupCache, cacheKey); - cached = yield* Cache.get(prLookupCache, cacheKey); + yield* prLookupCache.invalidate(cacheKey); + cached = yield* prLookupCache.get(cacheKey); const refreshedIdentity = yield* resolvePrLookupRepositoryIdentity( cacheCwd, branch, diff --git a/apps/server/src/git/LookupResultCache.memory.test.ts b/apps/server/src/git/LookupResultCache.memory.test.ts new file mode 100644 index 000000000000..7effc2c6a29b --- /dev/null +++ b/apps/server/src/git/LookupResultCache.memory.test.ts @@ -0,0 +1,23 @@ +// @effect-diagnostics nodeBuiltinImport:off - retention assertions need an isolated process with explicit GC. +import * as NodeChildProcess from "node:child_process"; +import * as NodeURL from "node:url"; +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)", + async (failure) => { + const fixture = NodeURL.fileURLToPath( + new URL("./testing/LookupRetention.fixture.mjs", import.meta.url), + ); + const { stdout } = await execFile(process.execPath, [ + "--expose-gc", + fixture, + ...(failure ? ["failure"] : []), + ]); + expect(JSON.parse(stdout)).toEqual({ + retained: false, + cachedResult: failure ? "unavailable" : 1, + }); + }, +); diff --git a/apps/server/src/git/LookupResultCache.test.ts b/apps/server/src/git/LookupResultCache.test.ts new file mode 100644 index 000000000000..d9b1c31f18b7 --- /dev/null +++ b/apps/server/src/git/LookupResultCache.test.ts @@ -0,0 +1,112 @@ +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 new file mode 100644 index 000000000000..9ce31e760e83 --- /dev/null +++ b/apps/server/src/git/LookupResultCache.ts @@ -0,0 +1,90 @@ +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/testing/LookupRetention.fixture.mjs b/apps/server/src/git/testing/LookupRetention.fixture.mjs new file mode 100644 index 000000000000..9ad6092f9ce9 --- /dev/null +++ b/apps/server/src/git/testing/LookupRetention.fixture.mjs @@ -0,0 +1,41 @@ +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 })); From 347c444ec3ecfe4f97a2fe210132e73268af0b32 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 23 Sep 2026 13:24:23 -0700 Subject: [PATCH 2/3] fix(server): defer history estimates until handoff delivery needs them --- .../ContextHandoffBudget.test.ts | 34 +++++++++++++++++++ .../ContextHandoffDelivery.ts | 13 +++---- .../ProviderTurnStartService.ts | 22 ++++++------ 3 files changed, 53 insertions(+), 16 deletions(-) diff --git a/apps/server/src/orchestration-v2/ContextHandoffBudget.test.ts b/apps/server/src/orchestration-v2/ContextHandoffBudget.test.ts index b65142a3cbf4..67fc4e35bafe 100644 --- a/apps/server/src/orchestration-v2/ContextHandoffBudget.test.ts +++ b/apps/server/src/orchestration-v2/ContextHandoffBudget.test.ts @@ -365,6 +365,40 @@ describe("handoff delivery", () => { }), ); } + it.effect("loads the history budget only when a handoff needs delivery", () => + Effect.gen(function* () { + let reads = 0; + const input = { + providerThread, + budget: Effect.sync(() => { + reads++; + return 16_000; + }), + alreadyDeliveredItemIds: new Set(), + persist: () => Effect.void, + }; + yield* deliverContextHandoffs({ ...input, handoffs: [] }); + yield* deliverContextHandoffs({ ...input, handoffs: [handoff], deferInline: true }); + yield* deliverContextHandoffs({ + ...input, + handoffs: [ + { + ...handoff, + delivery: { + nativeThreadId: providerThread.nativeThreadRef!.nativeId!, + status: "injected", + itemIds: [], + }, + }, + ], + }); + assert.equal(reads, 0); + const result = yield* deliverContextHandoffs({ ...input, handoffs: [handoff] }); + assert.equal(reads, 1); + assert.include(result.context, messages[0]!.text); + }), + ); + it.effect("persists successful injection before turn start and skips it on retry", () => Effect.gen(function* () { let durable = handoff; diff --git a/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts b/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts index c4ed32549d92..f3df022afd73 100644 --- a/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts +++ b/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts @@ -9,10 +9,10 @@ import { historyCost, renderHistory, selectHistory } from "./ContextHandoffBudge /** Persist before/after injection: an ambiguous pending delivery requires a fresh native thread. */ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextHandoffs")( - function* (input: { + function* (input: { readonly handoffs: ReadonlyArray; readonly providerThread: OrchestrationV2ProviderThread; - readonly budget: number; + readonly budget: number | Effect.Effect; readonly deferInline?: boolean; readonly alreadyDeliveredItemIds: ReadonlySet; readonly inject?: ( @@ -29,6 +29,7 @@ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextH ); if (pending.length === 0 || (input.deferInline && input.inject === undefined)) return { context: "", delivered: Effect.void }; + const budget = typeof input.budget === "number" ? input.budget : yield* input.budget; let coverage = pending .map( (handoff) => @@ -41,7 +42,7 @@ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextH // Repeated failures can accumulate many recovery markers. Keep a single // thread-level entry point when detailed coverage would crowd out history; // its activity includes the original handoff/fork source references. - if (historyCost([], coverage) > Math.min(4_000, input.budget / 2)) { + if (historyCost([], coverage) > Math.min(4_000, budget / 2)) { const strategies = Array.from(new Set(pending.map((handoff) => handoff.strategy))); coverage = `Context handoff (${strategies.join(", ")}). ${pending.length} handoff records; detailed coverage references omitted. Recover history with t3_thread_read({threadId:"${input.providerThread.appThreadId ?? pending[0]!.threadId}",view:"activity",limit:20,maxCharsPerItem:4000}); paginate with afterPosition=nextPosition. Follow fork/handoff source references in activity. For long items use itemId and textOffset=nextTextOffset until null.`; } @@ -60,16 +61,16 @@ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextH .map((handoff) => handoff.summaryText) .join("\n\n"); const fullCoverage = - oldContext && historyCost([], `${coverage}\n${oldContext}`) + 512 <= input.budget + oldContext && historyCost([], `${coverage}\n${oldContext}`) + 512 <= budget ? `${coverage}\n${oldContext}` : coverage; const selected = selectHistory({ messages, coverage: fullCoverage, omittedItems: pending.reduce((sum, handoff) => sum + (handoff.history?.omittedItems ?? 0), 0), - budget: input.budget, + budget, }); - if (historyCost(selected.messages, selected.context) > input.budget) { + if (historyCost(selected.messages, selected.context) > budget) { if (input.deferInline) return { context: "", delivered: Effect.void }; return yield* new ContextHandoffBudgetError(); } diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 9ea71cb4c9d5..9e0563620747 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -1042,16 +1042,18 @@ export const layer: Layer.Layer< handoffs: [...effectiveHandoffs, ...retryHandoff], deferInline: compact, providerThread: runningProviderThread, - budget: handoffBudget({ - tokenCap, - modelContextWindow, - userText, - attachments: message.attachments, - providerThread: budgetProviderThread, - nativeContextEstimate: - budgetProviderThread.contextUsage?.usedTokens === undefined - ? yield* nativeContextEstimate - : 0, + budget: Effect.gen(function* () { + return handoffBudget({ + tokenCap, + modelContextWindow, + userText, + attachments: message.attachments, + providerThread: budgetProviderThread, + nativeContextEstimate: + budgetProviderThread.contextUsage?.usedTokens === undefined + ? yield* nativeContextEstimate + : 0, + }); }), alreadyDeliveredItemIds: deliveredItemIds, ...(session.injectHistory === undefined From 5b09254ac5e8aa257640cf7583a288b693f3190e Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 23 Sep 2026 13:36:27 -0700 Subject: [PATCH 3/3] fix(server): infer absent handoff injector errors as never --- apps/server/src/orchestration-v2/ContextHandoffDelivery.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts b/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts index f3df022afd73..2ef3a5a5cd89 100644 --- a/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts +++ b/apps/server/src/orchestration-v2/ContextHandoffDelivery.ts @@ -9,7 +9,7 @@ import { historyCost, renderHistory, selectHistory } from "./ContextHandoffBudge /** Persist before/after injection: an ambiguous pending delivery requires a fresh native thread. */ export const deliverContextHandoffs = Effect.fn("orchestrationV2.deliverContextHandoffs")( - function* (input: { + function* (input: { readonly handoffs: ReadonlyArray; readonly providerThread: OrchestrationV2ProviderThread; readonly budget: number | Effect.Effect;