skip old item

This commit is contained in:
Hongbo Wu 2023-07-11 12:59:55 +08:00
parent 44473ba089
commit eb9a3eddd0
4 changed files with 30 additions and 90 deletions

View file

@ -1,12 +1,10 @@
/* 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 { createGCSFile } from '../../utils/uploads'
import { enqueueRssFeedFetch } from '../../utils/createTask'
export function rssFeedRouter() {
const router = express.Router()
@ -22,7 +20,6 @@ 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({
@ -34,33 +31,22 @@ export function rssFeedRouter() {
relations: ['user'],
})
// 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)
// 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)
}
})
)
subscriptions.forEach((sub) => {
stringifier.write([sub.id, sub.user.id, sub.url])
})
res.send('OK')
} catch (error) {
console.log('error fetching rss feeds', error)
return res.status(500).send('Internal Server Error')
} finally {
writeStream?.end()
res.status(500).send('Internal Server Error')
}
res.send('OK')
})
return router

View file

@ -21,7 +21,6 @@
},
"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

@ -9,10 +9,16 @@ interface RssFeedRequest {
subscriptionId: string
userId: string
feedUrl: string
lastFetchedAt: Date
}
function isRssFeedRequest(body: any): body is RssFeedRequest {
return 'subscriptionId' in body && 'userId' in body && 'feedUrl' in body
return (
'subscriptionId' in body &&
'userId' in body &&
'feedUrl' in body &&
'lastFetchedAt' in body
)
}
const sendSavePageMutation = async (userId: string, input: unknown) => {
@ -148,19 +154,24 @@ export const rssHandler = Sentry.GCPFunction.wrapHttpFunction(
return res.status(400).send('INVALID_REQUEST_BODY')
}
const { userId, feedUrl, subscriptionId } = req.body
const { userId, feedUrl, subscriptionId, lastFetchedAt } = req.body
// fetch feed
const feed = await parser.parseURL(feedUrl)
const lastFetchedAt = new Date()
console.log('Fetched feed', feed.title, lastFetchedAt)
const newFetchedAt = new Date()
console.log('Fetched feed', feed.title, newFetchedAt)
// save each item in the feed
for (const item of feed.items) {
if (!item.link || !item.title || !item.content) {
if (!item.link || !item.title || !item.content || !item.isoDate) {
console.log('Invalid feed item', item)
continue
}
if (new Date(item.isoDate) <= lastFetchedAt) {
console.log('Skipping old feed item', item.title)
continue
}
const input = {
source: 'rss-feeder',
url: item.link,
@ -184,7 +195,7 @@ export const rssHandler = Sentry.GCPFunction.wrapHttpFunction(
const updatedSubscription = await sendUpdateSubscriptionMutation(
userId,
subscriptionId,
lastFetchedAt
newFetchedAt
)
console.log('Updated subscription', updatedSubscription)

View file

@ -1,56 +0,0 @@
/* 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
})
}