From 25e374f6ff6ad6f909cc85f7a9a29398cbaf3bc7 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Thu, 1 Feb 2024 19:50:59 +0800 Subject: [PATCH 1/3] fix: replace webhook endpoint with a bullmq job --- packages/api/src/jobs/call_webhook.ts | 56 +++++++++++++++++++++++++++ packages/api/src/pubsub.ts | 16 +++++++- packages/api/src/queue-processor.ts | 3 ++ packages/api/src/services/webhook.ts | 18 ++++++++- packages/api/src/utils/createTask.ts | 15 +++++++ 5 files changed, 106 insertions(+), 2 deletions(-) create mode 100644 packages/api/src/jobs/call_webhook.ts diff --git a/packages/api/src/jobs/call_webhook.ts b/packages/api/src/jobs/call_webhook.ts new file mode 100644 index 000000000..de8e986cb --- /dev/null +++ b/packages/api/src/jobs/call_webhook.ts @@ -0,0 +1,56 @@ +import axios, { Method } from 'axios' +import { findWebhooksByEventTypes } from '../services/webhook' +import { logger } from '../utils/logger' + +export interface CallWebhookJobData { + data: unknown + userId: string + type: string + action: string +} + +export const CALL_WEBHOOK_JOB_NAME = 'call-webhook' +const TIMEOUT = 5000 // 5s + +export const callWebhook = async (jobData: CallWebhookJobData) => { + const { data, type, action, userId } = jobData + const eventType = `${type}_${action}`.toUpperCase() + const webhooks = await findWebhooksByEventTypes(userId, [eventType]) + + if (webhooks.length <= 0) { + return + } + + await Promise.all( + webhooks.map((webhook) => { + const url = webhook.url + const method = webhook.method as Method + const body = { + action, + userId, + [type]: data, + } + + logger.info('triggering webhook', { url, method }) + + return axios + .request({ + url, + method, + headers: { + 'Content-Type': webhook.contentType, + }, + data: body, + timeout: TIMEOUT, + }) + .then(() => logger.info('webhook triggered')) + .catch((error) => { + if (axios.isAxiosError(error)) { + logger.info('webhook failed', error.response) + } else { + logger.info('webhook failed', error) + } + }) + }) + ) +} diff --git a/packages/api/src/pubsub.ts b/packages/api/src/pubsub.ts index 02c35f3d4..2f5abf639 100644 --- a/packages/api/src/pubsub.ts +++ b/packages/api/src/pubsub.ts @@ -3,7 +3,7 @@ import express from 'express' import { RuleEventType } from './entity/rule' import { env } from './env' import { ReportType } from './generated/graphql' -import { enqueueTriggerRuleJob } from './utils/createTask' +import { enqueueTriggerRuleJob, enqueueWebhookJob } from './utils/createTask' import { deepDelete } from './utils/helpers' import { buildLogger } from './utils/logger' @@ -63,6 +63,13 @@ export const createPubSubClient = (): PubsubClient => { [...fieldsToDelete] ) + await enqueueWebhookJob({ + userId, + type, + action: 'created', + data, + }) + return publish( 'entityCreated', Buffer.from(JSON.stringify({ type, userId, ...cleanData })) @@ -88,6 +95,13 @@ export const createPubSubClient = (): PubsubClient => { [...fieldsToDelete] ) + await enqueueWebhookJob({ + userId, + type, + action: 'updated', + data, + }) + return publish( 'entityUpdated', Buffer.from(JSON.stringify({ type, userId, ...cleanData })) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 162611742..b93b327d9 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -15,6 +15,7 @@ import { SnakeNamingStrategy } from 'typeorm-naming-strategies' import { appDataSource } from './data_source' import { env } from './env' import { bulkAction, BULK_ACTION_JOB_NAME } from './jobs/bulk_action' +import { CALL_WEBHOOK_JOB_NAME, callWebhook } from './jobs/call_webhook' import { findThumbnail, THUMBNAIL_JOB } from './jobs/find_thumbnail' import { refreshAllFeeds } from './jobs/rss/refreshAllFeeds' import { refreshFeed } from './jobs/rss/refreshFeed' @@ -87,6 +88,8 @@ export const createWorker = (connection: ConnectionOptions) => return syncReadPositionsJob(job.data) case BULK_ACTION_JOB_NAME: return bulkAction(job.data) + case CALL_WEBHOOK_JOB_NAME: + return callWebhook(job.data) } }, { diff --git a/packages/api/src/services/webhook.ts b/packages/api/src/services/webhook.ts index 4414f44a1..149560fb5 100644 --- a/packages/api/src/services/webhook.ts +++ b/packages/api/src/services/webhook.ts @@ -1,4 +1,4 @@ -import { DeepPartial, EntityManager } from 'typeorm' +import { ArrayContainedBy, DeepPartial, EntityManager } from 'typeorm' import { Webhook } from '../entity/webhook' import { authTrx } from '../repository' @@ -34,6 +34,22 @@ export const findWebhooks = async (userId: string) => { ) } +export const findWebhooksByEventTypes = async ( + userId: string, + eventTypes: string[] +) => { + return authTrx( + (tx) => + tx.getRepository(Webhook).findBy({ + user: { id: userId }, + enabled: true, + eventTypes: ArrayContainedBy(eventTypes), + }), + undefined, + userId + ) +} + export const findWebhookById = async (id: string, userId: string) => { return authTrx( (tx) => tx.getRepository(Webhook).findOneBy({ id, user: { id: userId } }), diff --git a/packages/api/src/utils/createTask.ts b/packages/api/src/utils/createTask.ts index da1115200..53a15e041 100644 --- a/packages/api/src/utils/createTask.ts +++ b/packages/api/src/utils/createTask.ts @@ -15,6 +15,7 @@ import { CreateLabelInput, } from '../generated/graphql' import { BulkActionData, BULK_ACTION_JOB_NAME } from '../jobs/bulk_action' +import { CallWebhookJobData, CALL_WEBHOOK_JOB_NAME } from '../jobs/call_webhook' import { THUMBNAIL_JOB } from '../jobs/find_thumbnail' import { queueRSSRefreshFeedJob } from '../jobs/rss/refreshAllFeeds' import { TriggerRuleJobData, TRIGGER_RULE_JOB_NAME } from '../jobs/trigger_rule' @@ -670,6 +671,20 @@ export const enqueueTriggerRuleJob = async (data: TriggerRuleJobData) => { }) } +export const enqueueWebhookJob = async (data: CallWebhookJobData) => { + const queue = await getBackendQueue() + if (!queue) { + return undefined + } + + return queue.add(CALL_WEBHOOK_JOB_NAME, data, { + priority: 1, + attempts: 1, + removeOnComplete: true, + removeOnFail: true, + }) +} + export const bulkEnqueueUpdateLabels = async (data: UpdateLabelsData[]) => { const queue = await getBackendQueue() if (!queue) { From 1742ffda24b14df48a9313ce70c6f2833bc39d2a Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Thu, 1 Feb 2024 20:53:48 +0800 Subject: [PATCH 2/3] change priority to 5 --- packages/api/src/utils/createTask.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/api/src/utils/createTask.ts b/packages/api/src/utils/createTask.ts index 53a15e041..3a1145aec 100644 --- a/packages/api/src/utils/createTask.ts +++ b/packages/api/src/utils/createTask.ts @@ -664,7 +664,7 @@ export const enqueueTriggerRuleJob = async (data: TriggerRuleJobData) => { } return queue.add(TRIGGER_RULE_JOB_NAME, data, { - priority: 1, + priority: 5, attempts: 1, removeOnComplete: true, removeOnFail: true, @@ -678,7 +678,7 @@ export const enqueueWebhookJob = async (data: CallWebhookJobData) => { } return queue.add(CALL_WEBHOOK_JOB_NAME, data, { - priority: 1, + priority: 5, attempts: 1, removeOnComplete: true, removeOnFail: true, From 74783313da41cf2b91f019d1319474a6b9956aa7 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Fri, 2 Feb 2024 16:17:36 +0800 Subject: [PATCH 3/3] fix query --- packages/api/src/jobs/call_webhook.ts | 4 ++-- packages/api/src/services/rules.ts | 4 ++-- packages/api/src/services/webhook.ts | 8 ++++---- 3 files changed, 8 insertions(+), 8 deletions(-) diff --git a/packages/api/src/jobs/call_webhook.ts b/packages/api/src/jobs/call_webhook.ts index de8e986cb..a7f152c02 100644 --- a/packages/api/src/jobs/call_webhook.ts +++ b/packages/api/src/jobs/call_webhook.ts @@ -1,5 +1,5 @@ import axios, { Method } from 'axios' -import { findWebhooksByEventTypes } from '../services/webhook' +import { findWebhooksByEventType } from '../services/webhook' import { logger } from '../utils/logger' export interface CallWebhookJobData { @@ -15,7 +15,7 @@ const TIMEOUT = 5000 // 5s export const callWebhook = async (jobData: CallWebhookJobData) => { const { data, type, action, userId } = jobData const eventType = `${type}_${action}`.toUpperCase() - const webhooks = await findWebhooksByEventTypes(userId, [eventType]) + const webhooks = await findWebhooksByEventType(userId, eventType) if (webhooks.length <= 0) { return diff --git a/packages/api/src/services/rules.ts b/packages/api/src/services/rules.ts index 5f27a3839..bf3637be7 100644 --- a/packages/api/src/services/rules.ts +++ b/packages/api/src/services/rules.ts @@ -1,4 +1,4 @@ -import { ArrayContainedBy, ArrayContains, ILike } from 'typeorm' +import { ArrayContains, ILike } from 'typeorm' import { Rule, RuleAction, RuleEventType } from '../entity/rule' import { authTrx, getRepository } from '../repository' @@ -61,6 +61,6 @@ export const findEnabledRules = async ( return getRepository(Rule).findBy({ user: { id: userId }, enabled: true, - eventTypes: ArrayContainedBy([eventType]), + eventTypes: ArrayContains([eventType]), }) } diff --git a/packages/api/src/services/webhook.ts b/packages/api/src/services/webhook.ts index 149560fb5..736b6fc54 100644 --- a/packages/api/src/services/webhook.ts +++ b/packages/api/src/services/webhook.ts @@ -1,4 +1,4 @@ -import { ArrayContainedBy, DeepPartial, EntityManager } from 'typeorm' +import { ArrayContains, DeepPartial, EntityManager } from 'typeorm' import { Webhook } from '../entity/webhook' import { authTrx } from '../repository' @@ -34,16 +34,16 @@ export const findWebhooks = async (userId: string) => { ) } -export const findWebhooksByEventTypes = async ( +export const findWebhooksByEventType = async ( userId: string, - eventTypes: string[] + eventType: string ) => { return authTrx( (tx) => tx.getRepository(Webhook).findBy({ user: { id: userId }, enabled: true, - eventTypes: ArrayContainedBy(eventTypes), + eventTypes: ArrayContains([eventType]), }), undefined, userId