Create rule handler cloud run

This commit is contained in:
Hongbo Wu 2022-11-19 11:15:26 +08:00
parent 075c1b4d14
commit dc92fe9c52
21 changed files with 309 additions and 269 deletions

View file

@ -1,4 +1,4 @@
import { CreateSubscriptionOptions, PubSub } from '@google-cloud/pubsub'
import { PubSub } from '@google-cloud/pubsub'
import { env } from '../env'
import { ReportType } from '../generated/graphql'
import express from 'express'
@ -143,38 +143,3 @@ 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()
}
}

View file

@ -8,15 +8,6 @@ import {
import { getRepository } from '../../entity/utils'
import { User } from '../../entity/user'
import { Rule } from '../../entity/rule'
import {
getPubSubSubscriptionName,
getPubSubSubscriptionOptions,
getPubSubTopicName,
} from '../../services/rules'
import {
createPubSubSubscription,
deletePubSubSubscription,
} from '../../datalayer/pubsub'
export const setRuleResolver = authorized<
SetRuleSuccess,
@ -40,28 +31,6 @@ export const setRuleResolver = authorized<
}
}
// 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 = await 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,

View file

@ -1,94 +0,0 @@
import { RuleAction, RuleActionType } from '../generated/graphql'
import { CreateSubscriptionOptions } from '@google-cloud/pubsub'
import { env } from '../env'
import { getDeviceTokensByUserId } from './user_device_tokens'
enum RuleTrigger {
ON_PAGE_UPDATE,
CRON,
}
export const getRuleTrigger = (action: RuleAction): RuleTrigger => {
switch (action.type) {
case RuleActionType.AddLabel:
case RuleActionType.Archive:
case RuleActionType.MarkAsRead:
case RuleActionType.SendNotification:
return RuleTrigger.ON_PAGE_UPDATE
// TODO: Add more actions, e.g. RuleActionType.SendEmail
}
return RuleTrigger.ON_PAGE_UPDATE
}
export const getPubSubTopicName = (action: RuleAction): string => {
const trigger = getRuleTrigger(action)
switch (trigger) {
case RuleTrigger.ON_PAGE_UPDATE:
return 'entityUpdated'
// TODO: Add more triggers, e.g. RuleTrigger.CRON
}
return 'entityUpdated'
}
export const getPubSubSubscriptionName = (
topicName: string,
userId: string,
ruleName: string
): string => {
return `${topicName}-${userId}-rule-${ruleName}`
}
export const getPubSubSubscriptionOptions = async (
userId: string,
ruleName: string,
filter: string,
action: RuleAction
): Promise<CreateSubscriptionOptions> => {
const options: CreateSubscriptionOptions = {
messageRetentionDuration: 60 * 10, // 10 minutes
expirationPolicy: {
ttl: null, // never expire
},
ackDeadlineSeconds: 10,
retryPolicy: {
minimumBackoff: {
seconds: 10,
},
maximumBackoff: {
seconds: 600,
},
},
filter,
}
switch (action.type) {
case RuleActionType.SendNotification: {
const params = action.params
if (params.length === 0) {
throw new Error('Missing notification messages')
}
const deviceTokens = await getDeviceTokensByUserId(userId)
if (!deviceTokens || deviceTokens.length === 0) {
throw new Error('No device tokens found')
}
options.pushConfig = {
pushEndpoint: `${env.queue.notificationEndpoint}?token=${env.queue.verificationToken}`,
attributes: {
userId,
filter,
messages: JSON.stringify(params),
tokens: JSON.stringify(deviceTokens.map((t) => t.token)),
},
}
break
}
// TODO: Add more actions, e.g. RuleActionType.SendEmail
}
return options
}

View file

@ -1,97 +0,0 @@
import * as Sentry from '@sentry/serverless'
import { Request, Response } from 'express'
import { sendBatchPushNotifications } from './sendNotification'
import { Message } from 'firebase-admin/lib/messaging'
import * as dotenv from 'dotenv' // see https://github.com/motdotla/dotenv#how-do-i-use-dotenv-with-import
dotenv.config()
interface SubscriptionAttributes {
messages: string[]
tokens: string[]
}
interface SubscriptionData {
attributes?: SubscriptionAttributes
data: string
}
const readPushSubscription = (req: Request): SubscriptionData | null => {
console.debug('request query', req.body)
if (req.query.token !== process.env.PUBSUB_VERIFICATION_TOKEN) {
console.log('query does not include valid pubsub token')
return null
}
// GCP PubSub sends the request as a base64 encoded string
if (!('message' in req.body)) {
console.log('Invalid pubsub message: message not in body')
return null
}
const body = req.body as {
message: { data: string }
attributes?: SubscriptionAttributes
}
const data = Buffer.from(body.message.data, 'base64').toString('utf-8')
return {
data,
attributes: body.attributes,
}
}
const getBatchMessages = (messages: string[], tokens: string[]): Message[] => {
const batchMessages: Message[] = []
messages.forEach((message) => {
tokens.forEach((token) => {
batchMessages.push({
token,
notification: {
body: message,
},
})
})
})
return batchMessages
}
export const notification = Sentry.GCPFunction.wrapHttpFunction(
async (req: Request, res: Response) => {
try {
const subscriptionData = readPushSubscription(req)
if (!subscriptionData) {
res.status(400).send('Invalid request')
return
}
const { attributes } = subscriptionData
if (!attributes) {
res.status(400).send('Invalid request')
return
}
const { messages, tokens } = attributes
if (
!messages ||
messages.length === 0 ||
!tokens ||
tokens.length === 0
) {
res.status(400).send('Invalid request')
return
}
const batchMessages = getBatchMessages(messages, tokens)
await sendBatchPushNotifications(batchMessages)
res.status(200).send('OK')
} catch (error) {
console.error(error)
res.status(500).send('Internal server error')
}
}
)

