diff --git a/packages/twenty-server/src/engine/core-modules/cache-lock/cache-lock.service.ts b/packages/twenty-server/src/engine/core-modules/cache-lock/cache-lock.service.ts index 8e2123481f..e7c20f5916 100644 --- a/packages/twenty-server/src/engine/core-modules/cache-lock/cache-lock.service.ts +++ b/packages/twenty-server/src/engine/core-modules/cache-lock/cache-lock.service.ts @@ -1,5 +1,9 @@ import { Injectable, Logger } from '@nestjs/common'; +import { + CacheLockException, + CacheLockExceptionCode, +} from 'src/engine/core-modules/cache-lock/exceptions/cache-lock.exception'; import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decorators/cache-storage.decorator'; import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service'; @@ -50,6 +54,9 @@ export class CacheLockService { await this.delay(ms); } - throw new Error(`Failed to acquire lock for key: ${key}`); + throw new CacheLockException( + `Failed to acquire lock for key: ${key}`, + CacheLockExceptionCode.LOCK_ACQUISITION_TIMEOUT, + ); } } diff --git a/packages/twenty-server/src/engine/core-modules/cache-lock/exceptions/cache-lock.exception.ts b/packages/twenty-server/src/engine/core-modules/cache-lock/exceptions/cache-lock.exception.ts new file mode 100644 index 0000000000..d00417b4ef --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/cache-lock/exceptions/cache-lock.exception.ts @@ -0,0 +1,26 @@ +import { type MessageDescriptor } from '@lingui/core'; +import { msg } from '@lingui/core/macro'; + +import { + appendCommonExceptionCode, + CustomException, +} from 'src/utils/custom-exception'; + +export class CacheLockException extends CustomException< + keyof typeof CacheLockExceptionCode +> { + constructor( + message: string, + code: keyof typeof CacheLockExceptionCode, + { userFriendlyMessage }: { userFriendlyMessage?: MessageDescriptor } = {}, + ) { + super(message, code, { + userFriendlyMessage: + userFriendlyMessage ?? msg`A cache lock error occurred.`, + }); + } +} + +export const CacheLockExceptionCode = appendCommonExceptionCode({ + LOCK_ACQUISITION_TIMEOUT: 'LOCK_ACQUISITION_TIMEOUT', +} as const); diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda.driver.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda.driver.ts index 2151fc370c..67976ca211 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda.driver.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda.driver.ts @@ -72,6 +72,7 @@ export class LambdaDriver implements LogicFunctionDriver { this.layerManager, options.cacheLockService, options.logicFunctionResourceService, + options.workspaceCacheService, ); } diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/services/lambda-executor-manager.service.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/services/lambda-executor-manager.service.ts index 834d77792d..be7a5f9d8c 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/services/lambda-executor-manager.service.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/services/lambda-executor-manager.service.ts @@ -18,6 +18,10 @@ import { isDefined } from 'twenty-shared/utils'; import { type FlatApplication } from 'src/engine/core-modules/application/types/flat-application.type'; import { type CacheLockService } from 'src/engine/core-modules/cache-lock/cache-lock.service'; +import { + CacheLockException, + CacheLockExceptionCode, +} from 'src/engine/core-modules/cache-lock/exceptions/cache-lock.exception'; import { EXECUTOR_LAMBDA_MEMORY_MB, EXECUTOR_LAMBDA_TIMEOUT_SECONDS, @@ -40,6 +44,7 @@ import { LogicFunctionExceptionCode, } from 'src/engine/metadata-modules/logic-function/logic-function.exception'; import { type FlatLogicFunction } from 'src/engine/metadata-modules/logic-function/types/flat-logic-function.type'; +import { type WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; type ExecutorBuildContext = { flatLogicFunction: FlatLogicFunction; @@ -56,6 +61,7 @@ export class LambdaExecutorManagerService { private readonly layerManager: LambdaLayerManagerService, private readonly cacheLockService: CacheLockService, private readonly logicFunctionResourceService: LogicFunctionResourceService, + private readonly workspaceCacheService: WorkspaceCacheService, ) {} async getLambdaExecutor( @@ -100,26 +106,74 @@ export class LambdaExecutorManagerService { const buildLockTtlMs = 120_000; const buildLockRetryMs = 500; const buildLockMaxRetries = 240; + const lockKey = `lambda-build:${context.flatLogicFunction.id}`; - await this.cacheLockService.withLock( - async () => { - // Need to check again inside the lock in case lock was not acquired immediately. - const { canSkip: canSkipAfterLock, lambdaExecutor } = - await this.checkBuildStatus(context); + try { + await this.cacheLockService.withLock( + async () => { + // Refresh the application before re-checking: another process may + // have rebuilt the SDK layer (and cleared isSdkLayerStale) while we + // were waiting for the lock, and our request-time snapshot cannot + // see it. Without this, every queued waiter rebuilds again. + const refreshedContext = await this.refreshBuildContext(context); - if (canSkipAfterLock) { - return; - } + const { canSkip: canSkipAfterLock, lambdaExecutor } = + await this.checkBuildStatus(refreshedContext); - await this.ensureExecutor({ ...context, lambdaExecutor }); - }, - `lambda-build:${context.flatLogicFunction.id}`, - { - ttl: buildLockTtlMs, - ms: buildLockRetryMs, - maxRetries: buildLockMaxRetries, - }, - ); + if (canSkipAfterLock) { + return; + } + + await this.ensureExecutor({ ...refreshedContext, lambdaExecutor }); + }, + lockKey, + { + ttl: buildLockTtlMs, + ms: buildLockRetryMs, + maxRetries: buildLockMaxRetries, + }, + ); + } catch (error) { + const isLockAcquisitionTimeout = + error instanceof CacheLockException && + error.code === CacheLockExceptionCode.LOCK_ACQUISITION_TIMEOUT; + + if (!isLockAcquisitionTimeout) { + throw error; + } + + // Lock wait budget exhausted. If concurrent builds left the executor in + // a usable state, proceed with the invocation instead of failing it. + const { canSkip: isExecutorUsable } = await this.checkBuildStatus( + await this.refreshBuildContext(context), + ); + + if (!isExecutorUsable) { + throw error; + } + + this.logger.warn( + `Lock acquisition timed out for ${lockKey} but executor is up to date, proceeding`, + ); + } + } + + private async refreshBuildContext( + context: ExecutorBuildContext, + ): Promise { + const { flatApplicationMaps } = + await this.workspaceCacheService.getOrRecompute( + context.flatLogicFunction.workspaceId, + ['flatApplicationMaps'], + ); + + const refreshedFlatApplication = + flatApplicationMaps.byId[context.flatApplication.id]; + + return { + ...context, + flatApplication: refreshedFlatApplication ?? context.flatApplication, + }; } async installPrebuiltBundle(context: ExecutorBuildContext): Promise { diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/types/lambda-driver.type.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/types/lambda-driver.type.ts index 6f23c8d538..57c6c13a8d 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/types/lambda-driver.type.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/drivers/lambda/types/lambda-driver.type.ts @@ -3,6 +3,7 @@ import { type LambdaClientConfig } from '@aws-sdk/client-lambda'; import { type CacheLockService } from 'src/engine/core-modules/cache-lock/cache-lock.service'; import { type LogicFunctionResourceService } from 'src/engine/core-modules/logic-function/logic-function-resource/logic-function-resource.service'; import { type SdkClientArchiveService } from 'src/engine/core-modules/sdk-client/sdk-client-archive.service'; +import { type WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; export type LambdaDriverExecutorPayload = { code?: string; @@ -37,6 +38,7 @@ export interface LambdaDriverOptions extends LambdaClientConfig { logicFunctionResourceService: LogicFunctionResourceService; sdkClientArchiveService: SdkClientArchiveService; cacheLockService: CacheLockService; + workspaceCacheService: WorkspaceCacheService; region: string; lambdaRole: string; subhostingRole?: string; diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/logic-function-driver.factory.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/logic-function-driver.factory.ts index 53f7cd0a02..befc0ee0fa 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/logic-function-driver.factory.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-drivers/logic-function-driver.factory.ts @@ -86,6 +86,7 @@ export class LogicFunctionDriverFactory extends DriverFactoryBase