From ae4a93c993dd179288227d3ddde422905da1a63b Mon Sep 17 00:00:00 2001
From: Harshit Singh <73997189+harshit078@users.noreply.github.com>
Date: Wed, 1 Oct 2025 18:47:05 +0530
Subject: [PATCH] feat: add create or update workflow trigger (#14708)
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
### Description
This PR focuses on
https://github.com/twentyhq/core-team-issues/issues/1476
- introduces the handling of upserted events in the workflow trigger
system
## Visual Appearance
### Changes
- Added `UPSERTED` action to the `DatabaseEventAction` enum.
- Updated the `WorkflowEditTriggerDatabaseEventForm` to support upserted
events.
---------
Co-authored-by: Thomas Trompette
---
.../src/generated-metadata/graphql.ts | 3 +-
.../twenty-front/src/generated/graphql.ts | 3 +-
.../splitWorkflowTriggerEventName.test.ts | 11 ++++++
.../WorkflowEditTriggerDatabaseEventForm.tsx | 4 +-
.../constants/DatabaseTriggerDefaultLabel.ts | 1 +
.../constants/DatabaseTriggerTypes.ts | 6 +++
.../enums/database-event-action.ts | 1 +
.../create-audit-log-from-internal-event.ts | 7 ++++
.../core-modules/audit/types/events.type.ts | 6 +++
.../object-event/object-record-upserted.ts | 15 ++++++++
.../types/object-record-event.event.ts | 4 +-
.../object-record-non-destructive-event.ts | 4 +-
.../types/object-record-upsert.event.ts | 13 +++++++
.../workspace-insert-query-builder.ts | 8 ++++
.../workspace-update-query-builder.ts | 18 +++++++++
.../workspace-event-emitter.ts | 38 +++++++++++++++++++
.../generate-fake-object-record-event.spec.ts | 31 +++++++++++++++
.../generate-fake-object-record-event.ts | 1 +
.../create-record.workflow-action.ts | 5 ---
.../constants/automated-trigger-settings.ts | 7 +++-
...orkflow-database-event-trigger.listener.ts | 35 ++++++++++++++++-
.../schemas/database-event-trigger-schema.ts | 8 ++--
22 files changed, 213 insertions(+), 16 deletions(-)
create mode 100644 packages/twenty-server/src/engine/core-modules/audit/utils/events/object-event/object-record-upserted.ts
create mode 100644 packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-upsert.event.ts
diff --git a/packages/twenty-front/src/generated-metadata/graphql.ts b/packages/twenty-front/src/generated-metadata/graphql.ts
index 0c912abec4..e8d9eace4a 100644
--- a/packages/twenty-front/src/generated-metadata/graphql.ts
+++ b/packages/twenty-front/src/generated-metadata/graphql.ts
@@ -966,7 +966,8 @@ export enum DatabaseEventAction {
DELETED = 'DELETED',
DESTROYED = 'DESTROYED',
RESTORED = 'RESTORED',
- UPDATED = 'UPDATED'
+ UPDATED = 'UPDATED',
+ UPSERTED = 'UPSERTED'
}
export type DatabaseEventTrigger = {
diff --git a/packages/twenty-front/src/generated/graphql.ts b/packages/twenty-front/src/generated/graphql.ts
index 565e666294..fd430ec221 100644
--- a/packages/twenty-front/src/generated/graphql.ts
+++ b/packages/twenty-front/src/generated/graphql.ts
@@ -916,7 +916,8 @@ export enum DatabaseEventAction {
DELETED = 'DELETED',
DESTROYED = 'DESTROYED',
RESTORED = 'RESTORED',
- UPDATED = 'UPDATED'
+ UPDATED = 'UPDATED',
+ UPSERTED = 'UPSERTED'
}
export type DatabaseEventTrigger = {
diff --git a/packages/twenty-front/src/modules/workflow/utils/__tests__/splitWorkflowTriggerEventName.test.ts b/packages/twenty-front/src/modules/workflow/utils/__tests__/splitWorkflowTriggerEventName.test.ts
index 64bfd43240..8395a77646 100644
--- a/packages/twenty-front/src/modules/workflow/utils/__tests__/splitWorkflowTriggerEventName.test.ts
+++ b/packages/twenty-front/src/modules/workflow/utils/__tests__/splitWorkflowTriggerEventName.test.ts
@@ -99,4 +99,15 @@ describe('splitWorkflowTriggerEventName', () => {
event: '',
});
});
+
+ it('should split event name with upserted event', () => {
+ const eventName = 'company.upserted';
+
+ const result = splitWorkflowTriggerEventName(eventName);
+
+ expect(result).toEqual({
+ objectType: 'company',
+ event: 'upserted',
+ });
+ });
});
diff --git a/packages/twenty-front/src/modules/workflow/workflow-trigger/components/WorkflowEditTriggerDatabaseEventForm.tsx b/packages/twenty-front/src/modules/workflow/workflow-trigger/components/WorkflowEditTriggerDatabaseEventForm.tsx
index 6e14786e71..c184771891 100644
--- a/packages/twenty-front/src/modules/workflow/workflow-trigger/components/WorkflowEditTriggerDatabaseEventForm.tsx
+++ b/packages/twenty-front/src/modules/workflow/workflow-trigger/components/WorkflowEditTriggerDatabaseEventForm.tsx
@@ -82,6 +82,8 @@ export const WorkflowEditTriggerDatabaseEventForm = ({
trigger.settings.eventName,
);
const isUpdateEvent = triggerEvent.event === 'updated';
+ const isUpsertEvent = triggerEvent.event === 'upserted';
+ const isFieldFilteringSupported = isUpdateEvent || isUpsertEvent;
const regularObjects = objectMetadataItems
.filter((item) => item.isActive && !item.isSystem)
@@ -271,7 +273,7 @@ export const WorkflowEditTriggerDatabaseEventForm = ({
dropdownOffset={{ y: parseInt(theme.spacing(1), 10) }}
/>
- {isDefined(selectedObjectMetadataItem) && isUpdateEvent && (
+ {isDefined(selectedObjectMetadataItem) && isFieldFilteringSupported && (
=
diff --git a/packages/twenty-server/src/engine/core-modules/audit/utils/events/object-event/object-record-upserted.ts b/packages/twenty-server/src/engine/core-modules/audit/utils/events/object-event/object-record-upserted.ts
new file mode 100644
index 0000000000..85ab5e7966
--- /dev/null
+++ b/packages/twenty-server/src/engine/core-modules/audit/utils/events/object-event/object-record-upserted.ts
@@ -0,0 +1,15 @@
+import { z } from 'zod';
+
+import { registerEvent } from 'src/engine/core-modules/audit/utils/events/workspace-event/track';
+
+export const OBJECT_RECORD_UPSERTED_EVENT = 'Object Record Upserted' as const;
+export const objectRecordUpsertedSchema = z.object({
+ event: z.literal(OBJECT_RECORD_UPSERTED_EVENT),
+ properties: z.looseObject({}),
+});
+
+export type ObjectRecordUpsertedTrackEvent = z.infer<
+ typeof objectRecordUpsertedSchema
+>;
+
+registerEvent(OBJECT_RECORD_UPSERTED_EVENT, objectRecordUpsertedSchema);
diff --git a/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-event.event.ts b/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-event.event.ts
index 045cfe629f..b721f4d0e5 100644
--- a/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-event.event.ts
+++ b/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-event.event.ts
@@ -3,10 +3,12 @@ import { type ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emit
import { type ObjectRecordDestroyEvent } from 'src/engine/core-modules/event-emitter/types/object-record-destroy.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 { type ObjectRecordUpsertEvent } from 'src/engine/core-modules/event-emitter/types/object-record-upsert.event';
export type ObjectRecordEvent =
| ObjectRecordUpdateEvent
| ObjectRecordDeleteEvent
| ObjectRecordCreateEvent
| ObjectRecordDestroyEvent
- | ObjectRecordRestoreEvent;
+ | ObjectRecordRestoreEvent
+ | ObjectRecordUpsertEvent;
diff --git a/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-non-destructive-event.ts b/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-non-destructive-event.ts
index 948c176372..fd36fcf575 100644
--- a/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-non-destructive-event.ts
+++ b/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-non-destructive-event.ts
@@ -2,9 +2,11 @@ import { type ObjectRecordCreateEvent } from 'src/engine/core-modules/event-emit
import { type ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emitter/types/object-record-delete.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 { type ObjectRecordUpsertEvent } from 'src/engine/core-modules/event-emitter/types/object-record-upsert.event';
export type ObjectRecordNonDestructiveEvent =
| ObjectRecordCreateEvent
| ObjectRecordUpdateEvent
| ObjectRecordDeleteEvent
- | ObjectRecordRestoreEvent;
+ | ObjectRecordRestoreEvent
+ | ObjectRecordUpsertEvent;
diff --git a/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-upsert.event.ts b/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-upsert.event.ts
new file mode 100644
index 0000000000..e1d165b2b6
--- /dev/null
+++ b/packages/twenty-server/src/engine/core-modules/event-emitter/types/object-record-upsert.event.ts
@@ -0,0 +1,13 @@
+import { type ObjectRecordDiff } from 'src/engine/core-modules/event-emitter/types/object-record-diff';
+import { ObjectRecordBaseEvent } from 'src/engine/core-modules/event-emitter/types/object-record.base.event';
+
+export class ObjectRecordUpsertEvent<
+ T = object,
+> extends ObjectRecordBaseEvent {
+ properties: {
+ before?: T;
+ after: T;
+ diff?: Partial>;
+ updatedFields?: string[];
+ };
+}
diff --git a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts
index 18dc54fc5e..fb43a411c2 100644
--- a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts
+++ b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts
@@ -168,6 +168,14 @@ export class WorkspaceInsertQueryBuilder<
authContext: this.authContext,
});
+ await this.internalContext.eventEmitterService.emitMutationEvent({
+ 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)
? []
diff --git a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts
index adec41c277..4fe9c41101 100644
--- a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts
+++ b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts
@@ -161,6 +161,15 @@ export class WorkspaceUpdateQueryBuilder<
authContext: this.authContext,
});
+ await this.internalContext.eventEmitterService.emitMutationEvent({
+ action: DatabaseEventAction.UPSERTED,
+ objectMetadataItem: objectMetadata,
+ workspaceId: this.internalContext.workspaceId,
+ entities: formattedAfter,
+ beforeEntities: formattedBefore,
+ authContext: this.authContext,
+ });
+
const formattedResult = formatResult(
result.raw,
objectMetadata,
@@ -293,6 +302,15 @@ export class WorkspaceUpdateQueryBuilder<
authContext: this.authContext,
});
+ await this.internalContext.eventEmitterService.emitMutationEvent({
+ action: DatabaseEventAction.UPSERTED,
+ objectMetadataItem: objectMetadata,
+ workspaceId: this.internalContext.workspaceId,
+ entities: formattedAfter,
+ beforeEntities: formattedBefore,
+ authContext: this.authContext,
+ });
+
const formattedResults = formatResult(
results.flatMap((result) => result.raw),
objectMetadata,
diff --git a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.ts b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.ts
index 1178248be0..15d6bf522d 100644
--- a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.ts
+++ b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.ts
@@ -12,6 +12,7 @@ import { ObjectRecordDestroyEvent } from 'src/engine/core-modules/event-emitter/
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';
@@ -24,6 +25,7 @@ type ActionEventMap = {
[DatabaseEventAction.DELETED]: ObjectRecordDeleteEvent;
[DatabaseEventAction.DESTROYED]: ObjectRecordDestroyEvent;
[DatabaseEventAction.RESTORED]: ObjectRecordRestoreEvent;
+ [DatabaseEventAction.UPSERTED]: ObjectRecordUpsertEvent;
};
@Injectable()
@@ -62,6 +64,7 @@ export class WorkspaceEventEmitter {
| ObjectRecordCreateEvent
| ObjectRecordUpdateEvent
| ObjectRecordDeleteEvent
+ | ObjectRecordUpsertEvent
)[] = [];
switch (action) {
@@ -140,6 +143,41 @@ export class WorkspaceEventEmitter {
return event;
});
break;
+ case DatabaseEventAction.UPSERTED:
+ events = entityArray.map((after, index) => {
+ const event = new ObjectRecordUpsertEvent();
+
+ 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>;
+
+ updatedFields = Object.keys(diff);
+
+ event.properties = {
+ after,
+ ...(before && { before }),
+ ...(diff && { diff }),
+ ...(updatedFields && { updatedFields }),
+ };
+
+ return event;
+ });
+ break;
default:
return;
}
diff --git a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/__tests__/generate-fake-object-record-event.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/__tests__/generate-fake-object-record-event.spec.ts
index 38beaf878d..e28882fe59 100644
--- a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/__tests__/generate-fake-object-record-event.spec.ts
+++ b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/__tests__/generate-fake-object-record-event.spec.ts
@@ -173,6 +173,37 @@ describe('generateFakeObjectRecordEvent', () => {
});
});
+ it('should generate record with "after" prefix for UPSERTED action', () => {
+ const result = generateFakeObjectRecordEvent(
+ objectMetadataInfo,
+ DatabaseEventAction.UPSERTED,
+ );
+
+ expect(result).toEqual({
+ object: {
+ isLeaf: true,
+ icon: 'test-company-icon',
+ label: 'Company',
+ value: 'A company',
+ fieldIdName: 'properties.after.id',
+ objectMetadataId: '20202020-c03c-45d6-a4b0-04afe1357c5c',
+ },
+ fields: {
+ 'properties.after.field1': {
+ type: 'TEXT',
+ value: 'test',
+ fieldMetadataId: '123e4567-e89b-12d3-a456-426614174000',
+ },
+ 'properties.after.field2': {
+ type: 'NUMBER',
+ value: 123,
+ fieldMetadataId: '123e4567-e89b-12d3-a456-426614174001',
+ },
+ },
+ _outputSchemaType: 'RECORD',
+ });
+ });
+
it('should throw error for unknown action', () => {
expect(() => {
generateFakeObjectRecordEvent(
diff --git a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/generate-fake-object-record-event.ts b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/generate-fake-object-record-event.ts
index 8a243c3905..8b59e3bf46 100644
--- a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/generate-fake-object-record-event.ts
+++ b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/utils/generate-fake-object-record-event.ts
@@ -45,6 +45,7 @@ export const generateFakeObjectRecordEvent = (
switch (action) {
case DatabaseEventAction.CREATED:
case DatabaseEventAction.UPDATED:
+ case DatabaseEventAction.UPSERTED:
return generateFakeObjectRecordEventWithPrefix({
objectMetadataInfo,
prefix: 'properties.after',
diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts
index 23a32edbc2..ed2e5c018f 100644
--- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts
+++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts
@@ -1,17 +1,14 @@
import { Injectable } from '@nestjs/common';
-import { InjectRepository } from '@nestjs/typeorm';
import { isDefined } from 'class-validator';
import { resolveInput } from 'twenty-shared/utils';
import { canObjectBeManagedByWorkflow } from 'twenty-shared/workflow';
-import { Repository } from 'typeorm';
import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/interfaces/workflow-action.interface';
import { RecordPositionService } from 'src/engine/core-modules/record-position/services/record-position.service';
import { RecordInputTransformerService } from 'src/engine/core-modules/record-transformer/services/record-input-transformer.service';
import { FieldActorSource } from 'src/engine/metadata-modules/field-metadata/composite-types/actor.composite-type';
-import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { ScopedWorkspaceContextFactory } from 'src/engine/twenty-orm/factories/scoped-workspace-context.factory';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
@@ -28,8 +25,6 @@ import { type WorkflowCreateRecordActionInput } from 'src/modules/workflow/workf
export class CreateRecordWorkflowAction implements WorkflowAction {
constructor(
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
- @InjectRepository(ObjectMetadataEntity)
- private readonly objectMetadataRepository: Repository,
private readonly scopedWorkspaceContextFactory: ScopedWorkspaceContextFactory,
private readonly recordPositionService: RecordPositionService,
private readonly recordInputTransformerService: RecordInputTransformerService,
diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings.ts
index e60fb71e95..abcf1a2472 100644
--- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings.ts
+++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings.ts
@@ -4,12 +4,17 @@ export type BaseDatabaseEventTriggerSettings = {
export type DatabaseEventTriggerSettings =
| BaseDatabaseEventTriggerSettings
- | UpdateEventTriggerSettings;
+ | UpdateEventTriggerSettings
+ | UpsertEventTriggerSettings;
export type UpdateEventTriggerSettings = BaseDatabaseEventTriggerSettings & {
fields: string[];
};
+export type UpsertEventTriggerSettings = BaseDatabaseEventTriggerSettings & {
+ fields: string[];
+};
+
export type CronTriggerSettings = {
pattern: string;
};
diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts
index 89f0328e47..0d176075a8 100644
--- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts
+++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts
@@ -12,6 +12,7 @@ import { type ObjectRecordDeleteEvent } from 'src/engine/core-modules/event-emit
import { type ObjectRecordDestroyEvent } from 'src/engine/core-modules/event-emitter/types/object-record-destroy.event';
import { type ObjectRecordNonDestructiveEvent } from 'src/engine/core-modules/event-emitter/types/object-record-non-destructive-event';
import { type ObjectRecordUpdateEvent } from 'src/engine/core-modules/event-emitter/types/object-record-update.event';
+import { type ObjectRecordUpsertEvent } from 'src/engine/core-modules/event-emitter/types/object-record-upsert.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';
@@ -22,7 +23,10 @@ import {
type WorkflowAutomatedTriggerWorkspaceEntity,
} from 'src/modules/workflow/common/standard-objects/workflow-automated-trigger.workspace-entity';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
-import { type UpdateEventTriggerSettings } from 'src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings';
+import {
+ type UpdateEventTriggerSettings,
+ type UpsertEventTriggerSettings,
+} from 'src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings';
import {
WorkflowTriggerJob,
type WorkflowTriggerJobData,
@@ -110,6 +114,22 @@ export class WorkflowDatabaseEventTriggerListener {
});
}
+ @OnDatabaseBatchEvent('*', DatabaseEventAction.UPSERTED)
+ async handleObjectRecordUpsertEvent(
+ payload: WorkspaceEventBatch,
+ ) {
+ if (await this.shouldIgnoreEvent(payload)) {
+ return;
+ }
+
+ const clonedPayload = structuredClone(payload);
+
+ await this.handleEvent({
+ payload: clonedPayload,
+ action: DatabaseEventAction.UPSERTED,
+ });
+ }
+
private async enrichCreatedEvent(
payload: WorkspaceEventBatch,
) {
@@ -316,6 +336,19 @@ export class WorkflowDatabaseEventTriggerListener {
);
}
+ if (action === DatabaseEventAction.UPSERTED) {
+ const settings = eventListener.settings as UpsertEventTriggerSettings;
+ const upsertEventPayload = eventPayload as ObjectRecordUpsertEvent;
+
+ return (
+ !settings.fields ||
+ settings.fields.length === 0 ||
+ settings.fields.some((field) =>
+ upsertEventPayload?.properties?.updatedFields?.includes(field),
+ )
+ );
+ }
+
return true;
}
}
diff --git a/packages/twenty-shared/src/workflow/schemas/database-event-trigger-schema.ts b/packages/twenty-shared/src/workflow/schemas/database-event-trigger-schema.ts
index b9024b1e7e..be0dfa5b6a 100644
--- a/packages/twenty-shared/src/workflow/schemas/database-event-trigger-schema.ts
+++ b/packages/twenty-shared/src/workflow/schemas/database-event-trigger-schema.ts
@@ -8,11 +8,11 @@ export const workflowDatabaseEventTriggerSchema = baseTriggerSchema
eventName: z
.string()
.regex(
- /^[a-z][a-zA-Z0-9_]*\.(created|updated|deleted)$/,
- 'Event name must follow the pattern: objectName.action (e.g., "company.created", "person.updated")',
+ /^[a-z][a-zA-Z0-9_]*\.(created|updated|deleted|upserted)$/,
+ 'Event name must follow the pattern: objectName.action (e.g., "company.created", "person.updated", "company.upserted")',
)
.describe(
- 'Event name in format: objectName.action (e.g., "company.created", "person.updated", "task.deleted"). Use lowercase object names.',
+ 'Event name in format: objectName.action (e.g., "company.created", "person.updated", "task.deleted", "company.upserted"). Use lowercase object names.',
),
input: z.looseObject({}).optional(),
outputSchema: z
@@ -25,5 +25,5 @@ export const workflowDatabaseEventTriggerSchema = baseTriggerSchema
}),
})
.describe(
- 'Database event trigger that fires when a record is created, updated, or deleted. The triggered record is accessible in workflow steps via {{trigger.object.fieldName}}.',
+ 'Database event trigger that fires when a record is created, updated, deleted, or upserted. The triggered record is accessible in workflow steps via {{trigger.object.fieldName}}.',
);