View file

@ -8,19 +8,19 @@ COPY yarn.lock .
COPY tsconfig.json .
COPY .eslintrc .
COPY /packages/notification/package.json ./packages/notification/package.json
COPY /packages/rule-handler/package.json ./packages/rule-handler/package.json
RUN yarn install --pure-lockfile
ADD /packages/notification ./packages/notification
RUN yarn workspace @omnivore/notification build
ADD /packages/rule-handler ./packages/rule-handler
RUN yarn workspace @omnivore/rule-handler build
# After building, fetch the production dependencies
RUN rm -rf /app/packages/notification/node_modules
RUN rm -rf /app/packages/rule-handler/node_modules
RUN rm -rf /app/node_modules
RUN yarn install --pure-lockfile --production
EXPOSE 8080
CMD ["yarn", "workspace", "@omnivore/notification", "start"]
CMD ["yarn", "workspace", "@omnivore/rule-handler", "start"]

View file

@ -1,5 +1,5 @@
{
"name": "@omnivore/notification",
"name": "@omnivore/rule-handler",
"version": "1.0.0",
"main": "build/src/index.js",
"files": [
@ -11,9 +11,9 @@
"lint": "eslint src --ext ts,js,tsx,jsx",
"compile": "tsc",
"build": "tsc",
"start": "functions-framework --target=notification",
"start": "functions-framework --target=ruleHandler",
"dev": "concurrently \"tsc -w\" \"nodemon --watch ./build/ --exec npm run start\"",
"gcloud-deploy": "gcloud functions deploy notification --gen2 --entry-point=notification --trigger-http --allow-unauthenticated --region=us-west2 --runtime nodejs14",
"gcloud-deploy": "gcloud functions deploy rule-handler --gen2 --entry-point=ruleHandler --trigger-http --allow-unauthenticated --region=us-west2 --runtime nodejs14",
"deploy": "yarn build && yarn gcloud-deploy"
},
"devDependencies": {
@ -23,9 +23,10 @@
},
"dependencies": {
"@google-cloud/functions-framework": "3.1.2",
"@google-cloud/pubsub": "^3.2.1",
"dotenv": "^16.0.1",
"firebase-admin": "^10.0.2",
"@sentry/serverless": "^6.16.1"
"@sentry/serverless": "^6.16.1",
"typeorm": "^0.3.4",
"typeorm-naming-strategies": "^4.1.0"
}
}

View file

@ -0,0 +1,30 @@
import { DataSource, EntityTarget, Repository } from 'typeorm'
import { SnakeNamingStrategy } from 'typeorm-naming-strategies'
import * as dotenv from 'dotenv'
dotenv.config()
const AppDataSource = new DataSource({
type: 'postgres',
host: process.env.PG_HOST,
port: Number(process.env.PG_PORT),
schema: 'omnivore',
username: process.env.PG_USER,
password: process.env.PG_PASSWORD,
database: process.env.PG_DB,
logging: ['query', 'info'],
entities: [__dirname + '/entity/**/*{.js,.ts}'],
namingStrategy: new SnakeNamingStrategy(),
})
export const createDBConnection = async () => {
await AppDataSource.initialize()
}
export const closeDBConnection = async () => {
await AppDataSource.destroy()
}
export const getRepository = <T>(entity: EntityTarget<T>): Repository<T> => {
return AppDataSource.getRepository(entity)
}

View file

@ -0,0 +1,41 @@
import {
Column,
CreateDateColumn,
Entity,
JoinColumn,
ManyToOne,
PrimaryGeneratedColumn,
UpdateDateColumn,
} from 'typeorm'
import { User } from './user'
@Entity({ name: 'rules' })
export class Rules {
@PrimaryGeneratedColumn('uuid')
id!: string
@ManyToOne(() => User, { onDelete: 'CASCADE' })
@JoinColumn({ name: 'user_id' })
user!: User
@Column('text')
name!: string
@Column('text')
filter!: string
@Column('simple-json')
actions!: { type: string; params: string[] }[]
@Column('text', { nullable: true })
description?: string | null
@Column('boolean', { default: true })
enabled!: boolean
@CreateDateColumn({ default: () => 'CURRENT_TIMESTAMP' })
createdAt!: Date
@UpdateDateColumn({ default: () => 'CURRENT_TIMESTAMP' })
updatedAt!: Date
}

