diff --git a/packages/twenty-server/src/engine/core-entity-cache/services/__tests__/core-entity-cache.service.spec.ts b/packages/twenty-server/src/engine/core-entity-cache/services/__tests__/core-entity-cache.service.spec.ts new file mode 100644 index 0000000000..211f873f39 --- /dev/null +++ b/packages/twenty-server/src/engine/core-entity-cache/services/__tests__/core-entity-cache.service.spec.ts @@ -0,0 +1,48 @@ +import { type DiscoveryService, type Reflector } from '@nestjs/core'; + +import { CoreEntityCacheService } from 'src/engine/core-entity-cache/services/core-entity-cache.service'; +import { type CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service'; + +describe('CoreEntityCacheService', () => { + let service: CoreEntityCacheService; + + beforeEach(() => { + jest.useFakeTimers(); + + service = new CoreEntityCacheService( + {} as CacheStorageService, + {} as DiscoveryService, + {} as Reflector, + ); + }); + + afterEach(() => { + jest.useRealTimers(); + }); + + it('should run the local cache expiration sweep at most once per minute', async () => { + const expirationSweepSpy = jest.spyOn( + service as unknown as { + evictExpiredLocalEntries: (now: number) => void; + }, + 'evictExpiredLocalEntries', + ); + + await expect( + service.get('workspaceEntity', 'invalid-entity-id'), + ).resolves.toBeNull(); + await expect( + service.get('workspaceEntity', 'invalid-entity-id'), + ).resolves.toBeNull(); + + expect(expirationSweepSpy).toHaveBeenCalledTimes(1); + + jest.advanceTimersByTime(60_000); + + await expect( + service.get('workspaceEntity', 'invalid-entity-id'), + ).resolves.toBeNull(); + + expect(expirationSweepSpy).toHaveBeenCalledTimes(2); + }); +}); diff --git a/packages/twenty-server/src/engine/core-entity-cache/services/core-entity-cache.service.ts b/packages/twenty-server/src/engine/core-entity-cache/services/core-entity-cache.service.ts index 37449d4553..7b7daa824b 100644 --- a/packages/twenty-server/src/engine/core-entity-cache/services/core-entity-cache.service.ts +++ b/packages/twenty-server/src/engine/core-entity-cache/services/core-entity-cache.service.ts @@ -21,6 +21,7 @@ import { PromiseMemoizer } from 'src/engine/twenty-orm/storage/promise-memoizer. const LOCAL_TTL_MS = 100; // 100ms const LOCAL_ENTRY_TTL_MS = 30 * 60 * 1000; // 30 minutes +const LOCAL_CACHE_EXPIRATION_SWEEP_INTERVAL_MS = 60 * 1000; const MEMOIZER_TTL_MS = 10_000; // 10 seconds const STALE_VERSION_TTL_MS = 5_000; // 5 seconds const MAX_LOCAL_STALE_VERSIONS = 5; @@ -44,6 +45,7 @@ export class CoreEntityCacheService implements OnModuleInit { private readonly memoizer = new PromiseMemoizer( MEMOIZER_TTL_MS, ); + private lastLocalCacheExpirationSweepAt: number | undefined; private readonly logger = new Logger(CoreEntityCacheService.name); @@ -82,7 +84,7 @@ export class CoreEntityCacheService implements OnModuleInit { cacheKeyName: K, entityId: string, ): Promise { - this.evictExpiredLocalEntries(); + this.evictExpiredLocalEntriesIfNeeded(); if (!isDefined(entityId) || !isValidUuid(entityId)) { return null; @@ -288,9 +290,22 @@ export class CoreEntityCacheService implements OnModuleInit { } } - private evictExpiredLocalEntries(): void { + private evictExpiredLocalEntriesIfNeeded(): void { const now = Date.now(); + if ( + isDefined(this.lastLocalCacheExpirationSweepAt) && + now - this.lastLocalCacheExpirationSweepAt < + LOCAL_CACHE_EXPIRATION_SWEEP_INTERVAL_MS + ) { + return; + } + + this.evictExpiredLocalEntries(now); + this.lastLocalCacheExpirationSweepAt = now; + } + + private evictExpiredLocalEntries(now: number): void { for (const [localKey, entry] of this.localCache) { for (const [hash, version] of entry.versions) { if (now - version.lastReadAt > LOCAL_ENTRY_TTL_MS) { diff --git a/packages/twenty-server/src/engine/workspace-cache/services/__tests__/workspace-cache.service.spec.ts b/packages/twenty-server/src/engine/workspace-cache/services/__tests__/workspace-cache.service.spec.ts index ef391dfb34..ea710bc663 100644 --- a/packages/twenty-server/src/engine/workspace-cache/services/__tests__/workspace-cache.service.spec.ts +++ b/packages/twenty-server/src/engine/workspace-cache/services/__tests__/workspace-cache.service.spec.ts @@ -282,6 +282,34 @@ describe('WorkspaceCacheService', () => { }); }); + describe('local cache expiration sweep', () => { + it('should run at most once per minute', async () => { + const expirationSweepSpy = jest.spyOn( + service as unknown as { + evictExpiredLocalEntries: (now: number) => void; + }, + 'evictExpiredLocalEntries', + ); + + await expect( + service.getOrRecompute('invalid-workspace-id', ['featureFlagsMap']), + ).rejects.toThrow(); + await expect( + service.getOrRecompute('invalid-workspace-id', ['featureFlagsMap']), + ).rejects.toThrow(); + + expect(expirationSweepSpy).toHaveBeenCalledTimes(1); + + jest.advanceTimersByTime(60_000); + + await expect( + service.getOrRecompute('invalid-workspace-id', ['featureFlagsMap']), + ).rejects.toThrow(); + + expect(expirationSweepSpy).toHaveBeenCalledTimes(2); + }); + }); + describe('invalidateAndRecompute', () => { beforeEach(async () => { discoveryService.getProviders.mockReturnValue([ diff --git a/packages/twenty-server/src/engine/workspace-cache/services/workspace-cache.service.ts b/packages/twenty-server/src/engine/workspace-cache/services/workspace-cache.service.ts index 72f00b4960..c05f088baa 100644 --- a/packages/twenty-server/src/engine/workspace-cache/services/workspace-cache.service.ts +++ b/packages/twenty-server/src/engine/workspace-cache/services/workspace-cache.service.ts @@ -36,6 +36,7 @@ import { combineCacheHashes } from 'src/engine/workspace-cache/utils/combine-cac const LOCAL_TTL_MS = 100; // 100ms const LOCAL_ENTRY_TTL_MS = 30 * 60 * 1000; // 30 minutes +const LOCAL_CACHE_EXPIRATION_SWEEP_INTERVAL_MS = 60 * 1000; const MEMOIZER_TTL_MS = 10_000; // 10 seconds const STALE_VERSION_TTL_MS = 5_000; // 5 seconds const MAX_LOCAL_STALE_VERSIONS = 5; // 5 stale versions @@ -71,6 +72,7 @@ export class WorkspaceCacheService implements OnModuleInit { private readonly memoizer = new PromiseMemoizer( MEMOIZER_TTL_MS, ); + private lastLocalCacheExpirationSweepAt: number | undefined; private readonly logger = new Logger(WorkspaceCacheService.name); @@ -135,7 +137,7 @@ export class WorkspaceCacheService implements OnModuleInit { workspaceId: string, cacheKeyNames: K, ): Promise> { - this.evictExpiredLocalEntries(); + this.evictExpiredLocalEntriesIfNeeded(); this.assertValidCacheParameters(workspaceId, cacheKeyNames); const memoKey = @@ -640,9 +642,22 @@ export class WorkspaceCacheService implements OnModuleInit { } } - private evictExpiredLocalEntries(): void { + private evictExpiredLocalEntriesIfNeeded(): void { const now = Date.now(); + if ( + isDefined(this.lastLocalCacheExpirationSweepAt) && + now - this.lastLocalCacheExpirationSweepAt < + LOCAL_CACHE_EXPIRATION_SWEEP_INTERVAL_MS + ) { + return; + } + + this.evictExpiredLocalEntries(now); + this.lastLocalCacheExpirationSweepAt = now; + } + + private evictExpiredLocalEntries(now: number): void { for (const [localKey, entry] of this.localCache) { for (const [hash, version] of entry.versions) { if (now - version.lastReadAt > LOCAL_ENTRY_TTL_MS) {