[CleanUp] Post api keys and webhooks migration cleanup (#13576)

This commit is contained in:
nitin
2025-08-05 12:43:09 +05:30
committed by GitHub
parent df9136c6c0
commit 676cf838af
15 changed files with 7 additions and 169 deletions
@@ -0,0 +1,94 @@
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 { WebhookService } from 'src/engine/core-modules/webhook/webhook.service';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type';
import {
CallWebhookJob,
CallWebhookJobData,
} from 'src/engine/core-modules/webhook/jobs/call-webhook.job';
import { ObjectRecordEventForWebhook } from 'src/engine/core-modules/webhook/types/object-record-event-for-webhook.type';
import { removeSecretFromWebhookRecord } from 'src/utils/remove-secret-from-webhook-record';
@Processor(MessageQueue.webhookQueue)
export class CallWebhookJobsJob {
constructor(
@InjectMessageQueue(MessageQueue.webhookQueue)
private readonly messageQueueService: MessageQueueService,
private readonly webhookService: WebhookService,
) {}
@Process(CallWebhookJobsJob.name)
async handle(
workspaceEventBatch: WorkspaceEventBatch<ObjectRecordEventForWebhook>,
): Promise<void> {
// If you change that function, double check it does not break Zapier
// trigger in packages/twenty-zapier/src/triggers/trigger_record.ts
// Also change the openApi schema for webhooks
// packages/twenty-server/src/engine/core-modules/open-api/utils/computeWebhooks.utils.ts
const [nameSingular, operation] = workspaceEventBatch.name.split('.');
const webhooks = await this.webhookService.findByOperations(
workspaceEventBatch.workspaceId,
[
`${nameSingular}.${operation}`,
`*.${operation}`,
`${nameSingular}.*`,
'*.*',
],
);
for (const eventData of workspaceEventBatch.events) {
const eventName = workspaceEventBatch.name;
const objectMetadata: Pick<ObjectMetadataEntity, 'id' | 'nameSingular'> =
{
id: eventData.objectMetadata.id,
nameSingular: eventData.objectMetadata.nameSingular,
};
const workspaceId = workspaceEventBatch.workspaceId;
const record =
'after' in eventData.properties && isDefined(eventData.properties.after)
? eventData.properties.after
: 'before' in eventData.properties &&
isDefined(eventData.properties.before)
? eventData.properties.before
: {};
const updatedFields =
'updatedFields' in eventData.properties
? eventData.properties.updatedFields
: undefined;
const isWebhookEvent = nameSingular === 'webhook';
const sanitizedRecord = removeSecretFromWebhookRecord(
record,
isWebhookEvent,
);
webhooks.forEach((webhook) => {
const webhookData = {
targetUrl: webhook.targetUrl,
secret: webhook.secret,
eventName,
objectMetadata,
workspaceId,
webhookId: webhook.id,
eventDate: new Date(),
record: sanitizedRecord,
...(updatedFields && { updatedFields }),
};
this.messageQueueService.add<CallWebhookJobData>(
CallWebhookJob.name,
webhookData,
{ retryLimit: 3 },
);
});
}
}
}
@@ -0,0 +1,95 @@
import { HttpService } from '@nestjs/axios';
import crypto from 'crypto';
import { getAbsoluteUrl } from 'twenty-shared/utils';
import { AuditService } from 'src/engine/core-modules/audit/services/audit.service';
import { WEBHOOK_RESPONSE_EVENT } from 'src/engine/core-modules/audit/utils/events/workspace-event/webhook/webhook-response';
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';
export type CallWebhookJobData = {
targetUrl: string;
eventName: string;
objectMetadata: { id: string; nameSingular: string };
workspaceId: string;
webhookId: string;
eventDate: Date;
// eslint-disable-next-line @typescript-eslint/no-explicit-any
record: any;
updatedFields?: string[];
secret?: string;
};
@Processor(MessageQueue.webhookQueue)
export class CallWebhookJob {
constructor(
private readonly httpService: HttpService,
private readonly auditService: AuditService,
) {}
private generateSignature(
payload: CallWebhookJobData,
secret: string,
timestamp: string,
): string {
return crypto
.createHmac('sha256', secret)
.update(`${timestamp}:${JSON.stringify(payload)}`)
.digest('hex');
}
@Process(CallWebhookJob.name)
async handle(data: CallWebhookJobData): Promise<void> {
const commonPayload = {
url: data.targetUrl,
webhookId: data.webhookId,
eventName: data.eventName,
};
const analytics = this.auditService.createContext({
workspaceId: data.workspaceId,
});
try {
const headers: Record<string, string> = {
'Content-Type': 'application/json',
};
const { secret, ...payloadWithoutSecret } = data;
if (secret) {
headers['X-Twenty-Webhook-Timestamp'] = Date.now().toString();
headers['X-Twenty-Webhook-Signature'] = this.generateSignature(
payloadWithoutSecret,
secret,
headers['X-Twenty-Webhook-Timestamp'],
);
headers['X-Twenty-Webhook-Nonce'] = crypto
.randomBytes(16)
.toString('hex');
}
const response = await this.httpService.axiosRef.post(
getAbsoluteUrl(data.targetUrl),
payloadWithoutSecret,
{ headers },
);
const success = response.status >= 200 && response.status < 300;
analytics.insertWorkspaceEvent(WEBHOOK_RESPONSE_EVENT, {
status: response.status,
success,
...commonPayload,
});
} catch (err) {
analytics.insertWorkspaceEvent(WEBHOOK_RESPONSE_EVENT, {
success: false,
...commonPayload,
...(err.response && { status: err.response.status }),
});
}
}
}
@@ -0,0 +1,13 @@
import { HttpModule } from '@nestjs/axios';
import { Module } from '@nestjs/common';
import { AuditModule } from 'src/engine/core-modules/audit/audit.module';
import { WebhookModule } from 'src/engine/core-modules/webhook/webhook.module';
import { CallWebhookJobsJob } from 'src/engine/core-modules/webhook/jobs/call-webhook-jobs.job';
import { CallWebhookJob } from 'src/engine/core-modules/webhook/jobs/call-webhook.job';
@Module({
imports: [HttpModule, AuditModule, WebhookModule],
providers: [CallWebhookJobsJob, CallWebhookJob],
})
export class WebhookJobModule {}
@@ -0,0 +1,9 @@
import { ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
export type ObjectRecordEventForWebhook = Omit<
ObjectRecordEvent,
'objectMetadata'
> & {
objectMetadata: Pick<ObjectMetadataEntity, 'id' | 'nameSingular'>;
};