Migrate cron, databaseEventTrigger, httpRoute triggers to serverless functions (#17488)

## Summary

Migrates trigger entities (`CronTriggerEntity`,
`DatabaseEventTriggerEntity`, `RouteTriggerEntity`) into
`ServerlessFunctionEntity` by storing trigger settings as JSONB columns
directly on the serverless function. This simplifies the architecture
since these relationships were effectively one-to-one.

## Changes

### Schema Changes
- Added three new nullable JSONB columns to `ServerlessFunctionEntity`:
  - `cronTriggerSettings` - stores cron pattern
- `databaseEventTriggerSettings` - stores event name and updated fields
filter
- `httpRouteTriggerSettings` - stores path, HTTP method, auth
requirements, and forwarded headers

### Core Logic Updates
- `CronTriggerCronJob` - now queries `ServerlessFunctionEntity` directly
instead of `CronTriggerEntity`
- `CallDatabaseEventTriggerJobsJob` - now queries
`ServerlessFunctionEntity` directly
- `RouteTriggerService` - now queries `ServerlessFunctionEntity`
directly
- `ApplicationSyncService` - extracts trigger settings from manifest and
writes to serverless function
This commit is contained in:
Charles Bochet
2026-01-27 20:54:09 +01:00
committed by GitHub
parent e1f92bd951
commit a4499d21bb
110 changed files with 439 additions and 5193 deletions
@@ -7,9 +7,7 @@ import { ApplicationResolver } from 'src/engine/core-modules/application/applica
import { ApplicationVariableEntityModule } from 'src/engine/core-modules/applicationVariable/application-variable.module';
import { FileStorageModule } from 'src/engine/core-modules/file-storage/file-storage.module';
import { FileEntity } from 'src/engine/core-modules/file/entities/file.entity';
import { CronTriggerModule } from 'src/engine/metadata-modules/cron-trigger/cron-trigger.module';
import { DataSourceModule } from 'src/engine/metadata-modules/data-source/data-source.module';
import { DatabaseEventTriggerModule } from 'src/engine/metadata-modules/database-event-trigger/database-event-trigger.module';
import { FieldMetadataModule } from 'src/engine/metadata-modules/field-metadata/field-metadata.module';
import { WorkspaceManyOrAllFlatEntityMapsCacheModule } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.module';
import { ObjectMetadataModule } from 'src/engine/metadata-modules/object-metadata/object-metadata.module';
@@ -17,7 +15,6 @@ import { ObjectPermissionModule } from 'src/engine/metadata-modules/object-permi
import { PermissionFlagModule } from 'src/engine/metadata-modules/permission-flag/permission-flag.module';
import { PermissionsModule } from 'src/engine/metadata-modules/permissions/permissions.module';
import { RoleModule } from 'src/engine/metadata-modules/role/role.module';
import { RouteTriggerModule } from 'src/engine/metadata-modules/route-trigger/route-trigger.module';
import { ServerlessFunctionLayerModule } from 'src/engine/metadata-modules/serverless-function-layer/serverless-function-layer.module';
import { ServerlessFunctionModule } from 'src/engine/metadata-modules/serverless-function/serverless-function.module';
import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache.module';
@@ -37,9 +34,6 @@ import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-commo
DataSourceModule,
ServerlessFunctionLayerModule,
ServerlessFunctionModule,
DatabaseEventTriggerModule,
CronTriggerModule,
RouteTriggerModule,
WorkspaceMigrationModule,
PermissionsModule,
RoleModule,
@@ -23,11 +23,7 @@ import {
import { ApplicationService } from 'src/engine/core-modules/application/application.service';
import { ApplicationInput } from 'src/engine/core-modules/application/dtos/application.input';
import { ApplicationVariableEntityService } from 'src/engine/core-modules/applicationVariable/application-variable.service';
import { CronTriggerV2Service } from 'src/engine/metadata-modules/cron-trigger/services/cron-trigger-v2.service';
import { FlatCronTrigger } from 'src/engine/metadata-modules/cron-trigger/types/flat-cron-trigger.type';
import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service';
import { DatabaseEventTriggerV2Service } from 'src/engine/metadata-modules/database-event-trigger/services/database-event-trigger-v2.service';
import { FlatDatabaseEventTrigger } from 'src/engine/metadata-modules/database-event-trigger/types/flat-database-event-trigger.type';
import { CreateFieldInput } from 'src/engine/metadata-modules/field-metadata/dtos/create-field.input';
import { FieldMetadataService } from 'src/engine/metadata-modules/field-metadata/services/field-metadata.service';
import { WorkspaceManyOrAllFlatEntityMapsCacheService } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.service';
@@ -39,13 +35,16 @@ import { FieldPermissionService } from 'src/engine/metadata-modules/object-permi
import { ObjectPermissionService } from 'src/engine/metadata-modules/object-permission/object-permission.service';
import { PermissionFlagService } from 'src/engine/metadata-modules/permission-flag/permission-flag.service';
import { RoleService } from 'src/engine/metadata-modules/role/role.service';
import { RouteTriggerV2Service } from 'src/engine/metadata-modules/route-trigger/services/route-trigger-v2.service';
import { FlatRouteTrigger } from 'src/engine/metadata-modules/route-trigger/types/flat-route-trigger.type';
import { ServerlessFunctionLayerService } from 'src/engine/metadata-modules/serverless-function-layer/serverless-function-layer.service';
import { ServerlessFunctionV2Service } from 'src/engine/metadata-modules/serverless-function/services/serverless-function-v2.service';
import { FlatServerlessFunction } from 'src/engine/metadata-modules/serverless-function/types/flat-serverless-function.type';
import { computeMetadataNameFromLabelOrThrow } from 'src/engine/metadata-modules/utils/compute-metadata-name-from-label-or-throw.util';
import { WorkspaceMigrationValidateBuildAndRunService } from 'src/engine/workspace-manager/workspace-migration/services/workspace-migration-validate-build-and-run-service';
import {
CronTriggerSettings,
DatabaseEventTriggerSettings,
HttpRouteTriggerSettings,
} from 'src/engine/metadata-modules/serverless-function/serverless-function.entity';
@Injectable()
export class ApplicationSyncService {
@@ -60,9 +59,6 @@ export class ApplicationSyncService {
private readonly serverlessFunctionV2Service: ServerlessFunctionV2Service,
private readonly flatEntityMapsCacheService: WorkspaceManyOrAllFlatEntityMapsCacheService,
private readonly dataSourceService: DataSourceService,
private readonly databaseEventTriggerV2Service: DatabaseEventTriggerV2Service,
private readonly cronTriggerV2Service: CronTriggerV2Service,
private readonly routeTriggerV2Service: RouteTriggerV2Service,
private readonly workspaceMigrationValidateBuildAndRunService: WorkspaceMigrationValidateBuildAndRunService,
private readonly roleService: RoleService,
private readonly objectPermissionService: ObjectPermissionService,
@@ -960,26 +956,8 @@ export class ApplicationSyncService {
workspaceId,
);
await this.syncDatabaseEventTriggersForServerlessFunction({
serverlessFunctionId: serverlessFunctionToUpdate.id,
triggersToSync: serverlessFunctionToSync.triggers || [],
workspaceId,
applicationId,
});
await this.syncCronTriggersForServerlessFunction({
serverlessFunctionId: serverlessFunctionToUpdate.id,
triggersToSync: serverlessFunctionToSync.triggers || [],
workspaceId,
applicationId,
});
await this.syncRouteTriggersForServerlessFunction({
serverlessFunctionId: serverlessFunctionToUpdate.id,
triggersToSync: serverlessFunctionToSync.triggers || [],
workspaceId,
applicationId,
});
// Trigger settings are now embedded in the serverless function entity
// They are handled through the update input
}
for (const serverlessFunctionToCreate of serverlessFunctionsToCreate) {
@@ -1001,390 +979,52 @@ export class ApplicationSyncService {
isTool: serverlessFunctionToCreate.isTool,
};
const createdServerlessFunction =
await this.serverlessFunctionV2Service.createOne({
createServerlessFunctionInput,
workspaceId,
applicationId,
});
await this.syncDatabaseEventTriggersForServerlessFunction({
serverlessFunctionId: createdServerlessFunction.id,
triggersToSync: serverlessFunctionToCreate.triggers || [],
await this.serverlessFunctionV2Service.createOne({
createServerlessFunctionInput,
workspaceId,
applicationId,
});
await this.syncCronTriggersForServerlessFunction({
serverlessFunctionId: createdServerlessFunction.id,
triggersToSync: serverlessFunctionToCreate.triggers || [],
workspaceId,
applicationId,
});
await this.syncRouteTriggersForServerlessFunction({
serverlessFunctionId: createdServerlessFunction.id,
triggersToSync: serverlessFunctionToCreate.triggers || [],
workspaceId,
applicationId,
});
// Trigger settings are now embedded in the serverless function entity
// They are handled through the create input
}
}
private async syncDatabaseEventTriggersForServerlessFunction({
serverlessFunctionId,
triggersToSync,
workspaceId,
applicationId,
}: {
serverlessFunctionId: string;
triggersToSync: ServerlessFunctionTriggerManifest[];
workspaceId: string;
applicationId: string;
}) {
const databaseEventTriggersToSync = triggersToSync.filter(
(trigger) => trigger.type === 'databaseEvent',
);
private extractTriggerSettingsFromManifest(
triggers: ServerlessFunctionTriggerManifest[] = [],
): {
cronTriggerSettings: CronTriggerSettings | null;
databaseEventTriggerSettings: DatabaseEventTriggerSettings | null;
httpRouteTriggerSettings: HttpRouteTriggerSettings | null;
} {
let cronTriggerSettings: CronTriggerSettings | null = null;
let databaseEventTriggerSettings: DatabaseEventTriggerSettings | null =
null;
let httpRouteTriggerSettings: HttpRouteTriggerSettings | null = null;
const { flatDatabaseEventTriggerMaps } =
await this.flatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps(
{
workspaceId,
flatMapsKeys: ['flatDatabaseEventTriggerMaps'],
},
);
const existingDatabaseEventTriggers = Object.values(
flatDatabaseEventTriggerMaps.byId,
).filter(
(trigger) =>
isDefined(trigger) &&
trigger.serverlessFunctionId === serverlessFunctionId,
) as FlatDatabaseEventTrigger[];
const triggersToSyncUniversalIdentifiers = databaseEventTriggersToSync.map(
(trigger) => trigger.universalIdentifier,
);
const existingTriggersUniversalIdentifiers =
existingDatabaseEventTriggers.map(
(trigger) => trigger.universalIdentifier,
);
const triggersToDelete = existingDatabaseEventTriggers.filter(
(trigger) =>
isDefined(trigger.universalIdentifier) &&
!triggersToSyncUniversalIdentifiers.includes(
trigger.universalIdentifier,
),
);
const triggersToUpdate = existingDatabaseEventTriggers.filter(
(trigger) =>
isDefined(trigger.universalIdentifier) &&
triggersToSyncUniversalIdentifiers.includes(
trigger.universalIdentifier,
),
);
const triggersToCreate = databaseEventTriggersToSync.filter(
(triggerToSync) =>
!existingTriggersUniversalIdentifiers.includes(
triggerToSync.universalIdentifier,
),
);
for (const triggerToDelete of triggersToDelete) {
await this.databaseEventTriggerV2Service.destroyOne({
destroyDatabaseEventTriggerInput: { id: triggerToDelete.id },
workspaceId,
});
}
for (const triggerToUpdate of triggersToUpdate) {
const triggerToSync = databaseEventTriggersToSync.find(
(trigger) =>
trigger.universalIdentifier === triggerToUpdate.universalIdentifier,
);
if (!triggerToSync || triggerToSync.type !== 'databaseEvent') {
throw new ApplicationException(
`Failed to find database event trigger to sync with universalIdentifier ${triggerToUpdate.universalIdentifier}`,
ApplicationExceptionCode.ENTITY_NOT_FOUND,
);
for (const trigger of triggers) {
if (trigger.type === 'cron') {
cronTriggerSettings = { pattern: trigger.pattern };
} else if (trigger.type === 'databaseEvent') {
databaseEventTriggerSettings = {
eventName: trigger.eventName,
updatedFields: trigger.updatedFields,
};
} else if (trigger.type === 'route') {
httpRouteTriggerSettings = {
path: trigger.path,
httpMethod: trigger.httpMethod as HTTPMethod,
isAuthRequired: trigger.isAuthRequired,
forwardedRequestHeaders: trigger.forwardedRequestHeaders,
};
}
const updateDatabaseEventTriggerInput = {
id: triggerToUpdate.id,
update: {
settings: {
eventName: triggerToSync.eventName,
updatedFields: triggerToSync.updatedFields,
},
},
};
await this.databaseEventTriggerV2Service.updateOne(
updateDatabaseEventTriggerInput,
workspaceId,
);
}
for (const triggerToCreate of triggersToCreate) {
if (triggerToCreate.type !== 'databaseEvent') {
continue;
}
const createDatabaseEventTriggerInput = {
settings: {
eventName: triggerToCreate.eventName,
updatedFields: triggerToCreate.updatedFields,
},
universalIdentifier: triggerToCreate.universalIdentifier,
serverlessFunctionId,
};
await this.databaseEventTriggerV2Service.createOne(
createDatabaseEventTriggerInput,
workspaceId,
applicationId,
);
}
}
private async syncCronTriggersForServerlessFunction({
serverlessFunctionId,
triggersToSync,
workspaceId,
applicationId,
}: {
serverlessFunctionId: string;
triggersToSync: ServerlessFunctionTriggerManifest[];
workspaceId: string;
applicationId: string;
}) {
const cronTriggersToSync = triggersToSync.filter(
(trigger) => trigger.type === 'cron',
);
const { flatCronTriggerMaps } =
await this.flatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps(
{
workspaceId,
flatMapsKeys: ['flatCronTriggerMaps'],
},
);
const existingCronTriggers = Object.values(flatCronTriggerMaps.byId).filter(
(trigger) =>
isDefined(trigger) &&
trigger.serverlessFunctionId === serverlessFunctionId,
) as FlatCronTrigger[];
const triggersToSyncUniversalIdentifiers = cronTriggersToSync.map(
(trigger) => trigger.universalIdentifier,
);
const existingTriggersUniversalIdentifiers = existingCronTriggers.map(
(trigger) => trigger.universalIdentifier,
);
const triggersToDelete = existingCronTriggers.filter(
(trigger) =>
isDefined(trigger.universalIdentifier) &&
!triggersToSyncUniversalIdentifiers.includes(
trigger.universalIdentifier,
),
);
const triggersToUpdate = existingCronTriggers.filter(
(trigger) =>
isDefined(trigger.universalIdentifier) &&
triggersToSyncUniversalIdentifiers.includes(
trigger.universalIdentifier,
),
);
const triggersToCreate = cronTriggersToSync.filter(
(triggerToSync) =>
!existingTriggersUniversalIdentifiers.includes(
triggerToSync.universalIdentifier,
),
);
for (const triggerToDelete of triggersToDelete) {
await this.cronTriggerV2Service.destroyOne({
destroyCronTriggerInput: { id: triggerToDelete.id },
workspaceId,
});
}
for (const triggerToUpdate of triggersToUpdate) {
const triggerToSync = cronTriggersToSync.find(
(trigger) =>
trigger.universalIdentifier === triggerToUpdate.universalIdentifier,
);
if (!triggerToSync || triggerToSync.type !== 'cron') {
throw new ApplicationException(
`Failed to find cron trigger to sync with universalIdentifier ${triggerToUpdate.universalIdentifier}`,
ApplicationExceptionCode.ENTITY_NOT_FOUND,
);
}
const updateCronTriggerInput = {
id: triggerToUpdate.id,
update: {
settings: {
pattern: triggerToSync.pattern,
},
},
};
await this.cronTriggerV2Service.updateOne(
updateCronTriggerInput,
workspaceId,
);
}
for (const triggerToCreate of triggersToCreate) {
if (triggerToCreate.type !== 'cron') {
continue;
}
const createCronTriggerInput = {
settings: {
pattern: triggerToCreate.pattern,
},
universalIdentifier: triggerToCreate.universalIdentifier,
serverlessFunctionId,
};
await this.cronTriggerV2Service.createOne(
createCronTriggerInput,
workspaceId,
applicationId,
);
}
}
private async syncRouteTriggersForServerlessFunction({
serverlessFunctionId,
triggersToSync,
workspaceId,
applicationId,
}: {
serverlessFunctionId: string;
triggersToSync: ServerlessFunctionTriggerManifest[];
workspaceId: string;
applicationId: string;
}) {
const routeTriggersToSync = triggersToSync.filter(
(trigger) => trigger.type === 'route',
);
const { flatRouteTriggerMaps } =
await this.flatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps(
{
workspaceId,
flatMapsKeys: ['flatRouteTriggerMaps'],
},
);
const existingRouteTriggers = Object.values(
flatRouteTriggerMaps.byId,
).filter(
(trigger) =>
isDefined(trigger) &&
trigger.serverlessFunctionId === serverlessFunctionId,
) as FlatRouteTrigger[];
const triggersToSyncUniversalIdentifiers = routeTriggersToSync.map(
(trigger) => trigger.universalIdentifier,
);
const existingTriggersUniversalIdentifiers = existingRouteTriggers.map(
(trigger) => trigger.universalIdentifier,
);
const triggersToDelete = existingRouteTriggers.filter(
(trigger) =>
isDefined(trigger.universalIdentifier) &&
!triggersToSyncUniversalIdentifiers.includes(
trigger.universalIdentifier,
),
);
const triggersToUpdate = existingRouteTriggers.filter(
(trigger) =>
isDefined(trigger.universalIdentifier) &&
triggersToSyncUniversalIdentifiers.includes(
trigger.universalIdentifier,
),
);
const triggersToCreate = routeTriggersToSync.filter(
(triggerToSync) =>
!existingTriggersUniversalIdentifiers.includes(
triggerToSync.universalIdentifier,
),
);
for (const triggerToDelete of triggersToDelete) {
await this.routeTriggerV2Service.destroyOne({
destroyRouteTriggerInput: { id: triggerToDelete.id },
workspaceId,
});
}
for (const triggerToUpdate of triggersToUpdate) {
const triggerToSync = routeTriggersToSync.find(
(trigger) =>
trigger.universalIdentifier === triggerToUpdate.universalIdentifier,
);
if (!triggerToSync || triggerToSync.type !== 'route') {
throw new ApplicationException(
`Failed to find route trigger to sync with universalIdentifier ${triggerToUpdate.universalIdentifier}`,
ApplicationExceptionCode.ENTITY_NOT_FOUND,
);
}
const updateRouteTriggerInput = {
id: triggerToUpdate.id,
update: {
path: triggerToSync.path,
httpMethod: triggerToSync.httpMethod as HTTPMethod,
isAuthRequired: triggerToSync.isAuthRequired,
forwardedRequestHeaders: triggerToSync.forwardedRequestHeaders ?? [],
},
};
await this.routeTriggerV2Service.updateOne(
updateRouteTriggerInput,
workspaceId,
);
}
for (const triggerToCreate of triggersToCreate) {
if (triggerToCreate.type !== 'route') {
continue;
}
const createRouteTriggerInput = {
path: triggerToCreate.path,
httpMethod: triggerToCreate.httpMethod as HTTPMethod,
isAuthRequired: triggerToCreate.isAuthRequired,
forwardedRequestHeaders: triggerToCreate.forwardedRequestHeaders ?? [],
serverlessFunctionId,
};
await this.routeTriggerV2Service.createOne(
createRouteTriggerInput,
workspaceId,
applicationId,
);
}
return {
cronTriggerSettings,
databaseEventTriggerSettings,
httpRouteTriggerSettings,
};
}
public async uninstallApplication({