mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix: retire stale review lookups after cache invalidation (#23182)
* 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).
This commit is contained in:
@@ -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<HostedReviewInfo | null>
|
||||
/** Releases the callers and unpins the branch; idempotent. */
|
||||
expire: () => void
|
||||
}
|
||||
|
||||
const entries = new Map<string, CacheEntry>()
|
||||
const inflight = new Map<string, InflightRecord>()
|
||||
// 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
|
||||
}
|
||||
|
||||
@@ -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<HostedReviewInfo | null>()
|
||||
const freshResponse = Promise.withResolvers<HostedReviewInfo | null>()
|
||||
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<HostedReviewInfo | null>()
|
||||
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<HostedReviewInfo | null>()
|
||||
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<HostedReviewInfo | null>()
|
||||
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<HostedReviewInfo | null>()
|
||||
const freshResponse = Promise.withResolvers<HostedReviewInfo | null>()
|
||||
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<HostedReviewInfo | null>()
|
||||
const second = Promise.withResolvers<HostedReviewInfo | null>()
|
||||
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<HostedReviewInfo | null>()
|
||||
const freshResponse = Promise.withResolvers<HostedReviewInfo | null>()
|
||||
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<typeof Promise.withResolvers<HostedReviewInfo | null>>[] = []
|
||||
const readers: Promise<unknown>[] = []
|
||||
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<HostedReviewInfo | null>()
|
||||
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)
|
||||
})
|
||||
})
|
||||
@@ -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)
|
||||
}
|
||||
)
|
||||
@@ -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<HostedReviewInfo | null>
|
||||
/** Releases the callers and unpins the branch; idempotent. */
|
||||
expire: () => void
|
||||
}
|
||||
|
||||
const inflight = new Map<string, InflightRecord>()
|
||||
/**
|
||||
* 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<InflightToken, InflightRecord>()
|
||||
|
||||
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()
|
||||
}
|
||||
Reference in New Issue
Block a user