Function trigger updates 2 (#16608)
- Improves route trigger job performances - expose function params types ## Before <img width="938" height="271" alt="image" src="https://github.com/user-attachments/assets/5752ba64-f31d-44ed-974d-536e63458f2c" /> ## After <img width="1000" height="559" alt="image" src="https://github.com/user-attachments/assets/b1f4927a-5f43-49f0-a606-244c72356772" />
This commit is contained in:
+26
-20
@@ -1,8 +1,10 @@
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
|
||||
import { Repository } from 'typeorm';
|
||||
import chunk from 'lodash.chunk';
|
||||
|
||||
import type { ObjectRecordEvent } from 'twenty-shared/database-events';
|
||||
|
||||
import type { ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
|
||||
import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator';
|
||||
import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator';
|
||||
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
|
||||
@@ -14,6 +16,9 @@ import {
|
||||
ServerlessFunctionTriggerJobData,
|
||||
} from 'src/engine/metadata-modules/serverless-function/jobs/serverless-function-trigger.job';
|
||||
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
|
||||
import { transformEventBatchToEventPayloads } from 'src/engine/metadata-modules/database-event-trigger/utils/transform-event-batch-to-event-payloads';
|
||||
|
||||
const DATABASE_EVENT_JOBS_CHUNK_SIZE = 20;
|
||||
|
||||
@Processor(MessageQueue.triggerQueue)
|
||||
export class CallDatabaseEventTriggerJobsJob {
|
||||
@@ -35,29 +40,30 @@ export class CallDatabaseEventTriggerJobsJob {
|
||||
relations: ['serverlessFunction'],
|
||||
});
|
||||
|
||||
for (const databaseEventListener of databaseEventListeners) {
|
||||
if (
|
||||
!this.shouldTriggerJob({
|
||||
const databaseEventListenersToTrigger = databaseEventListeners.filter(
|
||||
(databaseEventListener) =>
|
||||
this.shouldTriggerJob({
|
||||
workspaceEventBatch,
|
||||
eventName: databaseEventListener.settings.eventName,
|
||||
})
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
}),
|
||||
);
|
||||
|
||||
const { events, ...batchEventInfo } = workspaceEventBatch;
|
||||
const serverlessFunctionPayloads = transformEventBatchToEventPayloads({
|
||||
databaseEventListeners: databaseEventListenersToTrigger,
|
||||
workspaceEventBatch,
|
||||
});
|
||||
|
||||
for (const event of events) {
|
||||
await this.messageQueueService.add<ServerlessFunctionTriggerJobData>(
|
||||
ServerlessFunctionTriggerJob.name,
|
||||
{
|
||||
serverlessFunctionId: databaseEventListener.serverlessFunction.id,
|
||||
workspaceId: databaseEventListener.workspaceId,
|
||||
payload: { ...batchEventInfo, ...event },
|
||||
},
|
||||
{ retryLimit: 3 },
|
||||
);
|
||||
}
|
||||
const serverlessFunctionPayloadsChunks = chunk(
|
||||
serverlessFunctionPayloads,
|
||||
DATABASE_EVENT_JOBS_CHUNK_SIZE,
|
||||
);
|
||||
|
||||
for (const serverlessFunctionPayloadsChunk of serverlessFunctionPayloadsChunks) {
|
||||
await this.messageQueueService.add<ServerlessFunctionTriggerJobData[]>(
|
||||
ServerlessFunctionTriggerJob.name,
|
||||
serverlessFunctionPayloadsChunk,
|
||||
{ retryLimit: 3 },
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+34
@@ -0,0 +1,34 @@
|
||||
import type {
|
||||
DatabaseEventPayload,
|
||||
ObjectRecordEvent,
|
||||
} from 'twenty-shared/database-events';
|
||||
|
||||
import type { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
|
||||
import { type DatabaseEventTriggerEntity } from 'src/engine/metadata-modules/database-event-trigger/entities/database-event-trigger.entity';
|
||||
import { type ServerlessFunctionTriggerJobData } from 'src/engine/metadata-modules/serverless-function/jobs/serverless-function-trigger.job';
|
||||
|
||||
export const transformEventBatchToEventPayloads = ({
|
||||
workspaceEventBatch,
|
||||
databaseEventListeners,
|
||||
}: {
|
||||
workspaceEventBatch: WorkspaceEventBatch<ObjectRecordEvent>;
|
||||
databaseEventListeners: DatabaseEventTriggerEntity[];
|
||||
}): ServerlessFunctionTriggerJobData[] => {
|
||||
const result: ServerlessFunctionTriggerJobData[] = [];
|
||||
|
||||
for (const databaseEventListener of databaseEventListeners) {
|
||||
const { events, ...batchEventInfo } = workspaceEventBatch;
|
||||
|
||||
for (const event of events) {
|
||||
const payload: DatabaseEventPayload = { ...batchEventInfo, ...event };
|
||||
|
||||
result.push({
|
||||
serverlessFunctionId: databaseEventListener.serverlessFunction.id,
|
||||
workspaceId: databaseEventListener.workspaceId,
|
||||
payload,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
return result;
|
||||
};
|
||||
+12
-7
@@ -21,12 +21,17 @@ export class ServerlessFunctionTriggerJob {
|
||||
) {}
|
||||
|
||||
@Process(ServerlessFunctionTriggerJob.name)
|
||||
async handle(data: ServerlessFunctionTriggerJobData) {
|
||||
await this.serverlessFunctionService.executeOneServerlessFunction({
|
||||
id: data.serverlessFunctionId,
|
||||
workspaceId: data.workspaceId,
|
||||
payload: data.payload || {},
|
||||
version: 'draft',
|
||||
});
|
||||
async handle(serverlessFunctionPayloads: ServerlessFunctionTriggerJobData[]) {
|
||||
await Promise.all(
|
||||
serverlessFunctionPayloads.map(
|
||||
async (serverlessFunctionPayload) =>
|
||||
await this.serverlessFunctionService.executeOneServerlessFunction({
|
||||
id: serverlessFunctionPayload.serverlessFunctionId,
|
||||
workspaceId: serverlessFunctionPayload.workspaceId,
|
||||
payload: serverlessFunctionPayload.payload || {},
|
||||
version: 'draft',
|
||||
}),
|
||||
),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user