From dc8f623a3a23e22e178ca604cee9c51df6bc18e2 Mon Sep 17 00:00:00 2001 From: Neil <4138956+nwparker@users.noreply.github.com> Date: Sat, 26 Sep 2026 23:06:57 -0700 Subject: [PATCH] fix: retire stale review lookups after cache invalidation (#23182) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix: retire stale review lookups after cache invalidation * test: pin the fence a retired review lookup answer has to clear The extraction widened canAdoptDetachedAnswer: a retired owner is no longer in the in-flight index, so with no replacement reader the scope generation is the only thing stopping its pre-invalidation "no review" from being adopted as current — which would short-circuit the lookup for the review Orca just opened. Cover that interleaving, and restore the invariant comments the extraction dropped (token identity, expire idempotency, and that a size-cap eviction forfeits the wall-clock sweep). --- .../hosted-review-branch-cache.ts | 83 +----- ...osted-review-inflight-invalidation.test.ts | 271 ++++++++++++++++++ .../hosted-review-inflight-lookups.test.ts | 44 +++ .../hosted-review-inflight-lookups.ts | 101 +++++++ 4 files changed, 431 insertions(+), 68 deletions(-) create mode 100644 src/main/source-control/hosted-review-inflight-invalidation.test.ts create mode 100644 src/main/source-control/hosted-review-inflight-lookups.test.ts create mode 100644 src/main/source-control/hosted-review-inflight-lookups.ts diff --git a/src/main/source-control/hosted-review-branch-cache.ts b/src/main/source-control/hosted-review-branch-cache.ts index 0e6f207ab9c..7a9f02e50c8 100644 --- a/src/main/source-control/hosted-review-branch-cache.ts +++ b/src/main/source-control/hosted-review-branch-cache.ts @@ -13,6 +13,15 @@ import { settleDetachedLookup, settleLookup } from './hosted-review-unsettled-lookups' +import { + __resetHostedReviewInflightLookupsForTests, + expireOverdueInflight, + getInflightLookup, + releaseInflight, + retireInflightWithPrefix, + trackInflight, + type InflightToken +} from './hosted-review-inflight-lookups' import { __resetHostedReviewScopeGenerationsForTests, bumpScopeGeneration, @@ -29,7 +38,6 @@ import { ACTIVE_REFRESH_INTERVAL_MS, HOSTED_REVIEW_LOOKUP_DEADLINE_MS, MAX_BRANCH_MAP_ENTRIES, - MAX_INFLIGHT_LOOKUPS, NO_REVIEW_REFRESH_INTERVAL_MS } from './hosted-review-refresh-pacing' @@ -59,22 +67,7 @@ type CacheEntry = { startedAt: number } -declare const inflightTokenBrand: unique symbol - -/** Identity token for one lookup; only ever compared by reference. */ -type InflightToken = { readonly [inflightTokenBrand]?: never } - -type InflightRecord = { - /** Identity, so a detached lookup can only ever clear its own entry. */ - token: InflightToken - startedAt: number - promise: Promise - /** Releases the callers and unpins the branch; idempotent. */ - expire: () => void -} - const entries = new Map() -const inflight = new Map() // Why: NUL is the one byte a repo path or branch name cannot contain, so a // scope prefix cannot straddle a component boundary — invalidating `/a/b` must // not also flush the unrelated repo at `/a/b c`. @@ -158,53 +151,6 @@ function storeEntry(key: string, entry: CacheEntry): void { } } -/** Clears the key's in-flight record only if it is still this lookup's. */ -function releaseInflight(key: string, token: InflightToken): boolean { - if (inflight.get(key)?.token !== token) { - return false - } - inflight.delete(key) - return true -} - -/** - * Expires records that outlived the deadline without their timer firing. Main's - * timers are suspended across a system sleep, so wall-clock age — not - * `setTimeout` alone — is what actually bounds how long a branch stays pinned. - * - * The guarantee covers tracked records only: one the size cap evicted is no - * longer reachable here and falls back to its own suspended timer. - */ -function expireOverdueInflight(now: number): void { - let overdue: InflightRecord[] | undefined - for (const record of inflight.values()) { - if (now - record.startedAt >= HOSTED_REVIEW_LOOKUP_DEADLINE_MS) { - overdue ??= [] - overdue.push(record) - } - } - // Expire after the walk: each one deletes its own entry from the map. - for (const record of overdue ?? []) { - record.expire() - } -} - -function trackInflight(key: string, record: InflightRecord): void { - inflight.set(key, record) - while (inflight.size > MAX_INFLIGHT_LOOKUPS) { - const oldest = inflight.keys().next().value - if (oldest === undefined) { - break - } - // Why: drop the record without expiring it — its own deadline still - // releases its callers, and evicting is about memory, not about failing. - // It does forfeit the sweep's wall-clock release, so the cap must stay far - // above realistic concurrency: below it, sleep-suspended timers are all an - // evicted record's callers have left. - inflight.delete(oldest) - } -} - /** * Drops every cached answer for a repo. Called when Orca itself opens a review, * so the new one is visible immediately instead of after the no-review interval. @@ -221,13 +167,14 @@ export function invalidateHostedReviewBranchCache( entries.delete(key) } } + retireInflightWithPrefix(prefix) dropFailuresWithPrefix(prefix) } /** @internal - exposed for tests only */ export function __resetHostedReviewBranchCacheForTests(): void { entries.clear() - inflight.clear() + __resetHostedReviewInflightLookupsForTests() __resetHostedReviewLookupBackoffForTests() __resetHostedReviewActiveClaimsForTests() __resetUnsettledHostedReviewLookupsForTests() @@ -244,7 +191,7 @@ export function __resetHostedReviewBranchCacheForTests(): void { * was about to give the real one. */ function canAdoptDetachedAnswer(key: string, startedAt: number): boolean { - if (inflight.has(key)) { + if (getInflightLookup(key) !== undefined) { return false } const current = entries.get(key) @@ -341,7 +288,7 @@ function startLookup( // straggler and has to prove it still outranks what is there. const stored = generation === scopeGeneration(scope) && - (inflight.get(key)?.token === token || canAdoptDetachedAnswer(key, startedAt)) + (getInflightLookup(key)?.token === token || canAdoptDetachedAnswer(key, startedAt)) if (stored) { storeEntry(key, { review, fetchedAt: Date.now(), headOid, startedAt }) } @@ -360,7 +307,7 @@ function startLookup( } // Why: a record the size cap dropped has a live successor, and backing the // branch off would slow the retry that is already running. - if (inflight.get(key)?.token === token) { + if (getInflightLookup(key)?.token === token) { noteFailure(key) } // Why: the last good review beats an error card here just as it does on @@ -417,7 +364,7 @@ export async function withHostedReviewBranchCache( return cached.review } - const pending = inflight.get(key) + const pending = getInflightLookup(key) if (pending) { return pending.promise } diff --git a/src/main/source-control/hosted-review-inflight-invalidation.test.ts b/src/main/source-control/hosted-review-inflight-invalidation.test.ts new file mode 100644 index 00000000000..6d8653ba76b --- /dev/null +++ b/src/main/source-control/hosted-review-inflight-invalidation.test.ts @@ -0,0 +1,271 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { HostedReviewInfo } from '../../shared/hosted-review' +import { + __resetHostedReviewBranchCacheForTests, + invalidateHostedReviewBranchCache, + withHostedReviewBranchCache +} from './hosted-review-branch-cache' +import { + HOSTED_REVIEW_LOOKUP_DEADLINE_MS, + MAX_UNSETTLED_LOOKUP_KEYS, + MAX_UNSETTLED_LOOKUPS_PER_KEY +} from './hosted-review-refresh-pacing' + +const identity = { + repoPath: '/repo', + executionHostId: 'ssh:host-a' as const, + branch: 'feature/review', + localGitExecOptions: { admissionTier: 'background' } +} +const options = { headOid: null } +const review: HostedReviewInfo = { + provider: 'gitlab', + number: 7, + title: 'Created review', + state: 'open', + url: 'https://git.example/team/repo/merge_requests/7', + status: 'neutral', + updatedAt: '', + mergeable: 'UNKNOWN' +} + +beforeEach(() => { + __resetHostedReviewBranchCacheForTests() + vi.useFakeTimers() + vi.setSystemTime(1_000_000) +}) + +afterEach(() => { + __resetHostedReviewBranchCacheForTests() + vi.useRealTimers() +}) + +describe('hosted review in-flight invalidation', () => { + it.each(['before', 'after'] as const)( + 'refreshes post-create readers when the old request finishes %s its replacement', + async (completionOrder) => { + const oldResponse = Promise.withResolvers() + const freshResponse = Promise.withResolvers() + const oldLookup = vi.fn(() => oldResponse.promise) + const freshLookup = vi.fn(() => freshResponse.promise) + const admitted = withHostedReviewBranchCache(identity, options, oldLookup) + expect( + await withHostedReviewBranchCache( + { ...identity, localGitExecOptions: { admissionTier: 'interactive' } }, + options, + async () => null + ) + ).toBeNull() + + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + const next = withHostedReviewBranchCache(identity, options, freshLookup) + const concurrent = withHostedReviewBranchCache(identity, options, freshLookup) + try { + expect(freshLookup).toHaveBeenCalledTimes(1) + if (completionOrder === 'before') { + oldResponse.resolve(null) + expect(await admitted).toBeNull() + } + freshResponse.resolve(review) + expect(await next).toEqual(review) + expect(await concurrent).toEqual(review) + if (completionOrder === 'after') { + oldResponse.resolve(null) + expect(await admitted).toBeNull() + } + expect(await withHostedReviewBranchCache(identity, options, freshLookup)).toEqual(review) + expect(oldLookup).toHaveBeenCalledTimes(1) + expect(freshLookup).toHaveBeenCalledTimes(1) + } finally { + oldResponse.resolve(null) + freshResponse.resolve(review) + await Promise.all([admitted, next, concurrent]) + } + } + ) + + it('does not make a post-create reader wait for a stale request deadline', async () => { + const oldResponse = Promise.withResolvers() + const admitted = withHostedReviewBranchCache(identity, options, () => oldResponse.promise) + const admittedResult = admitted.catch(() => null) + await vi.advanceTimersByTimeAsync(10_000) + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + const freshLookup = vi.fn(async () => review) + let settled = false + const next = withHostedReviewBranchCache(identity, options, freshLookup).then((value) => { + settled = true + return value + }) + try { + await vi.advanceTimersByTimeAsync(0) + expect(settled).toBe(true) + expect(await next).toEqual(review) + expect(freshLookup).toHaveBeenCalledTimes(1) + await vi.advanceTimersByTimeAsync(HOSTED_REVIEW_LOOKUP_DEADLINE_MS - 10_000) + expect(await admittedResult).toEqual(review) + expect(await withHostedReviewBranchCache(identity, options, freshLookup)).toEqual(review) + expect(freshLookup).toHaveBeenCalledTimes(2) + } finally { + oldResponse.resolve(null) + await Promise.all([admittedResult, next]) + } + }) + + it('discards a retired answer that landed with no replacement to outrank it', async () => { + const oldResponse = Promise.withResolvers() + const admitted = withHostedReviewBranchCache(identity, options, () => oldResponse.promise) + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + // Nothing replaced it, so only the scope generation stands between this + // pre-invalidation "no review" and the key it no longer owns. Adopting it + // would read as fresh and short-circuit the lookup for the review Orca just + // opened. + oldResponse.resolve(null) + expect(await admitted).toBeNull() + const freshLookup = vi.fn(async () => review) + expect(await withHostedReviewBranchCache(identity, options, freshLookup)).toEqual(review) + expect(freshLookup).toHaveBeenCalledTimes(1) + }) + + it('keeps other hosts, paths and their pending promises isolated', async () => { + const others = [ + { ...identity, executionHostId: 'local' as const }, + { ...identity, executionHostId: 'ssh:host-b' as const }, + { ...identity, repoPath: '/repo-other' } + ].map((key) => { + const response = Promise.withResolvers() + const lookup = vi.fn(() => response.promise) + return { key, response, lookup, admitted: withHostedReviewBranchCache(key, options, lookup) } + }) + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + const readers = others.map(({ key, lookup }) => + withHostedReviewBranchCache(key, options, lookup) + ) + for (const other of others) { + expect(other.lookup).toHaveBeenCalledTimes(1) + other.response.resolve(review) + } + expect(await Promise.all(readers)).toEqual([review, review, review]) + await Promise.all(others.map((other) => other.admitted)) + }) + + it('sweeps invalidated readers after sleep without expiring their replacement', async () => { + const oldResponse = Promise.withResolvers() + const freshResponse = Promise.withResolvers() + const admitted = withHostedReviewBranchCache(identity, options, () => oldResponse.promise) + const admittedResult = admitted.catch((error: unknown) => error) + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + vi.setSystemTime(Date.now() + 90_000) + const freshLookup = vi.fn(() => freshResponse.promise) + const replacement = withHostedReviewBranchCache(identity, options, freshLookup).catch( + (error: unknown) => error + ) + vi.setSystemTime(Date.now() + 40_000) + const joined = withHostedReviewBranchCache(identity, options, freshLookup).catch( + (error: unknown) => error + ) + try { + await vi.advanceTimersByTimeAsync(0) + expect(await admittedResult).toMatchObject({ + message: expect.stringContaining('timed out') + }) + expect(freshLookup).toHaveBeenCalledTimes(1) + freshResponse.resolve(review) + expect(await replacement).toEqual(review) + expect(await joined).toEqual(review) + oldResponse.resolve(null) + expect(await withHostedReviewBranchCache(identity, options, freshLookup)).toEqual(review) + expect(freshLookup).toHaveBeenCalledTimes(1) + } finally { + oldResponse.resolve(null) + freshResponse.resolve(review) + await Promise.all([admittedResult, replacement, joined]) + } + }) + + it('keeps the two-unsettled-request cap across repeated invalidation', async () => { + const first = Promise.withResolvers() + const second = Promise.withResolvers() + const firstRead = withHostedReviewBranchCache(identity, options, () => first.promise) + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + const secondRead = withHostedReviewBranchCache(identity, options, () => second.promise) + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + const thirdLookup = vi.fn(async () => review) + const third = withHostedReviewBranchCache(identity, options, thirdLookup) + const thirdResult = third.catch((error: unknown) => error) + try { + await vi.advanceTimersByTimeAsync(0) + expect(thirdLookup).not.toHaveBeenCalled() + first.resolve(null) + second.resolve(null) + await Promise.all([firstRead, secondRead]) + expect(await thirdResult).toMatchObject({ + message: expect.stringContaining('never answered') + }) + expect(await withHostedReviewBranchCache(identity, options, thirdLookup)).toEqual(review) + expect(thirdLookup).toHaveBeenCalledTimes(1) + } finally { + first.resolve(null) + second.resolve(null) + await Promise.all([firstRead, secondRead, thirdResult]) + } + }) + + it('does not let an invalidated rejection penalize the replacement', async () => { + const oldResponse = Promise.withResolvers() + const freshResponse = Promise.withResolvers() + const admitted = withHostedReviewBranchCache(identity, options, () => oldResponse.promise) + const rejected = expect(admitted).rejects.toThrow('old lookup failed') + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + const fresh = withHostedReviewBranchCache(identity, options, () => freshResponse.promise) + oldResponse.reject(new Error('old lookup failed')) + await rejected + freshResponse.resolve(review) + expect(await fresh).toEqual(review) + await vi.advanceTimersByTimeAsync(60_001) + const next = vi.fn(async () => review) + expect(await withHostedReviewBranchCache(identity, options, next)).toEqual(review) + expect(next).toHaveBeenCalledTimes(1) + }) + + it('bounds retired owners by admission and sweeps all of them after sleep', async () => { + const requests: ReturnType>[] = [] + const readers: Promise[] = [] + for (let index = 0; index < MAX_UNSETTLED_LOOKUP_KEYS; index += 1) { + for (let attempt = 0; attempt < MAX_UNSETTLED_LOOKUPS_PER_KEY; attempt += 1) { + const response = Promise.withResolvers() + requests.push(response) + readers.push( + withHostedReviewBranchCache( + { ...identity, branch: `branch-${index}` }, + options, + () => response.promise + ).catch((error: unknown) => error) + ) + invalidateHostedReviewBranchCache(identity.repoPath, identity.executionHostId) + } + } + const freshLookup = vi.fn(async () => review) + try { + await expect(withHostedReviewBranchCache(identity, options, freshLookup)).rejects.toThrow( + 'Too many hosted review lookups are already in progress' + ) + expect(freshLookup).not.toHaveBeenCalled() + expect(vi.getTimerCount()).toBe(MAX_UNSETTLED_LOOKUP_KEYS * MAX_UNSETTLED_LOOKUPS_PER_KEY) + vi.setSystemTime(Date.now() + HOSTED_REVIEW_LOOKUP_DEADLINE_MS) + await expect(withHostedReviewBranchCache(identity, options, freshLookup)).rejects.toThrow() + expect(vi.getTimerCount()).toBe(0) + expect(await Promise.all(readers)).toEqual( + requests.map(() => + expect.objectContaining({ message: expect.stringContaining('timed out') }) + ) + ) + } finally { + for (const response of requests) { + response.resolve(null) + } + await Promise.all(readers) + } + expect(await withHostedReviewBranchCache(identity, options, freshLookup)).toEqual(review) + expect(freshLookup).toHaveBeenCalledTimes(1) + }) +}) diff --git a/src/main/source-control/hosted-review-inflight-lookups.test.ts b/src/main/source-control/hosted-review-inflight-lookups.test.ts new file mode 100644 index 00000000000..149de43c097 --- /dev/null +++ b/src/main/source-control/hosted-review-inflight-lookups.test.ts @@ -0,0 +1,44 @@ +import { afterEach, expect, it, vi } from 'vitest' +import { + __resetHostedReviewInflightLookupsForTests, + expireOverdueInflight, + getInflightLookup, + releaseInflight, + retireInflightWithPrefix, + trackInflight +} from './hosted-review-inflight-lookups' +import { HOSTED_REVIEW_LOOKUP_DEADLINE_MS } from './hosted-review-refresh-pacing' + +afterEach(__resetHostedReviewInflightLookupsForTests) + +it.each(['completion', 'deadline'])( + 'releases a retired owner on %s while preserving its successor', + (outcome) => { + const oldToken = {} + const oldExpire = vi.fn(() => releaseInflight('repo\0branch', oldToken)) + trackInflight('repo\0branch', { + token: oldToken, + startedAt: 0, + promise: Promise.resolve(null), + expire: oldExpire + }) + retireInflightWithPrefix('repo\0') + const replacement = { + token: {}, + startedAt: HOSTED_REVIEW_LOOKUP_DEADLINE_MS, + promise: Promise.resolve(null), + expire: vi.fn() + } + trackInflight('repo\0branch', replacement) + if (outcome === 'completion') { + expect(releaseInflight('repo\0branch', oldToken)).toBe(false) + } + expireOverdueInflight(HOSTED_REVIEW_LOOKUP_DEADLINE_MS) + expect(oldExpire).toHaveBeenCalledTimes(outcome === 'deadline' ? 1 : 0) + expect(replacement.expire).not.toHaveBeenCalled() + expect(getInflightLookup('repo\0branch')).toBe(replacement) + expireOverdueInflight(HOSTED_REVIEW_LOOKUP_DEADLINE_MS + 1) + expect(oldExpire).toHaveBeenCalledTimes(outcome === 'deadline' ? 1 : 0) + expect(releaseInflight('repo\0branch', replacement.token)).toBe(true) + } +) diff --git a/src/main/source-control/hosted-review-inflight-lookups.ts b/src/main/source-control/hosted-review-inflight-lookups.ts new file mode 100644 index 00000000000..79d4f12ba13 --- /dev/null +++ b/src/main/source-control/hosted-review-inflight-lookups.ts @@ -0,0 +1,101 @@ +import type { HostedReviewInfo } from '../../shared/hosted-review' +import { + HOSTED_REVIEW_LOOKUP_DEADLINE_MS, + MAX_INFLIGHT_LOOKUPS +} from './hosted-review-refresh-pacing' + +declare const inflightTokenBrand: unique symbol + +/** Identity token for one lookup; only ever compared by reference. */ +export type InflightToken = { readonly [inflightTokenBrand]?: never } + +export type InflightRecord = { + /** Identity, so a detached lookup can only ever clear its own entry. */ + token: InflightToken + startedAt: number + promise: Promise + /** Releases the callers and unpins the branch; idempotent. */ + expire: () => void +} + +const inflight = new Map() +/** + * Owners an invalidation took off their key. They keep running — nothing here + * can cancel a lookup — but no new reader may join them, which is what stops a + * post-invalidation read from waiting out a stale request's deadline. + * + * Admission bounds these to two per key across at most 1,000 unsettled keys, so + * this map cannot outgrow the lookups already counted as in progress. + */ +const retired = new Map() + +export function getInflightLookup(key: string): InflightRecord | undefined { + return inflight.get(key) +} + +/** Clears only this owner's records; the return value identifies the current owner. */ +export function releaseInflight(key: string, token: InflightToken): boolean { + retired.delete(token) + if (inflight.get(key)?.token !== token) { + return false + } + inflight.delete(key) + return true +} + +/** Takes every owner under `prefix` off its key, without failing it. */ +export function retireInflightWithPrefix(prefix: string): void { + for (const [key, record] of inflight) { + if (key.startsWith(prefix)) { + retired.set(record.token, record) + inflight.delete(key) + } + } +} + +/** + * Expires records that outlived the deadline without their timer firing. Main's + * timers are suspended across a system sleep, so wall-clock age — not + * `setTimeout` alone — is what actually bounds how long a branch stays pinned. + * Retired owners are swept too: their readers are gone, but their own callers + * still need releasing. + * + * The guarantee covers tracked records only: one the size cap evicted is in + * neither map and falls back to its own suspended timer. + */ +export function expireOverdueInflight(now: number): void { + let overdue: InflightRecord[] | undefined + for (const records of [inflight, retired]) { + for (const record of records.values()) { + if (now - record.startedAt >= HOSTED_REVIEW_LOOKUP_DEADLINE_MS) { + overdue ??= [] + overdue.push(record) + } + } + } + for (const record of overdue ?? []) { + record.expire() + } +} + +export function trackInflight(key: string, record: InflightRecord): void { + inflight.set(key, record) + while (inflight.size > MAX_INFLIGHT_LOOKUPS) { + const oldest = inflight.keys().next().value + if (oldest === undefined) { + break + } + // Why: drop the record without expiring it — its own deadline still releases + // its callers, and evicting is about memory, not about failing. It does + // forfeit the sweep above, so the cap must stay far above realistic + // concurrency: below it, sleep-suspended timers are all an evicted record's + // callers have left. + inflight.delete(oldest) + } +} + +/** @internal - exposed for tests only */ +export function __resetHostedReviewInflightLookupsForTests(): void { + inflight.clear() + retired.clear() +}