import { Injectable } from '@nestjs/common'; import { Record as IRecord, RecordFilter, RecordOrderBy, } from 'src/engine/api/graphql/workspace-query-builder/interfaces/record.interface'; import { IConnection } from 'src/engine/api/graphql/workspace-query-runner/interfaces/connection.interface'; import { WorkspaceQueryRunnerOptions } from 'src/engine/api/graphql/workspace-query-runner/interfaces/query-runner-option.interface'; import { CreateManyResolverArgs, CreateOneResolverArgs, DestroyOneResolverArgs, FindManyResolverArgs, FindOneResolverArgs, ResolverArgsType, } from 'src/engine/api/graphql/workspace-resolver-builder/interfaces/workspace-resolvers-builder.interface'; import { ObjectMetadataInterface } from 'src/engine/metadata-modules/field-metadata/interfaces/object-metadata.interface'; import { GraphqlQueryCreateManyResolverService } from 'src/engine/api/graphql/graphql-query-runner/resolvers/graphql-query-create-many-resolver.service'; import { GraphqlQueryDestroyOneResolverService } from 'src/engine/api/graphql/graphql-query-runner/resolvers/graphql-query-destroy-one-resolver.service'; import { GraphqlQueryFindManyResolverService } from 'src/engine/api/graphql/graphql-query-runner/resolvers/graphql-query-find-many-resolver.service'; import { GraphqlQueryFindOneResolverService } from 'src/engine/api/graphql/graphql-query-runner/resolvers/graphql-query-find-one-resolver.service'; import { QueryRunnerArgsFactory } from 'src/engine/api/graphql/workspace-query-runner/factories/query-runner-args.factory'; import { CallWebhookJobsJob, CallWebhookJobsJobData, CallWebhookJobsJobOperation, } from 'src/engine/api/graphql/workspace-query-runner/jobs/call-webhook-jobs.job'; import { assertIsValidUuid } from 'src/engine/api/graphql/workspace-query-runner/utils/assert-is-valid-uuid.util'; import { WorkspaceQueryHookService } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-hook/workspace-query-hook.service'; import { WorkspaceQueryRunnerException, WorkspaceQueryRunnerExceptionCode, } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-runner.exception'; import { AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type'; import { ObjectRecordCreateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-create.event'; 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 { LogExecutionTime } from 'src/engine/decorators/observability/log-execution-time.decorator'; import { assertMutationNotOnRemoteObject } from 'src/engine/metadata-modules/object-metadata/utils/assert-mutation-not-on-remote-object.util'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter'; @Injectable() export class GraphqlQueryRunnerService { constructor( private readonly twentyORMGlobalManager: TwentyORMGlobalManager, private readonly workspaceQueryHookService: WorkspaceQueryHookService, private readonly queryRunnerArgsFactory: QueryRunnerArgsFactory, private readonly workspaceEventEmitter: WorkspaceEventEmitter, @InjectMessageQueue(MessageQueue.webhookQueue) private readonly messageQueueService: MessageQueueService, ) {} @LogExecutionTime() async findOne< ObjectRecord extends IRecord = IRecord, Filter extends RecordFilter = RecordFilter, >( args: FindOneResolverArgs, options: WorkspaceQueryRunnerOptions, ): Promise { const graphqlQueryFindOneResolverService = new GraphqlQueryFindOneResolverService(this.twentyORMGlobalManager); const { authContext, objectMetadataItem } = options; if (!args.filter || Object.keys(args.filter).length === 0) { throw new WorkspaceQueryRunnerException( 'Missing filter argument', WorkspaceQueryRunnerExceptionCode.INVALID_QUERY_INPUT, ); } const hookedArgs = await this.workspaceQueryHookService.executePreQueryHooks( authContext, objectMetadataItem.nameSingular, 'findOne', args, ); const computedArgs = (await this.queryRunnerArgsFactory.create( hookedArgs, options, ResolverArgsType.FindOne, )) as FindOneResolverArgs; return graphqlQueryFindOneResolverService.findOne(computedArgs, options); } @LogExecutionTime() async findMany< ObjectRecord extends IRecord = IRecord, Filter extends RecordFilter = RecordFilter, OrderBy extends RecordOrderBy = RecordOrderBy, >( args: FindManyResolverArgs, options: WorkspaceQueryRunnerOptions, ): Promise> { const graphqlQueryFindManyResolverService = new GraphqlQueryFindManyResolverService(this.twentyORMGlobalManager); const { authContext, objectMetadataItem } = options; const hookedArgs = await this.workspaceQueryHookService.executePreQueryHooks( authContext, objectMetadataItem.nameSingular, 'findMany', args, ); const computedArgs = (await this.queryRunnerArgsFactory.create( hookedArgs, options, ResolverArgsType.FindMany, )) as FindManyResolverArgs; return graphqlQueryFindManyResolverService.findMany(computedArgs, options); } @LogExecutionTime() async createOne( args: CreateOneResolverArgs>, options: WorkspaceQueryRunnerOptions, ): Promise { const graphqlQueryCreateManyResolverService = new GraphqlQueryCreateManyResolverService(this.twentyORMGlobalManager); const { authContext, objectMetadataItem } = options; assertMutationNotOnRemoteObject(objectMetadataItem); if (args.data.id) { assertIsValidUuid(args.data.id); } const createManyArgs = { data: [args.data], upsert: args.upsert, } as CreateManyResolverArgs; const hookedArgs = await this.workspaceQueryHookService.executePreQueryHooks( authContext, objectMetadataItem.nameSingular, 'createMany', createManyArgs, ); const computedArgs = (await this.queryRunnerArgsFactory.create( hookedArgs, options, ResolverArgsType.CreateMany, )) as CreateManyResolverArgs; const results = (await graphqlQueryCreateManyResolverService.createMany( computedArgs, options, )) as ObjectRecord[]; await this.triggerWebhooks( results, CallWebhookJobsJobOperation.create, options, ); this.emitCreateEvents( results, authContext, objectMetadataItem, ); return results?.[0] as ObjectRecord; } @LogExecutionTime() async createMany( args: CreateManyResolverArgs>, options: WorkspaceQueryRunnerOptions, ): Promise { const graphqlQueryCreateManyResolverService = new GraphqlQueryCreateManyResolverService(this.twentyORMGlobalManager); const { authContext, objectMetadataItem } = options; assertMutationNotOnRemoteObject(objectMetadataItem); args.data.forEach((record) => { if (record?.id) { assertIsValidUuid(record.id); } }); const hookedArgs = await this.workspaceQueryHookService.executePreQueryHooks( authContext, objectMetadataItem.nameSingular, 'createMany', args, ); const computedArgs = (await this.queryRunnerArgsFactory.create( hookedArgs, options, ResolverArgsType.CreateMany, )) as CreateManyResolverArgs; const results = (await graphqlQueryCreateManyResolverService.createMany( computedArgs, options, )) as ObjectRecord[]; await this.workspaceQueryHookService.executePostQueryHooks( authContext, objectMetadataItem.nameSingular, 'createMany', results, ); await this.triggerWebhooks( results, CallWebhookJobsJobOperation.create, options, ); this.emitCreateEvents( results, authContext, objectMetadataItem, ); return results; } private emitCreateEvents( records: BaseRecord[], authContext: AuthContext, objectMetadataItem: ObjectMetadataInterface, ) { this.workspaceEventEmitter.emit( `${objectMetadataItem.nameSingular}.created`, records.map( (record) => ({ userId: authContext.user?.id, recordId: record.id, objectMetadata: objectMetadataItem, properties: { after: record, }, }) satisfies ObjectRecordCreateEvent, ), authContext.workspace.id, ); } private async triggerWebhooks( jobsData: Record[] | undefined, operation: CallWebhookJobsJobOperation, options: WorkspaceQueryRunnerOptions, ) { if (!Array.isArray(jobsData)) { return; } jobsData.forEach((jobData) => { this.messageQueueService.add( CallWebhookJobsJob.name, { record: jobData, workspaceId: options.authContext.workspace.id, operation, objectMetadataItem: options.objectMetadataItem, }, { retryLimit: 3 }, ); }); } @LogExecutionTime() async destroyOne( args: DestroyOneResolverArgs, options: WorkspaceQueryRunnerOptions, ): Promise { const graphqlQueryDestroyOneResolverService = new GraphqlQueryDestroyOneResolverService(this.twentyORMGlobalManager); return graphqlQueryDestroyOneResolverService.destroyOne(args, options); } }