Add TwentyORM query read timeout exception (#13603)

In this PR:
- adding a try / catch around all ORM internal methods save, insert,
upsert, findOne, ...
- leveraging this error to prevent messageChannels to get FAILED
- optimizing messaging BATCH_SIZE and THROTTLE threshold according to
local tests
- 

<img width="1510" height="851" alt="image"
src="https://github.com/user-attachments/assets/802fd933-caac-4291-9cde-34a1ddf59c06"
/>
This commit is contained in:
Charles Bochet
2025-08-04 17:16:16 +02:00
committed by GitHub
parent d958447bb6
commit 2eeddcddc6
14 changed files with 656 additions and 553 deletions
@@ -54,6 +54,7 @@ import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.
import { formatData } from 'src/engine/twenty-orm/utils/format-data.util';
import { formatResult } from 'src/engine/twenty-orm/utils/format-result.util';
import { getObjectMetadataFromEntityTarget } from 'src/engine/twenty-orm/utils/get-object-metadata-from-entity-target.util';
import { computeTwentyORMException } from 'src/engine/twenty-orm/error-handling/compute-twenty-orm-exception';
type PermissionOptions = {
shouldBypassPermissionChecks?: boolean;
@@ -1057,165 +1058,170 @@ export class WorkspaceEntityManager extends EntityManager {
| (SaveOptions & { reload: false }),
permissionOptions?: PermissionOptions,
): Promise<(T & Entity) | (T & Entity)[] | Entity | Entity[]> {
const permissionOptionsFromArgs =
maybeOptionsOrMaybePermissionOptions &&
('shouldBypassPermissionChecks' in maybeOptionsOrMaybePermissionOptions ||
'objectRecordsPermissions' in maybeOptionsOrMaybePermissionOptions)
try {
const permissionOptionsFromArgs =
maybeOptionsOrMaybePermissionOptions &&
('shouldBypassPermissionChecks' in
maybeOptionsOrMaybePermissionOptions ||
'objectRecordsPermissions' in maybeOptionsOrMaybePermissionOptions)
? maybeOptionsOrMaybePermissionOptions
: permissionOptions;
let target =
arguments.length > 1 &&
(typeof targetOrEntity === 'function' ||
InstanceChecker.isEntitySchema(targetOrEntity) ||
typeof targetOrEntity === 'string')
? targetOrEntity
: undefined;
const entity = target ? entityOrMaybeOptions : targetOrEntity;
const options = target
? maybeOptionsOrMaybePermissionOptions
: permissionOptions;
: entityOrMaybeOptions;
let target =
arguments.length > 1 &&
(typeof targetOrEntity === 'function' ||
InstanceChecker.isEntitySchema(targetOrEntity) ||
typeof targetOrEntity === 'string')
? targetOrEntity
: undefined;
if (InstanceChecker.isEntitySchema(target)) target = target.options.name;
if (Array.isArray(entity) && entity.length === 0)
return Promise.resolve(entity as Entity[]);
const entity = target ? entityOrMaybeOptions : targetOrEntity;
const queryRunnerForEntityPersistExecutor =
this.connection.createQueryRunnerForEntityPersistExecutor();
const options = target
? maybeOptionsOrMaybePermissionOptions
: entityOrMaybeOptions;
const isEntityArray = Array.isArray(entity);
const entityTarget =
target ?? (isEntityArray ? entity[0]?.constructor : entity.constructor);
if (InstanceChecker.isEntitySchema(target)) target = target.options.name;
if (Array.isArray(entity) && entity.length === 0)
return Promise.resolve(entity as Entity[]);
const entityArray = isEntityArray ? entity : [entity];
const queryRunnerForEntityPersistExecutor =
this.connection.createQueryRunnerForEntityPersistExecutor();
const isEntityArray = Array.isArray(entity);
const entityTarget =
target ?? (isEntityArray ? entity[0]?.constructor : entity.constructor);
const entityArray = isEntityArray ? entity : [entity];
const relationNestedQueries = new RelationNestedQueries(
this.internalContext,
);
const relationNestedConfig =
relationNestedQueries.prepareNestedRelationQueries(
entityArray,
entityTarget,
const relationNestedQueries = new RelationNestedQueries(
this.internalContext,
);
const entityWithConnectedRelations = isDefined(relationNestedConfig)
? await relationNestedQueries.processRelationNestedQueries({
entities: entityArray,
relationNestedConfig,
queryBuilder: this.createQueryBuilder(
undefined,
undefined,
undefined,
permissionOptions,
),
})
: entityArray;
const relationNestedConfig =
relationNestedQueries.prepareNestedRelationQueries(
entityArray,
entityTarget,
);
const entityIds = entityArray
.map((entity) => (entity as { id: string }).id)
.filter(isDefined);
const beforeUpdate = await this.find(
entityTarget,
{
where: { id: In(entityIds) },
},
{ shouldBypassPermissionChecks: true }, // Bypass as this is for event emission
);
const entityWithConnectedRelations = isDefined(relationNestedConfig)
? await relationNestedQueries.processRelationNestedQueries({
entities: entityArray,
relationNestedConfig,
queryBuilder: this.createQueryBuilder(
undefined,
undefined,
undefined,
permissionOptions,
),
})
: entityArray;
const beforeUpdateMapById = beforeUpdate.reduce(
(acc, e: ObjectLiteral) => {
acc[e.id] = e;
const entityIds = entityArray
.map((entity) => (entity as { id: string }).id)
.filter(isDefined);
const beforeUpdate = await this.find(
entityTarget,
{
where: { id: In(entityIds) },
},
{ shouldBypassPermissionChecks: true }, // Bypass as this is for event emission
);
return acc;
},
{} as Record<string, ObjectLiteral>,
);
const beforeUpdateMapById = beforeUpdate.reduce(
(acc, e: ObjectLiteral) => {
acc[e.id] = e;
const objectMetadataItem = getObjectMetadataFromEntityTarget(
entityTarget,
this.internalContext,
);
return acc;
},
{} as Record<string, ObjectLiteral>,
);
const formattedEntityOrEntities = formatData(
entityWithConnectedRelations,
objectMetadataItem,
);
const objectMetadataItem = getObjectMetadataFromEntityTarget(
entityTarget,
this.internalContext,
);
const updatedColumns = formattedEntityOrEntities
.map((e) => Object.keys(e))
.flat();
this.validatePermissions({
target: targetOrEntity,
operationType: 'update',
permissionOptions: permissionOptionsFromArgs,
selectedColumns: [],
updatedColumns,
});
const result = await new EntityPersistExecutor(
this.connection,
queryRunnerForEntityPersistExecutor,
'save',
target,
formattedEntityOrEntities as ObjectLiteral[],
options as SaveOptions | (SaveOptions & { reload: false }),
)
.execute()
.then(() => formattedEntityOrEntities as Entity[])
.finally(() => queryRunnerForEntityPersistExecutor.release());
const resultArray = Array.isArray(result) ? result : [result];
let formattedResult = formatResult<Entity[]>(
resultArray,
objectMetadataItem,
this.internalContext.objectMetadataMaps,
);
const updatedEntities = formattedResult.filter(
(entity) => beforeUpdateMapById[entity.id],
);
const createdEntities = formattedResult.filter(
(entity) => !beforeUpdateMapById[entity.id],
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.UPDATED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: updatedEntities,
beforeEntities: updatedEntities.map(
(entity) => beforeUpdateMapById[entity.id],
),
});
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.CREATED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: createdEntities,
});
const isFieldPermissionsEnabled =
this.getFeatureFlagMap().IS_FIELDS_PERMISSIONS_ENABLED;
const permissionCheckApplies =
permissionOptionsFromArgs?.shouldBypassPermissionChecks !== true &&
objectMetadataItem.isSystem !== true;
if (isFieldPermissionsEnabled && permissionCheckApplies) {
formattedResult = this.getFormattedResultWithoutNonReadableFields({
formattedResult,
const formattedEntityOrEntities = formatData(
entityWithConnectedRelations,
objectMetadataItem,
permissionOptionsFromArgs,
});
}
);
return isEntityArray ? formattedResult : formattedResult[0];
const updatedColumns = formattedEntityOrEntities
.map((e) => Object.keys(e))
.flat();
this.validatePermissions({
target: targetOrEntity,
operationType: 'update',
permissionOptions: permissionOptionsFromArgs,
selectedColumns: [],
updatedColumns,
});
const result = await new EntityPersistExecutor(
this.connection,
queryRunnerForEntityPersistExecutor,
'save',
target,
formattedEntityOrEntities as ObjectLiteral[],
options as SaveOptions | (SaveOptions & { reload: false }),
)
.execute()
.then(() => formattedEntityOrEntities as Entity[])
.finally(() => queryRunnerForEntityPersistExecutor.release());
const resultArray = Array.isArray(result) ? result : [result];
let formattedResult = formatResult<Entity[]>(
resultArray,
objectMetadataItem,
this.internalContext.objectMetadataMaps,
);
const updatedEntities = formattedResult.filter(
(entity) => beforeUpdateMapById[entity.id],
);
const createdEntities = formattedResult.filter(
(entity) => !beforeUpdateMapById[entity.id],
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.UPDATED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: updatedEntities,
beforeEntities: updatedEntities.map(
(entity) => beforeUpdateMapById[entity.id],
),
});
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.CREATED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: createdEntities,
});
const isFieldPermissionsEnabled =
this.getFeatureFlagMap().IS_FIELDS_PERMISSIONS_ENABLED;
const permissionCheckApplies =
permissionOptionsFromArgs?.shouldBypassPermissionChecks !== true &&
objectMetadataItem.isSystem !== true;
if (isFieldPermissionsEnabled && permissionCheckApplies) {
formattedResult = this.getFormattedResultWithoutNonReadableFields({
formattedResult,
objectMetadataItem,
permissionOptionsFromArgs,
});
}
return isEntityArray ? formattedResult : formattedResult[0];
} catch (error) {
throw computeTwentyORMException(error);
}
}
private getFormattedResultWithoutNonReadableFields<