feat: mutualize CRUD tools between workflows and AI (#14996)

Closes [#1662](https://github.com/twentyhq/core-team-issues/issues/1662)
This commit is contained in:
Abdul Rahman
2025-10-10 00:55:34 +05:30
committed by GitHub
parent b508bd8f9d
commit 5adc9fe9b2
29 changed files with 899 additions and 1499 deletions
@@ -11,7 +11,7 @@ import { ToolAdapterService } from 'src/engine/core-modules/ai/services/tool-ada
import { ToolService } from 'src/engine/core-modules/ai/services/tool.service';
import { TokenModule } from 'src/engine/core-modules/auth/token/token.module';
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
import { RecordTransformerModule } from 'src/engine/core-modules/record-transformer/record-transformer.module';
import { RecordCrudModule } from 'src/engine/core-modules/record-crud/record-crud.module';
import { ToolRegistryService } from 'src/engine/core-modules/tool/services/tool-registry.service';
import { SendEmailTool } from 'src/engine/core-modules/tool/tools/send-email-tool/send-email-tool';
import { ObjectMetadataModule } from 'src/engine/metadata-modules/object-metadata/object-metadata.module';
@@ -29,7 +29,7 @@ import { MessagingModule } from 'src/modules/messaging/messaging.module';
TypeOrmModule.forFeature([RoleEntity]),
TokenModule,
FeatureFlagModule,
RecordTransformerModule,
RecordCrudModule,
ObjectMetadataModule,
WorkspacePermissionsCacheModule,
WorkspaceCacheStorageModule,
@@ -1,6 +1,10 @@
import { Test } from '@nestjs/testing';
import { ToolService } from 'src/engine/core-modules/ai/services/tool.service';
import { CreateRecordService } from 'src/engine/core-modules/record-crud/services/create-record.service';
import { DeleteRecordService } from 'src/engine/core-modules/record-crud/services/delete-record.service';
import { FindRecordsService } from 'src/engine/core-modules/record-crud/services/find-records.service';
import { UpdateRecordService } from 'src/engine/core-modules/record-crud/services/update-record.service';
import { RecordInputTransformerService } from 'src/engine/core-modules/record-transformer/services/record-input-transformer.service';
import { ObjectMetadataService } from 'src/engine/metadata-modules/object-metadata/object-metadata.service';
import { WorkspacePermissionsCacheService } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service';
@@ -24,10 +28,7 @@ describe('ToolService', () => {
const roleId = 'role_1';
let service: ToolService;
let ormManager: TwentyORMGlobalManager;
let permissionsCacheService: WorkspacePermissionsCacheService;
let transformer: RecordInputTransformerService;
let workspaceCache: WorkspaceCacheStorageService;
const testObject = getMockObjectMetadataEntity({
workspaceId: '',
@@ -102,14 +103,27 @@ describe('ToolService', () => {
}),
},
},
{
provide: CreateRecordService,
useValue: { execute: jest.fn() },
},
{
provide: UpdateRecordService,
useValue: { execute: jest.fn() },
},
{
provide: DeleteRecordService,
useValue: { execute: jest.fn() },
},
{
provide: FindRecordsService,
useValue: { execute: jest.fn() },
},
],
}).compile();
service = moduleRef.get(ToolService);
ormManager = moduleRef.get(TwentyORMGlobalManager);
permissionsCacheService = moduleRef.get(WorkspacePermissionsCacheService);
transformer = moduleRef.get(RecordInputTransformerService);
workspaceCache = moduleRef.get(WorkspaceCacheStorageService);
});
describe('listTools', () => {
@@ -133,68 +147,6 @@ describe('ToolService', () => {
});
});
describe('createRecord', () => {
it('should create a record successfully', async () => {
const record = { id: 'r1', name: 'Test' };
mockRepo.save.mockResolvedValue(record);
const result = await service.createRecord(
'testObject',
{ name: 'Test' },
workspaceId,
roleId,
);
expect(result.success).toBe(true);
expect(result.result).toEqual(record);
expect(ormManager.getRepositoryForWorkspace).toHaveBeenCalledWith(
workspaceId,
'testObject',
{ roleId },
);
expect(workspaceCache.getObjectMetadataMapsOrThrow).toHaveBeenCalledWith(
workspaceId,
);
expect(transformer.process).toHaveBeenCalled();
expect(mockRepo.save).toHaveBeenCalledWith({ name: 'Test' });
});
});
describe('updateRecord', () => {
it('should return error when id is missing', async () => {
const result = await (service as any).updateRecord(
'testObject',
{ name: 'No ID' },
workspaceId,
roleId,
);
expect(result.success).toBe(false);
expect(result.error).toBe('Record ID is required for update');
});
});
describe('findRecords', () => {
it('should return records and count', async () => {
const records = [{ id: 'a' }, { id: 'b' }];
mockRepo.find.mockResolvedValue(records);
const result = await (service as any).findRecords(
'testObject',
{},
workspaceId,
roleId,
);
expect(result.success).toBe(true);
expect(result.result.records).toEqual(records);
expect(result.result.count).toBe(2);
expect(mockRepo.find).toHaveBeenCalled();
});
});
describe('softDeleteManyRecords', () => {
it('should error when filter is invalid', async () => {
const result = await (service as any).softDeleteManyRecords(
@@ -2,8 +2,10 @@ import { Injectable } from '@nestjs/common';
import { type ToolSet } from 'ai';
import { buildWhereConditions } from 'src/engine/core-modules/ai/utils/find-records-filters.utils';
import { RecordInputTransformerService } from 'src/engine/core-modules/record-transformer/services/record-input-transformer.service';
import { CreateRecordService } from 'src/engine/core-modules/record-crud/services/create-record.service';
import { DeleteRecordService } from 'src/engine/core-modules/record-crud/services/delete-record.service';
import { FindRecordsService } from 'src/engine/core-modules/record-crud/services/find-records.service';
import { UpdateRecordService } from 'src/engine/core-modules/record-crud/services/update-record.service';
import {
generateBulkDeleteToolSchema,
generateFindOneToolSchema,
@@ -13,10 +15,8 @@ import {
} from 'src/engine/metadata-modules/agent/utils/agent-tool-schema.utils';
import { isWorkflowRunObject } from 'src/engine/metadata-modules/agent/utils/is-workflow-run-object.util';
import { ObjectMetadataService } from 'src/engine/metadata-modules/object-metadata/object-metadata.service';
import { getObjectMetadataMapItemByNameSingular } from 'src/engine/metadata-modules/utils/get-object-metadata-map-item-by-name-singular.util';
import { WorkspacePermissionsCacheService } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkspaceCacheStorageService } from 'src/engine/workspace-cache-storage/workspace-cache-storage.service';
@Injectable()
export class ToolService {
@@ -24,8 +24,10 @@ export class ToolService {
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
private readonly objectMetadataService: ObjectMetadataService,
protected readonly workspacePermissionsCacheService: WorkspacePermissionsCacheService,
private readonly recordInputTransformerService: RecordInputTransformerService,
private readonly workspaceCacheStorageService: WorkspaceCacheStorageService,
private readonly createRecordService: CreateRecordService,
private readonly updateRecordService: UpdateRecordService,
private readonly deleteRecordService: DeleteRecordService,
private readonly findRecordsService: FindRecordsService,
) {}
async listTools(roleId: string, workspaceId: string): Promise<ToolSet> {
@@ -63,12 +65,12 @@ export class ToolService {
description: `Create a new ${objectMetadata.labelSingular} record. Provide all required fields and any optional fields you want to set. The system will automatically handle timestamps and IDs. Returns the created record with all its data.`,
inputSchema: getRecordInputSchema(objectMetadata),
execute: async (parameters) => {
return this.createRecord(
objectMetadata.nameSingular,
parameters.input,
return this.createRecordService.execute({
objectName: objectMetadata.nameSingular,
objectRecord: parameters.input,
workspaceId,
roleId,
);
});
},
};
@@ -76,12 +78,15 @@ export class ToolService {
description: `Update an existing ${objectMetadata.labelSingular} record. Provide the record ID and only the fields you want to change. Unspecified fields will remain unchanged. Returns the updated record with all current data.`,
inputSchema: getRecordInputSchema(objectMetadata),
execute: async (parameters) => {
return this.updateRecord(
objectMetadata.nameSingular,
parameters.input,
const { id, ...objectRecord } = parameters.input;
return this.updateRecordService.execute({
objectName: objectMetadata.nameSingular,
objectRecordId: id,
objectRecord,
workspaceId,
roleId,
);
});
},
};
}
@@ -91,12 +96,16 @@ export class ToolService {
description: `Search for ${objectMetadata.labelSingular} records using flexible filtering criteria. Supports exact matches, pattern matching, ranges, and null checks. Use limit/offset for pagination. Returns an array of matching records with their full data.`,
inputSchema: generateFindToolSchema(objectMetadata),
execute: async (parameters) => {
return this.findRecords(
objectMetadata.nameSingular,
parameters.input,
const { limit, offset, ...filter } = parameters.input;
return this.findRecordsService.execute({
objectName: objectMetadata.nameSingular,
filter,
limit,
offset,
workspaceId,
roleId,
);
});
},
};
@@ -104,12 +113,13 @@ export class ToolService {
description: `Retrieve a single ${objectMetadata.labelSingular} record by its unique ID. Use this when you know the exact record ID and need the complete record data. Returns the full record or an error if not found.`,
inputSchema: generateFindOneToolSchema(),
execute: async (parameters) => {
return this.findOneRecord(
objectMetadata.nameSingular,
parameters.input,
return this.findRecordsService.execute({
objectName: objectMetadata.nameSingular,
filter: { id: { eq: parameters.input.id } },
limit: 1,
workspaceId,
roleId,
);
});
},
};
}
@@ -119,12 +129,13 @@ export class ToolService {
description: `Soft delete a ${objectMetadata.labelSingular} record by marking it as deleted. The record remains in the database but is hidden from normal queries. This is reversible and preserves all data. Use this for temporary removal.`,
inputSchema: generateSoftDeleteToolSchema(),
execute: async (parameters) => {
return this.softDeleteRecord(
objectMetadata.nameSingular,
parameters.input,
return this.deleteRecordService.execute({
objectName: objectMetadata.nameSingular,
objectRecordId: parameters.input.id,
workspaceId,
roleId,
);
soft: true,
});
},
};
@@ -146,340 +157,6 @@ export class ToolService {
return tools;
}
private async findRecords(
objectName: string,
parameters: Record<string, unknown>,
workspaceId: string,
roleId: string,
) {
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
{ roleId },
);
const { limit = 100, offset = 0, ...searchCriteria } = parameters;
const whereConditions = buildWhereConditions(searchCriteria);
const records = await repository.find({
where: whereConditions,
take: limit as number,
skip: offset as number,
order: { createdAt: 'DESC' },
});
return {
success: true,
message: `Found ${records.length} ${objectName} records`,
result: {
records,
count: records.length,
},
};
} catch (error) {
return {
success: false,
message: `Failed to find ${objectName} records`,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
}
private async findOneRecord(
objectName: string,
parameters: Record<string, unknown>,
workspaceId: string,
roleId: string,
) {
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
{ roleId },
);
const { id } = parameters;
if (!id || typeof id !== 'string') {
return {
success: false,
message: `Failed to find ${objectName}: Record ID is required`,
error: 'Record ID is required',
};
}
const record = await repository.findOne({
where: { id },
});
if (!record) {
return {
success: false,
message: `Failed to find ${objectName}: Record with ID ${id} not found`,
error: 'Record not found',
};
}
return {
success: true,
message: `Found ${objectName} record`,
result: record,
};
} catch (error) {
return {
success: false,
message: `Failed to find ${objectName} record`,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
}
async createRecord(
objectName: string,
parameters: Record<string, unknown>,
workspaceId: string,
roleId: string,
) {
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
{ roleId },
);
const objectMetadataMaps =
await this.workspaceCacheStorageService.getObjectMetadataMapsOrThrow(
workspaceId,
);
const objectMetadataItemWithFieldsMaps =
getObjectMetadataMapItemByNameSingular(objectMetadataMaps, objectName);
if (!objectMetadataItemWithFieldsMaps) {
return {
success: false,
message: `Failed to create ${objectName}: Object metadata not found`,
error: 'Object metadata not found',
};
}
const transformedCreateData =
await this.recordInputTransformerService.process({
recordInput: parameters,
objectMetadataMapItem: objectMetadataItemWithFieldsMaps,
});
const createdRecord = await repository.save(transformedCreateData);
return {
success: true,
message: `Successfully created ${objectName}`,
result: createdRecord,
};
} catch (error) {
return {
success: false,
message: `Failed to create ${objectName}`,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
}
private async updateRecord(
objectName: string,
parameters: Record<string, unknown>,
workspaceId: string,
roleId: string,
) {
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
{ roleId },
);
const { id, ...updateData } = parameters;
if (!id || typeof id !== 'string') {
return {
success: false,
message: `Failed to update ${objectName}: Record ID is required`,
error: 'Record ID is required for update',
};
}
const existingRecord = await repository.findOne({
where: { id },
});
if (!existingRecord) {
return {
success: false,
message: `Failed to update ${objectName}: Record with ID ${id} not found`,
error: 'Record not found',
};
}
const objectMetadataMaps =
await this.workspaceCacheStorageService.getObjectMetadataMapsOrThrow(
workspaceId,
);
const objectMetadataItemWithFieldsMaps =
getObjectMetadataMapItemByNameSingular(objectMetadataMaps, objectName);
if (!objectMetadataItemWithFieldsMaps) {
return {
success: false,
message: `Failed to update ${objectName}: Object metadata not found`,
error: 'Object metadata not found',
};
}
const transformedUpdateData =
await this.recordInputTransformerService.process({
recordInput: updateData,
objectMetadataMapItem: objectMetadataItemWithFieldsMaps,
});
await repository.update(id as string, transformedUpdateData);
const updatedRecord = await repository.findOne({
where: { id: id as string },
});
if (!updatedRecord) {
return {
success: false,
message: `Failed to update ${objectName}: Could not retrieve updated record`,
error: 'Failed to retrieve updated record',
};
}
return {
success: true,
message: `Successfully updated ${objectName}`,
result: updatedRecord,
};
} catch (error) {
return {
success: false,
message: `Failed to update ${objectName}`,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
}
private async softDeleteRecord(
objectName: string,
parameters: Record<string, unknown>,
workspaceId: string,
roleId: string,
) {
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
{ roleId },
);
const { id } = parameters;
if (!id || typeof id !== 'string') {
return {
success: false,
message: `Failed to soft delete ${objectName}: Record ID is required`,
error: 'Record ID is required for soft delete',
};
}
const existingRecord = await repository.findOne({
where: { id },
});
if (!existingRecord) {
return {
success: false,
message: `Failed to soft delete ${objectName}: Record with ID ${id} not found`,
error: 'Record not found',
};
}
await repository.softDelete(id);
return {
success: true,
message: `Successfully soft deleted ${objectName}`,
result: { id },
};
} catch (error) {
return {
success: false,
message: `Failed to soft delete ${objectName}`,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
}
private async _destroyRecord(
objectName: string,
parameters: Record<string, unknown>,
workspaceId: string,
roleId: string,
) {
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
{ roleId },
);
const { id } = parameters;
if (!id || typeof id !== 'string') {
return {
success: false,
message: `Failed to destroy ${objectName}: Record ID is required`,
error: 'Record ID is required for destroy',
};
}
const existingRecord = await repository.findOne({
where: { id },
});
if (!existingRecord) {
return {
success: false,
message: `Failed to destroy ${objectName}: Record with ID ${id} not found`,
error: 'Record not found',
};
}
await repository.remove(existingRecord);
return {
success: true,
message: `Successfully destroyed ${objectName}`,
result: { id },
};
} catch (error) {
return {
success: false,
message: `Failed to destroy ${objectName}`,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
}
private async softDeleteManyRecords(
objectName: string,
parameters: Record<string, unknown>,
@@ -545,70 +222,4 @@ export class ToolService {
};
}
}
private async _destroyManyRecords(
objectName: string,
parameters: Record<string, unknown>,
workspaceId: string,
roleId: string,
) {
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
{ roleId },
);
const { filter } = parameters;
if (!filter || typeof filter !== 'object' || !('id' in filter)) {
return {
success: false,
message: `Failed to destroy many ${objectName}: Filter with record IDs is required`,
error: 'Filter with record IDs is required for bulk destroy',
};
}
const idFilter = filter.id as Record<string, unknown>;
const recordIds = idFilter.in as string[];
if (!Array.isArray(recordIds) || recordIds.length === 0) {
return {
success: false,
message: `Failed to destroy many ${objectName}: At least one record ID is required`,
error: 'At least one record ID is required for bulk destroy',
};
}
const existingRecords = await repository.find({
where: { id: { in: recordIds } },
});
if (existingRecords.length === 0) {
return {
success: false,
message: `Failed to destroy many ${objectName}: No records found with the provided IDs`,
error: 'No records found to destroy',
};
}
await repository.delete({ id: { in: recordIds } });
return {
success: true,
message: `Successfully destroyed ${existingRecords.length} ${objectName} records`,
result: {
count: existingRecords.length,
destroyedIds: recordIds,
},
};
} catch (error) {
return {
success: false,
message: `Failed to destroy many ${objectName}`,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
}
}
@@ -1,6 +1,6 @@
import { faker } from '@faker-js/faker';
import { type EachTestingContext } from 'twenty-shared/testing';
import { FieldMetadataType } from 'twenty-shared/types';
import { faker } from '@faker-js/faker';
import { NumberDataType } from 'src/engine/metadata-modules/field-metadata/interfaces/field-metadata-settings.interface';
@@ -61,6 +61,7 @@ describe('computeSchemaComponents', () => {
"EMAIL",
"CALENDAR",
"WORKFLOW",
"AGENT",
"API",
"IMPORT",
"MANUAL",
@@ -292,6 +293,7 @@ describe('computeSchemaComponents', () => {
"EMAIL",
"CALENDAR",
"WORKFLOW",
"AGENT",
"API",
"IMPORT",
"MANUAL",
@@ -562,6 +564,7 @@ describe('computeSchemaComponents', () => {
"EMAIL",
"CALENDAR",
"WORKFLOW",
"AGENT",
"API",
"IMPORT",
"MANUAL",
@@ -0,0 +1,19 @@
import { CustomException } from 'src/utils/custom-exception';
export class RecordCrudException extends CustomException {
code: RecordCrudExceptionCode;
constructor(message: string, code: RecordCrudExceptionCode) {
super(message, code);
}
}
export enum RecordCrudExceptionCode {
INVALID_REQUEST = 'INVALID_REQUEST',
WORKSPACE_ID_NOT_FOUND = 'WORKSPACE_ID_NOT_FOUND',
OBJECT_NOT_FOUND = 'OBJECT_NOT_FOUND',
RECORD_NOT_FOUND = 'RECORD_NOT_FOUND',
RECORD_CREATION_FAILED = 'RECORD_CREATION_FAILED',
RECORD_UPDATE_FAILED = 'RECORD_UPDATE_FAILED',
RECORD_DELETION_FAILED = 'RECORD_DELETION_FAILED',
QUERY_FAILED = 'QUERY_FAILED',
}
@@ -0,0 +1,32 @@
import { Module } from '@nestjs/common';
import { CreateRecordService } from 'src/engine/core-modules/record-crud/services/create-record.service';
import { DeleteRecordService } from 'src/engine/core-modules/record-crud/services/delete-record.service';
import { FindRecordsService } from 'src/engine/core-modules/record-crud/services/find-records.service';
import { UpdateRecordService } from 'src/engine/core-modules/record-crud/services/update-record.service';
import { RecordPositionModule } from 'src/engine/core-modules/record-position/record-position.module';
import { RecordTransformerModule } from 'src/engine/core-modules/record-transformer/record-transformer.module';
import { TwentyORMModule } from 'src/engine/twenty-orm/twenty-orm.module';
import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module';
@Module({
imports: [
TwentyORMModule,
RecordPositionModule,
RecordTransformerModule,
WorkflowCommonModule,
],
providers: [
CreateRecordService,
UpdateRecordService,
DeleteRecordService,
FindRecordsService,
],
exports: [
CreateRecordService,
UpdateRecordService,
DeleteRecordService,
FindRecordsService,
],
})
export class RecordCrudModule {}
@@ -0,0 +1,122 @@
import { Injectable, Logger } from '@nestjs/common';
import { isDefined } from 'class-validator';
import { canObjectBeManagedByWorkflow } from 'twenty-shared/workflow';
import {
RecordCrudException,
RecordCrudExceptionCode,
} from 'src/engine/core-modules/record-crud/exceptions/record-crud.exception';
import { type CreateRecordParams } from 'src/engine/core-modules/record-crud/types/create-record-params.type';
import { RecordPositionService } from 'src/engine/core-modules/record-position/services/record-position.service';
import { RecordInputTransformerService } from 'src/engine/core-modules/record-transformer/services/record-input-transformer.service';
import { type ToolOutput } from 'src/engine/core-modules/tool/types/tool-output.type';
import { FieldActorSource } from 'src/engine/metadata-modules/field-metadata/composite-types/actor.composite-type';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
@Injectable()
// eslint-disable-next-line @nx/workspace-inject-workspace-repository
export class CreateRecordService {
private readonly logger = new Logger(CreateRecordService.name);
constructor(
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
private readonly recordPositionService: RecordPositionService,
private readonly recordInputTransformerService: RecordInputTransformerService,
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
) {}
async execute(params: CreateRecordParams): Promise<ToolOutput> {
const { objectName, objectRecord, workspaceId, roleId } = params;
if (!workspaceId) {
return {
success: false,
message: 'Failed to create record: Workspace ID is required',
error: 'Workspace ID not found',
};
}
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
roleId ? { roleId } : { shouldBypassPermissionChecks: true },
);
const { objectMetadataItemWithFieldsMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
objectName,
workspaceId,
);
if (
!canObjectBeManagedByWorkflow({
nameSingular: objectMetadataItemWithFieldsMaps.nameSingular,
isSystem: objectMetadataItemWithFieldsMaps.isSystem,
})
) {
throw new RecordCrudException(
'Failed to create: Object cannot be created by workflow',
RecordCrudExceptionCode.INVALID_REQUEST,
);
}
const position = await this.recordPositionService.buildRecordPosition({
value: 'first',
objectMetadata: objectMetadataItemWithFieldsMaps,
workspaceId,
});
const validObjectRecord = Object.fromEntries(
Object.entries(objectRecord).filter(([key]) =>
isDefined(objectMetadataItemWithFieldsMaps.fieldIdByName[key]),
),
);
const transformedObjectRecord =
await this.recordInputTransformerService.process({
recordInput: validObjectRecord,
objectMetadataMapItem: objectMetadataItemWithFieldsMaps,
});
const insertResult = await repository.insert({
...transformedObjectRecord,
position,
createdBy: {
source: roleId ? FieldActorSource.AGENT : FieldActorSource.WORKFLOW,
name: roleId ? 'Agent' : 'Workflow',
},
});
const [createdRecord] = insertResult.generatedMaps;
this.logger.log(`Record created successfully in ${objectName}`);
return {
success: true,
message: `Record created successfully in ${objectName}`,
result: createdRecord,
};
} catch (error) {
if (error instanceof RecordCrudException) {
return {
success: false,
message: `Failed to create record in ${objectName}`,
error: error.message,
};
}
this.logger.error(`Failed to create record: ${error}`);
return {
success: false,
message: `Failed to create record in ${objectName}`,
error:
error instanceof Error ? error.message : 'Failed to create record',
};
}
}
}
@@ -0,0 +1,137 @@
import { Injectable, Logger } from '@nestjs/common';
import { isDefined, isValidUuid } from 'twenty-shared/utils';
import { canObjectBeManagedByWorkflow } from 'twenty-shared/workflow';
import {
RecordCrudException,
RecordCrudExceptionCode,
} from 'src/engine/core-modules/record-crud/exceptions/record-crud.exception';
import { type DeleteRecordParams } from 'src/engine/core-modules/record-crud/types/delete-record-params.type';
import { type ToolOutput } from 'src/engine/core-modules/tool/types/tool-output.type';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
@Injectable()
// eslint-disable-next-line @nx/workspace-inject-workspace-repository
export class DeleteRecordService {
private readonly logger = new Logger(DeleteRecordService.name);
constructor(
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
) {}
async execute(params: DeleteRecordParams): Promise<ToolOutput> {
const {
objectName,
objectRecordId,
workspaceId,
roleId,
soft = true,
} = params;
if (!workspaceId) {
return {
success: false,
message: 'Failed to delete record: Workspace ID is required',
error: 'Workspace ID not found',
};
}
if (!isDefined(objectRecordId) || !isValidUuid(objectRecordId)) {
return {
success: false,
message: 'Failed to delete: Object record ID must be a valid UUID',
error: 'Invalid object record ID',
};
}
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
roleId ? { roleId } : { shouldBypassPermissionChecks: true },
);
const { objectMetadataItemWithFieldsMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
objectName,
workspaceId,
);
if (
!canObjectBeManagedByWorkflow({
nameSingular: objectMetadataItemWithFieldsMaps.nameSingular,
isSystem: objectMetadataItemWithFieldsMaps.isSystem,
})
) {
throw new RecordCrudException(
'Failed to delete: Object cannot be deleted by workflow',
RecordCrudExceptionCode.INVALID_REQUEST,
);
}
const objectRecord = await repository.findOne({
where: {
id: objectRecordId,
},
});
if (!objectRecord) {
throw new RecordCrudException(
`Failed to delete: Record ${objectName} with id ${objectRecordId} not found`,
RecordCrudExceptionCode.RECORD_NOT_FOUND,
);
}
if (soft) {
const columnsToReturnForSoftDelete: string[] = [];
await repository.softDelete(
objectRecordId,
undefined,
columnsToReturnForSoftDelete,
);
this.logger.log(`Record soft deleted successfully from ${objectName}`);
return {
success: true,
message: `Record soft deleted successfully from ${objectName}`,
result: objectRecord,
};
} else {
await repository.remove(objectRecord);
this.logger.log(
`Record permanently deleted successfully from ${objectName}`,
);
return {
success: true,
message: `Record permanently deleted successfully from ${objectName}`,
result: { id: objectRecordId },
};
}
} catch (error) {
if (error instanceof RecordCrudException) {
return {
success: false,
message: `Failed to delete record from ${objectName}`,
error: error.message,
};
}
this.logger.error(`Failed to delete record: ${error}`);
return {
success: false,
message: `Failed to delete record from ${objectName}`,
error:
error instanceof Error ? error.message : 'Failed to delete record',
};
}
}
}
@@ -0,0 +1,185 @@
import { Injectable, Logger } from '@nestjs/common';
import { QUERY_MAX_RECORDS } from 'twenty-shared/constants';
import { type ObjectLiteral } from 'typeorm';
import {
type ObjectRecordFilter,
type ObjectRecordOrderBy,
OrderByDirection,
} from 'src/engine/api/graphql/workspace-query-builder/interfaces/object-record.interface';
import { GraphqlQueryParser } from 'src/engine/api/graphql/graphql-query-runner/graphql-query-parsers/graphql-query.parser';
import { type FindRecordsParams } from 'src/engine/core-modules/record-crud/types/find-records-params.type';
import { FindRecordsResult } from 'src/engine/core-modules/record-crud/types/find-records-result.type';
import { type ToolOutput } from 'src/engine/core-modules/tool/types/tool-output.type';
import { type WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
@Injectable()
// eslint-disable-next-line @nx/workspace-inject-workspace-repository
export class FindRecordsService {
private readonly logger = new Logger(FindRecordsService.name);
constructor(
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
) {}
async execute(
params: FindRecordsParams,
): Promise<ToolOutput<FindRecordsResult>> {
const {
objectName,
filter,
orderBy,
limit,
offset = 0,
workspaceId,
roleId,
} = params;
if (!workspaceId) {
return {
success: false,
message: 'Failed to find records: Workspace ID is required',
error: 'Workspace ID not found',
};
}
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
roleId ? { roleId } : { shouldBypassPermissionChecks: true },
);
const { objectMetadataItemWithFieldsMaps, objectMetadataMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
objectName,
workspaceId,
);
const graphqlQueryParser = new GraphqlQueryParser(
objectMetadataItemWithFieldsMaps,
objectMetadataMaps,
);
const records = await this.getObjectRecords({
objectName,
filter,
orderBy,
limit,
offset,
repository,
graphqlQueryParser,
});
const totalCount = await this.getTotalCount({
objectName,
filter,
repository,
graphqlQueryParser,
});
this.logger.log(`Found ${records.length} records in ${objectName}`);
return {
success: true,
message: `Found ${records.length} ${objectName} records`,
result: {
records,
count: totalCount,
},
};
} catch (error) {
this.logger.error(`Failed to find records: ${error}`);
return {
success: false,
message: `Failed to find ${objectName} records`,
error:
error instanceof Error ? error.message : 'Failed to find records',
};
}
}
private async getObjectRecords<T extends ObjectLiteral>({
objectName,
filter,
orderBy,
limit,
offset,
repository,
graphqlQueryParser,
}: {
objectName: string;
filter:
| Partial<ObjectRecordFilter>
| Partial<ObjectRecordFilter>[]
| undefined;
orderBy: Partial<ObjectRecordOrderBy> | undefined;
limit: number | undefined;
offset: number;
repository: WorkspaceRepository<T>;
graphqlQueryParser: GraphqlQueryParser;
}): Promise<T[]> {
const queryBuilder = repository.createQueryBuilder(objectName);
const withFilterQueryBuilder = graphqlQueryParser.applyFilterToBuilder(
queryBuilder,
objectName,
filter ?? {},
);
const orderByWithIdCondition: ObjectRecordOrderBy = [
...(orderBy ?? []).filter((item) => item !== undefined),
{ id: OrderByDirection.AscNullsFirst },
];
const withOrderByQueryBuilder = graphqlQueryParser.applyOrderToBuilder(
withFilterQueryBuilder,
orderByWithIdCondition,
objectName,
false,
);
return withOrderByQueryBuilder
.skip(offset)
.take(limit ? Math.min(limit, QUERY_MAX_RECORDS) : QUERY_MAX_RECORDS)
.getMany();
}
private async getTotalCount({
objectName,
filter,
repository,
graphqlQueryParser,
}: {
objectName: string;
filter:
| Partial<ObjectRecordFilter>
| Partial<ObjectRecordFilter>[]
| undefined;
repository: WorkspaceRepository<ObjectLiteral>;
graphqlQueryParser: GraphqlQueryParser;
}): Promise<number> {
const countQueryBuilder = repository.createQueryBuilder(objectName);
const withFilterCountQueryBuilder = graphqlQueryParser.applyFilterToBuilder(
countQueryBuilder,
objectName,
filter ?? {},
);
const withDeletedCountQueryBuilder =
graphqlQueryParser.applyDeletedAtToBuilder(
withFilterCountQueryBuilder,
filter ?? {},
);
return withDeletedCountQueryBuilder.getCount();
}
}
@@ -0,0 +1,160 @@
import { Injectable, Logger } from '@nestjs/common';
import deepEqual from 'deep-equal';
import { isDefined, isValidUuid } from 'twenty-shared/utils';
import { canObjectBeManagedByWorkflow } from 'twenty-shared/workflow';
import {
RecordCrudException,
RecordCrudExceptionCode,
} from 'src/engine/core-modules/record-crud/exceptions/record-crud.exception';
import { type UpdateRecordParams } from 'src/engine/core-modules/record-crud/types/update-record-params.type';
import { RecordInputTransformerService } from 'src/engine/core-modules/record-transformer/services/record-input-transformer.service';
import { type ToolOutput } from 'src/engine/core-modules/tool/types/tool-output.type';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
@Injectable()
// eslint-disable-next-line @nx/workspace-inject-workspace-repository
export class UpdateRecordService {
private readonly logger = new Logger(UpdateRecordService.name);
constructor(
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
private readonly recordInputTransformerService: RecordInputTransformerService,
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
) {}
async execute(params: UpdateRecordParams): Promise<ToolOutput> {
const {
objectName,
objectRecordId,
objectRecord,
fieldsToUpdate,
workspaceId,
roleId,
} = params;
if (!workspaceId) {
return {
success: false,
message: 'Failed to update record: Workspace ID is required',
error: 'Workspace ID not found',
};
}
if (!isDefined(objectRecordId) || !isValidUuid(objectRecordId)) {
return {
success: false,
message: 'Failed to update: Object record ID must be a valid UUID',
error: 'Invalid object record ID',
};
}
try {
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
objectName,
roleId ? { roleId } : { shouldBypassPermissionChecks: true },
);
const previousObjectRecord = await repository.findOne({
where: {
id: objectRecordId,
},
});
if (!previousObjectRecord) {
throw new RecordCrudException(
`Failed to update: Record ${objectName} with id ${objectRecordId} not found`,
RecordCrudExceptionCode.RECORD_NOT_FOUND,
);
}
const fieldsToUpdateArray = fieldsToUpdate || Object.keys(objectRecord);
if (fieldsToUpdateArray.length === 0) {
return {
success: true,
message: 'No fields to update',
result: previousObjectRecord,
};
}
const { objectMetadataItemWithFieldsMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
objectName,
workspaceId,
);
if (
!canObjectBeManagedByWorkflow({
nameSingular: objectMetadataItemWithFieldsMaps.nameSingular,
isSystem: objectMetadataItemWithFieldsMaps.isSystem,
})
) {
throw new RecordCrudException(
'Failed to update: Object cannot be updated by workflow',
RecordCrudExceptionCode.INVALID_REQUEST,
);
}
const objectRecordWithFilteredFields = Object.keys(objectRecord).reduce(
(acc, key) => {
if (fieldsToUpdateArray.includes(key)) {
return {
...acc,
[key]: objectRecord[key],
};
}
return acc;
},
{},
);
const transformedObjectRecord =
await this.recordInputTransformerService.process({
recordInput: objectRecordWithFilteredFields,
objectMetadataMapItem: objectMetadataItemWithFieldsMaps,
});
const updatedObjectRecord = {
...previousObjectRecord,
...objectRecordWithFilteredFields,
};
if (!deepEqual(updatedObjectRecord, previousObjectRecord)) {
await repository.update(objectRecordId, {
...transformedObjectRecord,
});
}
this.logger.log(`Record updated successfully in ${objectName}`);
return {
success: true,
message: `Record updated successfully in ${objectName}`,
result: updatedObjectRecord,
};
} catch (error) {
if (error instanceof RecordCrudException) {
return {
success: false,
message: `Failed to update record in ${objectName}`,
error: error.message,
};
}
this.logger.error(`Failed to update record: ${error}`);
return {
success: false,
message: `Failed to update record in ${objectName}`,
error:
error instanceof Error ? error.message : 'Failed to update record',
};
}
}
}
@@ -0,0 +1,8 @@
import { type ObjectRecordProperties } from 'src/engine/core-modules/record-crud/types/object-record-properties.type';
export type CreateRecordParams = {
objectName: string;
objectRecord: ObjectRecordProperties;
workspaceId: string;
roleId?: string;
};
@@ -0,0 +1,7 @@
export type DeleteRecordParams = {
objectName: string;
objectRecordId: string;
workspaceId: string;
roleId?: string;
soft?: boolean;
};
@@ -0,0 +1,18 @@
import {
type ObjectRecordFilter,
type ObjectRecordOrderBy,
} from 'src/engine/api/graphql/workspace-query-builder/interfaces/object-record.interface';
export type FindRecordsParams = {
objectName: string;
filter?:
| Record<string, unknown>
| Record<string, unknown>[]
| Partial<ObjectRecordFilter>
| Partial<ObjectRecordFilter>[];
orderBy?: Partial<ObjectRecordOrderBy>;
limit?: number;
offset?: number;
workspaceId: string;
roleId?: string;
};
@@ -0,0 +1,4 @@
export type FindRecordsResult = {
records: unknown[];
count: number;
};
@@ -0,0 +1,2 @@
// eslint-disable-next-line @typescript-eslint/no-explicit-any
export type ObjectRecordProperties = Record<string, any>;
@@ -0,0 +1,10 @@
import { type ObjectRecordProperties } from 'src/engine/core-modules/record-crud/types/object-record-properties.type';
export type UpdateRecordParams = {
objectName: string;
objectRecordId: string;
objectRecord: ObjectRecordProperties;
fieldsToUpdate?: string[];
workspaceId: string;
roleId?: string;
};
@@ -1,6 +1,6 @@
export type ToolOutput = {
export type ToolOutput<T = object> = {
success: boolean;
message: string;
error?: string;
result?: unknown;
result?: T;
};