From e53d5036dc61729b91d050d9fc0417ed96fc334f Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Fri, 18 Nov 2022 14:33:27 +0800 Subject: [PATCH] Add notification endpoint --- packages/api/src/datalayer/pubsub.ts | 37 ++++++++++- packages/api/src/entity/rule.ts | 2 +- packages/api/src/generated/graphql.ts | 6 +- packages/api/src/generated/schema.graphql | 4 +- packages/api/src/resolvers/rules/index.ts | 77 ++++++++++++++--------- packages/api/src/schema.ts | 4 +- packages/api/src/services/rules.ts | 25 +++++--- packages/api/src/util.ts | 3 + packages/api/src/utils/pubsub.ts | 39 ------------ packages/db/migrations/0099.do.rules.sql | 4 +- 10 files changed, 111 insertions(+), 90 deletions(-) delete mode 100644 packages/api/src/utils/pubsub.ts diff --git a/packages/api/src/datalayer/pubsub.ts b/packages/api/src/datalayer/pubsub.ts index 8ef12ba83..f81b3210f 100644 --- a/packages/api/src/datalayer/pubsub.ts +++ b/packages/api/src/datalayer/pubsub.ts @@ -1,4 +1,4 @@ -import { PubSub } from '@google-cloud/pubsub' +import { CreateSubscriptionOptions, PubSub } from '@google-cloud/pubsub' import { env } from '../env' import { ReportType } from '../generated/graphql' import express from 'express' @@ -143,3 +143,38 @@ export const readPushSubscription = ( return { message: message, expired: expired(body) } } + +export const createPubSubSubscription = async ( + topicName: string, + subscriptionName: string, + options?: CreateSubscriptionOptions +) => { + const topic = client.topic(topicName) + const [exists] = await topic.exists() + if (!exists) { + await topic.create() + } + + const subscription = topic.subscription(subscriptionName) + const [subscriptionExists] = await subscription.exists() + if (!subscriptionExists) { + await subscription.create(options) + } +} + +export const deletePubSubSubscription = async ( + topicName: string, + subscriptionName: string +) => { + const topic = client.topic(topicName) + const [exists] = await topic.exists() + if (!exists) { + return + } + + const subscription = topic.subscription(subscriptionName) + const [subscriptionExists] = await subscription.exists() + if (subscriptionExists) { + await subscription.delete() + } +} diff --git a/packages/api/src/entity/rule.ts b/packages/api/src/entity/rule.ts index a61121230..bee6b0ed0 100644 --- a/packages/api/src/entity/rule.ts +++ b/packages/api/src/entity/rule.ts @@ -22,7 +22,7 @@ export class Rule { name!: string @Column('text') - query!: string + filter!: string @Column('simple-json') actions!: { type: string; params: string[] }[] diff --git a/packages/api/src/generated/graphql.ts b/packages/api/src/generated/graphql.ts index 932e82d4e..b5e08b5aa 100644 --- a/packages/api/src/generated/graphql.ts +++ b/packages/api/src/generated/graphql.ts @@ -1652,9 +1652,9 @@ export type Rule = { actions: Array; createdAt: Scalars['Date']; enabled: Scalars['Boolean']; + filter: Scalars['String']; id: Scalars['ID']; name: Scalars['String']; - query: Scalars['String']; updatedAt: Scalars['Date']; }; @@ -1973,9 +1973,9 @@ export type SetRuleInput = { actions: Array; description?: InputMaybe; enabled: Scalars['Boolean']; + filter: Scalars['String']; id?: InputMaybe; name: Scalars['String']; - query: Scalars['String']; }; export type SetRuleResult = SetRuleError | SetRuleSuccess; @@ -4406,9 +4406,9 @@ export type RuleResolvers, ParentType, ContextType>; createdAt?: Resolver; enabled?: Resolver; + filter?: Resolver; id?: Resolver; name?: Resolver; - query?: Resolver; updatedAt?: Resolver; __isTypeOf?: IsTypeOfResolverFn; }; diff --git a/packages/api/src/generated/schema.graphql b/packages/api/src/generated/schema.graphql index e8bae6a58..44c004f0c 100644 --- a/packages/api/src/generated/schema.graphql +++ b/packages/api/src/generated/schema.graphql @@ -1160,9 +1160,9 @@ type Rule { actions: [RuleAction!]! createdAt: Date! enabled: Boolean! + filter: String! id: ID! name: String! - query: String! updatedAt: Date! } @@ -1457,9 +1457,9 @@ input SetRuleInput { actions: [RuleActionInput!]! description: String enabled: Boolean! + filter: String! id: ID name: String! - query: String! } union SetRuleResult = SetRuleError | SetRuleSuccess diff --git a/packages/api/src/resolvers/rules/index.ts b/packages/api/src/resolvers/rules/index.ts index 18b09ef25..10acb18da 100644 --- a/packages/api/src/resolvers/rules/index.ts +++ b/packages/api/src/resolvers/rules/index.ts @@ -16,7 +16,7 @@ import { import { createPubSubSubscription, deletePubSubSubscription, -} from '../../utils/pubsub' +} from '../../datalayer/pubsub' export const setRuleResolver = authorized< SetRuleSuccess, @@ -32,42 +32,57 @@ export const setRuleResolver = authorized< }, }) - const user = await getRepository(User).findOneBy({ id: claims.uid }) - if (!user) { - return { - errorCodes: [SetRuleErrorCode.Unauthorized], + try { + const user = await getRepository(User).findOneBy({ id: claims.uid }) + if (!user) { + return { + errorCodes: [SetRuleErrorCode.Unauthorized], + } } - } - // create or delete pubsub subscription based on action and enabled state - for (const action of input.actions) { - const topicName = getPubSubTopicName(action) - const subscriptionName = getPubSubSubscriptionName( - topicName, - user.id, - input.name - ) - - if (input.enabled) { - const options = getPubSubSubscriptionOptions( + // create or delete pubsub subscription based on action and enabled state + for (const action of input.actions) { + const topicName = getPubSubTopicName(action) + const subscriptionName = getPubSubSubscriptionName( + topicName, user.id, - input.name, - input.query, - action + input.name ) - await createPubSubSubscription(topicName, subscriptionName, options) - } else { - await deletePubSubSubscription(topicName, subscriptionName) + + if (input.enabled) { + const options = getPubSubSubscriptionOptions( + user.id, + input.name, + input.filter, + action + ) + await createPubSubSubscription(topicName, subscriptionName, options) + } else { + await deletePubSubSubscription(topicName, subscriptionName) + } } - } - const rule = await getRepository(Rule).save({ - ...input, - id: input.id || undefined, - user: { id: claims.uid }, - }) + const rule = await getRepository(Rule).save({ + ...input, + id: input.id || undefined, + user: { id: claims.uid }, + }) - return { - rule, + return { + rule, + } + } catch (error) { + log.error('Error setting rules', { + error, + labels: { + source: 'resolver', + resolver: 'setRulesResolver', + uid: claims.uid, + }, + }) + + return { + errorCodes: [SetRuleErrorCode.BadRequest], + } } }) diff --git a/packages/api/src/schema.ts b/packages/api/src/schema.ts index d2ba1eaf3..1cf55942b 100755 --- a/packages/api/src/schema.ts +++ b/packages/api/src/schema.ts @@ -1959,7 +1959,7 @@ const schema = gql` type Rule { id: ID! name: String! - query: String! + filter: String! actions: [RuleAction!]! enabled: Boolean! createdAt: Date! @@ -1991,7 +1991,7 @@ const schema = gql` id: ID name: String! description: String - query: String! + filter: String! actions: [RuleActionInput!]! enabled: Boolean! } diff --git a/packages/api/src/services/rules.ts b/packages/api/src/services/rules.ts index c60411ca4..c42e4cbad 100644 --- a/packages/api/src/services/rules.ts +++ b/packages/api/src/services/rules.ts @@ -1,5 +1,6 @@ import { RuleAction, RuleActionType } from '../generated/graphql' import { CreateSubscriptionOptions } from '@google-cloud/pubsub' +import { env } from '../env' enum RuleTrigger { ON_PAGE_UPDATE, @@ -42,14 +43,10 @@ export const getPubSubSubscriptionName = ( export const getPubSubSubscriptionOptions = ( userId: string, ruleName: string, - query: string, + filter: string, action: RuleAction ): CreateSubscriptionOptions => { - const topic = getPubSubTopicName(action) - const name = getPubSubSubscriptionName(topic, userId, ruleName) const options: CreateSubscriptionOptions = { - name, - topic, messageRetentionDuration: 60 * 10, // 10 minutes expirationPolicy: { ttl: null, // never expire @@ -63,14 +60,24 @@ export const getPubSubSubscriptionOptions = ( seconds: 600, }, }, - filter: query, + filter, } switch (action.type) { - case RuleActionType.SendNotification: - options.pushEndpoint = `${process.env - .PUSH_NOTIFICATION_ENDPOINT!}?message=${action.params[0]}` + case RuleActionType.SendNotification: { + const params = action.params + if (params.length === 0) { + throw new Error('Missing notification messages') + } + + options.pushConfig = { + pushEndpoint: `${env.queue.notificationEndpoint}/${userId}`, + attributes: { + messages: JSON.stringify(params), + }, + } break + } // TODO: Add more actions, e.g. RuleActionType.SendEmail } diff --git a/packages/api/src/util.ts b/packages/api/src/util.ts index 0ca9ad8b8..6317207fb 100755 --- a/packages/api/src/util.ts +++ b/packages/api/src/util.ts @@ -65,6 +65,7 @@ interface BackendEnv { reminderTaskHanderUrl: string integrationTaskHandlerUrl: string textToSpeechTaskHandlerUrl: string + notificationEndpoint: string } fileUpload: { gcsUploadBucket: string @@ -152,6 +153,7 @@ const nullableEnvVars = [ 'AZURE_SPEECH_KEY', 'AZURE_SPEECH_REGION', 'GCP_LOCATION', + 'NOTIFICATION_ENDPOINT', ] // Allow some vars to be null/empty /* If not in GAE and Prod/QA/Demo env (f.e. on localhost/dev env), allow following env vars to be null */ @@ -237,6 +239,7 @@ export function getEnv(): BackendEnv { reminderTaskHanderUrl: parse('REMINDER_TASK_HANDLER_URL'), integrationTaskHandlerUrl: parse('INTEGRATION_TASK_HANDLER_URL'), textToSpeechTaskHandlerUrl: parse('TEXT_TO_SPEECH_TASK_HANDLER_URL'), + notificationEndpoint: parse('NOTIFICATION_ENDPOINT'), } const imageProxy = { url: parse('IMAGE_PROXY_URL'), diff --git a/packages/api/src/utils/pubsub.ts b/packages/api/src/utils/pubsub.ts deleted file mode 100644 index 8dfb91e59..000000000 --- a/packages/api/src/utils/pubsub.ts +++ /dev/null @@ -1,39 +0,0 @@ -import { CreateSubscriptionOptions, PubSub } from '@google-cloud/pubsub' - -// init pubsub client -const pubsub = new PubSub() - -export const createPubSubSubscription = async ( - topicName: string, - subscriptionName: string, - options?: CreateSubscriptionOptions -) => { - const topic = pubsub.topic(topicName) - const [exists] = await topic.exists() - if (!exists) { - await topic.create() - } - - const subscription = topic.subscription(subscriptionName) - const [subscriptionExists] = await subscription.exists() - if (!subscriptionExists) { - await subscription.create(options) - } -} - -export const deletePubSubSubscription = async ( - topicName: string, - subscriptionName: string -) => { - const topic = pubsub.topic(topicName) - const [exists] = await topic.exists() - if (!exists) { - return - } - - const subscription = topic.subscription(subscriptionName) - const [subscriptionExists] = await subscription.exists() - if (subscriptionExists) { - await subscription.delete() - } -} diff --git a/packages/db/migrations/0099.do.rules.sql b/packages/db/migrations/0099.do.rules.sql index 0c89f7991..406bc912b 100755 --- a/packages/db/migrations/0099.do.rules.sql +++ b/packages/db/migrations/0099.do.rules.sql @@ -9,12 +9,12 @@ CREATE TABLE omnivore.rules ( user_id uuid NOT NULL REFERENCES omnivore.user ON DELETE CASCADE, name text NOT NULL, description text, - query text NOT NULL, + filter text NOT NULL, actions json NOT NULL, -- array of actions of type {type: 'action_type', params: [action_params]} enabled boolean NOT NULL DEFAULT true, created_at timestamptz NOT NULL DEFAULT current_timestamp, updated_at timestamptz NOT NULL DEFAULT current_timestamp, - UNIQUE (user_id, name) + UNIQUE (user_id, filter) ); CREATE TRIGGER rules_modtime BEFORE UPDATE ON omnivore.rules