1658 post mortem 0710 send batch events in webhook (#15022)

Use concurrent calls

---------

Co-authored-by: Charles Bochet <charles@twenty.com>
This commit is contained in:
martmull
2025-10-11 12:48:40 +02:00
committed by GitHub
co-authored by Charles Bochet
parent 68c86871dd
commit c681fb7fb6
58 changed files with 920 additions and 521 deletions
@@ -10,16 +10,15 @@ import { type ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/ty
import { type ObjectRecordNonDestructiveEvent } from 'src/engine/core-modules/event-emitter/types/object-record-non-destructive-event';
import { type ObjectRecordRestoreEvent } from 'src/engine/core-modules/event-emitter/types/object-record-restore.event';
import { type ObjectRecordUpdateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-update.event';
import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service';
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 { CallWebhookJobsJob } from 'src/engine/core-modules/webhook/jobs/call-webhook-jobs.job';
import { type ObjectRecordEventForWebhook } from 'src/engine/core-modules/webhook/types/object-record-event-for-webhook.type';
import { CallDatabaseEventTriggerJobsJob } from 'src/engine/metadata-modules/database-event-trigger/jobs/call-database-event-trigger-jobs.job';
import { SubscriptionsService } from 'src/engine/subscriptions/subscriptions.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { UpsertTimelineActivityFromInternalEvent } from 'src/modules/timeline/jobs/upsert-timeline-activity-from-internal-event.job';
import { WorkspaceEventBatchForWebhook } from 'src/engine/core-modules/webhook/types/workspace-event-batch-for-webhook.type';
@Injectable()
export class EntityEventsToDbListener {
@@ -31,7 +30,6 @@ export class EntityEventsToDbListener {
@InjectMessageQueue(MessageQueue.triggerQueue)
private readonly triggerQueueService: MessageQueueService,
private readonly subscriptionsService: SubscriptionsService,
private readonly featureFlagService: FeatureFlagService,
) {}
@OnDatabaseBatchEvent('*', DatabaseEventAction.CREATED)
@@ -67,26 +65,21 @@ export class EntityEventsToDbListener {
batchEvent: WorkspaceEventBatch<T>,
action: DatabaseEventAction,
) {
const auditLogsEvents = batchEvent.events.filter(
(event) => event.objectMetadata?.isAuditLogged,
);
const isAuditLogBatchEvent = batchEvent.objectMetadata?.isAuditLogged;
const batchEventEventsForWebhook: ObjectRecordEventForWebhook[] =
batchEvent.events.map((event) => ({
...event,
objectMetadata: {
id: event.objectMetadata.id,
nameSingular: event.objectMetadata.nameSingular,
},
}));
const batchEventForWebhook = {
...batchEvent,
objectMetadata: {
id: batchEvent.objectMetadata.id,
nameSingular: batchEvent.objectMetadata.nameSingular,
},
};
const promises = [
this.subscriptionsService.publish(batchEvent),
this.webhookQueueService.add<
WorkspaceEventBatch<ObjectRecordEventForWebhook>
>(
this.webhookQueueService.add<WorkspaceEventBatchForWebhook<T>>(
CallWebhookJobsJob.name,
{ ...batchEvent, events: batchEventEventsForWebhook },
batchEventForWebhook,
{
retryLimit: 3,
},
@@ -101,14 +94,11 @@ export class EntityEventsToDbListener {
),
);
if (auditLogsEvents.length > 0) {
if (isAuditLogBatchEvent) {
promises.push(
this.entityEventsToDbQueueService.add<WorkspaceEventBatch<T>>(
CreateAuditLogFromInternalEvent.name,
{
...batchEvent,
events: auditLogsEvents,
},
batchEvent,
),
);
@@ -116,10 +106,7 @@ export class EntityEventsToDbListener {
promises.push(
this.entityEventsToDbQueueService.add<
WorkspaceEventBatch<ObjectRecordNonDestructiveEvent>
>(UpsertTimelineActivityFromInternalEvent.name, {
...batchEvent,
events: auditLogsEvents,
}),
>(UpsertTimelineActivityFromInternalEvent.name, batchEvent),
);
}
}
@@ -6,7 +6,7 @@ import { AuditService } from 'src/engine/core-modules/audit/services/audit.servi
import { USER_SIGNUP_EVENT } from 'src/engine/core-modules/audit/utils/events/workspace-event/user/user-signup';
import { type ObjectRecordCreateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-create.event';
import { TelemetryService } from 'src/engine/core-modules/telemetry/telemetry.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
@Injectable()
export class TelemetryListener {
@@ -17,7 +17,7 @@ export class TelemetryListener {
@OnCustomBatchEvent(USER_SIGNUP_EVENT_NAME)
async handleUserSignup(
payload: WorkspaceEventBatch<ObjectRecordCreateEvent>,
payload: CustomWorkspaceEventBatch<ObjectRecordCreateEvent>,
) {
await Promise.all(
payload.events.map(async (eventPayload) => {
@@ -8,7 +8,6 @@ import { WorkspaceQueryHookModule } from 'src/engine/api/graphql/workspace-query
import { AuditModule } from 'src/engine/core-modules/audit/audit.module';
import { AuthModule } from 'src/engine/core-modules/auth/auth.module';
import { FeatureFlag } from 'src/engine/core-modules/feature-flag/feature-flag.entity';
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
import { FileModule } from 'src/engine/core-modules/file/file.module';
import { RecordPositionModule } from 'src/engine/core-modules/record-position/record-position.module';
import { RecordTransformerModule } from 'src/engine/core-modules/record-transformer/record-transformer.module';
@@ -28,7 +27,6 @@ import { EntityEventsToDbListener } from './listeners/entity-events-to-db.listen
AuditModule,
TelemetryModule,
FileModule,
FeatureFlagModule,
RecordTransformerModule,
RecordPositionModule,
SubscriptionsModule,
@@ -7,7 +7,7 @@ import { type ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/ty
import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator';
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
@Processor(MessageQueue.entityEventsToDbQueue)
export class CreateAuditLogFromInternalEvent {
@@ -37,25 +37,25 @@ export class CreateAuditLogFromInternalEvent {
auditService.createObjectEvent(OBJECT_RECORD_UPDATED_EVENT, {
...eventProperties,
recordId: eventData.recordId,
objectMetadataId: eventData.objectMetadata.id,
objectMetadataId: workspaceEventBatch.objectMetadata.id,
});
} else if (workspaceEventBatch.name.endsWith('.created')) {
auditService.createObjectEvent(OBJECT_RECORD_CREATED_EVENT, {
...eventProperties,
recordId: eventData.recordId,
objectMetadataId: eventData.objectMetadata.id,
objectMetadataId: workspaceEventBatch.objectMetadata.id,
});
} else if (workspaceEventBatch.name.endsWith('.deleted')) {
auditService.createObjectEvent(OBJECT_RECORD_DELETED_EVENT, {
...eventProperties,
recordId: eventData.recordId,
objectMetadataId: eventData.objectMetadata.id,
objectMetadataId: workspaceEventBatch.objectMetadata.id,
});
} else if (workspaceEventBatch.name.endsWith('.upserted')) {
auditService.createObjectEvent(OBJECT_RECORD_UPSERTED_EVENT, {
...eventProperties,
recordId: eventData.recordId,
objectMetadataId: eventData.objectMetadata.id,
objectMetadataId: workspaceEventBatch.objectMetadata.id,
});
}
}
@@ -2,12 +2,14 @@
import { Injectable } from '@nestjs/common';
import { isDefined } from 'twenty-shared/utils';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { BILLING_FEATURE_USED } from 'src/engine/core-modules/billing/constants/billing-feature-used.constant';
import { BillingUsageService } from 'src/engine/core-modules/billing/services/billing-usage.service';
import { type BillingUsageEvent } from 'src/engine/core-modules/billing/types/billing-usage-event.type';
import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
@Injectable()
export class BillingFeatureUsedListener {
@@ -18,8 +20,12 @@ export class BillingFeatureUsedListener {
@OnCustomBatchEvent(BILLING_FEATURE_USED)
async handleBillingFeatureUsedEvent(
payload: WorkspaceEventBatch<BillingUsageEvent>,
payload: CustomWorkspaceEventBatch<BillingUsageEvent>,
) {
if (!isDefined(payload.workspaceId)) {
return;
}
if (!this.twentyConfigService.get('IS_BILLING_ENABLED')) {
return;
}
@@ -13,7 +13,7 @@ import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decora
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 { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type WorkspaceMemberWorkspaceEntity } from 'src/modules/workspace-member/standard-objects/workspace-member.workspace-entity';
@Injectable()
@@ -1,5 +1,4 @@
import { type ObjectRecordDiff } from 'src/engine/core-modules/event-emitter/types/object-record-diff';
import { type ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
type Properties<T> = {
updatedFields?: string[];
@@ -12,6 +11,5 @@ export class ObjectRecordBaseEvent<T = object> {
recordId: string;
userId?: string;
workspaceMemberId?: string;
objectMetadata: Omit<ObjectMetadataEntity, 'indexMetadatas'>;
properties: Properties<T>;
}
@@ -10,7 +10,7 @@ import {
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type AttachmentWorkspaceEntity } from 'src/modules/attachment/standard-objects/attachment.workspace-entity';
@Injectable()
@@ -10,7 +10,7 @@ import {
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type WorkspaceMemberWorkspaceEntity } from 'src/modules/workspace-member/standard-objects/workspace-member.workspace-entity';
@Injectable()
@@ -1,6 +1,6 @@
import { Logger } from '@nestjs/common';
import { isDefined } from 'twenty-shared/utils';
import chunk from 'lodash.chunk';
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';
@@ -11,11 +11,12 @@ import {
CallWebhookJob,
type CallWebhookJobData,
} from 'src/engine/core-modules/webhook/jobs/call-webhook.job';
import { type ObjectRecordEventForWebhook } from 'src/engine/core-modules/webhook/types/object-record-event-for-webhook.type';
import { WebhookService } from 'src/engine/core-modules/webhook/webhook.service';
import { type ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { removeSecretFromWebhookRecord } from 'src/utils/remove-secret-from-webhook-record';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import type { ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import { transformEventBatchToWebhookEvents } from 'src/engine/core-modules/webhook/utils/transform-event-batch-to-webhook-events';
const WEBHOOK_JOBS_CHUNK_SIZE = 20;
@Processor(MessageQueue.webhookQueue)
export class CallWebhookJobsJob {
@@ -28,7 +29,7 @@ export class CallWebhookJobsJob {
@Process(CallWebhookJobsJob.name)
async handle(
workspaceEventBatch: WorkspaceEventBatch<ObjectRecordEventForWebhook>,
workspaceEventBatch: WorkspaceEventBatch<ObjectRecordEvent>,
): Promise<void> {
// If you change that function, double check it does not break Zapier
// trigger in packages/twenty-zapier/src/triggers/trigger_record.ts
@@ -47,51 +48,19 @@ export class CallWebhookJobsJob {
],
);
for (const eventData of workspaceEventBatch.events) {
const eventName = workspaceEventBatch.name;
const objectMetadata: Pick<ObjectMetadataEntity, 'id' | 'nameSingular'> =
{
id: eventData.objectMetadata.id,
nameSingular: eventData.objectMetadata.nameSingular,
};
const workspaceId = workspaceEventBatch.workspaceId;
const record =
'after' in eventData.properties && isDefined(eventData.properties.after)
? eventData.properties.after
: 'before' in eventData.properties &&
isDefined(eventData.properties.before)
? eventData.properties.before
: {};
const updatedFields =
'updatedFields' in eventData.properties
? eventData.properties.updatedFields
: undefined;
const webhookEvents = transformEventBatchToWebhookEvents({
workspaceEventBatch,
webhooks,
});
const isWebhookEvent = nameSingular === 'webhook';
const sanitizedRecord = removeSecretFromWebhookRecord(
record,
isWebhookEvent,
const webhookEventsChunks = chunk(webhookEvents, WEBHOOK_JOBS_CHUNK_SIZE);
for (const webhookEventsChunk of webhookEventsChunks) {
await this.messageQueueService.add<CallWebhookJobData[]>(
CallWebhookJob.name,
webhookEventsChunk,
{ retryLimit: 3 },
);
webhooks.forEach((webhook) => {
const webhookData = {
targetUrl: webhook.targetUrl,
secret: webhook.secret,
eventName,
objectMetadata,
workspaceId,
webhookId: webhook.id,
eventDate: new Date(),
record: sanitizedRecord,
...(updatedFields && { updatedFields }),
};
this.messageQueueService.add<CallWebhookJobData>(
CallWebhookJob.name,
webhookData,
{ retryLimit: 3 },
);
});
}
}
}
@@ -42,7 +42,15 @@ export class CallWebhookJob {
}
@Process(CallWebhookJob.name)
async handle(data: CallWebhookJobData): Promise<void> {
async handle(webhookJobEvents: CallWebhookJobData[]): Promise<void> {
await Promise.all(
webhookJobEvents.map(
async (webhookJobEvent) => await this.callWebhook(webhookJobEvent),
),
);
}
private async callWebhook(data: CallWebhookJobData): Promise<void> {
const commonPayload = {
url: data.targetUrl,
webhookId: data.webhookId,
@@ -1,9 +0,0 @@
import { type ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import { type ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
export type ObjectRecordEventForWebhook = Omit<
ObjectRecordEvent,
'objectMetadata'
> & {
objectMetadata: Pick<ObjectMetadataEntity, 'id' | 'nameSingular'>;
};
@@ -0,0 +1,9 @@
import { type ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
export type WorkspaceEventBatchForWebhook<WorkspaceEvent> = Omit<
WorkspaceEventBatch<WorkspaceEvent>,
'objectMetadata'
> & {
objectMetadata: Pick<ObjectMetadataEntity, 'id' | 'nameSingular'>;
};
@@ -0,0 +1,235 @@
import type { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import type { ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import type { Webhook } from 'src/engine/core-modules/webhook/webhook.entity';
import { transformEventBatchToWebhookEvents } from 'src/engine/core-modules/webhook/utils/transform-event-batch-to-webhook-events';
import { getMockObjectMetadataEntity } from 'src/utils/__test__/get-object-metadata-entity.mock';
const mockObjectMetadata = getMockObjectMetadataEntity({
id: 'id',
nameSingular: 'nameSingular',
namePlural: 'namePlural',
workspaceId: 'workspaceId',
});
describe('transformEventBatchToWebhookEvents', () => {
it('should transform properly', () => {
const workspaceEventBatch = {
workspaceId: 'workspaceId',
objectMetadata: mockObjectMetadata,
name: 'objectNameSingular.created',
events: [
{
recordId: 'recordId-1',
properties: {
after: {
id: 'id-1',
nameSingular: 'nameSingular-1',
},
},
},
{
recordId: 'recordId-2',
properties: {
before: {
id: 'id-2',
nameSingular: 'nameSingular-2',
},
},
},
{
recordId: 'recordId-3',
properties: {
after: {
id: 'id-3',
nameSingular: 'nameSingular-3',
secret: 'secret-3',
},
updatedFields: ['nameSingular'],
},
},
],
} as WorkspaceEventBatch<ObjectRecordEvent>;
const webhooks = [
{
id: 'webhook-id',
targetUrl: 'targetUrl',
secret: 'secret',
},
{
id: 'webhook-id-2',
targetUrl: 'targetUrl-2',
secret: 'secret-2',
},
] as Webhook[];
const result = transformEventBatchToWebhookEvents({
workspaceEventBatch,
webhooks,
});
const expectedResultWithoutEventDate = [
{
targetUrl: 'targetUrl',
eventName: 'objectNameSingular.created',
objectMetadata: {
id: mockObjectMetadata.id,
nameSingular: mockObjectMetadata.nameSingular,
},
workspaceId: 'workspaceId',
webhookId: 'webhook-id',
record: {
id: 'id-1',
nameSingular: 'nameSingular-1',
},
secret: 'secret',
},
{
targetUrl: 'targetUrl',
eventName: 'objectNameSingular.created',
objectMetadata: {
id: mockObjectMetadata.id,
nameSingular: mockObjectMetadata.nameSingular,
},
workspaceId: 'workspaceId',
webhookId: 'webhook-id',
record: {
id: 'id-2',
nameSingular: 'nameSingular-2',
},
secret: 'secret',
},
{
targetUrl: 'targetUrl',
eventName: 'objectNameSingular.created',
objectMetadata: {
id: mockObjectMetadata.id,
nameSingular: mockObjectMetadata.nameSingular,
},
workspaceId: 'workspaceId',
webhookId: 'webhook-id',
record: {
id: 'id-3',
nameSingular: 'nameSingular-3',
secret: 'secret-3',
},
updatedFields: ['nameSingular'],
secret: 'secret',
},
{
targetUrl: 'targetUrl-2',
eventName: 'objectNameSingular.created',
objectMetadata: {
id: mockObjectMetadata.id,
nameSingular: mockObjectMetadata.nameSingular,
},
workspaceId: 'workspaceId',
webhookId: 'webhook-id-2',
record: {
id: 'id-1',
nameSingular: 'nameSingular-1',
},
secret: 'secret-2',
},
{
targetUrl: 'targetUrl-2',
eventName: 'objectNameSingular.created',
objectMetadata: {
id: mockObjectMetadata.id,
nameSingular: mockObjectMetadata.nameSingular,
},
workspaceId: 'workspaceId',
webhookId: 'webhook-id-2',
record: {
id: 'id-2',
nameSingular: 'nameSingular-2',
},
secret: 'secret-2',
},
{
targetUrl: 'targetUrl-2',
eventName: 'objectNameSingular.created',
objectMetadata: {
id: mockObjectMetadata.id,
nameSingular: mockObjectMetadata.nameSingular,
},
workspaceId: 'workspaceId',
webhookId: 'webhook-id-2',
record: {
id: 'id-3',
nameSingular: 'nameSingular-3',
secret: 'secret-3',
},
updatedFields: ['nameSingular'],
secret: 'secret-2',
},
];
const resultWithoutEventDate = result.map((event) => {
const { eventDate: _, ...eventWithoutEventDate } = event;
return eventWithoutEventDate;
});
expect(resultWithoutEventDate).toEqual(expectedResultWithoutEventDate);
});
it('should sanitize records properly', () => {
const workspaceEventBatch = {
workspaceId: 'workspaceId',
objectMetadata: mockObjectMetadata,
name: 'webhook.created',
events: [
{
recordId: 'recordId-1',
properties: {
after: {
id: 'id-1',
targetUrl: 'targetUrl-1',
secret: 'secret-1',
},
},
},
],
} as WorkspaceEventBatch<ObjectRecordEvent>;
const webhooks = [
{
id: 'webhook-id',
targetUrl: 'targetUrl',
secret: 'secret',
},
] as Webhook[];
const result = transformEventBatchToWebhookEvents({
workspaceEventBatch,
webhooks,
});
const expectedResultWithoutEventDate = [
{
targetUrl: 'targetUrl',
eventName: 'webhook.created',
objectMetadata: {
id: mockObjectMetadata.id,
nameSingular: mockObjectMetadata.nameSingular,
},
workspaceId: 'workspaceId',
webhookId: 'webhook-id',
record: {
id: 'id-1',
targetUrl: 'targetUrl-1',
// No secret
},
secret: 'secret',
},
];
const resultWithoutEventDate = result.map((event) => {
const { eventDate: _, ...eventWithoutEventDate } = event;
return eventWithoutEventDate;
});
expect(resultWithoutEventDate).toEqual(expectedResultWithoutEventDate);
});
});
@@ -0,0 +1,34 @@
import { transformEventToWebhookEvent } from 'src/engine/core-modules/webhook/utils/transform-event-to-webhook-event';
import type { ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
describe('transformEventToWebhookEvent', () => {
it('should transform event to webhook event', () => {
const record = {
recordId: 'recordId',
properties: {
after: {
id: 'id',
nameSingular: 'nameSingular',
secret: 'secret',
},
updatedFields: ['nameSingular'],
},
} as ObjectRecordEvent;
const expectedResult = {
record: {
id: 'id',
nameSingular: 'nameSingular',
secret: 'secret',
},
updatedFields: ['nameSingular'],
};
expect(
transformEventToWebhookEvent({
eventName: 'nameSingular.created',
event: record,
}),
).toEqual(expectedResult);
});
});
@@ -0,0 +1,49 @@
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type CallWebhookJobData } from 'src/engine/core-modules/webhook/jobs/call-webhook.job';
import { type Webhook } from 'src/engine/core-modules/webhook/webhook.entity';
import type { ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import { transformEventToWebhookEvent } from 'src/engine/core-modules/webhook/utils/transform-event-to-webhook-event';
export const transformEventBatchToWebhookEvents = ({
workspaceEventBatch,
webhooks,
}: {
workspaceEventBatch: WorkspaceEventBatch<ObjectRecordEvent>;
webhooks: Webhook[];
}): CallWebhookJobData[] => {
const result: CallWebhookJobData[] = [];
for (const webhook of webhooks) {
const targetUrl = webhook.targetUrl;
const eventName = workspaceEventBatch.name;
const objectMetadataForWebhook = {
id: workspaceEventBatch.objectMetadata.id,
nameSingular: workspaceEventBatch.objectMetadata.nameSingular,
};
const workspaceId = workspaceEventBatch.workspaceId;
const webhookId = webhook.id;
const eventDate = new Date();
const secret = webhook.secret;
for (const eventData of workspaceEventBatch.events) {
const { record, updatedFields } = transformEventToWebhookEvent({
eventName: workspaceEventBatch.name,
event: eventData,
});
result.push({
targetUrl,
eventName,
objectMetadata: objectMetadataForWebhook,
workspaceId,
webhookId,
eventDate,
record,
...(updatedFields && { updatedFields }),
secret,
});
}
}
return result;
};
@@ -0,0 +1,34 @@
import { isDefined } from 'twenty-shared/utils';
import type { ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import { removeSecretFromWebhookRecord } from 'src/utils/remove-secret-from-webhook-record';
export const transformEventToWebhookEvent = ({
eventName,
event,
}: {
eventName: string;
event: ObjectRecordEvent;
}) => {
const [nameSingular, _] = eventName.split('.');
const record =
'after' in event.properties && isDefined(event.properties.after)
? event.properties.after
: 'before' in event.properties && isDefined(event.properties.before)
? event.properties.before
: {};
const updatedFields =
'updatedFields' in event.properties
? event.properties.updatedFields
: undefined;
const isWebhookEvent = nameSingular === 'webhook';
const sanitizedRecord = removeSecretFromWebhookRecord(record, isWebhookEvent);
return {
record: sanitizedRecord,
...(updatedFields && { updatedFields }),
};
};
@@ -10,7 +10,7 @@ import { type ObjectRecordUpdateEvent } 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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type WorkspaceMemberWorkspaceEntity } from 'src/modules/workspace-member/standard-objects/workspace-member.workspace-entity';
import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
@@ -13,7 +13,7 @@ import {
ServerlessFunctionTriggerJob,
ServerlessFunctionTriggerJobData,
} from 'src/engine/metadata-modules/serverless-function/jobs/serverless-function-trigger.job';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
@Processor(MessageQueue.triggerQueue)
export class CallDatabaseEventTriggerJobsJob {
@@ -1,11 +1,10 @@
import { Inject, Injectable } from '@nestjs/common';
import { RedisPubSub } from 'graphql-redis-subscriptions';
import { isDefined } from 'twenty-shared/utils';
import { type ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { removeSecretFromWebhookRecord } from 'src/utils/remove-secret-from-webhook-record';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { transformEventToWebhookEvent } from 'src/engine/core-modules/webhook/utils/transform-event-to-webhook-event';
@Injectable()
export class SubscriptionsService {
@@ -14,32 +13,20 @@ export class SubscriptionsService {
async publish(
workspaceEventBatch: WorkspaceEventBatch<ObjectRecordEvent>,
): Promise<void> {
for (const eventData of workspaceEventBatch.events) {
const [nameSingular, operation] = workspaceEventBatch.name.split('.');
const record =
'after' in eventData.properties && isDefined(eventData.properties.after)
? eventData.properties.after
: 'before' in eventData.properties &&
isDefined(eventData.properties.before)
? eventData.properties.before
: {};
const updatedFields =
'updatedFields' in eventData.properties
? eventData.properties.updatedFields
: undefined;
const [nameSingular, operation] = workspaceEventBatch.name.split('.');
const isWebhookEvent = nameSingular === 'webhook';
const sanitizedRecord = removeSecretFromWebhookRecord(
record,
isWebhookEvent,
);
for (const eventData of workspaceEventBatch.events) {
const { record, updatedFields } = transformEventToWebhookEvent({
eventName: workspaceEventBatch.name,
event: eventData,
});
await this.pubSub.publish('onDbEvent', {
onDbEvent: {
action: operation,
objectNameSingular: nameSingular,
eventDate: new Date(),
record: sanitizedRecord,
record,
...(updatedFields && { updatedFields }),
},
});
@@ -56,6 +56,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 { formatTwentyOrmEventToDatabaseBatchEvent } from 'src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util';
type PermissionOptions = {
shouldBypassPermissionChecks?: boolean;
@@ -1195,22 +1196,26 @@ export class WorkspaceEntityManager extends EntityManager {
(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],
),
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
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,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.CREATED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: createdEntities,
}),
);
const permissionCheckApplies =
permissionOptionsFromArgs?.shouldBypassPermissionChecks !== true &&
@@ -1387,12 +1392,14 @@ export class WorkspaceEntityManager extends EntityManager {
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.DESTROYED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: formattedResult,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.DESTROYED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: formattedResult,
}),
);
return isEntityArray ? formattedResult : formattedResult[0];
}
@@ -1500,12 +1507,14 @@ export class WorkspaceEntityManager extends EntityManager {
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.DELETED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: formattedResult,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.DELETED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: formattedResult,
}),
);
return isEntityArray ? formattedResult : formattedResult[0];
}
@@ -1609,12 +1618,14 @@ export class WorkspaceEntityManager extends EntityManager {
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.RESTORED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: formattedResult,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.RESTORED,
objectMetadataItem,
workspaceId: this.internalContext.workspaceId,
entities: formattedResult,
}),
);
return isEntityArray ? formattedResult : formattedResult[0];
}
@@ -27,6 +27,7 @@ import { applyTableAliasOnWhereCondition } from 'src/engine/twenty-orm/utils/app
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 { computeTableName } from 'src/engine/utils/compute-table-name.util';
import { formatTwentyOrmEventToDatabaseBatchEvent } from 'src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util';
export class WorkspaceDeleteQueryBuilder<
T extends ObjectLiteral,
@@ -122,13 +123,15 @@ export class WorkspaceDeleteQueryBuilder<
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.DESTROYED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedBefore,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.DESTROYED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedBefore,
authContext: this.authContext,
}),
);
return {
raw: result.raw,
@@ -30,6 +30,7 @@ import { type WorkspaceUpdateQueryBuilder } from 'src/engine/twenty-orm/reposito
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 { formatTwentyOrmEventToDatabaseBatchEvent } from 'src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util';
export class WorkspaceInsertQueryBuilder<
T extends ObjectLiteral,
@@ -160,21 +161,25 @@ export class WorkspaceInsertQueryBuilder<
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.CREATED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedResultForEvent,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.CREATED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedResultForEvent,
authContext: this.authContext,
}),
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.UPSERTED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedResultForEvent,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.UPSERTED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedResultForEvent,
authContext: this.authContext,
}),
);
// TypeORM returns all entity columns for insertions
const resultWithoutInsertionExtraColumns = !isDefined(result.raw)
@@ -23,6 +23,7 @@ import { type WorkspaceSelectQueryBuilder } from 'src/engine/twenty-orm/reposito
import { type WorkspaceUpdateQueryBuilder } from 'src/engine/twenty-orm/repository/workspace-update-query-builder';
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 { formatTwentyOrmEventToDatabaseBatchEvent } from 'src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util';
export class WorkspaceSoftDeleteQueryBuilder<
T extends ObjectLiteral,
@@ -85,13 +86,15 @@ export class WorkspaceSoftDeleteQueryBuilder<
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.DELETED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.DELETED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
authContext: this.authContext,
}),
);
return {
raw: after.raw,
@@ -29,6 +29,7 @@ import { type WorkspaceSoftDeleteQueryBuilder } from 'src/engine/twenty-orm/repo
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 { formatTwentyOrmEventToDatabaseBatchEvent } from 'src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util';
export class WorkspaceUpdateQueryBuilder<
T extends ObjectLiteral,
@@ -152,23 +153,27 @@ export class WorkspaceUpdateQueryBuilder<
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.UPDATED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.UPDATED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
}),
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.UPSERTED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.UPSERTED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
}),
);
const formattedResult = formatResult<T[]>(
result.raw,
@@ -293,23 +298,27 @@ export class WorkspaceUpdateQueryBuilder<
this.internalContext.objectMetadataMaps,
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.UPDATED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.UPDATED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
}),
);
await this.internalContext.eventEmitterService.emitMutationEvent({
action: DatabaseEventAction.UPSERTED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
});
this.internalContext.eventEmitterService.emitDatabaseBatchEvent(
formatTwentyOrmEventToDatabaseBatchEvent({
action: DatabaseEventAction.UPSERTED,
objectMetadataItem: objectMetadata,
workspaceId: this.internalContext.workspaceId,
entities: formattedAfter,
beforeEntities: formattedBefore,
authContext: this.authContext,
}),
);
const formattedResults = formatResult<T[]>(
results.flatMap((result) => result.raw),
@@ -0,0 +1,174 @@
import { isDefined } from 'twenty-shared/utils';
import type { ObjectLiteral } from 'typeorm';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import type { ObjectMetadataItemWithFieldMaps } from 'src/engine/metadata-modules/types/object-metadata-item-with-field-maps';
import type { AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type';
import { STANDARD_OBJECT_IDS } from 'src/engine/workspace-manager/workspace-sync-metadata/constants/standard-object-ids';
import { ObjectRecordCreateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-create.event';
import { ObjectRecordUpdateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-update.event';
import { ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emitter/types/object-record-delete.event';
import { ObjectRecordUpsertEvent } from 'src/engine/core-modules/event-emitter/types/object-record-upsert.event';
import { objectRecordChangedValues } from 'src/engine/core-modules/event-emitter/utils/object-record-changed-values';
import type { ObjectRecordDiff } from 'src/engine/core-modules/event-emitter/types/object-record-diff';
import { ObjectRecordDestroyEvent } from 'src/engine/core-modules/event-emitter/types/object-record-destroy.event';
import { type DatabaseBatchEventInput } from 'src/engine/workspace-event-emitter/workspace-event-emitter';
export const formatTwentyOrmEventToDatabaseBatchEvent = <
T extends ObjectLiteral,
>({
action,
objectMetadataItem,
workspaceId,
authContext,
entities,
beforeEntities,
}: {
action: DatabaseEventAction;
objectMetadataItem: ObjectMetadataItemWithFieldMaps;
workspaceId: string;
authContext?: AuthContext;
entities: T | T[];
beforeEntities?: T | T[];
}): DatabaseBatchEventInput<T, DatabaseEventAction> | undefined => {
if (objectMetadataItem.standardId === STANDARD_OBJECT_IDS.timelineActivity) {
return;
}
const objectMetadataNameSingular = objectMetadataItem.nameSingular;
const fields = Object.values(objectMetadataItem.fieldsById ?? {});
const entityArray = isDefined(entities)
? Array.isArray(entities)
? entities
: [entities]
: [];
let events: (
| ObjectRecordCreateEvent<T>
| ObjectRecordUpdateEvent<T>
| ObjectRecordDeleteEvent<T>
| ObjectRecordUpsertEvent<T>
)[] = [];
switch (action) {
case DatabaseEventAction.CREATED:
events = entityArray.map((after) => {
const event = new ObjectRecordCreateEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = after.id;
event.properties = { after };
return event;
});
break;
case DatabaseEventAction.UPDATED:
events = entityArray
.map((after, idx) => {
if (!beforeEntities) {
throw new Error('beforeEntities is required for UPDATED action');
}
const before = Array.isArray(beforeEntities)
? beforeEntities?.[idx]
: beforeEntities;
const diff = objectRecordChangedValues(
before,
after,
objectMetadataItem,
) as Partial<ObjectRecordDiff<T>>;
const updatedFields = Object.keys(diff);
if (updatedFields.length === 0) {
return;
}
const event = new ObjectRecordUpdateEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = after.id;
event.properties = {
before,
after,
updatedFields,
diff,
};
return event;
})
.filter(isDefined);
break;
case DatabaseEventAction.DELETED:
events = entityArray.map((before) => {
const event = new ObjectRecordDeleteEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = before.id;
event.properties = { before };
return event;
});
break;
case DatabaseEventAction.DESTROYED:
events = entityArray.map((before) => {
const event = new ObjectRecordDestroyEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = before.id;
event.properties = { before };
return event;
});
break;
case DatabaseEventAction.UPSERTED:
events = entityArray.map((after, index) => {
const event = new ObjectRecordUpsertEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = after.id;
const before = beforeEntities
? Array.isArray(beforeEntities)
? beforeEntities[index]
: beforeEntities
: undefined;
let updatedFields;
let diff;
diff = objectRecordChangedValues(
before ?? {},
after,
objectMetadataItem,
) as Partial<ObjectRecordDiff<T>>;
updatedFields = Object.keys(diff);
event.properties = {
after,
...(before && { before }),
...(diff && { diff }),
...(updatedFields && { updatedFields }),
};
return event;
});
break;
default:
return;
}
if (!events.length) {
return;
}
return {
objectMetadataNameSingular,
action,
events,
objectMetadata: { ...objectMetadataItem, fields },
workspaceId,
};
};
@@ -0,0 +1,5 @@
export type CustomWorkspaceEventBatch<WorkspaceEvent> = {
name: string;
workspaceId?: string;
events: WorkspaceEvent[];
};
@@ -0,0 +1,8 @@
import type { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
export type WorkspaceEventBatch<WorkspaceEvent> = {
name: string;
workspaceId: string;
objectMetadata: Omit<ObjectMetadataEntity, 'indexMetadatas'>;
events: WorkspaceEvent[];
};
@@ -1,5 +0,0 @@
export type WorkspaceEventBatch<WorkspaceEvent> = {
name: string;
workspaceId: string;
events: WorkspaceEvent[];
};
@@ -2,22 +2,19 @@ import { Injectable } from '@nestjs/common';
import { EventEmitter2 } from '@nestjs/event-emitter';
import { isDefined } from 'twenty-shared/utils';
import { type ObjectLiteral } from 'typeorm';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import { type 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 { ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emitter/types/object-record-delete.event';
import { ObjectRecordDestroyEvent } from 'src/engine/core-modules/event-emitter/types/object-record-destroy.event';
import { type ObjectRecordDiff } from 'src/engine/core-modules/event-emitter/types/object-record-diff';
import { type ObjectRecordRestoreEvent } from 'src/engine/core-modules/event-emitter/types/object-record-restore.event';
import { ObjectRecordUpdateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-update.event';
import { ObjectRecordUpsertEvent } from 'src/engine/core-modules/event-emitter/types/object-record-upsert.event';
import { objectRecordChangedValues } from 'src/engine/core-modules/event-emitter/utils/object-record-changed-values';
import { type ObjectMetadataItemWithFieldMaps } from 'src/engine/metadata-modules/types/object-metadata-item-with-field-maps';
import { type CustomEventName } from 'src/engine/workspace-event-emitter/types/custom-event-name.type';
import { computeEventName } from 'src/engine/workspace-event-emitter/utils/compute-event-name';
import { STANDARD_OBJECT_IDS } from 'src/engine/workspace-manager/workspace-sync-metadata/constants/standard-object-ids';
import type { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
type ActionEventMap<T> = {
[DatabaseEventAction.CREATED]: ObjectRecordCreateEvent<T>;
@@ -28,159 +25,32 @@ type ActionEventMap<T> = {
[DatabaseEventAction.UPSERTED]: ObjectRecordUpsertEvent<T>;
};
export type DatabaseBatchEventInput<T, A extends keyof ActionEventMap<T>> = {
objectMetadataNameSingular: string;
action: A;
events: ActionEventMap<T>[A][];
objectMetadata: Omit<ObjectMetadataEntity, 'indexMetadatas'>;
workspaceId: string;
};
@Injectable()
export class WorkspaceEventEmitter {
constructor(private readonly eventEmitter: EventEmitter2) {}
async emitMutationEvent<T extends ObjectLiteral>({
action,
objectMetadataItem,
workspaceId,
authContext,
entities,
beforeEntities,
}: {
action: DatabaseEventAction;
objectMetadataItem: ObjectMetadataItemWithFieldMaps;
workspaceId: string;
authContext?: AuthContext;
entities: T | T[];
beforeEntities?: T | T[];
}) {
if (
objectMetadataItem.standardId === STANDARD_OBJECT_IDS.timelineActivity
) {
public emitDatabaseBatchEvent<T, A extends keyof ActionEventMap<T>>(
databaseBatchEventInput: DatabaseBatchEventInput<T, A> | undefined,
) {
if (!isDefined(databaseBatchEventInput)) {
return;
}
const objectMetadataNameSingular = objectMetadataItem.nameSingular;
const fields = Object.values(objectMetadataItem.fieldsById ?? {});
const entityArray = isDefined(entities)
? Array.isArray(entities)
? entities
: [entities]
: [];
let events: (
| ObjectRecordCreateEvent<T>
| ObjectRecordUpdateEvent<T>
| ObjectRecordDeleteEvent<T>
| ObjectRecordUpsertEvent<T>
)[] = [];
switch (action) {
case DatabaseEventAction.CREATED:
events = entityArray.map((after) => {
const event = new ObjectRecordCreateEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = after.id;
event.objectMetadata = { ...objectMetadataItem, fields };
event.properties = { after };
return event;
});
break;
case DatabaseEventAction.UPDATED:
events = entityArray
.map((after, idx) => {
if (!beforeEntities) {
throw new Error('beforeEntities is required for UPDATED action');
}
const before = Array.isArray(beforeEntities)
? beforeEntities?.[idx]
: beforeEntities;
const diff = objectRecordChangedValues(
before,
after,
objectMetadataItem,
) as Partial<ObjectRecordDiff<T>>;
const updatedFields = Object.keys(diff);
if (updatedFields.length === 0) {
return;
}
const event = new ObjectRecordUpdateEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = after.id;
event.objectMetadata = { ...objectMetadataItem, fields };
event.properties = {
before,
after,
updatedFields,
diff,
};
return event;
})
.filter(isDefined);
break;
case DatabaseEventAction.DELETED:
events = entityArray.map((before) => {
const event = new ObjectRecordDeleteEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = before.id;
event.objectMetadata = { ...objectMetadataItem, fields };
event.properties = { before };
return event;
});
break;
case DatabaseEventAction.DESTROYED:
events = entityArray.map((before) => {
const event = new ObjectRecordDestroyEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = before.id;
event.objectMetadata = { ...objectMetadataItem, fields };
event.properties = { before };
return event;
});
break;
case DatabaseEventAction.UPSERTED:
events = entityArray.map((after, index) => {
const event = new ObjectRecordUpsertEvent<T>();
event.userId = authContext?.user?.id;
event.recordId = after.id;
event.objectMetadata = { ...objectMetadataItem, fields };
const before = beforeEntities
? Array.isArray(beforeEntities)
? beforeEntities[index]
: beforeEntities
: undefined;
let updatedFields;
let diff;
diff = objectRecordChangedValues(
before ?? {},
after,
objectMetadataItem,
) as Partial<ObjectRecordDiff<T>>;
updatedFields = Object.keys(diff);
event.properties = {
after,
...(before && { before }),
...(diff && { diff }),
...(updatedFields && { updatedFields }),
};
return event;
});
break;
default:
return;
}
const {
objectMetadataNameSingular,
action,
events,
objectMetadata,
workspaceId,
} = databaseBatchEventInput;
if (!events.length) {
return;
@@ -188,35 +58,14 @@ export class WorkspaceEventEmitter {
const eventName = computeEventName(objectMetadataNameSingular, action);
this.eventEmitter.emit(eventName, {
const workspaceEventBatch: WorkspaceEventBatch<ActionEventMap<T>[A]> = {
name: eventName,
workspaceId,
objectMetadata,
events,
});
}
};
public emitDatabaseBatchEvent<T, A extends keyof ActionEventMap<T>>({
objectMetadataNameSingular,
action,
events,
workspaceId,
}: {
objectMetadataNameSingular: string;
action: A;
events: ActionEventMap<T>[A][];
workspaceId: string | undefined;
}) {
if (!events.length) {
return;
}
const eventName = `${objectMetadataNameSingular}.${action}`;
this.eventEmitter.emit(eventName, {
name: eventName,
workspaceId,
events,
});
this.eventEmitter.emit(eventName, workspaceEventBatch);
}
public emitCustomBatchEvent<T extends object>(
@@ -228,10 +77,12 @@ export class WorkspaceEventEmitter {
return;
}
this.eventEmitter.emit(eventName, {
const customWorkspaceEventBatch: CustomWorkspaceEventBatch<T> = {
name: eventName,
workspaceId,
events,
});
};
this.eventEmitter.emit(eventName, customWorkspaceEventBatch);
}
}
@@ -7,7 +7,7 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type BlocklistWorkspaceEntity } from 'src/modules/blocklist/standard-objects/blocklist.workspace-entity';
import { CalendarEventCleanerService } from 'src/modules/calendar/calendar-event-cleaner/services/calendar-event-cleaner.service';
import { type CalendarChannelEventAssociationWorkspaceEntity } from 'src/modules/calendar/common/standard-objects/calendar-channel-event-association.workspace-entity';
@@ -7,7 +7,7 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type BlocklistWorkspaceEntity } from 'src/modules/blocklist/standard-objects/blocklist.workspace-entity';
import { CalendarChannelSyncStatusService } from 'src/modules/calendar/common/services/calendar-channel-sync-status.service';
import {
@@ -6,7 +6,7 @@ import { type ObjectRecordUpdateEvent } 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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type BlocklistWorkspaceEntity } from 'src/modules/blocklist/standard-objects/blocklist.workspace-entity';
import {
BlocklistItemDeleteCalendarEventsJob,
@@ -4,7 +4,7 @@ import { type ObjectRecordDeleteEvent } 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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
DeleteConnectedAccountAssociatedCalendarDataJob,
type DeleteConnectedAccountAssociatedCalendarDataJobData,
@@ -11,7 +11,7 @@ import { objectRecordChangedProperties as objectRecordUpdateEventChangedProperti
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
CalendarEventParticipantMatchParticipantJob,
type CalendarEventParticipantMatchParticipantJobData,
@@ -8,7 +8,7 @@ import { objectRecordChangedProperties as objectRecordUpdateEventChangedProperti
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
CalendarEventParticipantMatchParticipantJob,
type CalendarEventParticipantMatchParticipantJobData,
@@ -7,10 +7,10 @@ import { Repository } from 'typeorm';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { InjectObjectMetadataRepository } from 'src/engine/object-metadata-repository/object-metadata-repository.decorator';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { type CalendarEventParticipantWorkspaceEntity } from 'src/modules/calendar/common/standard-objects/calendar-event-participant.workspace-entity';
import { TimelineActivityRepository } from 'src/modules/timeline/repositories/timeline-activity.repository';
import { TimelineActivityWorkspaceEntity } from 'src/modules/timeline/standard-objects/timeline-activity.workspace-entity';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
@Injectable()
export class CalendarEventParticipantListener {
@@ -23,11 +23,15 @@ export class CalendarEventParticipantListener {
@OnCustomBatchEvent('calendarEventParticipant_matched')
public async handleCalendarEventParticipantMatchedEvent(
batchEvent: WorkspaceEventBatch<{
batchEvent: CustomWorkspaceEventBatch<{
workspaceMemberId: string;
participants: CalendarEventParticipantWorkspaceEntity[];
}>,
): Promise<void> {
if (!isDefined(batchEvent.workspaceId)) {
return;
}
const calendarEventObjectMetadata =
await this.objectMetadataRepository.findOneOrFail({
where: {
@@ -7,7 +7,7 @@ import { type ObjectRecordDestroyEvent } from 'src/engine/core-modules/event-emi
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
DeleteWorkspaceMemberConnectedAccountsCleanupJob,
type DeleteWorkspaceMemberConnectedAccountsCleanupJobData,
@@ -4,7 +4,7 @@ import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runne
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import { type ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emitter/types/object-record-delete.event';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { AccountsToReconnectService } from 'src/modules/connected-account/services/accounts-to-reconnect.service';
import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity';
import { type WorkspaceMemberWorkspaceEntity } from 'src/modules/workspace-member/standard-objects/workspace-member.workspace-entity';
@@ -57,9 +57,9 @@ export class ConnectedAccountDeleteOnePreQueryHook
this.workspaceEventEmitter.emitDatabaseBatchEvent({
objectMetadataNameSingular: 'messageChannel',
action: DatabaseEventAction.DESTROYED,
objectMetadata,
events: messageChannels.map((messageChannel) => ({
recordId: messageChannel.id,
objectMetadata,
properties: {
before: messageChannel,
},
@@ -5,7 +5,7 @@ import { objectRecordChangedProperties } from 'src/engine/core-modules/event-emi
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
CalendarCreateCompanyAndContactAfterSyncJob,
type CalendarCreateCompanyAndContactAfterSyncJobData,
@@ -5,7 +5,7 @@ import { objectRecordChangedProperties } from 'src/engine/core-modules/event-emi
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity';
import {
MessagingCreateCompanyAndContactAfterSyncJob,
@@ -4,7 +4,7 @@ import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runne
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import { type ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emitter/types/object-record-delete.event';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type FavoriteFolderWorkspaceEntity } from 'src/modules/favorite-folder/standard-objects/favorite-folder.workspace-entity';
import { type FavoriteWorkspaceEntity } from 'src/modules/favorite/standard-objects/favorite.workspace-entity';
@@ -6,7 +6,7 @@ import { type ObjectRecordDeleteEvent } 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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
FavoriteDeletionJob,
type FavoriteDeletionJobData,
@@ -7,7 +7,7 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type BlocklistWorkspaceEntity } from 'src/modules/blocklist/standard-objects/blocklist.workspace-entity';
import { type MessageChannelMessageAssociationWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel-message-association.workspace-entity';
import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity';
@@ -7,7 +7,7 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type BlocklistWorkspaceEntity } from 'src/modules/blocklist/standard-objects/blocklist.workspace-entity';
import { MessageChannelSyncStatusService } from 'src/modules/messaging/common/services/message-channel-sync-status.service';
import {
@@ -6,7 +6,7 @@ import { type ObjectRecordUpdateEvent } 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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type BlocklistWorkspaceEntity } from 'src/modules/blocklist/standard-objects/blocklist.workspace-entity';
import {
BlocklistItemDeleteMessagesJob,
@@ -4,7 +4,7 @@ import { type ObjectRecordDeleteEvent } 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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity';
import {
MessagingConnectedAccountDeletionCleanupJob,
@@ -4,7 +4,7 @@ import { type ObjectRecordDeleteEvent } 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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity';
import {
MessagingCleanCacheJob,
@@ -11,7 +11,7 @@ import { objectRecordChangedProperties as objectRecordUpdateEventChangedProperti
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
MessageParticipantMatchParticipantJob,
type MessageParticipantMatchParticipantJobData,
@@ -13,7 +13,7 @@ import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decora
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 { Workspace } from 'src/engine/core-modules/workspace/workspace.entity';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
MessageParticipantMatchParticipantJob,
type MessageParticipantMatchParticipantJobData,
@@ -7,10 +7,10 @@ import { Repository } from 'typeorm';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { InjectObjectMetadataRepository } from 'src/engine/object-metadata-repository/object-metadata-repository.decorator';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { type MessageParticipantWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-participant.workspace-entity';
import { TimelineActivityRepository } from 'src/modules/timeline/repositories/timeline-activity.repository';
import { TimelineActivityWorkspaceEntity } from 'src/modules/timeline/standard-objects/timeline-activity.workspace-entity';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
@Injectable()
export class MessageParticipantListener {
@@ -23,11 +23,15 @@ export class MessageParticipantListener {
@OnCustomBatchEvent('messageParticipant_matched')
public async handleMessageParticipantMatched(
batchEvent: WorkspaceEventBatch<{
batchEvent: CustomWorkspaceEventBatch<{
workspaceMemberId: string;
participants: MessageParticipantWorkspaceEntity[];
}>,
): Promise<void> {
if (!isDefined(batchEvent.workspaceId)) {
return;
}
const messageObjectMetadata =
await this.objectMetadataRepository.findOneOrFail({
where: {
@@ -6,7 +6,7 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { TimelineActivityService } from 'src/modules/timeline/services/timeline-activity.service';
import { WorkspaceMemberWorkspaceEntity } from 'src/modules/workspace-member/standard-objects/workspace-member.workspace-entity';
@@ -21,6 +21,18 @@ export class UpsertTimelineActivityFromInternalEvent {
async handle(
workspaceEventBatch: WorkspaceEventBatch<ObjectRecordNonDestructiveEvent>,
): Promise<void> {
if (workspaceEventBatch.events.length === 0) {
return;
}
if (
workspaceEventBatch.objectMetadata.isSystem &&
workspaceEventBatch.objectMetadata.nameSingular !== 'noteTarget' &&
workspaceEventBatch.objectMetadata.nameSingular !== 'taskTarget'
) {
return;
}
const workspaceMemberRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceEventBatch.workspaceId,
@@ -48,35 +60,22 @@ export class UpsertTimelineActivityFromInternalEvent {
}
}
const filteredEvents = workspaceEventBatch.events
.filter((event) => {
return (
!event.objectMetadata.isSystem ||
event.objectMetadata.nameSingular === 'noteTarget' ||
event.objectMetadata.nameSingular === 'taskTarget'
);
})
.map((event) => {
if ('diff' in event.properties && event.properties.diff) {
return {
...event,
properties: {
diff: event.properties.diff,
},
};
}
const formattedEvents = workspaceEventBatch.events.map((event) => {
if ('diff' in event.properties && event.properties.diff) {
return {
...event,
properties: {
diff: event.properties.diff,
},
};
}
return event;
});
if (filteredEvents.length === 0) {
return;
}
return event;
});
await this.timelineActivityService.upsertEvents({
events: filteredEvents,
eventName: workspaceEventBatch.name,
workspaceId: workspaceEventBatch.workspaceId,
...workspaceEventBatch,
events: formattedEvents,
});
}
}
@@ -12,15 +12,10 @@ import { parseEventNameOrThrow } from 'src/engine/workspace-event-emitter/utils/
import { TimelineActivityRepository } from 'src/modules/timeline/repositories/timeline-activity.repository';
import { TimelineActivityWorkspaceEntity } from 'src/modules/timeline/standard-objects/timeline-activity.workspace-entity';
import { type TimelineActivityPayload } from 'src/modules/timeline/types/timeline-activity-payload';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
type ActivityType = 'note' | 'task';
type EventsWithNameAndWorkspaceId = {
events: ObjectRecordBaseEvent[];
eventName: string;
workspaceId: string;
};
@Injectable()
export class TimelineActivityService {
constructor(
@@ -36,16 +31,22 @@ export class TimelineActivityService {
async upsertEvents({
events,
eventName,
name,
objectMetadata,
workspaceId,
}: EventsWithNameAndWorkspaceId) {
const { objectSingularName } = parseEventNameOrThrow(eventName);
}: WorkspaceEventBatch<ObjectRecordBaseEvent>) {
if (!isDefined(workspaceId)) {
return;
}
const { objectSingularName } = parseEventNameOrThrow(name);
const timelineActivitiesPayloads =
await this.transformEventsToTimelineActivityPayloads({
events,
objectMetadata,
workspaceId,
eventName,
name,
});
if (
@@ -82,11 +83,12 @@ export class TimelineActivityService {
private async transformEventsToTimelineActivityPayloads({
events,
workspaceId,
eventName,
}: EventsWithNameAndWorkspaceId): Promise<
objectMetadata,
name,
}: WorkspaceEventBatch<ObjectRecordBaseEvent>): Promise<
TimelineActivityPayload[] | undefined
> {
const { objectSingularName } = parseEventNameOrThrow(eventName);
const { objectSingularName } = parseEventNameOrThrow(name);
if (objectSingularName === 'note') {
const noteEventsTimelineActivities =
@@ -94,13 +96,14 @@ export class TimelineActivityService {
events,
activityType: 'note',
workspaceId,
eventName,
objectMetadata,
name,
});
return [
...noteEventsTimelineActivities,
...(events.map((event) => ({
name: eventName,
name,
objectSingularName,
recordId: event.recordId,
workspaceMemberId: event.workspaceMemberId,
@@ -115,13 +118,14 @@ export class TimelineActivityService {
events,
activityType: 'task',
workspaceId,
eventName,
objectMetadata,
name,
});
return [
...taskEventsTimelineActivities,
...(events.map((event) => ({
name: eventName,
name,
objectSingularName,
recordId: event.recordId,
workspaceMemberId: event.workspaceMemberId,
@@ -138,12 +142,13 @@ export class TimelineActivityService {
events,
activityType: objectSingularName === 'noteTarget' ? 'note' : 'task',
workspaceId,
eventName,
objectMetadata,
name,
});
}
return events.map((event) => ({
name: eventName,
name,
objectSingularName,
recordId: event.recordId,
workspaceMemberId: event.workspaceMemberId,
@@ -154,12 +159,17 @@ export class TimelineActivityService {
private async computeTimelineActivityPayloadsForActivities({
events,
activityType,
eventName,
name,
workspaceId,
}: EventsWithNameAndWorkspaceId & { activityType: ActivityType }): Promise<
TimelineActivityPayload[]
> {
const { action } = parseEventNameOrThrow(eventName);
objectMetadata,
}: WorkspaceEventBatch<ObjectRecordBaseEvent> & {
activityType: ActivityType;
}): Promise<TimelineActivityPayload[]> {
if (!isDefined(workspaceId)) {
return [];
}
const { action } = parseEventNameOrThrow(name);
const activityTargetRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
@@ -218,9 +228,9 @@ export class TimelineActivityService {
recordId: activityTarget[targetColumn.replace(/Id$/, '')],
linkedRecordCachedName: activityTitle,
linkedRecordId: activityId,
linkedObjectMetadataId: event.objectMetadata.id,
linkedObjectMetadataId: objectMetadata.id,
properties: event.properties,
overrideObjectSingularName: event.objectMetadata.nameSingular,
overrideObjectSingularName: objectMetadata.nameSingular,
} satisfies TimelineActivityPayload;
});
})
@@ -230,12 +240,13 @@ export class TimelineActivityService {
private async computeTimelineActivityPayloadsForActivityTargets({
events,
activityType,
eventName,
name,
objectMetadata,
workspaceId,
}: EventsWithNameAndWorkspaceId & { activityType: ActivityType }): Promise<
TimelineActivityPayload[]
> {
const { action } = parseEventNameOrThrow(eventName);
}: WorkspaceEventBatch<ObjectRecordBaseEvent> & {
activityType: ActivityType;
}): Promise<TimelineActivityPayload[]> {
const { action } = parseEventNameOrThrow(name);
const activityRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
@@ -277,7 +288,7 @@ export class TimelineActivityService {
return;
}
const activityObjectMetadataId = event.objectMetadata.fields.find(
const activityObjectMetadataId = objectMetadata.fields.find(
(field) => field.name === activityType,
)?.relationTargetObjectMetadataId;
@@ -1,11 +1,12 @@
import { Injectable } from '@nestjs/common';
import { isDefined } from 'twenty-shared/utils';
import { type ObjectRecordCreateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-create.event';
import { type ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emitter/types/object-record-delete.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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import {
WorkflowVersionStatus,
type WorkflowVersionWorkspaceEntity,
@@ -20,6 +21,7 @@ import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runne
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import { WORKFLOW_VERSION_STATUS_UPDATED } from 'src/modules/workflow/workflow-status/constants/workflow-version-status-updated.constants';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
@Injectable()
export class WorkflowVersionStatusListener {
@@ -30,10 +32,14 @@ export class WorkflowVersionStatusListener {
@OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.CREATED)
async handleWorkflowVersionCreated(
batchEvent: WorkspaceEventBatch<
batchEvent: CustomWorkspaceEventBatch<
ObjectRecordCreateEvent<WorkflowVersionWorkspaceEntity>
>,
): Promise<void> {
if (!isDefined(batchEvent.workspaceId)) {
return;
}
const workflowIds = batchEvent.events
.filter(
(event) =>
@@ -58,8 +64,12 @@ export class WorkflowVersionStatusListener {
@OnCustomBatchEvent(WORKFLOW_VERSION_STATUS_UPDATED)
async handleWorkflowVersionUpdated(
batchEvent: WorkspaceEventBatch<WorkflowVersionStatusUpdate>,
batchEvent: CustomWorkspaceEventBatch<WorkflowVersionStatusUpdate>,
): Promise<void> {
if (!isDefined(batchEvent.workspaceId)) {
return;
}
await this.messageQueueService.add<WorkflowVersionBatchEvent>(
WorkflowStatusesUpdateJob.name,
{
@@ -72,10 +82,14 @@ export class WorkflowVersionStatusListener {
@OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.DELETED)
async handleWorkflowVersionDeleted(
batchEvent: WorkspaceEventBatch<
batchEvent: CustomWorkspaceEventBatch<
ObjectRecordDeleteEvent<WorkflowVersionWorkspaceEntity>
>,
): Promise<void> {
if (!isDefined(batchEvent.workspaceId)) {
return;
}
const workflowIds = batchEvent.events
.filter(
(event) =>
@@ -96,30 +96,30 @@ describe('WorkflowDatabaseEventTriggerListener', () => {
const mockPayload = {
workspaceId,
name: databaseEventName,
objectMetadata: getMockObjectMetadataEntity({
id: 'test-object-metadata',
workspaceId,
nameSingular: 'testObject',
namePlural: 'testObjects',
labelSingular: 'Test Object',
labelPlural: 'Test Objects',
description: 'Test object for testing',
targetTableName: 'test_objects',
isSystem: false,
isCustom: false,
isActive: true,
isRemote: false,
isAuditLogged: true,
isSearchable: true,
createdAt: new Date(),
updatedAt: new Date(),
fields: [],
indexMetadatas: [],
icon: 'Icon123',
}),
events: [
{
recordId: 'test-record',
objectMetadata: getMockObjectMetadataEntity({
id: 'test-object-metadata',
workspaceId,
nameSingular: 'testObject',
namePlural: 'testObjects',
labelSingular: 'Test Object',
labelPlural: 'Test Objects',
description: 'Test object for testing',
targetTableName: 'test_objects',
isSystem: false,
isCustom: false,
isActive: true,
isRemote: false,
isAuditLogged: true,
isSearchable: true,
createdAt: new Date(),
updatedAt: new Date(),
fields: [],
indexMetadatas: [],
icon: 'Icon123',
}),
properties: {
updatedFields: ['field1', 'field2'],
before: { field1: 'old', field2: 'old' },
@@ -19,7 +19,7 @@ import { MessageQueueService } from 'src/engine/core-modules/message-queue/servi
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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import {
AutomatedTriggerType,
type WorkflowAutomatedTriggerWorkspaceEntity,
@@ -138,7 +138,7 @@ export class WorkflowDatabaseEventTriggerListener {
const workspaceId = payload.workspaceId;
const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
payload.events[0].objectMetadata.nameSingular,
payload.objectMetadata.nameSingular,
workspaceId,
);
@@ -156,7 +156,7 @@ export class WorkflowDatabaseEventTriggerListener {
const workspaceId = payload.workspaceId;
const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
payload.events[0].objectMetadata.nameSingular,
payload.objectMetadata.nameSingular,
workspaceId,
);
@@ -180,7 +180,7 @@ export class WorkflowDatabaseEventTriggerListener {
const workspaceId = payload.workspaceId;
const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
payload.events[0].objectMetadata.nameSingular,
payload.objectMetadata.nameSingular,
workspaceId,
);
@@ -198,7 +198,7 @@ export class WorkflowDatabaseEventTriggerListener {
const workspaceId = payload.workspaceId;
const { objectMetadataMaps, objectMetadataItemWithFieldsMaps } =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
payload.events[0].objectMetadata.nameSingular,
payload.objectMetadata.nameSingular,
workspaceId,
);
@@ -1,4 +1,4 @@
import { Bundle, ZObject } from 'zapier-platform-core';
import { type Bundle, type ZObject } from 'zapier-platform-core';
import handleQueryParams from '../../utils/handleQueryParams';
import requestDb, {
@@ -93,9 +93,7 @@ export const performList = async (
return results.map((result) => ({
record: result,
...(bundle.inputData.operation === DatabaseEventAction.UPDATED && {
updatedFields: [
Object.keys(result).filter((key) => key !== 'id')?.[0],
] || ['updatedField'],
updatedFields: [Object.keys(result).filter((key) => key !== 'id')?.[0]],
}),
}));
};