View file

@ -0,0 +1,31 @@
import {
Column,
CreateDateColumn,
Entity,
PrimaryGeneratedColumn,
UpdateDateColumn,
} from 'typeorm'
@Entity()
export class User {
@PrimaryGeneratedColumn('uuid')
id!: string
@Column('text')
name!: string
@Column('text')
email!: string
@Column('text')
sourceUserId!: string
@CreateDateColumn()
createdAt!: Date
@UpdateDateColumn()
updatedAt!: Date
@Column('varchar', { length: 255, nullable: true })
password?: string
}

View file

@ -0,0 +1,25 @@
import {
Column,
CreateDateColumn,
Entity,
JoinColumn,
ManyToOne,
PrimaryGeneratedColumn,
} from 'typeorm'
import { User } from './user'
@Entity({ name: 'user_device_tokens' })
export class UserDeviceToken {
@PrimaryGeneratedColumn('uuid')
id!: string
@Column('text')
token!: string
@ManyToOne(() => User)
@JoinColumn({ name: 'user_id' })
user!: User
@CreateDateColumn()
createdAt!: Date
}

View file

@ -0,0 +1,102 @@
import * as Sentry from '@sentry/serverless'
import express, { Request, Response } from 'express'
import * as dotenv from 'dotenv'
import { closeDBConnection, createDBConnection, getRepository } from './db'
import { Rules } from './entity/rules'
import { triggerActions } from './rule' // see https://github.com/motdotla/dotenv#how-do-i-use-dotenv-with-import
dotenv.config()
interface PubSubRequestMessage {
data: string
publishTime: string
}
interface PubSubRequestBody {
message: PubSubRequestMessage
}
enum EntityType {
PAGE = 'page',
HIGHLIGHT = 'highlight',
LABEL = 'label',
}
const expired = (body: PubSubRequestBody): boolean => {
const now = new Date()
const expiredTime = new Date(body.message.publishTime)
expiredTime.setHours(expiredTime.getHours() + 1)
return now > expiredTime
}
const readPushSubscription = (
req: express.Request
): { message: string | undefined; expired: boolean } => {
console.debug('request query', req.body)
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')
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) }
}
export const ruleHandler = Sentry.GCPFunction.wrapHttpFunction(
async (req: Request, res: Response) => {
const { message: msgStr, expired } = readPushSubscription(req)
if (!msgStr) {
res.status(400).send('Bad Request')
return
}
if (expired) {
console.log('discarding expired message')
res.status(200).send('Expired')
return
}
try {
const data = JSON.parse(msgStr) as { userId: string; type: EntityType }
const { userId, type } = data
if (!userId || !type) {
console.log('No userId or type found in message')
res.status(400).send('Bad Request')
return
}
if (type !== EntityType.PAGE) {
console.log('Not a page update')
res.status(200).send('Not Page')
return
}
await createDBConnection()
const rules = await getRepository(Rules).findBy({
user: { id: userId },
enabled: true,
})
await triggerActions(userId, rules, data)
await closeDBConnection()
res.status(200).send('OK')
} catch (error) {
console.error(error)
res.status(500).send('Internal server error')
}
}
)

View file

@ -0,0 +1,48 @@
import { Rules } from './entity/rules'
import {
getBatchMessages,
sendBatchPushNotifications,
} from './sendNotification'
import { getRepository } from './db'
import { UserDeviceToken } from './entity/user_device_tokens'
enum RuleActionType {
AddLabel = 'ADD_LABEL',
Archive = 'ARCHIVE',
MarkAsRead = 'MARK_AS_READ',
SendNotification = 'SEND_NOTIFICATION',
}
export const triggerActions = async (
userId: string,
rules: Rules[],
data: any
) => {
for (const rule of rules) {
// TODO: filter out rules that don't match the trigger
for (const action of rule.actions) {
switch (action.type) {
case RuleActionType.AddLabel:
case RuleActionType.Archive:
case RuleActionType.MarkAsRead:
continue
case RuleActionType.SendNotification:
await sendNotification(userId, action.params)
}
}
}
}
export const sendNotification = async (userId: string, messages: string[]) => {
const tokens = await getRepository(UserDeviceToken).findBy({
user: { id: userId },
})
const batchMessages = getBatchMessages(
messages,
tokens.map((t) => t.token)
)
return sendBatchPushNotifications(batchMessages)
}

View file

@ -11,6 +11,25 @@ initializeApp({
credential: applicationDefault(),
})
export const getBatchMessages = (
messages: string[],
tokens: string[]
): Message[] => {
const batchMessages: Message[] = []
messages.forEach((message) => {
tokens.forEach((token) => {
batchMessages.push({
token,
notification: {
body: message,
},
})
})
})
return batchMessages
}
export const sendPushNotification = async (
message: Message
): Promise<string | undefined> => {

View file

@ -1,5 +1,5 @@
{
"extends": "@tsconfig/node14/tsconfig.json",
"extends": "./../../tsconfig.json",
"compilerOptions": {
"outDir": "build",
"rootDir": ".",