feat: emit metadata events for schema changes with actor context for webhooks (#17622)
## Summary
This PR adds **metadata eventing**: when schema metadata
(objectMetadata, fieldMetadata, view, viewField, etc.) is created,
updated, or deleted, we now emit events that can trigger webhooks and
future audit logs. It also adds **actor context** (`userId`,
`workspaceMemberId`) to those events so subscribers can attribute
changes to a user or API key.
## What changed
### 1. Metadata eventing (first commit)
- **MetadataEventEmitter**
New service that emits batch events after successful workspace
migrations. Event names follow `metadata.{entity}.{action}` (e.g.
`metadata.objectMetadata.created`, `metadata.fieldMetadata.updated`).
- **MetadataEventsToDbListener**
Listens for metadata events and enqueues webhook delivery via
`CallWebhookJobsForMetadataJob`.
- **Event types** (twenty-shared)
`MetadataEventAction`, `MetadataEventBatch`, and record event types for
create/update/delete.
- **WorkspaceMigrationValidateBuildAndRunService**
Calls the metadata event emitter after running migrations so all
metadata changes (from any module) emit events from a single place.
- **Create events**
Sourced from the create action payload (`flatEntity` /
`flatFieldMetadatas`) because `fromToAllFlatEntityMaps` does not provide
a before/after diff for creates. Update/delete events still use the
fromToAllFlatEntityMaps comparison.
### 2. Actor context (second commit)
- **MetadataEventEmitter**
Accepts optional `actorContext` (`userId`, `workspaceMemberId`) and
includes it on emitted batch events.
- **WorkspaceMigrationValidateBuildAndRunService**
Passes `actorContext` from the request into the metadata event emitter.
- **Metadata resolvers & services**
All metadata modules resolve `@AuthUser({ allowUndefined: true })` and
`@AuthUserWorkspaceId()` and pass `userId` and `workspaceMemberId`
through to the migration/event pipeline. Both are optional so
API-key–authenticated requests (no user) still emit events without a
user identity.
Shared some questions on Discord about the PR.
---------
Co-authored-by: Félix Malfait <felix.malfait@gmail.com>
Co-authored-by: prastoin <paul@twenty.com>
This commit is contained in:
+1
-1
@@ -1,8 +1,8 @@
|
||||
import { UseFilters, UseGuards, UsePipes } from '@nestjs/common';
|
||||
import { Args, Context, Mutation, Parent, ResolveField } from '@nestjs/graphql';
|
||||
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { PermissionFlagType } from 'twenty-shared/constants';
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
|
||||
import { PreventNestToAutoLogGraphqlErrorsFilter } from 'src/engine/core-modules/graphql/filters/prevent-nest-to-auto-log-graphql-errors.filter';
|
||||
import { ResolverValidationPipe } from 'src/engine/core-modules/graphql/pipes/resolver-validation.pipe';
|
||||
|
||||
+99
@@ -0,0 +1,99 @@
|
||||
import chunk from 'lodash.chunk';
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
|
||||
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';
|
||||
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 { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
|
||||
import { type MetadataEventBatch } from 'src/engine/metadata-event-emitter/types/metadata-event-batch.type';
|
||||
import { type FlatWebhook } from 'src/engine/metadata-modules/flat-webhook/types/flat-webhook.type';
|
||||
import { CallWebhookJob } from 'src/engine/metadata-modules/webhook/jobs/call-webhook.job';
|
||||
import { type CallMetadataWebhookJobData } from 'src/engine/metadata-modules/webhook/types/webhook-job-data.type';
|
||||
import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
|
||||
|
||||
const WEBHOOK_JOBS_CHUNK_SIZE = 20;
|
||||
|
||||
@Processor(MessageQueue.webhookQueue)
|
||||
export class CallWebhookJobsForMetadataJob {
|
||||
constructor(
|
||||
@InjectMessageQueue(MessageQueue.webhookQueue)
|
||||
private readonly messageQueueService: MessageQueueService,
|
||||
private readonly workspaceCacheService: WorkspaceCacheService,
|
||||
) {}
|
||||
|
||||
@Process(CallWebhookJobsForMetadataJob.name)
|
||||
async handle(metadataEventBatch: MetadataEventBatch): Promise<void> {
|
||||
const eventName = metadataEventBatch.name;
|
||||
const metadataName = metadataEventBatch.metadataName;
|
||||
const operation = metadataEventBatch.type;
|
||||
|
||||
const operationsToMatch = [
|
||||
eventName,
|
||||
`metadata.${metadataName}.*`,
|
||||
`metadata.*.${operation}`,
|
||||
`metadata.*.*`,
|
||||
'*.*',
|
||||
];
|
||||
|
||||
const { flatWebhookMaps } = await this.workspaceCacheService.getOrRecompute(
|
||||
metadataEventBatch.workspaceId,
|
||||
['flatWebhookMaps'],
|
||||
);
|
||||
|
||||
const webhooks = Object.values(flatWebhookMaps.byUniversalIdentifier)
|
||||
.filter(isDefined)
|
||||
.filter((webhook) =>
|
||||
operationsToMatch.some((operationToMatch) =>
|
||||
webhook.operations.includes(operationToMatch),
|
||||
),
|
||||
);
|
||||
|
||||
if (webhooks.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
const webhookEvents = this.transformMetadataEventBatchToWebhookEvents({
|
||||
metadataEventBatch,
|
||||
webhooks,
|
||||
});
|
||||
|
||||
const webhookEventsChunks = chunk(webhookEvents, WEBHOOK_JOBS_CHUNK_SIZE);
|
||||
|
||||
for (const webhookEventsChunk of webhookEventsChunks) {
|
||||
await this.messageQueueService.add<CallMetadataWebhookJobData[]>(
|
||||
CallWebhookJob.name,
|
||||
webhookEventsChunk,
|
||||
{ retryLimit: 3 },
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private transformMetadataEventBatchToWebhookEvents({
|
||||
metadataEventBatch,
|
||||
webhooks,
|
||||
}: {
|
||||
metadataEventBatch: MetadataEventBatch;
|
||||
webhooks: FlatWebhook[];
|
||||
}): CallMetadataWebhookJobData[] {
|
||||
const result: CallMetadataWebhookJobData[] = [];
|
||||
|
||||
for (const webhook of webhooks) {
|
||||
for (const event of metadataEventBatch.events) {
|
||||
result.push({
|
||||
targetUrl: webhook.targetUrl,
|
||||
eventName: metadataEventBatch.name,
|
||||
workspaceId: metadataEventBatch.workspaceId,
|
||||
webhookId: webhook.id,
|
||||
eventDate: new Date(),
|
||||
userId: metadataEventBatch.userId,
|
||||
apiKeyId: metadataEventBatch.apiKeyId,
|
||||
secret: webhook.secret,
|
||||
event,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
}
|
||||
+2
-4
@@ -10,10 +10,8 @@ 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 { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
|
||||
import {
|
||||
CallWebhookJob,
|
||||
type CallWebhookJobData,
|
||||
} from 'src/engine/metadata-modules/webhook/jobs/call-webhook.job';
|
||||
import { CallWebhookJob } from 'src/engine/metadata-modules/webhook/jobs/call-webhook.job';
|
||||
import { type CallWebhookJobData } from 'src/engine/metadata-modules/webhook/types/webhook-job-data.type';
|
||||
import { transformEventBatchToWebhookEvents } from 'src/engine/metadata-modules/webhook/utils/transform-event-batch-to-webhook-events';
|
||||
import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
|
||||
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
|
||||
|
||||
+4
-18
@@ -10,21 +10,7 @@ import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queu
|
||||
import { MetricsService } from 'src/engine/core-modules/metrics/metrics.service';
|
||||
import { MetricsKeys } from 'src/engine/core-modules/metrics/types/metrics-keys.type';
|
||||
import { SecureHttpClientService } from 'src/engine/core-modules/secure-http-client/secure-http-client.service';
|
||||
|
||||
export type CallWebhookJobData = {
|
||||
targetUrl: string;
|
||||
eventName: string;
|
||||
objectMetadata: { id: string; nameSingular: string };
|
||||
workspaceId: string;
|
||||
webhookId: string;
|
||||
eventDate: Date;
|
||||
userId?: string;
|
||||
workspaceMemberId?: string;
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
record: any;
|
||||
updatedFields?: string[];
|
||||
secret?: string;
|
||||
};
|
||||
import { type WebhookJobData } from 'src/engine/metadata-modules/webhook/types/webhook-job-data.type';
|
||||
|
||||
@Processor(MessageQueue.webhookQueue)
|
||||
export class CallWebhookJob {
|
||||
@@ -35,7 +21,7 @@ export class CallWebhookJob {
|
||||
) {}
|
||||
|
||||
private generateSignature(
|
||||
payload: CallWebhookJobData,
|
||||
payload: Record<string, unknown>,
|
||||
secret: string,
|
||||
timestamp: string,
|
||||
): string {
|
||||
@@ -46,7 +32,7 @@ export class CallWebhookJob {
|
||||
}
|
||||
|
||||
@Process(CallWebhookJob.name)
|
||||
async handle(webhookJobEvents: CallWebhookJobData[]): Promise<void> {
|
||||
async handle(webhookJobEvents: WebhookJobData[]): Promise<void> {
|
||||
await Promise.all(
|
||||
webhookJobEvents.map(
|
||||
async (webhookJobEvent) => await this.callWebhook(webhookJobEvent),
|
||||
@@ -54,7 +40,7 @@ export class CallWebhookJob {
|
||||
);
|
||||
}
|
||||
|
||||
private async callWebhook(data: CallWebhookJobData): Promise<void> {
|
||||
private async callWebhook(data: WebhookJobData): Promise<void> {
|
||||
const commonPayload = {
|
||||
url: data.targetUrl,
|
||||
webhookId: data.webhookId,
|
||||
|
||||
+6
-1
@@ -3,6 +3,7 @@ import { Module } from '@nestjs/common';
|
||||
import { AuditModule } from 'src/engine/core-modules/audit/audit.module';
|
||||
import { MetricsModule } from 'src/engine/core-modules/metrics/metrics.module';
|
||||
import { SecureHttpClientModule } from 'src/engine/core-modules/secure-http-client/secure-http-client.module';
|
||||
import { CallWebhookJobsForMetadataJob } from 'src/engine/metadata-modules/webhook/jobs/call-webhook-jobs-for-metadata.job';
|
||||
import { CallWebhookJobsJob } from 'src/engine/metadata-modules/webhook/jobs/call-webhook-jobs.job';
|
||||
import { CallWebhookJob } from 'src/engine/metadata-modules/webhook/jobs/call-webhook.job';
|
||||
import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache.module';
|
||||
@@ -14,6 +15,10 @@ import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache
|
||||
SecureHttpClientModule,
|
||||
WorkspaceCacheModule,
|
||||
],
|
||||
providers: [CallWebhookJobsJob, CallWebhookJob],
|
||||
providers: [
|
||||
CallWebhookJobsJob,
|
||||
CallWebhookJobsForMetadataJob,
|
||||
CallWebhookJob,
|
||||
],
|
||||
})
|
||||
export class WebhookJobModule {}
|
||||
|
||||
+27
@@ -0,0 +1,27 @@
|
||||
import { type MetadataEvent } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/metadata-event';
|
||||
|
||||
type WebhookJobBase = {
|
||||
targetUrl: string;
|
||||
eventName: string;
|
||||
workspaceId: string;
|
||||
webhookId: string;
|
||||
eventDate: Date;
|
||||
userId?: string;
|
||||
apiKeyId?: string;
|
||||
secret?: string;
|
||||
};
|
||||
|
||||
export type CallWebhookJobData = WebhookJobBase & {
|
||||
objectMetadata: { id: string; nameSingular: string };
|
||||
workspaceMemberId?: string;
|
||||
applicationId?: string;
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
record: any;
|
||||
updatedFields?: string[];
|
||||
};
|
||||
|
||||
export type CallMetadataWebhookJobData = WebhookJobBase & {
|
||||
event: MetadataEvent;
|
||||
};
|
||||
|
||||
export type WebhookJobData = CallWebhookJobData | CallMetadataWebhookJobData;
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
import type { ObjectRecordEvent } from 'twenty-shared/database-events';
|
||||
|
||||
import { type WebhookEntity } from 'src/engine/metadata-modules/webhook/entities/webhook.entity';
|
||||
import { type CallWebhookJobData } from 'src/engine/metadata-modules/webhook/jobs/call-webhook.job';
|
||||
import { type CallWebhookJobData } from 'src/engine/metadata-modules/webhook/types/webhook-job-data.type';
|
||||
import { transformEventToWebhookEvent } from 'src/engine/metadata-modules/webhook/utils/transform-event-to-webhook-event';
|
||||
import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
|
||||
|
||||
|
||||
@@ -4,8 +4,8 @@ import { InjectRepository } from '@nestjs/typeorm';
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { IsNull, Repository } from 'typeorm';
|
||||
|
||||
import { ApplicationService } from 'src/engine/core-modules/application/services/application.service';
|
||||
import { ApplicationEntity } from 'src/engine/core-modules/application/application.entity';
|
||||
import { ApplicationService } from 'src/engine/core-modules/application/services/application.service';
|
||||
import { WorkspaceManyOrAllFlatEntityMapsCacheService } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.service';
|
||||
import { findFlatEntityByIdInFlatEntityMapsOrThrow } from 'src/engine/metadata-modules/flat-entity/utils/find-flat-entity-by-id-in-flat-entity-maps-or-throw.util';
|
||||
import { fromCreateWebhookInputToFlatWebhookToCreate } from 'src/engine/metadata-modules/flat-webhook/utils/from-create-webhook-input-to-flat-webhook-to-create.util';
|
||||
|
||||
Reference in New Issue
Block a user