diff --git a/packages/twenty-front/src/modules/action-menu/actions/record-actions/run-workflow-actions/hooks/useRunWorkflowRecordActions.tsx b/packages/twenty-front/src/modules/action-menu/actions/record-actions/run-workflow-actions/hooks/useRunWorkflowRecordActions.tsx index 010110e58c5..a35d66f485a 100644 --- a/packages/twenty-front/src/modules/action-menu/actions/record-actions/run-workflow-actions/hooks/useRunWorkflowRecordActions.tsx +++ b/packages/twenty-front/src/modules/action-menu/actions/record-actions/run-workflow-actions/hooks/useRunWorkflowRecordActions.tsx @@ -49,7 +49,10 @@ export const useRunWorkflowRecordActions = ({ ({ snapshot }) => async ( selectedRecordIds: string[], - activeWorkflowVersion: WorkflowVersion, + activeWorkflowVersion: Pick< + WorkflowVersion, + 'id' | 'workflowId' | 'trigger' + >, ) => { if ( isIteratorEnabled && diff --git a/packages/twenty-front/src/modules/action-menu/actions/record-actions/single-record/workflow-actions/components/SeeActiveVersionWorkflowSingleRecordAction.tsx b/packages/twenty-front/src/modules/action-menu/actions/record-actions/single-record/workflow-actions/components/SeeActiveVersionWorkflowSingleRecordAction.tsx index f0a62dd69d4..e16b28fec5c 100644 --- a/packages/twenty-front/src/modules/action-menu/actions/record-actions/single-record/workflow-actions/components/SeeActiveVersionWorkflowSingleRecordAction.tsx +++ b/packages/twenty-front/src/modules/action-menu/actions/record-actions/single-record/workflow-actions/components/SeeActiveVersionWorkflowSingleRecordAction.tsx @@ -8,7 +8,9 @@ import { AppPath } from 'twenty-shared/types'; export const SeeActiveVersionWorkflowSingleRecordAction = () => { const recordId = useSelectedRecordIdOrThrow(); - const { workflowVersion, loading } = useActiveWorkflowVersion(recordId); + const { workflowVersion, loading } = useActiveWorkflowVersion({ + workflowId: recordId, + }); if (loading) { return ; diff --git a/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersion.ts b/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersion.ts index 5064ca387cd..f1df077bd92 100644 --- a/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersion.ts +++ b/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersion.ts @@ -1,14 +1,16 @@ import { CoreObjectNameSingular } from '@/object-metadata/types/CoreObjectNameSingular'; import { useFindManyRecords } from '@/object-record/hooks/useFindManyRecords'; -import { type Workflow, type WorkflowVersion } from '@/workflow/types/Workflow'; +import { type WorkflowVersion } from '@/workflow/types/Workflow'; -export const useActiveWorkflowVersion = (workflowId: string) => { +type UseActiveWorkflowVersionProps = { + workflowId: string; +}; + +export const useActiveWorkflowVersion = ({ + workflowId, +}: UseActiveWorkflowVersionProps) => { const { records: workflowVersions, loading } = useFindManyRecords< - WorkflowVersion & { - workflow: Omit & { - versions: Array<{ __typename: string }>; - }; - } + Pick >({ objectNameSingular: CoreObjectNameSingular.WorkflowVersion, filter: { @@ -21,21 +23,6 @@ export const useActiveWorkflowVersion = (workflowId: string) => { }, recordGqlFields: { id: true, - name: true, - createdAt: true, - updatedAt: true, - workflowId: true, - trigger: true, - steps: true, - status: true, - workflow: { - id: true, - name: true, - statuses: true, - versions: { - totalCount: true, - }, - }, }, }); diff --git a/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersionsWithManualTrigger.ts b/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersionsWithManualTrigger.ts index 63b98a83428..66f21f4455f 100644 --- a/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersionsWithManualTrigger.ts +++ b/packages/twenty-front/src/modules/workflow/hooks/useActiveWorkflowVersionsWithManualTrigger.ts @@ -1,8 +1,6 @@ import { isGlobalManualTrigger } from '@/action-menu/actions/record-actions/utils/isGlobalManualTrigger'; -import { useObjectMetadataItem } from '@/object-metadata/hooks/useObjectMetadataItem'; import { CoreObjectNameSingular } from '@/object-metadata/types/CoreObjectNameSingular'; import { type ObjectMetadataItem } from '@/object-metadata/types/ObjectMetadataItem'; -import { generateDepthOneRecordGqlFields } from '@/object-record/graphql/utils/generateDepthOneRecordGqlFields'; import { useFindManyRecords } from '@/object-record/hooks/useFindManyRecords'; import { type ManualTriggerWorkflowVersion, @@ -52,23 +50,24 @@ export const useActiveWorkflowVersionsWithManualTrigger = ({ filters.push(objectTypeFilter); } - const { objectMetadataItem: workflowVersionObjectMetadataItem } = - useObjectMetadataItem({ - objectNameSingular: CoreObjectNameSingular.WorkflowVersion, - }); - const { records } = useFindManyRecords< - ManualTriggerWorkflowVersion & { workflow: Workflow } + Pick< + ManualTriggerWorkflowVersion, + 'id' | '__typename' | 'trigger' | 'status' | 'workflowId' + > & { + workflow: Workflow; + } >({ objectNameSingular: CoreObjectNameSingular.WorkflowVersion, filter: { and: filters, }, recordGqlFields: { - ...generateDepthOneRecordGqlFields({ - objectMetadataItem: workflowVersionObjectMetadataItem, - }), + id: true, + trigger: true, + workflowId: true, workflow: true, + status: true, }, skip, }); diff --git a/packages/twenty-server/src/engine/api/graphql/graphql-config/hooks/use-cached-metadata.ts b/packages/twenty-server/src/engine/api/graphql/graphql-config/hooks/use-cached-metadata.ts index 637de7388b9..48193b5793e 100644 --- a/packages/twenty-server/src/engine/api/graphql/graphql-config/hooks/use-cached-metadata.ts +++ b/packages/twenty-server/src/engine/api/graphql/graphql-config/hooks/use-cached-metadata.ts @@ -43,14 +43,20 @@ export function useCachedMetadata(config: CacheMetadataPluginConfig): Plugin { return { onRequest: async ({ endResponse, serverContext }) => { + // TODO: we should probably override the graphql-yoga request type to include the workspace and locale + const request = (serverContext as unknown as { req: Request }).req; + + if (!request.workspace?.id) { + return; + } + if (!config.operationsToCache.includes(getOperationName(serverContext))) { return; } const cacheKey = computeCacheKey({ operationName: getOperationName(serverContext), - // TODO: we should probably override the graphql-yoga request type to include the workspace and locale - request: (serverContext as unknown as { req: Request }).req, + request, }); const cachedResponse = await config.cacheGetter(cacheKey); @@ -61,13 +67,19 @@ export function useCachedMetadata(config: CacheMetadataPluginConfig): Plugin { } }, onResponse: async ({ response, serverContext }) => { + const request = (serverContext as unknown as { req: Request }).req; + + if (!request.workspace?.id) { + return; + } + if (!config.operationsToCache.includes(getOperationName(serverContext))) { return; } const cacheKey = computeCacheKey({ operationName: getOperationName(serverContext), - request: (serverContext as unknown as { req: Request }).req, + request, }); const cachedResponse = await config.cacheGetter(cacheKey); diff --git a/packages/twenty-server/src/engine/api/rest/core/controllers/rest-api-core.controller.ts b/packages/twenty-server/src/engine/api/rest/core/controllers/rest-api-core.controller.ts index 2ebb8790c83..d40b76ad8cd 100644 --- a/packages/twenty-server/src/engine/api/rest/core/controllers/rest-api-core.controller.ts +++ b/packages/twenty-server/src/engine/api/rest/core/controllers/rest-api-core.controller.ts @@ -2,6 +2,7 @@ import { Controller, Delete, Get, + Logger, Patch, Post, Put, @@ -22,10 +23,14 @@ import { WorkspaceAuthGuard } from 'src/engine/guards/workspace-auth.guard'; @UseGuards(JwtAuthGuard, WorkspaceAuthGuard) @UseFilters(RestApiExceptionFilter) export class RestApiCoreController { + private readonly logger = new Logger(RestApiCoreController.name); constructor(private readonly restApiCoreService: RestApiCoreService) {} @Post('batch/*') async handleApiPostBatch(@Req() request: Request, @Res() res: Response) { + this.logger.log( + `[REST API] Processing BATCH request to ${request.path} on workspace ${request.workspaceId}`, + ); const result = await this.restApiCoreService.createMany(request); res.status(201).send(result); @@ -33,6 +38,9 @@ export class RestApiCoreController { @Post('*/duplicates') async handleApiFindDuplicates(@Req() request: Request, @Res() res: Response) { + this.logger.log( + `[REST API] Processing DUPLICATES request to ${request.path} on workspace ${request.workspaceId}`, + ); const result = await this.restApiCoreService.findDuplicates(request); res.status(200).send(result); @@ -40,6 +48,9 @@ export class RestApiCoreController { @Post('*') async handleApiPost(@Req() request: Request, @Res() res: Response) { + this.logger.log( + `[REST API] Processing POST request to ${request.path} on workspace ${request.workspaceId}`, + ); const result = await this.restApiCoreService.createOne(request); res.status(201).send(result); @@ -47,6 +58,9 @@ export class RestApiCoreController { @Get('*') async handleApiGet(@Req() request: Request, @Res() res: Response) { + this.logger.log( + `[REST API] Processing GET request to ${request.path} on workspace ${request.workspaceId}`, + ); const result = await this.restApiCoreService.get(request); res.status(200).send(result); @@ -54,6 +68,9 @@ export class RestApiCoreController { @Delete('*') async handleApiDelete(@Req() request: Request, @Res() res: Response) { + this.logger.log( + `[REST API] Processing DELETE request to ${request.path} on workspace ${request.workspaceId}`, + ); const result = await this.restApiCoreService.delete(request); res.status(200).send(result); @@ -61,6 +78,9 @@ export class RestApiCoreController { @Patch('*') async handleApiPatch(@Req() request: Request, @Res() res: Response) { + this.logger.log( + `[REST API] Processing PATCH request to ${request.path} on workspace ${request.workspaceId}`, + ); const result = await this.restApiCoreService.update(request); res.status(200).send(result); @@ -71,6 +91,9 @@ export class RestApiCoreController { // of PATCH, and because the PUT verb is often used as a PATCH. @Put('*') async handleApiPut(@Req() request: Request, @Res() res: Response) { + this.logger.log( + `[REST API] Processing PUT request to ${request.path} on workspace ${request.workspaceId}`, + ); const result = await this.restApiCoreService.update(request); res.status(200).send(result); 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 10ad62e3df1..ace12db201d 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 @@ -15,7 +15,9 @@ export class CacheStorageService { ) {} async get(key: string): Promise { - return this.cache.get(this.getKey(key)); + const value = await this.cache.get(this.getKey(key)); + + return value; } async set(key: string, value: T, ttl?: Milliseconds) { diff --git a/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts b/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts index bd43061cbcc..e8164090cdd 100644 --- a/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts +++ b/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts @@ -95,6 +95,11 @@ export const useGraphQLErrorHandlerHook = < 'Anonymous Operation'; const workspaceInfo = extractWorkspaceInfo(args.contextValue.req); + // eslint-disable-next-line no-console + console.log( + `[GQL Execute] Processing GQL query ${opName} on workspace ${workspaceInfo?.id}`, + ); + return { onExecuteDone(payload) { const handleResult: OnExecuteDoneHookResultOnNextHook = ({ diff --git a/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.module.ts b/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.module.ts index f2b0128a9db..42f66ac25e8 100644 --- a/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.module.ts +++ b/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.module.ts @@ -10,6 +10,7 @@ import { WorkspaceCacheStorageModule } from 'src/engine/workspace-cache-storage/ imports: [ TypeOrmModule.forFeature([Workspace, FeatureFlag]), WorkspaceCacheStorageModule, + WorkspaceCacheStorageModule, ], providers: [WorkspaceFeatureFlagsMapCacheService], exports: [WorkspaceFeatureFlagsMapCacheService], diff --git a/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.service.ts b/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.service.ts index 4e2a80ab695..f9d27366e71 100644 --- a/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.service.ts @@ -7,7 +7,7 @@ import { type FeatureFlagMap } from 'src/engine/core-modules/feature-flag/interf import { FeatureFlag } from 'src/engine/core-modules/feature-flag/feature-flag.entity'; import { TwentyORMExceptionCode } from 'src/engine/twenty-orm/exceptions/twenty-orm.exception'; -import { getFromCacheWithRecompute } from 'src/engine/utils/get-data-from-cache-with-recompute.util'; +import { GetDataFromCacheWithRecomputeService } from 'src/engine/workspace-cache-storage/services/get-data-from-cache-with-recompute.service'; import { WorkspaceCacheStorageService } from 'src/engine/workspace-cache-storage/workspace-cache-storage.service'; const FEATURE_FLAG_MAP = 'Feature flag map'; @@ -20,6 +20,10 @@ export class WorkspaceFeatureFlagsMapCacheService { private readonly workspaceCacheStorageService: WorkspaceCacheStorageService, @InjectRepository(FeatureFlag) private readonly featureFlagRepository: Repository, + private readonly getFromCacheWithRecomputeService: GetDataFromCacheWithRecomputeService< + string, + FeatureFlagMap + >, ) {} async getWorkspaceFeatureFlagsMap({ @@ -38,7 +42,7 @@ export class WorkspaceFeatureFlagsMapCacheService { }: { workspaceId: string; }) { - return getFromCacheWithRecompute({ + return this.getFromCacheWithRecomputeService.getFromCacheWithRecompute({ workspaceId, getCacheData: () => this.workspaceCacheStorageService.getFeatureFlagsMap(workspaceId), @@ -49,7 +53,6 @@ export class WorkspaceFeatureFlagsMapCacheService { recomputeCache: (params) => this.recomputeFeatureFlagsMapCache(params), cachedEntityName: FEATURE_FLAG_MAP, exceptionCode: TwentyORMExceptionCode.FEATURE_FLAG_MAP_VERSION_NOT_FOUND, - logger: this.logger, }); } diff --git a/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache-storage.service.ts b/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache-storage.service.ts index 27f37b15c6b..a0e1e0326e5 100644 --- a/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache-storage.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache-storage.service.ts @@ -98,6 +98,14 @@ export class WorkspacePermissionsCacheStorageService { ); } + getUserWorkspaceRoleMapVersion( + workspaceId: string, + ): Promise { + return this.cacheStorageService.get( + `${WorkspaceCacheKeys.MetadataPermissionsUserWorkspaceRoleMapVersion}:${workspaceId}`, + ); + } + removeUserWorkspaceRoleMap(workspaceId: string) { return this.cacheStorageService.del( `${WorkspaceCacheKeys.MetadataPermissionsUserWorkspaceRoleMap}:${workspaceId}`, @@ -126,6 +134,14 @@ export class WorkspacePermissionsCacheStorageService { ); } + async getApiKeyRoleMapVersion( + workspaceId: string, + ): Promise { + return this.cacheStorageService.get( + `${WorkspaceCacheKeys.MetadataPermissionsApiKeyRoleMapVersion}:${workspaceId}`, + ); + } + async removeApiKeyRoleMap(workspaceId: string): Promise { await Promise.all([ this.cacheStorageService.del( diff --git a/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service.ts b/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service.ts index 26ccb42c436..64585ac2e37 100644 --- a/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service.ts @@ -1,6 +1,7 @@ import { Injectable, Logger } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; +import { Record } from 'cloudflare/core'; import { ObjectsPermissions, type ObjectsPermissionsByRoleIdDeprecated, @@ -14,10 +15,10 @@ import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadat import { PermissionFlagType } from 'src/engine/metadata-modules/permissions/constants/permission-flag-type.constants'; import { RoleTargetsEntity } from 'src/engine/metadata-modules/role/role-targets.entity'; import { RoleEntity } from 'src/engine/metadata-modules/role/role.entity'; -import { type UserWorkspaceRoleMap } from 'src/engine/metadata-modules/workspace-permissions-cache/types/user-workspace-role-map.type'; +import { UserWorkspaceRoleMap } from 'src/engine/metadata-modules/workspace-permissions-cache/types/user-workspace-role-map.type'; import { WorkspacePermissionsCacheStorageService } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache-storage.service'; import { TwentyORMExceptionCode } from 'src/engine/twenty-orm/exceptions/twenty-orm.exception'; -import { getFromCacheWithRecompute } from 'src/engine/utils/get-data-from-cache-with-recompute.util'; +import { GetDataFromCacheWithRecomputeService } from 'src/engine/workspace-cache-storage/services/get-data-from-cache-with-recompute.service'; import { STANDARD_OBJECT_IDS } from 'src/engine/workspace-manager/workspace-sync-metadata/constants/standard-object-ids'; type CacheResult = { @@ -45,6 +46,18 @@ export class WorkspacePermissionsCacheService { @InjectRepository(RoleTargetsEntity) private readonly roleTargetsRepository: Repository, private readonly workspacePermissionsCacheStorageService: WorkspacePermissionsCacheStorageService, + private readonly getRolesPermissionsFromCacheWithRecomputeService: GetDataFromCacheWithRecomputeService< + string, + ObjectsPermissionsByRoleIdDeprecated + >, + private readonly getUserWorkspaceRoleMapFromCacheWithRecomputeService: GetDataFromCacheWithRecomputeService< + string, + UserWorkspaceRoleMap + >, + private readonly getApiKeyRoleMapFromCacheWithRecomputeService: GetDataFromCacheWithRecomputeService< + string, + Record + >, ) {} async recomputeRolesPermissionsCache({ @@ -109,44 +122,48 @@ export class WorkspacePermissionsCacheService { }: { workspaceId: string; }): Promise> { - return getFromCacheWithRecompute< - string, - ObjectsPermissionsByRoleIdDeprecated - >({ - workspaceId, - getCacheData: () => - this.workspacePermissionsCacheStorageService.getRolesPermissions( - workspaceId, - ), - getCacheVersion: () => - this.workspacePermissionsCacheStorageService.getRolesPermissionsVersion( - workspaceId, - ), - recomputeCache: (params) => this.recomputeRolesPermissionsCache(params), - cachedEntityName: ROLES_PERMISSIONS, - exceptionCode: TwentyORMExceptionCode.ROLES_PERMISSIONS_VERSION_NOT_FOUND, - logger: this.logger, - }); + return this.getRolesPermissionsFromCacheWithRecomputeService.getFromCacheWithRecompute( + { + workspaceId, + getCacheData: () => + this.workspacePermissionsCacheStorageService.getRolesPermissions( + workspaceId, + ), + getCacheVersion: () => + this.workspacePermissionsCacheStorageService.getRolesPermissionsVersion( + workspaceId, + ), + recomputeCache: (params) => this.recomputeRolesPermissionsCache(params), + cachedEntityName: ROLES_PERMISSIONS, + exceptionCode: + TwentyORMExceptionCode.ROLES_PERMISSIONS_VERSION_NOT_FOUND, + }, + ); } async getUserWorkspaceRoleMapFromCache({ workspaceId, }: { workspaceId: string; - }): Promise> { - return getFromCacheWithRecompute({ - workspaceId, - getCacheData: () => - this.workspacePermissionsCacheStorageService.getUserWorkspaceRoleMap( - workspaceId, - ), - recomputeCache: (params) => - this.recomputeUserWorkspaceRoleMapCache(params), - cachedEntityName: USER_WORKSPACE_ROLE_MAP, - exceptionCode: - TwentyORMExceptionCode.USER_WORKSPACE_ROLE_MAP_VERSION_NOT_FOUND, - logger: this.logger, - }); + }): Promise> { + return this.getUserWorkspaceRoleMapFromCacheWithRecomputeService.getFromCacheWithRecompute( + { + workspaceId, + getCacheData: () => + this.workspacePermissionsCacheStorageService.getUserWorkspaceRoleMap( + workspaceId, + ), + getCacheVersion: () => + this.workspacePermissionsCacheStorageService.getUserWorkspaceRoleMapVersion( + workspaceId, + ), + recomputeCache: (params) => + this.recomputeUserWorkspaceRoleMapCache(params), + cachedEntityName: USER_WORKSPACE_ROLE_MAP, + exceptionCode: + TwentyORMExceptionCode.USER_WORKSPACE_ROLE_MAP_VERSION_NOT_FOUND, + }, + ); } async getRoleIdFromUserWorkspaceId({ @@ -364,18 +381,24 @@ export class WorkspacePermissionsCacheService { workspaceId, }: { workspaceId: string; - }): Promise>> { - return getFromCacheWithRecompute>({ - workspaceId, - getCacheData: () => - this.workspacePermissionsCacheStorageService.getApiKeyRoleMap( - workspaceId, - ), - recomputeCache: (params) => this.recomputeApiKeyRoleMapCache(params), - cachedEntityName: 'API_KEY_ROLE_MAP', - exceptionCode: TwentyORMExceptionCode.API_KEY_ROLE_MAP_VERSION_NOT_FOUND, - logger: this.logger, - }); + }): Promise>> { + return this.getApiKeyRoleMapFromCacheWithRecomputeService.getFromCacheWithRecompute( + { + workspaceId, + getCacheData: () => + this.workspacePermissionsCacheStorageService.getApiKeyRoleMap( + workspaceId, + ), + getCacheVersion: () => + this.workspacePermissionsCacheStorageService.getApiKeyRoleMapVersion( + workspaceId, + ), + recomputeCache: (params) => this.recomputeApiKeyRoleMapCache(params), + cachedEntityName: 'API_KEY_ROLE_MAP', + exceptionCode: + TwentyORMExceptionCode.API_KEY_ROLE_MAP_VERSION_NOT_FOUND, + }, + ); } private async getApiKeyRoleMapFromDatabase({ diff --git a/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts b/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts index 23f2b513206..476accd244f 100644 --- a/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts +++ b/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts @@ -25,7 +25,7 @@ import { import { EntitySchemaFactory } from 'src/engine/twenty-orm/factories/entity-schema.factory'; import { PromiseMemoizer } from 'src/engine/twenty-orm/storage/promise-memoizer.storage'; import { type CacheKey } from 'src/engine/twenty-orm/storage/types/cache-key.type'; -import { getFromCacheWithRecompute } from 'src/engine/utils/get-data-from-cache-with-recompute.util'; +import { GetDataFromCacheWithRecomputeService } from 'src/engine/workspace-cache-storage/services/get-data-from-cache-with-recompute.service'; import { WorkspaceCacheStorageService } from 'src/engine/workspace-cache-storage/workspace-cache-storage.service'; import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter'; @@ -53,6 +53,10 @@ export class WorkspaceDatasourceFactory { @InjectRepository(Workspace) private readonly workspaceRepository: Repository, private readonly workspaceEventEmitter: WorkspaceEventEmitter, + private readonly getFromCacheWithRecomputeService: GetDataFromCacheWithRecomputeService< + string, + ObjectsPermissionsByRoleIdDeprecated + >, ) {} private async safelyDestroyDataSource( @@ -236,10 +240,7 @@ export class WorkspaceDatasourceFactory { }: { workspaceId: string; }): Promise> { - return getFromCacheWithRecompute< - string, - ObjectsPermissionsByRoleIdDeprecated - >({ + return this.getFromCacheWithRecomputeService.getFromCacheWithRecompute({ workspaceId, getCacheData: () => this.workspacePermissionsCacheStorageService.getRolesPermissions( @@ -255,7 +256,6 @@ export class WorkspaceDatasourceFactory { }), cachedEntityName: ROLES_PERMISSIONS, exceptionCode: TwentyORMExceptionCode.ROLES_PERMISSIONS_VERSION_NOT_FOUND, - logger: this.logger, }); } diff --git a/packages/twenty-server/src/engine/utils/get-data-from-cache-with-recompute.util.ts b/packages/twenty-server/src/engine/utils/get-data-from-cache-with-recompute.util.ts deleted file mode 100644 index 5a72e97f8d3..00000000000 --- a/packages/twenty-server/src/engine/utils/get-data-from-cache-with-recompute.util.ts +++ /dev/null @@ -1,85 +0,0 @@ -import { type Logger } from '@nestjs/common'; - -import { isDefined } from 'twenty-shared/utils'; - -import { - TwentyORMException, - type TwentyORMExceptionCode, -} from 'src/engine/twenty-orm/exceptions/twenty-orm.exception'; - -type CacheResult = { - version: T; - data: U; -}; - -const getFromCacheWithRecompute = async ({ - workspaceId, - getCacheData, - getCacheVersion, - recomputeCache, - cachedEntityName, - exceptionCode, - logger, -}: { - workspaceId: string; - getCacheData: (workspaceId: string) => Promise; - getCacheVersion?: (workspaceId: string) => Promise; - recomputeCache: (params: { workspaceId: string }) => Promise; - cachedEntityName: string; - exceptionCode: TwentyORMExceptionCode; - logger: Logger; -}): Promise> => { - let cachedVersion: T | undefined; - let cachedData: U | undefined; - - const expectCacheVersion = isDefined(getCacheVersion); - - if (expectCacheVersion) { - cachedVersion = await getCacheVersion(workspaceId); - } - - cachedData = await getCacheData(workspaceId); - - if ( - !isDefined(cachedData) || - (expectCacheVersion && !isDefined(cachedVersion)) - ) { - logger.warn( - `Triggering cache recompute for ${cachedEntityName} (workspace ${workspaceId})`, - { - cachedVersion, - cachedData, - }, - ); - await recomputeCache({ workspaceId }); - - cachedData = await getCacheData(workspaceId); - if (expectCacheVersion) { - cachedVersion = await getCacheVersion(workspaceId); - } - - if ( - !isDefined(cachedData) || - (expectCacheVersion && !isDefined(cachedVersion)) - ) { - logger.warn( - `Data still missing after recompute for ${cachedEntityName} (workspace ${workspaceId})`, - { - cachedVersion, - cachedData, - }, - ); - throw new TwentyORMException( - `${cachedEntityName} not found after recompute for workspace ${workspaceId} (missingData: ${!isDefined(cachedData)}, missingVersion: ${expectCacheVersion && !isDefined(cachedVersion)})`, - exceptionCode, - ); - } - } - - return { - version: cachedVersion as T, - data: cachedData, - }; -}; - -export { CacheResult, getFromCacheWithRecompute }; diff --git a/packages/twenty-server/src/engine/workspace-cache-storage/services/get-data-from-cache-with-recompute.service.ts b/packages/twenty-server/src/engine/workspace-cache-storage/services/get-data-from-cache-with-recompute.service.ts new file mode 100644 index 00000000000..bd39f80bae1 --- /dev/null +++ b/packages/twenty-server/src/engine/workspace-cache-storage/services/get-data-from-cache-with-recompute.service.ts @@ -0,0 +1,90 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import { logger } from '@sentry/node'; +import { isDefined } from 'twenty-shared/utils'; + +import { + TwentyORMException, + TwentyORMExceptionCode, +} from 'src/engine/twenty-orm/exceptions/twenty-orm.exception'; + +type CacheResult = { + version: T; + data: U; +}; + +@Injectable() +export class GetDataFromCacheWithRecomputeService { + private cache = new Map>(); + + logger = new Logger(GetDataFromCacheWithRecomputeService.name); + constructor() {} + + getFromCacheWithRecompute = async ({ + workspaceId, + getCacheData, + getCacheVersion, + recomputeCache, + cachedEntityName, + exceptionCode, + }: { + workspaceId: string; + getCacheData: (workspaceId: string) => Promise; + getCacheVersion: (workspaceId: string) => Promise; + recomputeCache: (params: { workspaceId: string }) => Promise; + cachedEntityName: string; + exceptionCode: TwentyORMExceptionCode; + }): Promise> => { + let cachedVersion: T | undefined; + let cachedData: U | undefined; + + cachedVersion = await getCacheVersion(workspaceId); + + const cacheKey = `${workspaceId}-${cachedVersion}`; + const cachedValue = this.cache.get(cacheKey); + + if (cachedValue) { + return cachedValue; + } + + cachedData = await getCacheData(workspaceId); + + if (!isDefined(cachedData) || !isDefined(cachedVersion)) { + logger.warn( + `Triggering cache recompute for ${cachedEntityName} (workspace ${workspaceId})`, + { + cachedVersion, + cachedData, + }, + ); + await recomputeCache({ workspaceId }); + + cachedData = await getCacheData(workspaceId); + cachedVersion = await getCacheVersion(workspaceId); + + if (!isDefined(cachedData) || !isDefined(cachedVersion)) { + logger.warn( + `Data still missing after recompute for ${cachedEntityName} (workspace ${workspaceId})`, + { + cachedVersion, + cachedData, + }, + ); + throw new TwentyORMException( + `${cachedEntityName} not found after recompute for workspace ${workspaceId} (missingData: ${!isDefined(cachedData)}, missingVersion: ${!isDefined(cachedVersion)})`, + exceptionCode, + ); + } + } + + this.cache.set(cacheKey, { + version: cachedVersion, + data: cachedData, + }); + + return { + version: cachedVersion, + data: cachedData, + }; + }; +} diff --git a/packages/twenty-server/src/engine/workspace-cache-storage/workspace-cache-storage.module.ts b/packages/twenty-server/src/engine/workspace-cache-storage/workspace-cache-storage.module.ts index fec5215284b..bf4ecb08c52 100644 --- a/packages/twenty-server/src/engine/workspace-cache-storage/workspace-cache-storage.module.ts +++ b/packages/twenty-server/src/engine/workspace-cache-storage/workspace-cache-storage.module.ts @@ -1,9 +1,13 @@ import { Module } from '@nestjs/common'; +import { GetDataFromCacheWithRecomputeService } from 'src/engine/workspace-cache-storage/services/get-data-from-cache-with-recompute.service'; import { WorkspaceCacheStorageService } from 'src/engine/workspace-cache-storage/workspace-cache-storage.service'; @Module({ - providers: [WorkspaceCacheStorageService], - exports: [WorkspaceCacheStorageService], + providers: [ + WorkspaceCacheStorageService, + GetDataFromCacheWithRecomputeService, + ], + exports: [WorkspaceCacheStorageService, GetDataFromCacheWithRecomputeService], }) export class WorkspaceCacheStorageModule {} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts index b3301d79e91..1127a9894b9 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts @@ -18,7 +18,7 @@ import { type MessagingMessageListFetchJobData, } from 'src/modules/messaging/message-import-manager/jobs/messaging-message-list-fetch.job'; -export const MESSAGING_MESSAGE_LIST_FETCH_CRON_PATTERN = '*/5 * * * *'; +export const MESSAGING_MESSAGE_LIST_FETCH_CRON_PATTERN = '2-59/5 * * * *'; @Processor(MessageQueue.cronQueue) export class MessagingMessageListFetchCronJob { diff --git a/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts b/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts index a06f47613a8..767f3d4f27a 100644 --- a/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts @@ -17,15 +17,15 @@ import { WorkflowVersionStatus, type WorkflowVersionWorkspaceEntity, } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity'; +import { + WorkflowStatus, + type WorkflowWorkspaceEntity, +} from 'src/modules/workflow/common/standard-objects/workflow.workspace-entity'; import { WorkflowActionType } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; import { WorkflowTriggerException, WorkflowTriggerExceptionCode, } from 'src/modules/workflow/workflow-trigger/exceptions/workflow-trigger.exception'; -import { - WorkflowStatus, - type WorkflowWorkspaceEntity, -} from 'src/modules/workflow/common/standard-objects/workflow.workspace-entity'; export type ObjectMetadataInfo = { objectMetadataItemWithFieldsMaps: ObjectMetadataItemWithFieldMaps; @@ -94,6 +94,7 @@ export class WorkflowCommonWorkspaceService { async getObjectMetadataMaps( workspaceId: string, ): Promise { + // TODO: replace this with the new cache service const objectMetadataMaps = await this.workspaceCacheStorageService.getObjectMetadataMapsOrThrow( workspaceId, diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts index 0d176075a84..9c85d1b977a 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts @@ -16,6 +16,8 @@ import { type ObjectRecordUpsertEvent } from 'src/engine/core-modules/event-emit import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; +import { ObjectMetadataItemWithFieldMaps } from 'src/engine/metadata-modules/types/object-metadata-item-with-field-maps'; +import { ObjectMetadataMaps } from 'src/engine/metadata-modules/types/object-metadata-maps'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type'; import { @@ -134,11 +136,17 @@ export class WorkflowDatabaseEventTriggerListener { payload: WorkspaceEventBatch, ) { const workspaceId = payload.workspaceId; + const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } = + await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps( + payload.events[0].objectMetadata.nameSingular, + workspaceId, + ); await this.enrichRecordsWithRelations({ records: payload.events.map((event) => event.properties.after), - objectMetadataNameSingular: payload.events[0].objectMetadata.nameSingular, workspaceId, + objectMetadataMaps, + objectMetadataItemWithFieldsMaps, }); } @@ -146,16 +154,23 @@ export class WorkflowDatabaseEventTriggerListener { payload: WorkspaceEventBatch, ) { const workspaceId = payload.workspaceId; + const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } = + await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps( + payload.events[0].objectMetadata.nameSingular, + workspaceId, + ); await this.enrichRecordsWithRelations({ records: payload.events.map((event) => event.properties.before), - objectMetadataNameSingular: payload.events[0].objectMetadata.nameSingular, workspaceId, + objectMetadataMaps, + objectMetadataItemWithFieldsMaps, }); await this.enrichRecordsWithRelations({ records: payload.events.map((event) => event.properties.after), - objectMetadataNameSingular: payload.events[0].objectMetadata.nameSingular, workspaceId, + objectMetadataMaps, + objectMetadataItemWithFieldsMaps, }); } @@ -163,11 +178,17 @@ export class WorkflowDatabaseEventTriggerListener { payload: WorkspaceEventBatch, ) { const workspaceId = payload.workspaceId; + const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } = + await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps( + payload.events[0].objectMetadata.nameSingular, + workspaceId, + ); await this.enrichRecordsWithRelations({ records: payload.events.map((event) => event.properties.before), - objectMetadataNameSingular: payload.events[0].objectMetadata.nameSingular, workspaceId, + objectMetadataMaps, + objectMetadataItemWithFieldsMaps, }); } @@ -175,29 +196,31 @@ export class WorkflowDatabaseEventTriggerListener { payload: WorkspaceEventBatch, ) { const workspaceId = payload.workspaceId; + const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } = + await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps( + payload.events[0].objectMetadata.nameSingular, + workspaceId, + ); await this.enrichRecordsWithRelations({ records: payload.events.map((event) => event.properties.before), - objectMetadataNameSingular: payload.events[0].objectMetadata.nameSingular, + objectMetadataMaps, workspaceId, + objectMetadataItemWithFieldsMaps, }); } private async enrichRecordsWithRelations({ records, - objectMetadataNameSingular, workspaceId, + objectMetadataMaps, + objectMetadataItemWithFieldsMaps, }: { records: Partial[]; - objectMetadataNameSingular: string; workspaceId: string; + objectMetadataMaps: ObjectMetadataMaps; + objectMetadataItemWithFieldsMaps: ObjectMetadataItemWithFieldMaps; }) { - const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } = - await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps( - objectMetadataNameSingular, - workspaceId, - ); - for (const [joinColumnName, joinFieldId] of Object.entries( objectMetadataItemWithFieldsMaps.fieldIdByJoinColumnName, )) {