diff --git a/packages/rule-handler/src/index.ts b/packages/rule-handler/src/index.ts index 029efcd26..9bc6b2e08 100644 --- a/packages/rule-handler/src/index.ts +++ b/packages/rule-handler/src/index.ts @@ -16,6 +16,7 @@ interface PubSubRequestMessage { interface PubSubRequestBody { message: PubSubRequestMessage + subscription: string } export interface PubSubData { @@ -45,24 +46,26 @@ const expired = (body: PubSubRequestBody): boolean => { const readPushSubscription = ( req: express.Request -): { message: string | undefined; expired: boolean } => { - console.debug('request query', req.body) - +): { message: string | undefined; expired: boolean; subscription?: string } => { if (req.query.token !== process.env.PUBSUB_VERIFICATION_TOKEN) { console.log('query does not include valid pubsub token') return { message: undefined, expired: false } } // GCP PubSub sends the request as a base64 encoded string - if (!('message' in req.body)) { - console.log('Invalid pubsub message: message not in body') + if (!('message' in req.body) || !('subscription' in req.body)) { + console.log('Invalid pubsub message: message or subscription not in body') return { message: undefined, expired: false } } const body = req.body as PubSubRequestBody const message = Buffer.from(body.message.data, 'base64').toString('utf-8') - return { message: message, expired: expired(body) } + return { + message: message, + expired: expired(body), + subscription: body.subscription, + } } export const getAuthToken = async ( @@ -81,9 +84,9 @@ export const ruleHandler = Sentry.GCPFunction.wrapHttpFunction( throw new Error('REST_BACKEND_ENDPOINT or JWT_SECRET not set') } - const { message: msgStr, expired } = readPushSubscription(req) + const { message: msgStr, expired, subscription } = readPushSubscription(req) - if (!msgStr) { + if (!msgStr || !subscription) { res.status(400).send('Bad Request') return } @@ -117,12 +120,21 @@ export const ruleHandler = Sentry.GCPFunction.wrapHttpFunction( return } + // subscription: 'projects/omnivore-demo/subscriptions/entityUpdated-rule' + const eventType = subscription.match(/entity(.*)-rule/)?.[1] + if (!eventType) { + console.log('No event type found') + res.status(200).send('No Event Type') + return + } + const triggeredActions = await triggerActions( userId, rules, data, apiEndpoint, - jwtSecret + jwtSecret, + eventType ) if (triggeredActions.length === 0) { console.log('No actions triggered') diff --git a/packages/rule-handler/src/rule.ts b/packages/rule-handler/src/rule.ts index b03a3420a..a1a2f898b 100644 --- a/packages/rule-handler/src/rule.ts +++ b/packages/rule-handler/src/rule.ts @@ -29,6 +29,8 @@ export interface Rule { updatedAt: Date } +const EVENT_FILTERS = ['event:created', 'event:updated'] + export const getEnabledRules = async ( userId: string, apiEndpoint: string, @@ -73,17 +75,34 @@ export const triggerActions = async ( rules: Rule[], data: PubSubData, apiEndpoint: string, - jwtSecret: string + jwtSecret: string, + eventType: string ) => { const authToken = await getAuthToken(userId, jwtSecret) const actionPromises: Promise | undefined>[] = [] for (const rule of rules) { + let filter = rule.filter + const filters = filter.split(' ') + + // Check if the rule is enabled for the event type + const eventFilterIndex = filters.findIndex((f) => EVENT_FILTERS.includes(f)) + if (eventFilterIndex !== -1) { + const eventFilter = filters[eventFilterIndex] + if (eventFilter !== `event:${eventType}`.toLowerCase()) { + continue + } + + // Remove the event filter from the filter string + filters.splice(eventFilterIndex, 1) + filter = filters.join(' ') + } + const filteredPage = await filterPage( userId, apiEndpoint, authToken, - rule.filter, + filter, data.id ) if (!filteredPage) {