upload feed subscriptions in cloud storage in cronjob

This commit is contained in:
Hongbo Wu 2023-07-11 12:22:49 +08:00
parent b88cf7a4e8
commit 44473ba089
5 changed files with 113 additions and 40 deletions

View file

@ -162,6 +162,7 @@ export interface Page {
listenedAt?: Date
wordsCount?: number
recommendations?: Recommendation[]
rssFeedUrl?: string
}
export interface SearchItem {

View file

@ -1,10 +1,12 @@
/* eslint-disable @typescript-eslint/no-misused-promises */
import { stringify } from 'csv-stringify/.'
import express from 'express'
import { DateTime } from 'luxon'
import { readPushSubscription } from '../../datalayer/pubsub'
import { Subscription } from '../../entity/subscription'
import { getRepository } from '../../entity/utils'
import { SubscriptionStatus, SubscriptionType } from '../../generated/graphql'
import { enqueueRssFeedFetch } from '../../utils/createTask'
import { createGCSFile } from '../../utils/uploads'
export function rssFeedRouter() {
const router = express.Router()
@ -20,9 +22,11 @@ export function rssFeedRouter() {
return res.status(200).send('Expired')
}
let writeStream: NodeJS.WritableStream | undefined
try {
// get all active rss feed subscriptions
const subscriptions = await getRepository(Subscription).find({
select: ['id', 'url', 'user'],
where: {
type: SubscriptionType.Rss,
status: SubscriptionStatus.Active,
@ -30,22 +34,33 @@ export function rssFeedRouter() {
relations: ['user'],
})
// create a cloud taks to fetch rss feed item for each subscription
await Promise.all(
subscriptions.map((subscription) => {
try {
return enqueueRssFeedFetch(subscription)
} catch (error) {
console.log('error creating rss feed fetch task', error)
}
})
)
// write the list of subscriptions to a csv file and upload it to gcs
// path style: rss/<date>.csv
const dateStr = DateTime.now().toISODate()
const fullPath = `rss/${dateStr}.csv`
// open a write_stream to the file
const file = createGCSFile(fullPath)
writeStream = file.createWriteStream({
contentType: 'text/csv',
})
// stringify the data and pipe it to the write_stream
const stringifier = stringify({
header: false,
columns: ['subscriptionId', 'userId', 'feedUrl'],
})
stringifier.pipe(writeStream)
res.send('OK')
subscriptions.forEach((sub) => {
stringifier.write([sub.id, sub.user.id, sub.url])
})
} catch (error) {
console.log('error fetching rss feeds', error)
res.status(500).send('Internal Server Error')
return res.status(500).send('Internal Server Error')
} finally {
writeStream?.end()
}
res.send('OK')
})
return router

View file

@ -21,6 +21,7 @@
},
"dependencies": {
"@google-cloud/functions-framework": "3.1.2",
"@google-cloud/tasks": "^3.0.5",
"@sentry/serverless": "^6.16.1",
"axios": "^1.4.0",
"dotenv": "^16.0.1",

View file

@ -148,43 +148,43 @@ export const rssHandler = Sentry.GCPFunction.wrapHttpFunction(
return res.status(400).send('INVALID_REQUEST_BODY')
}
const { userId, feedUrl } = req.body
const { userId, feedUrl, subscriptionId } = req.body
// fetch feed
const feed = await parser.parseURL(feedUrl)
console.log('Fetched feed', feed.title)
const lastFetchedAt = new Date()
console.log('Fetched feed', feed.title, lastFetchedAt)
// save each item in the feed
await Promise.all(
feed.items.map((item) => {
if (!item.link || !item.title || !item.content) {
console.log('Invalid feed item', item)
return
}
for (const item of feed.items) {
if (!item.link || !item.title || !item.content) {
console.log('Invalid feed item', item)
continue
}
const input = {
source: 'rss-feeder',
url: item.link,
saveRequestId: '',
labels: [{ name: 'RSS' }],
title: item.title,
originalContent: item.content,
}
const input = {
source: 'rss-feeder',
url: item.link,
saveRequestId: '',
labels: [{ name: 'RSS' }],
title: item.title,
originalContent: item.content,
}
try {
console.log('Saving page', input.title)
// save page
return sendSavePageMutation(userId, input)
} catch (error) {
console.error('Error while saving page', error)
}
})
)
try {
console.log('Saving page', input.title)
// save page
const result = await sendSavePageMutation(userId, input)
console.log('Saved page', result)
} catch (error) {
console.error('Error while saving page', error)
}
}
// update subscription lastFetchedAt
const updatedSubscription = await sendUpdateSubscriptionMutation(
userId,
req.body.subscriptionId,
new Date()
subscriptionId,
lastFetchedAt
)
console.log('Updated subscription', updatedSubscription)

View file

@ -0,0 +1,56 @@
/* eslint-disable @typescript-eslint/restrict-template-expressions */
import { CloudTasksClient, protos } from '@google-cloud/tasks'
const cloudTask = new CloudTasksClient()
export const emailUserUrl = () => {
const envar = process.env.INTERNAL_SVC_ENDPOINT
if (envar) {
return envar + 'api/user/email'
}
throw 'INTERNAL_SVC_ENDPOINT not set'
}
export const CONTENT_FETCH_URL = process.env.CONTENT_FETCH_GCF_URL
export const createCloudTask = async (
taskHandlerUrl: string | undefined,
payload: unknown,
requestHeaders?: Record<string, string>,
queue = 'omnivore-import-queue'
) => {
const location = process.env.GCP_LOCATION
const project = process.env.GCP_PROJECT_ID
if (!project || !location || !queue || !taskHandlerUrl) {
throw `Environment not configured: ${project}, ${location}, ${queue}, ${taskHandlerUrl}`
}
const serviceAccountEmail = `${project}@appspot.gserviceaccount.com`
const parent = cloudTask.queuePath(project, location, queue)
const convertedPayload = JSON.stringify(payload)
const body = Buffer.from(convertedPayload).toString('base64')
const task: protos.google.cloud.tasks.v2.ITask = {
httpRequest: {
httpMethod: 'POST',
url: taskHandlerUrl,
headers: {
'Content-Type': 'application/json',
...requestHeaders,
},
body,
...(serviceAccountEmail
? {
oidcToken: {
serviceAccountEmail,
},
}
: null),
},
}
return cloudTask.createTask({ parent, task }).then((result) => {
return result[0].name ?? undefined
})
}