Workflow run over a list of records + backfill availability on manual triggers (#14761)

- Allow to run workflow on a list of records
- Use new availability in manual trigger when iterator feature flag is
enabled
- Build a command to backfill availability in workflow trigger



https://github.com/user-attachments/assets/d685c01f-4059-4647-92a1-f5f529b560cf
This commit is contained in:
Thomas Trompette
2025-09-29 15:45:51 +02:00
committed by GitHub
parent 59be3aab19
commit 4329501743
17 changed files with 285 additions and 85 deletions
@@ -1,4 +1,5 @@
import { Action } from '@/action-menu/actions/components/Action';
import { isBulkRecordsManualTrigger } from '@/action-menu/actions/record-actions/utils/isBulkRecordsManualTrigger';
import { ActionScope } from '@/action-menu/actions/types/ActionScope';
import { ActionType } from '@/action-menu/actions/types/ActionType';
import { contextStoreTargetedRecordsRuleComponentState } from '@/context-store/states/contextStoreTargetedRecordsRuleComponentState';
@@ -10,9 +11,11 @@ import { useRunWorkflowVersion } from '@/workflow/hooks/useRunWorkflowVersion';
import { type WorkflowVersion } from '@/workflow/types/Workflow';
import { COMMAND_MENU_DEFAULT_ICON } from '@/workflow/workflow-trigger/constants/CommandMenuDefaultIcon';
import { useIsFeatureEnabled } from '@/workspace/hooks/useIsFeatureEnabled';
import { useRecoilCallback } from 'recoil';
import { capitalize, isDefined } from 'twenty-shared/utils';
import { useIcons } from 'twenty-ui/display';
import { FeatureFlagKey } from '~/generated/graphql';
export const useRunWorkflowRecordActions = ({
objectMetadataItem,
@@ -22,6 +25,9 @@ export const useRunWorkflowRecordActions = ({
skip?: boolean;
}) => {
const { getIcon } = useIcons();
const isIteratorEnabled = useIsFeatureEnabled(
FeatureFlagKey.IS_WORKFLOW_ITERATOR_ENABLED,
);
const contextStoreTargetedRecordsRule = useRecoilComponentValue(
contextStoreTargetedRecordsRuleComponentState,
);
@@ -45,23 +51,44 @@ export const useRunWorkflowRecordActions = ({
selectedRecordIds: string[],
activeWorkflowVersion: WorkflowVersion,
) => {
for (const selectedRecordId of selectedRecordIds) {
const selectedRecord = snapshot
.getLoadable(recordStoreFamilyState(selectedRecordId))
.getValue();
if (!isDefined(selectedRecord)) {
continue;
}
if (
isIteratorEnabled &&
isDefined(activeWorkflowVersion?.trigger) &&
isBulkRecordsManualTrigger(activeWorkflowVersion.trigger)
) {
const objectNamePlural = objectMetadataItem.namePlural;
const selectedRecords = selectedRecordIds
.map((recordId) =>
snapshot.getLoadable(recordStoreFamilyState(recordId)).getValue(),
)
.filter(isDefined);
await runWorkflowVersion({
workflowId: activeWorkflowVersion.workflowId,
workflowVersionId: activeWorkflowVersion.id,
payload: selectedRecord,
payload: {
[objectNamePlural]: selectedRecords,
},
});
} else {
for (const selectedRecordId of selectedRecordIds) {
const selectedRecord = snapshot
.getLoadable(recordStoreFamilyState(selectedRecordId))
.getValue();
if (!isDefined(selectedRecord)) {
continue;
}
await runWorkflowVersion({
workflowId: activeWorkflowVersion.workflowId,
workflowVersionId: activeWorkflowVersion.id,
payload: selectedRecord,
});
}
}
},
[runWorkflowVersion],
[runWorkflowVersion, isIteratorEnabled, objectMetadataItem],
);
return activeWorkflowVersions
@@ -0,0 +1,8 @@
import { type WorkflowTrigger } from '@/workflow/types/Workflow';
export const isBulkRecordsManualTrigger = (trigger: WorkflowTrigger) => {
return (
trigger.type === 'MANUAL' &&
trigger?.settings?.availability?.type === 'BULK_RECORDS'
);
};
@@ -0,0 +1,17 @@
import { type WorkflowTrigger } from '@/workflow/types/Workflow';
import { isDefined } from 'twenty-shared/utils';
export const isGlobalManualTrigger = (
trigger: WorkflowTrigger,
isIteratorEnabled: boolean,
) => {
if (trigger.type !== 'MANUAL') {
return false;
}
if (isIteratorEnabled && isDefined(trigger.settings?.availability)) {
return trigger.settings.availability.type === 'GLOBAL';
}
return !isDefined(trigger.settings.objectType);
};
@@ -1,3 +1,4 @@
import { isGlobalManualTrigger } from '@/action-menu/actions/record-actions/utils/isGlobalManualTrigger';
import { useObjectMetadataItem } from '@/object-metadata/hooks/useObjectMetadataItem';
import { CoreObjectNameSingular } from '@/object-metadata/types/CoreObjectNameSingular';
import { type ObjectMetadataItem } from '@/object-metadata/types/ObjectMetadataItem';
@@ -7,7 +8,9 @@ import {
type ManualTriggerWorkflowVersion,
type Workflow,
} from '@/workflow/types/Workflow';
import { useIsFeatureEnabled } from '@/workspace/hooks/useIsFeatureEnabled';
import { isDefined } from 'twenty-shared/utils';
import { FeatureFlagKey } from '~/generated/graphql';
export const useActiveWorkflowVersionsWithManualTrigger = ({
objectMetadataItem,
@@ -16,6 +19,10 @@ export const useActiveWorkflowVersionsWithManualTrigger = ({
objectMetadataItem?: ObjectMetadataItem;
skip?: boolean;
}) => {
const isIteratorEnabled = useIsFeatureEnabled(
FeatureFlagKey.IS_WORKFLOW_ITERATOR_ENABLED,
);
const filters = [
{
status: {
@@ -29,12 +36,20 @@ export const useActiveWorkflowVersionsWithManualTrigger = ({
},
];
const objectTypeFilter = isIteratorEnabled
? {
trigger: {
like: `%"objectNameSingular": "${objectMetadataItem?.nameSingular}"%`,
},
}
: {
trigger: {
like: `%"objectType": "${objectMetadataItem?.nameSingular}"%`,
},
};
if (isDefined(objectMetadataItem)) {
filters.push({
trigger: {
like: `%"objectType": "${objectMetadataItem.nameSingular}"%`,
},
});
filters.push(objectTypeFilter);
}
const { objectMetadataItem: workflowVersionObjectMetadataItem } =
@@ -64,8 +79,8 @@ export const useActiveWorkflowVersionsWithManualTrigger = ({
records: records.filter(
(record) =>
record.status === 'ACTIVE' &&
record.trigger?.type === 'MANUAL' &&
!isDefined(record.trigger?.settings.objectType),
isDefined(record.trigger) &&
isGlobalManualTrigger(record.trigger, isIteratorEnabled),
),
};
}
@@ -158,6 +158,7 @@ export const WorkflowEditTriggerManual = ({
type: availability.type,
objectNameSingular,
},
objectType: objectNameSingular,
outputSchema: {},
},
});
@@ -197,6 +197,10 @@ export const WorkflowEditTriggerManualDeprecated = ({
...trigger,
settings: {
...trigger.settings,
availability: {
objectNameSingular: updatedObject,
type: 'SINGLE_RECORD',
},
objectType: updatedObject,
outputSchema: {},
},
@@ -3,57 +3,71 @@ import { COMMAND_MENU_DEFAULT_ICON } from '@/workflow/workflow-trigger/constants
import { generatedMockObjectMetadataItems } from '~/testing/utils/generatedMockObjectMetadataItems';
import { getManualTriggerDefaultSettingsDeprecated } from '../getManualTriggerDefaultSettingsDeprecated';
it('returns settings for a manual trigger that can be activated from any where', () => {
expect(
getManualTriggerDefaultSettingsDeprecated({
availability: 'EVERYWHERE',
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
}),
).toStrictEqual({
objectType: undefined,
outputSchema: {},
icon: COMMAND_MENU_DEFAULT_ICON,
isPinned: false,
describe('getManualTriggerDefaultSettingsDeprecated', () => {
it('returns settings for a manual trigger that can be activated from any where', () => {
expect(
getManualTriggerDefaultSettingsDeprecated({
availability: 'EVERYWHERE',
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
}),
).toStrictEqual({
objectType: undefined,
outputSchema: {},
icon: COMMAND_MENU_DEFAULT_ICON,
isPinned: false,
availability: {
type: 'GLOBAL',
locations: [],
},
});
});
});
it('returns settings for a manual trigger that can be activated from any where', () => {
expect(
getManualTriggerDefaultSettingsDeprecated({
availability: 'WHEN_RECORD_SELECTED',
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
it('returns settings for a manual trigger that can be activated from any where', () => {
expect(
getManualTriggerDefaultSettingsDeprecated({
availability: 'WHEN_RECORD_SELECTED',
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
icon: 'IconTest',
}),
).toStrictEqual({
objectType: generatedMockObjectMetadataItems[0].nameSingular,
outputSchema: {},
icon: 'IconTest',
}),
).toStrictEqual({
objectType: generatedMockObjectMetadataItems[0].nameSingular,
outputSchema: {},
icon: 'IconTest',
isPinned: false,
isPinned: false,
availability: {
type: 'SINGLE_RECORD',
objectNameSingular: generatedMockObjectMetadataItems[0].nameSingular,
},
});
});
it('returns settings for WHEN_RECORD_SELECTED with default icon when no custom icon provided', () => {
expect(
getManualTriggerDefaultSettingsDeprecated({
availability: 'WHEN_RECORD_SELECTED',
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
}),
).toStrictEqual({
objectType: generatedMockObjectMetadataItems[0].nameSingular,
outputSchema: {},
icon: COMMAND_MENU_DEFAULT_ICON,
isPinned: false,
availability: {
type: 'SINGLE_RECORD',
objectNameSingular: generatedMockObjectMetadataItems[0].nameSingular,
},
});
});
it('throws error for unsupported availability type', () => {
const invalidAvailability =
'INVALID_AVAILABILITY' as WorkflowManualTriggerAvailability;
expect(() =>
getManualTriggerDefaultSettingsDeprecated({
availability: invalidAvailability,
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
}),
).toThrow("Didn't expect to get here.");
});
});
it('returns settings for WHEN_RECORD_SELECTED with default icon when no custom icon provided', () => {
expect(
getManualTriggerDefaultSettingsDeprecated({
availability: 'WHEN_RECORD_SELECTED',
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
}),
).toStrictEqual({
objectType: generatedMockObjectMetadataItems[0].nameSingular,
outputSchema: {},
icon: COMMAND_MENU_DEFAULT_ICON,
isPinned: false,
});
});
it('throws error for unsupported availability type', () => {
const invalidAvailability =
'INVALID_AVAILABILITY' as WorkflowManualTriggerAvailability;
expect(() =>
getManualTriggerDefaultSettingsDeprecated({
availability: invalidAvailability,
activeNonSystemObjectMetadataItems: generatedMockObjectMetadataItems,
}),
).toThrow("Didn't expect to get here.");
});
@@ -116,6 +116,10 @@ describe('getTriggerDefaultDefinition', () => {
name: 'Launch manually',
settings: {
objectType: generatedMockObjectMetadataItems[0].nameSingular,
availability: {
objectNameSingular: generatedMockObjectMetadataItems[0].nameSingular,
type: 'SINGLE_RECORD',
},
outputSchema: {},
icon: COMMAND_MENU_DEFAULT_ICON,
isPinned: false,
@@ -17,6 +17,7 @@ export const getManualTriggerDefaultSettings = ({
switch (availabilityType) {
case 'GLOBAL': {
return {
objectType: undefined,
availability: {
type: 'GLOBAL',
locations: undefined,
@@ -28,6 +29,7 @@ export const getManualTriggerDefaultSettings = ({
}
case 'SINGLE_RECORD': {
return {
objectType: activeNonSystemObjectMetadataItems[0].nameSingular,
availability: {
type: 'SINGLE_RECORD',
objectNameSingular:
@@ -40,6 +42,7 @@ export const getManualTriggerDefaultSettings = ({
}
case 'BULK_RECORDS': {
return {
objectType: activeNonSystemObjectMetadataItems[0].nameSingular,
availability: {
type: 'BULK_RECORDS',
objectNameSingular:
@@ -24,6 +24,10 @@ export const getManualTriggerDefaultSettingsDeprecated = ({
outputSchema: {},
icon: icon || COMMAND_MENU_DEFAULT_ICON,
isPinned: isPinned || false,
availability: {
type: 'GLOBAL',
locations: [],
},
};
}
case 'WHEN_RECORD_SELECTED': {
@@ -32,6 +36,11 @@ export const getManualTriggerDefaultSettingsDeprecated = ({
outputSchema: {},
icon: icon || COMMAND_MENU_DEFAULT_ICON,
isPinned: isPinned || false,
availability: {
type: 'SINGLE_RECORD',
objectNameSingular:
activeNonSystemObjectMetadataItems[0].nameSingular,
},
};
}
}
@@ -0,0 +1,85 @@
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
import { Command } from 'nest-commander';
import { isDefined } from 'twenty-shared/utils';
import { DataSource, Repository } from 'typeorm';
import {
ActiveOrSuspendedWorkspacesMigrationCommandRunner,
type RunOnWorkspaceArgs,
} from 'src/database/commands/command-runners/active-or-suspended-workspaces-migration.command-runner';
import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util';
import { WorkflowTriggerType } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
@Command({
name: 'upgrade:1-7:backfill-workflow-manual-trigger-availability',
description:
'Backfill workflow manual trigger availability based on objectType',
})
export class BackfillWorkflowManualTriggerAvailabilityCommand extends ActiveOrSuspendedWorkspacesMigrationCommandRunner {
constructor(
@InjectRepository(Workspace)
protected readonly workspaceRepository: Repository<Workspace>,
protected readonly twentyORMGlobalManager: TwentyORMGlobalManager,
@InjectDataSource()
private readonly coreDataSource: DataSource,
) {
super(workspaceRepository, twentyORMGlobalManager);
}
override async runOnWorkspace({
workspaceId,
}: RunOnWorkspaceArgs): Promise<void> {
const schemaName = getWorkspaceSchemaName(workspaceId);
const workflowVersions = await this.coreDataSource.query(
`SELECT * FROM ${schemaName}."workflowVersion"`,
);
for (const workflowVersion of workflowVersions) {
const { trigger } = workflowVersion;
if (trigger.type !== WorkflowTriggerType.MANUAL) {
continue;
}
const availability = trigger.settings.availability;
const objectType = trigger.settings.objectType;
if (isDefined(availability)) {
continue;
}
const newAvailability = objectType
? {
type: 'SINGLE_RECORD',
objectNameSingular: objectType,
}
: {
type: 'GLOBAL',
locations: [],
};
const updatedTrigger = {
...trigger,
settings: {
...trigger.settings,
availability: newAvailability,
},
};
this.logger.log(
`Updating workflow version ${workflowVersion.id} with new availability ${JSON.stringify(
newAvailability,
)}`,
);
await this.coreDataSource.query(
`UPDATE ${schemaName}."workflowVersion" SET trigger = $1 WHERE id = $2`,
[updatedTrigger, workflowVersion.id],
);
}
}
}
@@ -1,13 +1,20 @@
import { Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { BackfillWorkflowManualTriggerAvailabilityCommand } from 'src/database/commands/upgrade-version-command/1-7/1-7-backfill-workflow-manual-trigger-availability.command';
import { RegeneratePersonSearchVectorWithPhonesCommand } from 'src/database/commands/upgrade-version-command/1-7/1-7-regenerate-person-search-vector-with-phones.command';
import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity';
import { WorkspaceDataSourceModule } from 'src/engine/workspace-datasource/workspace-datasource.module';
@Module({
imports: [TypeOrmModule.forFeature([Workspace]), WorkspaceDataSourceModule],
providers: [RegeneratePersonSearchVectorWithPhonesCommand],
exports: [RegeneratePersonSearchVectorWithPhonesCommand],
providers: [
RegeneratePersonSearchVectorWithPhonesCommand,
BackfillWorkflowManualTriggerAvailabilityCommand,
],
exports: [
RegeneratePersonSearchVectorWithPhonesCommand,
BackfillWorkflowManualTriggerAvailabilityCommand,
],
})
export class V1_7_UpgradeVersionCommandModule {}
@@ -44,7 +44,7 @@ export class WorkflowTriggerResolver {
@Args('workflowVersionId', { type: () => UUIDScalarType })
workflowVersionId: string,
) {
return await this.workflowTriggerWorkspaceService.activateWorkflowVersion(
return this.workflowTriggerWorkspaceService.activateWorkflowVersion(
workflowVersionId,
);
}
@@ -54,7 +54,7 @@ export class WorkflowTriggerResolver {
@Args('workflowVersionId', { type: () => UUIDScalarType })
workflowVersionId: string,
) {
return await this.workflowTriggerWorkspaceService.deactivateWorkflowVersion(
return this.workflowTriggerWorkspaceService.deactivateWorkflowVersion(
workflowVersionId,
);
}
@@ -78,7 +78,7 @@ export class WorkflowTriggerResolver {
},
});
return await this.workflowTriggerWorkspaceService.runWorkflowVersion({
return this.workflowTriggerWorkspaceService.runWorkflowVersion({
workflowVersionId,
workflowRunId: workflowRunId ?? undefined,
payload: payload ?? {},
@@ -1,10 +1,11 @@
import { Module } from '@nestjs/common';
import { WorkflowSchemaWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service';
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module';
import { WorkflowSchemaWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service';
@Module({
imports: [WorkflowCommonModule],
imports: [WorkflowCommonModule, FeatureFlagModule],
providers: [WorkflowSchemaWorkspaceService],
exports: [WorkflowSchemaWorkspaceService],
})
@@ -12,6 +12,8 @@ import {
import { type DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import { checkStringIsDatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/utils/check-string-is-database-event-action';
import { FeatureFlagKey } from 'src/engine/core-modules/feature-flag/enums/feature-flag-key.enum';
import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service';
import { generateFakeValue } from 'src/engine/utils/generate-fake-value';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
import { DEFAULT_ITERATOR_CURRENT_ITEM } from 'src/modules/workflow/workflow-builder/workflow-schema/constants/default-iterator-current-item.const';
@@ -38,6 +40,7 @@ import {
export class WorkflowSchemaWorkspaceService {
constructor(
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
private readonly featureFlagService: FeatureFlagService,
) {}
async computeStepOutputSchema({
@@ -59,8 +62,21 @@ export class WorkflowSchemaWorkspaceService {
});
}
case WorkflowTriggerType.MANUAL: {
const isIteratorEnabled =
await this.featureFlagService.isFeatureEnabled(
FeatureFlagKey.IS_WORKFLOW_ITERATOR_ENABLED,
workspaceId,
);
const { objectType, availability } = step.settings;
if (isDefined(availability) && isIteratorEnabled) {
return this.computeTriggerOutputSchemaFromAvailability({
availability,
workspaceId,
});
}
// TODO: to be deprecated once all triggers are migrated to the new availability type
if (isDefined(objectType)) {
return this.computeRecordOutputSchema({
@@ -69,13 +85,6 @@ export class WorkflowSchemaWorkspaceService {
});
}
if (isDefined(availability)) {
return this.computeTriggerOutputSchemaFromAvailability({
availability,
workspaceId,
});
}
return {};
}
case WorkflowTriggerType.WEBHOOK:
@@ -291,7 +300,7 @@ export class WorkflowSchemaWorkspaceService {
);
return {
[availability.objectNameSingular]: {
[objectMetadataInfo.objectMetadataItemWithFieldsMaps.namePlural]: {
label:
objectMetadataInfo.objectMetadataItemWithFieldsMaps.labelPlural,
isLeaf: true,
@@ -2,7 +2,6 @@ import { Scope } from '@nestjs/common';
import { isDefined } from 'twenty-shared/utils';
import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service';
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';
@@ -37,7 +36,6 @@ export class RunWorkflowJob {
private readonly twentyConfigService: TwentyConfigService,
private readonly metricsService: MetricsService,
private readonly workflowRunQueueWorkspaceService: WorkflowRunQueueWorkspaceService,
private readonly featureFlagService: FeatureFlagService,
) {}
@Process(RunWorkflowJob.name)
@@ -1,7 +1,6 @@
import { Module } from '@nestjs/common';
import { BillingModule } from 'src/engine/core-modules/billing/billing.module';
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
import { MetricsModule } from 'src/engine/core-modules/metrics/metrics.module';
import { ThrottlerModule } from 'src/engine/core-modules/throttler/throttler.module';
import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module';
@@ -21,7 +20,6 @@ import { WorkflowRunnerWorkspaceService } from 'src/modules/workflow/workflow-ru
WorkflowRunModule,
MetricsModule,
WorkflowRunQueueModule,
FeatureFlagModule,
WorkflowVersionStepModule,
],
providers: [WorkflowRunnerWorkspaceService, RunWorkflowJob],