feat: add create or update workflow trigger (#14708)

### 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
<img width="1792" height="1039" alt="Screenshot 2025-09-21 at 11 40
28 PM"
src="https://github.com/user-attachments/assets/d926b566-f7d3-412b-b1ab-ed2064425f3d"
/>

<img width="501" height="1025" alt="Screenshot 2025-09-21 at 11 40
39 PM"
src="https://github.com/user-attachments/assets/13c3031a-0762-4c80-8df1-fa74ca70a540"
/>


### Changes
- Added `UPSERTED` action to the `DatabaseEventAction` enum.
- Updated the `WorkflowEditTriggerDatabaseEventForm` to support upserted
events.

---------

Co-authored-by: Thomas Trompette <thomas.trompette@sfr.fr>
This commit is contained in:
Harshit Singh
2025-10-01 18:47:05 +05:30
committed by GitHub
parent b5e6703f35
commit ae4a93c993
22 changed files with 213 additions and 16 deletions
@@ -966,7 +966,8 @@ export enum DatabaseEventAction {
DELETED = 'DELETED',
DESTROYED = 'DESTROYED',
RESTORED = 'RESTORED',
UPDATED = 'UPDATED'
UPDATED = 'UPDATED',
UPSERTED = 'UPSERTED'
}
export type DatabaseEventTrigger = {
@@ -916,7 +916,8 @@ export enum DatabaseEventAction {
DELETED = 'DELETED',
DESTROYED = 'DESTROYED',
RESTORED = 'RESTORED',
UPDATED = 'UPDATED'
UPDATED = 'UPDATED',
UPSERTED = 'UPSERTED'
}
export type DatabaseEventTrigger = {
@@ -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',
});
});
});
@@ -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) }}
/>
</StyledRecordTypeSelectContainer>
{isDefined(selectedObjectMetadataItem) && isUpdateEvent && (
{isDefined(selectedObjectMetadataItem) && isFieldFilteringSupported && (
<WorkflowFieldsMultiSelect
label="Fields (Optional)"
placeholder="Select specific fields to listen to"
@@ -2,4 +2,5 @@ export enum DatabaseTriggerDefaultLabel {
RECORD_IS_CREATED = 'Record is created',
RECORD_IS_UPDATED = 'Record is updated',
RECORD_IS_DELETED = 'Record is deleted',
RECORD_UPSERTED = 'Record is created or updated',
}
@@ -25,4 +25,10 @@ export const DATABASE_TRIGGER_TYPES: Array<{
icon: 'IconTrash',
event: 'deleted',
},
{
defaultLabel: DatabaseTriggerDefaultLabel.RECORD_UPSERTED,
type: 'DATABASE_EVENT',
icon: 'IconPencilPlus',
event: 'upserted',
},
];
@@ -4,4 +4,5 @@ export enum DatabaseEventAction {
DELETED = 'deleted',
DESTROYED = 'destroyed',
RESTORED = 'restored',
UPSERTED = 'upserted',
}
@@ -2,6 +2,7 @@ import { AuditService } from 'src/engine/core-modules/audit/services/audit.servi
import { OBJECT_RECORD_CREATED_EVENT } from 'src/engine/core-modules/audit/utils/events/object-event/object-record-created';
import { OBJECT_RECORD_DELETED_EVENT } from 'src/engine/core-modules/audit/utils/events/object-event/object-record-delete';
import { OBJECT_RECORD_UPDATED_EVENT } from 'src/engine/core-modules/audit/utils/events/object-event/object-record-updated';
import { OBJECT_RECORD_UPSERTED_EVENT } from 'src/engine/core-modules/audit/utils/events/object-event/object-record-upserted';
import { type ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator';
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
@@ -50,6 +51,12 @@ export class CreateAuditLogFromInternalEvent {
recordId: eventData.recordId,
objectMetadataId: eventData.objectMetadata.id,
});
} else if (workspaceEventBatch.name.endsWith('.upserted')) {
auditService.createObjectEvent(OBJECT_RECORD_UPSERTED_EVENT, {
...eventProperties,
recordId: eventData.recordId,
objectMetadataId: eventData.objectMetadata.id,
});
}
}
}
@@ -10,6 +10,10 @@ import {
type OBJECT_RECORD_UPDATED_EVENT,
type ObjectRecordUpdatedTrackEvent,
} from 'src/engine/core-modules/audit/utils/events/object-event/object-record-updated';
import {
type OBJECT_RECORD_UPSERTED_EVENT,
type ObjectRecordUpsertedTrackEvent,
} from 'src/engine/core-modules/audit/utils/events/object-event/object-record-upserted';
import {
type CUSTOM_DOMAIN_ACTIVATED_EVENT,
type CustomDomainActivatedTrackEvent,
@@ -50,6 +54,7 @@ export type TrackEventName =
| typeof OBJECT_RECORD_CREATED_EVENT
| typeof OBJECT_RECORD_UPDATED_EVENT
| typeof OBJECT_RECORD_DELETED_EVENT
| typeof OBJECT_RECORD_UPSERTED_EVENT
| typeof USER_SIGNUP_EVENT;
// Map event names to their corresponding event types
@@ -64,6 +69,7 @@ export interface TrackEvents {
[OBJECT_RECORD_DELETED_EVENT]: ObjectRecordDeletedTrackEvent;
[OBJECT_RECORD_CREATED_EVENT]: ObjectRecordCreatedTrackEvent;
[OBJECT_RECORD_UPDATED_EVENT]: ObjectRecordUpdatedTrackEvent;
[OBJECT_RECORD_UPSERTED_EVENT]: ObjectRecordUpsertedTrackEvent;
}
export type TrackEventProperties<T extends TrackEventName> =
@@ -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);
@@ -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<T = object> =
| ObjectRecordUpdateEvent<T>
| ObjectRecordDeleteEvent<T>
| ObjectRecordCreateEvent<T>
| ObjectRecordDestroyEvent<T>
| ObjectRecordRestoreEvent<T>;
| ObjectRecordRestoreEvent<T>
| ObjectRecordUpsertEvent<T>;
@@ -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;
@@ -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<T> {
properties: {
before?: T;
after: T;
diff?: Partial<ObjectRecordDiff<T>>;
updatedFields?: string[];
};
}
@@ -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)
? []
@@ -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<T[]>(
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<T[]>(
results.flatMap((result) => result.raw),
objectMetadata,
@@ -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<T> = {
[DatabaseEventAction.DELETED]: ObjectRecordDeleteEvent<T>;
[DatabaseEventAction.DESTROYED]: ObjectRecordDestroyEvent<T>;
[DatabaseEventAction.RESTORED]: ObjectRecordRestoreEvent<T>;
[DatabaseEventAction.UPSERTED]: ObjectRecordUpsertEvent<T>;
};
@Injectable()
@@ -62,6 +64,7 @@ export class WorkspaceEventEmitter {
| ObjectRecordCreateEvent<T>
| ObjectRecordUpdateEvent<T>
| ObjectRecordDeleteEvent<T>
| ObjectRecordUpsertEvent<T>
)[] = [];
switch (action) {
@@ -140,6 +143,41 @@ export class WorkspaceEventEmitter {
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;
}
@@ -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(
@@ -45,6 +45,7 @@ export const generateFakeObjectRecordEvent = (
switch (action) {
case DatabaseEventAction.CREATED:
case DatabaseEventAction.UPDATED:
case DatabaseEventAction.UPSERTED:
return generateFakeObjectRecordEventWithPrefix({
objectMetadataInfo,
prefix: 'properties.after',
@@ -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<ObjectMetadataEntity>,
private readonly scopedWorkspaceContextFactory: ScopedWorkspaceContextFactory,
private readonly recordPositionService: RecordPositionService,
private readonly recordInputTransformerService: RecordInputTransformerService,
@@ -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;
};
@@ -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<ObjectRecordUpsertEvent>,
) {
if (await this.shouldIgnoreEvent(payload)) {
return;
}
const clonedPayload = structuredClone(payload);
await this.handleEvent({
payload: clonedPayload,
action: DatabaseEventAction.UPSERTED,
});
}
private async enrichCreatedEvent(
payload: WorkspaceEventBatch<ObjectRecordCreateEvent>,
) {
@@ -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;
}
}
@@ -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}}.',
);