mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(cache): time out a getOrSet produce that never settles (#2297)
This commit is contained in:
+53
-4
@@ -2,6 +2,8 @@
|
||||
|
||||
const CACHE_INFLIGHT_MAX_ENTRIES = 10000;
|
||||
const CACHE_INFLIGHT_JOIN_RETRIES = 1;
|
||||
const CACHE_PRODUCE_TIMEOUT_MS = 15000;
|
||||
const CACHE_PRODUCE_TIMEOUT_MESSAGE = 'Cache produce timed out';
|
||||
|
||||
interface CacheMSetEntry<T> {
|
||||
key: string;
|
||||
@@ -75,7 +77,12 @@ export abstract class ICacheService {
|
||||
return entry.hit ? entry.value : null;
|
||||
}
|
||||
|
||||
async getOrSet<T>(key: string, valueFactory: () => Promise<T>, ttlSeconds?: CacheTtlSeconds<T>): Promise<T> {
|
||||
async getOrSet<T>(
|
||||
key: string,
|
||||
valueFactory: () => Promise<T>,
|
||||
ttlSeconds?: CacheTtlSeconds<T>,
|
||||
produceTimeoutMs: number = CACHE_PRODUCE_TIMEOUT_MS,
|
||||
): Promise<T> {
|
||||
let generation = this.trackProduce(key);
|
||||
try {
|
||||
for (let attempt = 0; ; attempt++) {
|
||||
@@ -85,7 +92,7 @@ export abstract class ICacheService {
|
||||
}
|
||||
const inflight = this.inflightValues.get(key);
|
||||
if (!inflight) {
|
||||
return await this.produceSingleFlight(key, valueFactory, ttlSeconds, generation);
|
||||
return await this.produceSingleFlight(key, valueFactory, ttlSeconds, generation, produceTimeoutMs);
|
||||
}
|
||||
const joined = await this.joinInflight<T>(inflight);
|
||||
if (joined.joined) {
|
||||
@@ -115,6 +122,14 @@ export abstract class ICacheService {
|
||||
return this.produceInvalidations.get(key)?.generation ?? 0;
|
||||
}
|
||||
|
||||
private abandonProduce(key: string, generation: number): void {
|
||||
this.trackProduce(key);
|
||||
const tracked = this.produceInvalidations.get(key);
|
||||
if (tracked?.generation === generation) {
|
||||
tracked.generation += 1;
|
||||
}
|
||||
}
|
||||
|
||||
private releaseProduce(key: string): void {
|
||||
const tracked = this.produceInvalidations.get(key);
|
||||
if (!tracked) {
|
||||
@@ -139,17 +154,51 @@ export abstract class ICacheService {
|
||||
valueFactory: () => Promise<T>,
|
||||
ttlSeconds: CacheTtlSeconds<T> | undefined,
|
||||
generation: number,
|
||||
produceTimeoutMs: number,
|
||||
): Promise<T> {
|
||||
const produced = this.boundProduce(
|
||||
this.produceAndStore(key, valueFactory, ttlSeconds, generation),
|
||||
key,
|
||||
generation,
|
||||
produceTimeoutMs,
|
||||
);
|
||||
if (this.inflightValues.size >= CACHE_INFLIGHT_MAX_ENTRIES) {
|
||||
return await this.produceAndStore(key, valueFactory, ttlSeconds, generation);
|
||||
return await produced;
|
||||
}
|
||||
const pending = this.produceAndStore(key, valueFactory, ttlSeconds, generation).finally(() => {
|
||||
const pending = produced.finally(() => {
|
||||
this.inflightValues.delete(key);
|
||||
});
|
||||
this.inflightValues.set(key, pending);
|
||||
return await pending;
|
||||
}
|
||||
|
||||
private boundProduce<T>(produced: Promise<T>, key: string, generation: number, produceTimeoutMs: number): Promise<T> {
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
let abandoned = false;
|
||||
const timer = setTimeout(() => {
|
||||
abandoned = true;
|
||||
this.abandonProduce(key, generation);
|
||||
reject(new Error(CACHE_PRODUCE_TIMEOUT_MESSAGE));
|
||||
}, produceTimeoutMs);
|
||||
const settle = () => {
|
||||
clearTimeout(timer);
|
||||
if (abandoned) {
|
||||
this.releaseProduce(key);
|
||||
}
|
||||
};
|
||||
produced.then(
|
||||
(value) => {
|
||||
settle();
|
||||
resolve(value);
|
||||
},
|
||||
(error: unknown) => {
|
||||
settle();
|
||||
reject(error);
|
||||
},
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
private async produceAndStore<T>(
|
||||
key: string,
|
||||
valueFactory: () => Promise<T>,
|
||||
|
||||
@@ -4,6 +4,8 @@ import {InMemoryProvider} from '@pkgs/cache/src/providers/InMemoryProvider';
|
||||
import {describe, expect, it} from 'vitest';
|
||||
|
||||
const INFLIGHT_OVERFLOW_ENTRIES = 10000;
|
||||
const PRODUCE_TIMEOUT_MS = 50;
|
||||
const PRODUCE_TIMEOUT_MESSAGE = 'Cache produce timed out';
|
||||
|
||||
function deferred<T>(): {promise: Promise<T>; resolve: (value: T) => void; reject: (error: Error) => void} {
|
||||
let resolve!: (value: T) => void;
|
||||
@@ -19,6 +21,16 @@ function flush(): Promise<void> {
|
||||
return new Promise((resolve) => setTimeout(resolve, 0));
|
||||
}
|
||||
|
||||
function settleWithin<T>(pending: Promise<T>, ms: number): Promise<T | 'pinned' | 'rejected'> {
|
||||
return Promise.race([
|
||||
pending.then(
|
||||
(value) => value,
|
||||
() => 'rejected' as const,
|
||||
),
|
||||
new Promise<'pinned'>((resolve) => setTimeout(() => resolve('pinned'), ms)),
|
||||
]);
|
||||
}
|
||||
|
||||
function trackedProduceKeys(cache: InMemoryProvider): Array<string> {
|
||||
return [...(cache as unknown as {produceInvalidations: Map<string, unknown>}).produceInvalidations.keys()];
|
||||
}
|
||||
@@ -150,6 +162,51 @@ describe('cache invalidation during an in-flight produce', () => {
|
||||
expect(trackedProduceKeys(cache)).toEqual([]);
|
||||
});
|
||||
|
||||
it('does not pin a key forever when the factory never settles', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const stuck = deferred<string>();
|
||||
const pinned = cache.getOrSet('session', async () => await stuck.promise, 30, PRODUCE_TIMEOUT_MS);
|
||||
await expect(settleWithin(pinned, 500)).resolves.toBe('rejected');
|
||||
await expect(pinned).rejects.toThrow(PRODUCE_TIMEOUT_MESSAGE);
|
||||
const recovered = cache.getOrSet('session', async () => 'recovered', 30, PRODUCE_TIMEOUT_MS);
|
||||
await expect(settleWithin(recovered, 500)).resolves.toBe('recovered');
|
||||
stuck.resolve('never-settled');
|
||||
await flush();
|
||||
expect(await cache.get('session')).toBe('recovered');
|
||||
});
|
||||
|
||||
it('does not store a value produced by a factory that settled after the timeout', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const stuck = deferred<string>();
|
||||
const pending = cache.getOrSet('session', async () => await stuck.promise, 30, PRODUCE_TIMEOUT_MS);
|
||||
await expect(pending).rejects.toThrow(PRODUCE_TIMEOUT_MESSAGE);
|
||||
stuck.resolve('late-produce');
|
||||
await flush();
|
||||
expect(await cache.get('session')).toBeNull();
|
||||
expect(trackedProduceKeys(cache)).toEqual([]);
|
||||
});
|
||||
|
||||
it('retries once for the joiners when the producer times out', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const gates: Array<ReturnType<typeof deferred<string>>> = [];
|
||||
const factory = async () => {
|
||||
const gate = deferred<string>();
|
||||
gates.push(gate);
|
||||
return await gate.promise;
|
||||
};
|
||||
const producer = cache.getOrSet('session', factory, 30, PRODUCE_TIMEOUT_MS);
|
||||
const joiner = cache.getOrSet('session', factory, 30, PRODUCE_TIMEOUT_MS);
|
||||
await expect(producer).rejects.toThrow(PRODUCE_TIMEOUT_MESSAGE);
|
||||
await flush();
|
||||
expect(gates).toHaveLength(2);
|
||||
gates[1].resolve('retried-session');
|
||||
await expect(joiner).resolves.toBe('retried-session');
|
||||
expect(await cache.get('session')).toBe('retried-session');
|
||||
gates[0].resolve('abandoned-produce');
|
||||
await flush();
|
||||
expect(await cache.get('session')).toBe('retried-session');
|
||||
});
|
||||
|
||||
it('keeps a later produce cacheable after an earlier one was invalidated', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const first = deferred<string>();
|
||||
|
||||
@@ -1055,6 +1055,7 @@ export class StripeCheckoutService {
|
||||
private static readonly LOCALIZED_CARD_PREAPPROVAL_CONTINUE_LOCK_TTL_SECONDS = seconds('30 seconds');
|
||||
private static readonly LOCALIZED_CARD_PREAPPROVAL_TTL_SECONDS = seconds('1 day');
|
||||
private static readonly PRICE_CACHE_TTL_SECONDS = seconds('1 hour');
|
||||
private static readonly PRICE_CACHE_PRODUCE_TIMEOUT_MS = 90000;
|
||||
|
||||
private resolveConfiguredPriceIds(countryCode?: string, pricingMode: PricingMode = 'localized'): ResolvedPriceIds {
|
||||
const recurringCurrencyPreferences =
|
||||
@@ -1247,6 +1248,7 @@ export class StripeCheckoutService {
|
||||
};
|
||||
},
|
||||
StripeCheckoutService.PRICE_CACHE_TTL_SECONDS,
|
||||
StripeCheckoutService.PRICE_CACHE_PRODUCE_TIMEOUT_MS,
|
||||
);
|
||||
} catch (error: unknown) {
|
||||
Logger.warn({error, priceId}, 'Failed to retrieve Stripe price summary');
|
||||
|
||||
@@ -790,6 +790,7 @@ export class StripeSubscriptionService {
|
||||
cacheKey,
|
||||
async () => this.loadCurrentSubscriptionPrice(user.stripeSubscriptionId!),
|
||||
StripeSubscriptionService.CURRENT_PRICE_CACHE_TTL_SECONDS,
|
||||
StripeSubscriptionService.PRICE_CACHE_PRODUCE_TIMEOUT_MS,
|
||||
);
|
||||
} catch (error) {
|
||||
Logger.warn(
|
||||
@@ -843,6 +844,7 @@ export class StripeSubscriptionService {
|
||||
return price.unit_amount ?? null;
|
||||
},
|
||||
StripeSubscriptionService.LIST_PRICE_CACHE_TTL_SECONDS,
|
||||
StripeSubscriptionService.PRICE_CACHE_PRODUCE_TIMEOUT_MS,
|
||||
);
|
||||
} catch (error) {
|
||||
Logger.warn({error, priceId}, 'Failed to retrieve Stripe list price amount');
|
||||
@@ -1023,6 +1025,7 @@ export class StripeSubscriptionService {
|
||||
|
||||
private static readonly CURRENT_PRICE_CACHE_TTL_SECONDS = seconds('5 minutes');
|
||||
private static readonly LIST_PRICE_CACHE_TTL_SECONDS = seconds('1 hour');
|
||||
private static readonly PRICE_CACHE_PRODUCE_TIMEOUT_MS = 90000;
|
||||
private static readonly USER_TRIAL_LOCK_TTL_SECONDS = seconds('30 seconds');
|
||||
private static readonly USER_TRIAL_LOCK_MAX_WAIT_MS = 15000;
|
||||
private static readonly USER_TRIAL_LOCK_RETRY_DELAY_MS = 100;
|
||||
|
||||
Reference in New Issue
Block a user