diff --git a/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts b/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts index 12e2d6d53b..0f5f49d350 100644 --- a/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts +++ b/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts @@ -24,6 +24,32 @@ export class CacheStorageService { return this.cache.set(this.getKey(key), value, ttl); } + async setIfAbsent( + key: string, + value: T, + ttl: Milliseconds, + ): Promise { + if (this.isRedisCache()) { + const result = await (this.cache as RedisCache).store.client.set( + this.getKey(key), + JSON.stringify(value), + ttl > 0 ? { NX: true, PX: ttl } : { NX: true }, + ); + + return result === 'OK'; + } + + const existingValue = await this.get(key); + + if (existingValue !== undefined) { + return false; + } + + await this.set(key, value, ttl); + + return true; + } + async del(key: string) { return this.cache.del(this.getKey(key)); } 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 51af77831b..cd0462c65d 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 @@ -6,7 +6,11 @@ import { WorkspaceCacheProvider } from 'src/engine/workspace-cache/interfaces/wo import { type CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service'; import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; import { MetricsService } from 'src/engine/core-modules/metrics/metrics.service'; -import { WORKSPACE_CACHE_KEY } from 'src/engine/workspace-cache/decorators/workspace-cache.decorator'; +import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; +import { + WORKSPACE_CACHE_KEY, + WORKSPACE_CACHE_OPTIONS, +} from 'src/engine/workspace-cache/decorators/workspace-cache.decorator'; import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; const WORKSPACE_ID = '20202020-0000-4000-8000-000000000000'; @@ -27,6 +31,14 @@ class MockRolesPermissionsCacheProvider extends WorkspaceCacheProvider<{ } } +class MockOrmEntityMetadatasCacheProvider extends WorkspaceCacheProvider<{ + testData: string; +}> { + async computeForCache(_workspaceId: string) { + return { testData: 'orm-computed-value' }; + } +} + describe('WorkspaceCacheService', () => { let service: WorkspaceCacheService; let cacheStorageService: jest.Mocked; @@ -48,6 +60,7 @@ describe('WorkspaceCacheService', () => { mget: jest.fn(), mset: jest.fn(), mdel: jest.fn(), + setIfAbsent: jest.fn(), }, }, { @@ -68,6 +81,12 @@ describe('WorkspaceCacheService', () => { incrementCounterBy: jest.fn(), }, }, + { + provide: TwentyConfigService, + useValue: { + get: jest.fn().mockReturnValue(604800), + }, + }, ], }).compile(); @@ -510,4 +529,145 @@ describe('WorkspaceCacheService', () => { expect(result).toEqual({ featureFlagsMap: { testData: 'latest-value' } }); }); }); + + describe('localDataOnly hash recovery', () => { + let localDataOnlyProvider: MockOrmEntityMetadatasCacheProvider; + const ormHashKey = `orm:entity-metadatas:${WORKSPACE_ID}:hash`; + const ormDataKey = `orm:entity-metadatas:${WORKSPACE_ID}:data`; + + beforeEach(async () => { + localDataOnlyProvider = new MockOrmEntityMetadatasCacheProvider(); + + discoveryService.getProviders.mockReturnValue([ + { instance: localDataOnlyProvider }, + ] as any); + + reflector.get.mockImplementation((key, target) => { + if (target !== MockOrmEntityMetadatasCacheProvider) { + return undefined; + } + + if (key === WORKSPACE_CACHE_KEY) { + return 'ORMEntityMetadatas'; + } + + if (key === WORKSPACE_CACHE_OPTIONS) { + return { localDataOnly: true }; + } + + return undefined; + }); + + await service.onModuleInit(); + }); + + it('should adopt the existing redis hash when recomputing instead of minting a new one', async () => { + cacheStorageService.mget.mockResolvedValue(['hash-from-another-pod']); + cacheStorageService.mset.mockResolvedValue(undefined); + + const computeSpy = jest.spyOn(localDataOnlyProvider, 'computeForCache'); + + const result = await service.getOrRecompute(WORKSPACE_ID, [ + 'ORMEntityMetadatas', + ]); + + expect(result).toEqual({ + ORMEntityMetadatas: { testData: 'orm-computed-value' }, + }); + expect(computeSpy).toHaveBeenCalledTimes(1); + expect(cacheStorageService.mset).not.toHaveBeenCalled(); + }); + + it('should treat the local copy as valid on later reads after adopting the redis hash', async () => { + cacheStorageService.mget.mockResolvedValue(['hash-from-another-pod']); + cacheStorageService.mset.mockResolvedValue(undefined); + + const computeSpy = jest.spyOn(localDataOnlyProvider, 'computeForCache'); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + jest.advanceTimersByTime(15_000); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + expect(computeSpy).toHaveBeenCalledTimes(1); + expect(cacheStorageService.mset).not.toHaveBeenCalled(); + }); + + it('should recompute and adopt the new hash when another pod invalidated the key', async () => { + cacheStorageService.mget.mockResolvedValue(['hash-v1']); + cacheStorageService.mset.mockResolvedValue(undefined); + + const computeSpy = jest.spyOn(localDataOnlyProvider, 'computeForCache'); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + jest.advanceTimersByTime(15_000); + + cacheStorageService.mget.mockResolvedValue(['hash-v2']); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + expect(computeSpy).toHaveBeenCalledTimes(2); + expect(cacheStorageService.mset).not.toHaveBeenCalled(); + }); + + it('should mint a hash and persist it only if still absent when redis has none', async () => { + cacheStorageService.mget.mockResolvedValue([undefined]); + cacheStorageService.setIfAbsent.mockResolvedValue(true); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + expect(cacheStorageService.setIfAbsent).toHaveBeenCalledWith( + ormHashKey, + expect.any(String), + 604800000, + ); + expect(cacheStorageService.mset).not.toHaveBeenCalled(); + }); + + it('should converge on a hash written concurrently by another pod instead of overwriting it', async () => { + cacheStorageService.mget.mockResolvedValue([undefined]); + cacheStorageService.setIfAbsent.mockResolvedValue(false); + + const computeSpy = jest.spyOn(localDataOnlyProvider, 'computeForCache'); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + expect(cacheStorageService.setIfAbsent).toHaveBeenCalledTimes(1); + expect(cacheStorageService.mset).not.toHaveBeenCalled(); + + jest.advanceTimersByTime(15_000); + + cacheStorageService.mget.mockResolvedValue(['winner-hash']); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + expect(computeSpy).toHaveBeenCalledTimes(2); + expect(cacheStorageService.setIfAbsent).toHaveBeenCalledTimes(1); + + jest.advanceTimersByTime(15_000); + + await service.getOrRecompute(WORKSPACE_ID, ['ORMEntityMetadatas']); + + expect(computeSpy).toHaveBeenCalledTimes(2); + }); + + it('should mint a new hash on invalidateAndRecompute', async () => { + cacheStorageService.mdel.mockResolvedValue(undefined); + cacheStorageService.mset.mockResolvedValue(undefined); + + await service.invalidateAndRecompute(WORKSPACE_ID, [ + 'ORMEntityMetadatas', + ]); + + expect(cacheStorageService.mdel).toHaveBeenCalledWith([ + ormDataKey, + ormHashKey, + ]); + expect(cacheStorageService.mset).toHaveBeenCalledWith([ + { key: ormHashKey, value: expect.any(String) }, + ]); + }); + }); }); 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 e79b55030b..1344bbf48b 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 @@ -12,6 +12,7 @@ import { CacheStorageService } from 'src/engine/core-modules/cache-storage/servi import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; import { MetricsService } from 'src/engine/core-modules/metrics/metrics.service'; import { MetricsKeys } from 'src/engine/core-modules/metrics/types/metrics-keys.type'; +import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; import { PromiseMemoizer } from 'src/engine/twenty-orm/storage/promise-memoizer.storage'; import { WORKSPACE_CACHE_KEY, @@ -41,6 +42,13 @@ const MIN_EVICT_KEYS = 100; type CacheDataType = WorkspaceCacheDataMap[WorkspaceCacheKeyName]; +type RecomputeHashResolution = + | { strategy: 'mint' } + | { + strategy: 'recover'; + adoptableHashes: Partial>; + }; + @Injectable() export class WorkspaceCacheService implements OnModuleInit { private readonly localCache = new Map< @@ -64,6 +72,7 @@ export class WorkspaceCacheService implements OnModuleInit { private readonly discoveryService: DiscoveryService, private readonly reflector: Reflector, private readonly metricsService: MetricsService, + private readonly twentyConfigService: TwentyConfigService, ) {} async onModuleInit() { @@ -135,8 +144,15 @@ export class WorkspaceCacheService implements OnModuleInit { } // Stage 2: Validate ttl stale keys against Redis hash - const { validKeys, keysNeedingDataFromRedis, keysNeedingRecompute } = - await this.validateLocalHashAgainstRedisHash(workspaceId, staleKeys); + const { + validKeys, + keysNeedingDataFromRedis, + keysNeedingRecompute, + adoptableHashes, + } = await this.validateLocalHashAgainstRedisHash( + workspaceId, + staleKeys, + ); const validatedData = this.getFromLocalCache(workspaceId, validKeys); // Stage 3: Fetch data from Redis @@ -150,6 +166,7 @@ export class WorkspaceCacheService implements OnModuleInit { const recomputedData = await this.recomputeDataFromProvider( workspaceId, keysToRecompute, + { strategy: 'recover', adoptableHashes }, ); return { @@ -171,7 +188,9 @@ export class WorkspaceCacheService implements OnModuleInit { await this.memoizer.clearKeys(`${workspaceId}-`); await this.flush(workspaceId, cacheKeyNames); - await this.recomputeDataFromProvider(workspaceId, cacheKeyNames); + await this.recomputeDataFromProvider(workspaceId, cacheKeyNames, { + strategy: 'mint', + }); // Clear memoizer again after recomputation to evict any stale entries // cached by concurrent getOrRecompute calls during the flush window. @@ -241,13 +260,20 @@ export class WorkspaceCacheService implements OnModuleInit { validKeys: WorkspaceCacheKeyName[]; keysNeedingDataFromRedis: WorkspaceCacheKeyName[]; keysNeedingRecompute: WorkspaceCacheKeyName[]; + adoptableHashes: Partial>; }> { const validKeys: WorkspaceCacheKeyName[] = []; const keysNeedingDataFromRedis: WorkspaceCacheKeyName[] = []; const keysNeedingRecompute: WorkspaceCacheKeyName[] = []; + const adoptableHashes: Partial> = {}; if (cacheKeyNames.length === 0) { - return { validKeys, keysNeedingDataFromRedis, keysNeedingRecompute }; + return { + validKeys, + keysNeedingDataFromRedis, + keysNeedingRecompute, + adoptableHashes, + }; } const hashKeys = cacheKeyNames.map( @@ -270,12 +296,21 @@ export class WorkspaceCacheService implements OnModuleInit { validKeys.push(keyName); } else if (this.localDataOnlyKeys.has(keyName)) { keysNeedingRecompute.push(keyName); + + if (isDefined(redisHash)) { + adoptableHashes[keyName] = redisHash; + } } else { keysNeedingDataFromRedis.push(keyName); } } - return { validKeys, keysNeedingDataFromRedis, keysNeedingRecompute }; + return { + validKeys, + keysNeedingDataFromRedis, + keysNeedingRecompute, + adoptableHashes, + }; } private async fetchDataFromRedis( @@ -321,6 +356,7 @@ export class WorkspaceCacheService implements OnModuleInit { private async recomputeDataFromProvider( workspaceId: string, cacheKeyNames: WorkspaceCacheKeyName[], + hashResolution: RecomputeHashResolution, ): Promise> { const result: Partial = {}; @@ -331,23 +367,41 @@ export class WorkspaceCacheService implements OnModuleInit { const computePromises = cacheKeyNames.map(async (keyName) => { const provider = this.getProviderOrThrow(keyName); const data = await provider.computeForCache(workspaceId); - const hash = crypto.randomUUID(); - return { keyName, data, hash }; + if (hashResolution.strategy === 'mint') { + return { keyName, data, hash: crypto.randomUUID(), isAdopted: false }; + } + + const adoptableHash = hashResolution.adoptableHashes[keyName]; + + return { + keyName, + data, + hash: adoptableHash ?? crypto.randomUUID(), + isAdopted: isDefined(adoptableHash), + }; }); const computed = await Promise.all(computePromises); const redisEntries: Array<{ key: string; value: unknown }> = []; + const bootstrapHashEntries: Array<{ key: string; value: string }> = []; - for (const { keyName, data, hash } of computed) { + for (const { keyName, data, hash, isAdopted } of computed) { Object.assign(result, { [keyName]: data }); const baseKey = this.buildCacheKey(workspaceId, keyName); + const isLocalDataOnly = this.localDataOnlyKeys.has(keyName); + const isRecoveryBootstrap = + hashResolution.strategy === 'recover' && !isAdopted && isLocalDataOnly; - redisEntries.push({ key: `${baseKey}:hash`, value: hash }); + if (isRecoveryBootstrap) { + bootstrapHashEntries.push({ key: `${baseKey}:hash`, value: hash }); + } else if (!isAdopted) { + redisEntries.push({ key: `${baseKey}:hash`, value: hash }); + } - if (!this.localDataOnlyKeys.has(keyName)) { + if (!isLocalDataOnly) { redisEntries.push({ key: `${baseKey}:data`, value: data }); } @@ -358,6 +412,17 @@ export class WorkspaceCacheService implements OnModuleInit { await this.cacheStorage.mset(redisEntries); } + if (bootstrapHashEntries.length > 0) { + const bootstrapHashTtlMs = + this.twentyConfigService.get('CACHE_STORAGE_TTL') * 1000; + + await Promise.all( + bootstrapHashEntries.map(({ key, value }) => + this.cacheStorage.setIfAbsent(key, value, bootstrapHashTtlMs), + ), + ); + } + return result; }