Avoid retaining accounting summary promises

This commit is contained in:
Ben
2026-08-04 12:30:29 -04:00
parent dd39c6edcf
commit 7d96c06bbf
@@ -102,6 +102,7 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
private rollupTimer: NodeJS.Timeout | null = null; private rollupTimer: NodeJS.Timeout | null = null;
private activeRollup: Promise<void> | null = null; private activeRollup: Promise<void> | null = null;
private summaryCache = new Map<string, SummaryCacheEntry>(); private summaryCache = new Map<string, SummaryCacheEntry>();
private summaryInFlight = new Map<string, Promise<ShareAccountingSummary>>();
private readonly poolSummaryCacheKey = 'accounting:pool-summary'; private readonly poolSummaryCacheKey = 'accounting:pool-summary';
private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE); private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE);
private readonly flushIntervalMs = this.readPositiveInt('SHARE_ACCOUNTING_FLUSH_INTERVAL_MS', DEFAULT_FLUSH_INTERVAL_MS); private readonly flushIntervalMs = this.readPositiveInt('SHARE_ACCOUNTING_FLUSH_INTERVAL_MS', DEFAULT_FLUSH_INTERVAL_MS);
@@ -565,18 +566,31 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
return cached.value; return cached.value;
} }
const value = this.loadCachedSummary(filter, cacheKey).catch(error => { if (cached != null) {
this.summaryCache.delete(cacheKey); this.summaryCache.delete(cacheKey);
throw error; }
});
const inFlight = this.summaryInFlight.get(cacheKey);
if (inFlight != null) {
return inFlight;
}
const value = this.loadCachedSummary(filter, cacheKey)
.then(summary => {
if (this.summaryCacheMs > 0) { if (this.summaryCacheMs > 0) {
this.summaryCache.set(cacheKey, { this.summaryCache.set(cacheKey, {
expiresAt: now + this.summaryCacheMs, expiresAt: Date.now() + this.summaryCacheMs,
value, value: summary,
}); });
this.trimSummaryCache(); this.trimSummaryCache();
} }
return summary;
})
.finally(() => {
this.summaryInFlight.delete(cacheKey);
});
this.summaryInFlight.set(cacheKey, value);
return value; return value;
} }
@@ -865,6 +879,13 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
} }
private trimSummaryCache(): void { private trimSummaryCache(): void {
const now = Date.now();
for (const [key, entry] of this.summaryCache) {
if (entry.expiresAt <= now) {
this.summaryCache.delete(key);
}
}
while (this.summaryCache.size > this.summaryCacheMax) { while (this.summaryCache.size > this.summaryCacheMax) {
const firstKey = this.summaryCache.keys().next().value; const firstKey = this.summaryCache.keys().next().value;
if (firstKey == null) { if (firstKey == null) {
@@ -901,5 +922,5 @@ interface PendingShare {
interface SummaryCacheEntry { interface SummaryCacheEntry {
expiresAt: number; expiresAt: number;
value: Promise<ShareAccountingSummary>; value: ShareAccountingSummary;
} }