From e9f22bba43da48bcd5e3ce172e303573c1329797 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Fri, 19 Jan 2024 19:26:36 +0800 Subject: [PATCH 01/29] Some debug --- packages/api/src/redis_data_source.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/packages/api/src/redis_data_source.ts b/packages/api/src/redis_data_source.ts index ef7cf55e3..f6bedf6c0 100644 --- a/packages/api/src/redis_data_source.ts +++ b/packages/api/src/redis_data_source.ts @@ -22,6 +22,8 @@ export class RedisDataSource { async initialize(): Promise { if (this.isInitialized) throw 'Error already initialized' + console.trace('initializing with options: ', { options: this.options }) + this.redisClient = createIORedisClient('app', this.options) this.workerRedisClient = createIORedisClient('worker', this.options) this.isInitialized = true From a9eeda9369d0e3acc2327489b43f486be5a0f216 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 09:23:45 +0800 Subject: [PATCH 02/29] Remove non-useful logging line --- packages/api/src/jobs/rss/refreshFeed.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index a6e46f6bb..cd6eb6e89 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -530,7 +530,6 @@ const processSubscription = async ( continue } - console.log('Fetching feed item', link) const feedItem = { ...item, isoDate, From 2961f69f02549df4dc277c56d3fbdfa64da878f1 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 09:47:46 +0800 Subject: [PATCH 03/29] Add a refresh context to make debugging rss jobs easier --- packages/api/src/jobs/rss/refreshAllFeeds.ts | 25 +++++++++++- packages/api/src/utils/createTask.ts | 41 +++----------------- 2 files changed, 29 insertions(+), 37 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshAllFeeds.ts b/packages/api/src/jobs/rss/refreshAllFeeds.ts index 43f4f38a3..05a98ca9e 100644 --- a/packages/api/src/jobs/rss/refreshAllFeeds.ts +++ b/packages/api/src/jobs/rss/refreshAllFeeds.ts @@ -5,8 +5,20 @@ import { redisDataSource } from '../../redis_data_source' import { RssSubscriptionGroup } from '../../utils/createTask' import { stringToHash } from '../../utils/helpers' import { validateUrl } from '../../services/create_page_save_request' +import { v4 as uuid } from 'uuid' + +type RSSRefreshContext = { + type: 'all' | 'user-added' + refreshID: string + startedAt: string +} export const refreshAllFeeds = async (db: DataSource): Promise => { + const refreshContext = { + type: 'all', + refreshID: uuid(), + startedAt: new Date().toISOString(), + } as RSSRefreshContext const subscriptionGroups = (await db.createEntityManager().query( ` SELECT @@ -30,19 +42,27 @@ export const refreshAllFeeds = async (db: DataSource): Promise => { ['RSS', 'ACTIVE', 'following'] )) as RssSubscriptionGroup[] + console.log(`rss: checking ${subscriptionGroups.length}`, { refreshContext }) + for (const group of subscriptionGroups) { try { - await updateSubscriptionGroup(group) + await updateSubscriptionGroup(group, refreshContext) } catch (err) { // we don't want to fail the whole job if one subscription group fails console.error('error updating subscription group') } } + console.log(`rss: finished queuing subscription groups at ${new Date()}`, { + refreshContext, + }) return true } -const updateSubscriptionGroup = async (group: RssSubscriptionGroup) => { +const updateSubscriptionGroup = async ( + group: RssSubscriptionGroup, + refreshContext: RSSRefreshContext +) => { let feedURL = group.url const userList = JSON.stringify(group.userIds.sort()) if (!feedURL) { @@ -63,6 +83,7 @@ const updateSubscriptionGroup = async (group: RssSubscriptionGroup) => { userList )}` const payload = { + refreshContext, subscriptionIds: group.subscriptionIds, feedUrl: group.url, lastFetchedTimestamps: group.fetchedDates.map( diff --git a/packages/api/src/utils/createTask.ts b/packages/api/src/utils/createTask.ts index 3d3b9f3ce..410fe7fc7 100644 --- a/packages/api/src/utils/createTask.ts +++ b/packages/api/src/utils/createTask.ts @@ -22,6 +22,7 @@ import View = google.cloud.tasks.v2.Task.View import { stringToHash } from './helpers' import { queueRSSRefreshFeedJob } from '../jobs/rss/refreshAllFeeds' import { redisDataSource } from '../redis_data_source' +import { v4 as uuid } from 'uuid' // Instantiates a client. const client = new CloudTasksClient() @@ -639,8 +640,12 @@ export interface RssSubscriptionGroup { export const enqueueRssFeedFetch = async ( subscriptionGroup: RssSubscriptionGroup ): Promise => { - const { GOOGLE_CLOUD_PROJECT, PUBSUB_VERIFICATION_TOKEN } = process.env const payload = { + refreshContext: { + type: 'user-added', + refreshID: uuid(), + startedAt: new Date().toISOString(), + }, subscriptionIds: subscriptionGroup.subscriptionIds, feedUrl: subscriptionGroup.url, lastFetchedTimestamps: subscriptionGroup.fetchedDates.map( @@ -670,40 +675,6 @@ export const enqueueRssFeedFetch = async ( } else { throw 'unable to queue rss-refresh-feed-job, redis is not configured' } - - // // If there is no Google Cloud Project Id exposed, it means that we are in local environment - // if (env.dev.isLocal || !GOOGLE_CLOUD_PROJECT) { - // if (env.queue.rssFeedTaskHandlerUrl) { - // // Calling the handler function directly. - // setTimeout(() => { - // axios - // .post( - // `${env.queue.rssFeedTaskHandlerUrl}?token=${PUBSUB_VERIFICATION_TOKEN}`, - // payload - // ) - // .catch((error) => { - // logError(error) - // }) - // }, 0) - // } - // return nanoid() - // } - - // const createdTasks = await createHttpTaskWithToken({ - // project: GOOGLE_CLOUD_PROJECT, - // queue: 'omnivore-rss-queue', - // payload, - // taskHandlerUrl: `${env.queue.rssFeedTaskHandlerUrl}?token=${PUBSUB_VERIFICATION_TOKEN}`, - // }) - - // if (!createdTasks || !createdTasks[0].name) { - // logger.error(`Unable to get the name of the task`, { - // payload, - // createdTasks, - // }) - // throw new CreateTaskError(`Unable to get the name of the task`) - // } - //return createdTasks[0].name } export default createHttpTaskWithToken From 4c1a182c7f03691686d8372bdc6a47517db399ed Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 09:50:21 +0800 Subject: [PATCH 04/29] Add a refresh context to make debugging refresh runs easier --- packages/api/src/jobs/rss/refreshAllFeeds.ts | 2 +- packages/api/src/jobs/rss/refreshFeed.ts | 5 ++++- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshAllFeeds.ts b/packages/api/src/jobs/rss/refreshAllFeeds.ts index 05a98ca9e..a18b7cafb 100644 --- a/packages/api/src/jobs/rss/refreshAllFeeds.ts +++ b/packages/api/src/jobs/rss/refreshAllFeeds.ts @@ -7,7 +7,7 @@ import { stringToHash } from '../../utils/helpers' import { validateUrl } from '../../services/create_page_save_request' import { v4 as uuid } from 'uuid' -type RSSRefreshContext = { +export type RSSRefreshContext = { type: 'all' | 'user-added' refreshID: string startedAt: string diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index cd6eb6e89..702764921 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -7,6 +7,7 @@ import { promisify } from 'util' import { env } from '../../env' import { redisDataSource } from '../../redis_data_source' import createHttpTaskWithToken from '../../utils/createTask' +import { RSSRefreshContext } from './refreshAllFeeds' type FolderType = 'following' | 'inbox' @@ -19,6 +20,7 @@ interface RefreshFeedRequest { userIds: string[] fetchContents: boolean[] folders: FolderType[] + refreshContext?: RSSRefreshContext } export const isRefreshFeedRequest = (data: any): data is RefreshFeedRequest => { @@ -638,8 +640,9 @@ export const _refreshFeed = async (request: RefreshFeedRequest) => { lastFetchedChecksums, fetchContents, folders, + refreshContext, } = request - console.log('Processing feed', feedUrl) + console.log('Processing feed', feedUrl, { refreshContext: refreshContext }) const isBlocked = await isFeedBlocked(feedUrl) if (isBlocked) { From f4af4593c25345863f0fb40602556acb172e9761 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 10:04:54 +0800 Subject: [PATCH 05/29] Fix debug line --- packages/api/src/jobs/rss/refreshAllFeeds.ts | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshAllFeeds.ts b/packages/api/src/jobs/rss/refreshAllFeeds.ts index a18b7cafb..7a87f4990 100644 --- a/packages/api/src/jobs/rss/refreshAllFeeds.ts +++ b/packages/api/src/jobs/rss/refreshAllFeeds.ts @@ -52,9 +52,13 @@ export const refreshAllFeeds = async (db: DataSource): Promise => { console.error('error updating subscription group') } } - console.log(`rss: finished queuing subscription groups at ${new Date()}`, { - refreshContext, - }) + const finishTime = new Date() + console.log( + `rss: finished queuing subscription groups at ${finishTime.toISOString()}`, + { + refreshContext, + } + ) return true } From 52050728baf70366d47a5dbe6d6a8ea37face245 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 10:29:07 +0800 Subject: [PATCH 06/29] Add more debugging logs --- packages/api/src/jobs/save_page.ts | 1 + packages/content-fetch/src/job.ts | 1 + 2 files changed, 2 insertions(+) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index bd28840be..61e6fb12b 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -353,6 +353,7 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { try { const url = encodeURI(data.url) + console.log(`savePageJob: ${userId} ${url}`) // get the fetch result from cache const { title, content, contentType, readabilityResult } = diff --git a/packages/content-fetch/src/job.ts b/packages/content-fetch/src/job.ts index b03ed9469..3d3946b18 100644 --- a/packages/content-fetch/src/job.ts +++ b/packages/content-fetch/src/job.ts @@ -57,6 +57,7 @@ export const queueSavePageJob = async (savePageJobs: savePageJob[]) => { data: job.data, opts: getOpts(job), })) + console.log('queue save page jobs:', { jobs }) return queue.addBulk(jobs) } From 2b536c6086cf010213d62635d5eced7fdf0355ad Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 10:48:57 +0800 Subject: [PATCH 07/29] Add more debug --- packages/api/src/jobs/rss/refreshFeed.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index 702764921..9e60d4cc5 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -323,6 +323,7 @@ const createTask = async ( return createItemWithPreviewContent(userId, feedUrl, item) } + console.log(`adding fetch content task ${userId} ${item.link.trim()}`) return addFetchContentTask(userId, folder, item) } From 4eeb012b787199efb01e081532296ca5b45bc337 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 11:36:10 +0800 Subject: [PATCH 08/29] Pass lits of fetchContentTasks into rss handlers This prevents the list from growing on each run --- packages/api/src/jobs/rss/refreshFeed.ts | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index 9e60d4cc5..771e4ab9c 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -70,8 +70,6 @@ interface FetchContentTask { item: RssFeedItem } -const fetchContentTasks = new Map() // url -> FetchContentTask - export const isOldItem = (item: RssFeedItem, lastFetchedAt: number) => { // existing items and items that were published before 24h const publishedAt = item.isoDate ? new Date(item.isoDate) : new Date() @@ -288,6 +286,7 @@ const isItemRecentlySaved = async (userId: string, url: string) => { } const addFetchContentTask = ( + fetchContentTasks: Map, userId: string, folder: FolderType, item: RssFeedItem @@ -307,6 +306,7 @@ const addFetchContentTask = ( } const createTask = async ( + fetchContentTasks: Map, userId: string, feedUrl: string, item: RssFeedItem, @@ -324,7 +324,7 @@ const createTask = async ( } console.log(`adding fetch content task ${userId} ${item.link.trim()}`) - return addFetchContentTask(userId, folder, item) + return addFetchContentTask(fetchContentTasks, userId, folder, item) } const fetchContentAndCreateItem = async ( @@ -482,6 +482,7 @@ const getLink = (links: RssFeedItemLink[]): string | undefined => { } const processSubscription = async ( + fetchContentTasks: Map, subscriptionId: string, userId: string, feedUrl: string, @@ -562,6 +563,7 @@ const processSubscription = async ( } const created = await createTask( + fetchContentTasks, userId, feedUrl, feedItem, @@ -591,6 +593,7 @@ const processSubscription = async ( // the feed has never been fetched, save at least the last valid item const created = await createTask( + fetchContentTasks, userId, feedUrl, lastValidItem, @@ -673,9 +676,11 @@ export const _refreshFeed = async (request: RefreshFeedRequest) => { console.log('Fetched feed', feed.title, new Date()) + const fetchContentTasks = new Map() // url -> FetchContentTask // process each subscription sequentially for (let i = 0; i < subscriptionIds.length; i++) { await processSubscription( + fetchContentTasks, subscriptionIds[i], userIds[i], feedUrl, From e5dec228f347d1602ee7d77aa4ad2730b1aac9d5 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 12:46:22 +0800 Subject: [PATCH 09/29] Remove some debug --- packages/api/src/jobs/rss/refreshFeed.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index 771e4ab9c..deda89220 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -521,7 +521,7 @@ const processSubscription = async ( // use published or updated if isoDate is not available for atom feeds const isoDate = item.isoDate || item.published || item.updated || item.created - console.log('Processing feed item', item.links, item.isoDate, feed.feedUrl) + console.log('Processing feed item', item.links, item.isoDate, feedUrl) if (!item.links || item.links.length === 0) { console.log('Invalid feed item', item) From 1898e326079e7c3e2b7ca52fd752e32aa0b0ee8d Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 13:27:54 +0800 Subject: [PATCH 10/29] Use service instead of API to update subscription --- packages/api/src/jobs/rss/refreshFeed.ts | 14 +++--- .../api/src/resolvers/subscriptions/index.ts | 30 +----------- .../api/src/services/update_subscription.ts | 47 +++++++++++++++++++ 3 files changed, 56 insertions(+), 35 deletions(-) create mode 100644 packages/api/src/services/update_subscription.ts diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index deda89220..dba8ed535 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -8,6 +8,7 @@ import { env } from '../../env' import { redisDataSource } from '../../redis_data_source' import createHttpTaskWithToken from '../../utils/createTask' import { RSSRefreshContext } from './refreshAllFeeds' +import { updateSubscription } from '../../services/update_subscription' type FolderType = 'following' | 'inbox' @@ -615,13 +616,12 @@ const processSubscription = async ( const nextScheduledAt = scheduledAt + updatePeriodInMs * updateFrequency // update subscription lastFetchedAt - const updatedSubscription = await sendUpdateSubscriptionMutation( - userId, - subscriptionId, - lastItemFetchedAt, - updatedLastFetchedChecksum, - new Date(nextScheduledAt) - ) + const updatedSubscription = await updateSubscription(userId, subscriptionId, { + id: subscriptionId, + lastFetchedAt: lastItemFetchedAt, + lastFetchedChecksum: updatedLastFetchedChecksum, + scheduledAt: new Date(nextScheduledAt), + }) console.log('Updated subscription', updatedSubscription) } diff --git a/packages/api/src/resolvers/subscriptions/index.ts b/packages/api/src/resolvers/subscriptions/index.ts index ab785b73a..2b0e0b12f 100644 --- a/packages/api/src/resolvers/subscriptions/index.ts +++ b/packages/api/src/resolvers/subscriptions/index.ts @@ -50,6 +50,7 @@ import { keysToCamelCase, } from '../../utils/helpers' import { parseFeed, parseOpml, RSS_PARSER_CONFIG } from '../../utils/parser' +import { updateSubscription } from '../../services/update_subscription' type PartialSubscription = Omit @@ -332,34 +333,7 @@ export const updateSubscriptionResolver = authorized< }, }) - const updatedSubscription = await authTrx(async (t) => { - const repo = t.getRepository(Subscription) - - // update subscription - await t.getRepository(Subscription).save({ - id: input.id, - name: input.name || undefined, - description: input.description || undefined, - lastFetchedAt: input.lastFetchedAt - ? new Date(input.lastFetchedAt) - : undefined, - lastFetchedChecksum: input.lastFetchedChecksum || undefined, - status: input.status || undefined, - scheduledAt: input.scheduledAt - ? new Date(input.scheduledAt) - : undefined, - autoAddToLibrary: input.autoAddToLibrary ?? undefined, - isPrivate: input.isPrivate ?? undefined, - fetchContent: input.fetchContent ?? undefined, - folder: input.folder ?? undefined, - }) - - return repo.findOneByOrFail({ - id: input.id, - user: { id: uid }, - }) - }) - + const updatedSubscription = await updateSubscription(uid, input.id, input) return { subscription: updatedSubscription, } diff --git a/packages/api/src/services/update_subscription.ts b/packages/api/src/services/update_subscription.ts new file mode 100644 index 000000000..33f9cdfb0 --- /dev/null +++ b/packages/api/src/services/update_subscription.ts @@ -0,0 +1,47 @@ +import { Subscription } from '../entity/subscription' +import { UpdateSubscriptionInput } from '../generated/graphql' +import { getRepository } from '../repository' + +const ensureOwns = async (userId: string, subscriptionId: string) => { + const repo = getRepository(Subscription) + + const existing = repo.findOneByOrFail({ + id: subscriptionId, + user: { id: userId }, + }) + if (!existing) { + throw new Error('Can not find subscription being updated.') + } +} + +export const updateSubscription = async ( + userId: string, + subscriptionId: string, + newData: UpdateSubscriptionInput +): Promise => { + ensureOwns(userId, subscriptionId) + + const repo = getRepository(Subscription) + await repo.save({ + id: subscriptionId, + name: newData.name || undefined, + description: newData.description || undefined, + lastFetchedAt: newData.lastFetchedAt + ? new Date(newData.lastFetchedAt) + : undefined, + lastFetchedChecksum: newData.lastFetchedChecksum || undefined, + status: newData.status || undefined, + scheduledAt: newData.scheduledAt + ? new Date(newData.scheduledAt) + : undefined, + autoAddToLibrary: newData.autoAddToLibrary ?? undefined, + isPrivate: newData.isPrivate ?? undefined, + fetchContent: newData.fetchContent ?? undefined, + folder: newData.folder ?? undefined, + }) + + return (await getRepository(Subscription).findOneByOrFail({ + id: subscriptionId, + user: { id: userId }, + })) as Subscription +} From 89537c13deda7fc97f1b3d615c5cf4580a372f8f Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 13:32:11 +0800 Subject: [PATCH 11/29] Simplify function input --- packages/api/src/jobs/rss/refreshFeed.ts | 1 - .../api/src/services/update_subscription.ts | 20 +++++++++++++++++-- 2 files changed, 18 insertions(+), 3 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index dba8ed535..ccefcd0d5 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -617,7 +617,6 @@ const processSubscription = async ( // update subscription lastFetchedAt const updatedSubscription = await updateSubscription(userId, subscriptionId, { - id: subscriptionId, lastFetchedAt: lastItemFetchedAt, lastFetchedChecksum: updatedLastFetchedChecksum, scheduledAt: new Date(nextScheduledAt), diff --git a/packages/api/src/services/update_subscription.ts b/packages/api/src/services/update_subscription.ts index 33f9cdfb0..70da2cfe0 100644 --- a/packages/api/src/services/update_subscription.ts +++ b/packages/api/src/services/update_subscription.ts @@ -1,5 +1,8 @@ import { Subscription } from '../entity/subscription' -import { UpdateSubscriptionInput } from '../generated/graphql' +import { + SubscriptionStatus, + UpdateSubscriptionInput, +} from '../generated/graphql' import { getRepository } from '../repository' const ensureOwns = async (userId: string, subscriptionId: string) => { @@ -14,10 +17,23 @@ const ensureOwns = async (userId: string, subscriptionId: string) => { } } +type UpdateSubscriptionData = { + autoAddToLibrary?: boolean | null + description?: string | null + fetchContent?: boolean | null + folder?: string | null + isPrivate?: boolean | null + lastFetchedAt?: Date | null + lastFetchedChecksum?: string | null + name?: string | null + scheduledAt?: Date | null + status?: SubscriptionStatus | null +} + export const updateSubscription = async ( userId: string, subscriptionId: string, - newData: UpdateSubscriptionInput + newData: UpdateSubscriptionData ): Promise => { ensureOwns(userId, subscriptionId) From 7786221fbc449daebc942da13206555d2b546b4f Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 13:38:25 +0800 Subject: [PATCH 12/29] Fix update_subscription awaits --- packages/api/src/services/update_subscription.ts | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/packages/api/src/services/update_subscription.ts b/packages/api/src/services/update_subscription.ts index 70da2cfe0..1106e3043 100644 --- a/packages/api/src/services/update_subscription.ts +++ b/packages/api/src/services/update_subscription.ts @@ -8,7 +8,7 @@ import { getRepository } from '../repository' const ensureOwns = async (userId: string, subscriptionId: string) => { const repo = getRepository(Subscription) - const existing = repo.findOneByOrFail({ + const existing = await repo.findOneByOrFail({ id: subscriptionId, user: { id: userId }, }) @@ -35,7 +35,7 @@ export const updateSubscription = async ( subscriptionId: string, newData: UpdateSubscriptionData ): Promise => { - ensureOwns(userId, subscriptionId) + await ensureOwns(userId, subscriptionId) const repo = getRepository(Subscription) await repo.save({ @@ -56,8 +56,8 @@ export const updateSubscription = async ( folder: newData.folder ?? undefined, }) - return (await getRepository(Subscription).findOneByOrFail({ + return await getRepository(Subscription).findOneByOrFail({ id: subscriptionId, user: { id: userId }, - })) as Subscription + }) } From 3a183272305f46753cdc0450f4407ae0f288d8e1 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 13:43:17 +0800 Subject: [PATCH 13/29] Remove unused --- packages/api/src/jobs/rss/refreshFeed.ts | 69 ------------------------ 1 file changed, 69 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshFeed.ts b/packages/api/src/jobs/rss/refreshFeed.ts index ccefcd0d5..f9339d6cc 100644 --- a/packages/api/src/jobs/rss/refreshFeed.ts +++ b/packages/api/src/jobs/rss/refreshFeed.ts @@ -205,75 +205,6 @@ const parseFeed = async (url: string, content: string) => { } } -const sendUpdateSubscriptionMutation = async ( - userId: string, - subscriptionId: string, - lastFetchedAt: Date, - lastFetchedChecksum: string, - scheduledAt: Date -) => { - if (!process.env.INTERNAL_API_URL || !env.server.jwtSecret) { - throw new Error( - 'Can not send update subscription, environment not configured.' - ) - } - const JWT_SECRET = env.server.jwtSecret - const REST_BACKEND_ENDPOINT = `${process.env.INTERNAL_API_URL}/api` - - if (!JWT_SECRET || !REST_BACKEND_ENDPOINT) { - throw 'Environment not configured correctly' - } - - const data = JSON.stringify({ - query: `mutation UpdateSubscription($input: UpdateSubscriptionInput!){ - updateSubscription(input:$input){ - ... on UpdateSubscriptionSuccess{ - subscription{ - id - lastFetchedAt - } - } - ... on UpdateSubscriptionError{ - errorCodes - } - } - }`, - variables: { - input: { - id: subscriptionId, - lastFetchedAt, - lastFetchedChecksum, - scheduledAt, - }, - }, - }) - - const auth = (await signToken({ uid: userId }, JWT_SECRET)) as string - try { - const response = await axios.post( - `${REST_BACKEND_ENDPOINT}/graphql`, - data, - { - headers: { - Cookie: `auth=${auth};`, - 'Content-Type': 'application/json', - }, - timeout: 30000, // 30s - } - ) - - /* eslint-disable @typescript-eslint/no-unsafe-member-access */ - return !!response.data.data.updateSubscription.subscription - } catch (error) { - if (axios.isAxiosError(error)) { - console.error('update subscription mutation error', error.message) - } else { - console.error(error) - } - return false - } -} - const isItemRecentlySaved = async (userId: string, url: string) => { const key = `recent-saved-item:${userId}:${url}` try { From 2f405ebed0b44ed4606646c7bf37448b44746cb8 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 14:07:22 +0800 Subject: [PATCH 14/29] Use savePage service instead of API call --- packages/api/src/jobs/save_page.ts | 63 ++++++++++++++++++++---------- 1 file changed, 43 insertions(+), 20 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index 61e6fb12b..2fa2d0daf 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -3,6 +3,10 @@ import jwt from 'jsonwebtoken' import { promisify } from 'util' import { env } from '../env' import { redisDataSource } from '../redis_data_source' +import { savePage } from '../services/save_page' +import { userRepository } from '../repository/user' +import { ArticleSavingRequestStatus, ParseResult } from '../generated/graphql' +import { logger } from '../utils/logger' const signToken = promisify(jwt.sign) @@ -359,6 +363,12 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { const { title, content, contentType, readabilityResult } = await getCachedFetchResult(url) + if (!title || !content || !contentType || !readabilityResult) { + throw new Error( + 'Invalid SavePage job, fetch result missing required data' + ) + } + // for pdf content, we need to upload the pdf if (contentType === 'application/pdf') { const uploadFileId = await uploadPdf(url, userId, articleSavingRequestId) @@ -383,28 +393,41 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { return true } - // for non-pdf content, we need to save the page - const apiResponse = await sendSavePageMutation(userId, { - url, - clientRequestId: articleSavingRequestId, - title, - originalContent: content, - parseResult: readabilityResult, - state, - labels, - rssFeedUrl, - savedAt, - publishedAt, - source, - folder, - }) - if (!apiResponse) { - throw new Error('error while saving page') + const user = await userRepository.findById(userId) + if (!user) { + logger.error('Unable to save job, user can not be found.', { + userId, + url, + }) + throw new Error('Unable to save job, user can not be found.') } - if ('error' in apiResponse && apiResponse.error === 'UNAUTHORIZED') { - console.log('user is deleted', userId) - return false + // for non-pdf content, we need to save the page + const result = await savePage( + { + url, + clientRequestId: articleSavingRequestId, + title, + originalContent: content, + parseResult: readabilityResult as ParseResult, + state: state as ArticleSavingRequestStatus, + labels: labels?.map((name) => { + return { + name, + } + }), + rssFeedUrl, + savedAt, + publishedAt, + source, + folder, + }, + user + ) + + if (result.__typename == 'SaveError') { + logger.error('Error saving page', { userId, url, result }) + throw new Error('Error saving page') } // if the readability result is not parsed, the import is failed From a4e707075fbd44713093a38bf12d558f954fb535 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sat, 20 Jan 2024 21:07:02 +0800 Subject: [PATCH 15/29] Refactor savePage a bit to use it as a service from the jobs --- packages/api/src/jobs/save_page.ts | 76 +++---------------- packages/api/src/resolvers/article/index.ts | 2 +- .../resolvers/article_saving_request/index.ts | 3 +- packages/api/src/resolvers/highlight/index.ts | 3 +- .../src/resolvers/recommendations/index.ts | 3 +- .../api/src/resolvers/subscriptions/index.ts | 7 +- packages/api/src/resolvers/update/index.ts | 3 +- .../api/src/resolvers/upload_files/index.ts | 4 +- packages/api/src/resolvers/user/index.ts | 3 +- packages/api/src/services/save_page.ts | 37 ++++++--- 10 files changed, 53 insertions(+), 88 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index 2fa2d0daf..cdeddc969 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -3,10 +3,10 @@ import jwt from 'jsonwebtoken' import { promisify } from 'util' import { env } from '../env' import { redisDataSource } from '../redis_data_source' -import { savePage } from '../services/save_page' +import { savePage, stringToRequestStatus } from '../services/save_page' import { userRepository } from '../repository/user' -import { ArticleSavingRequestStatus, ParseResult } from '../generated/graphql' import { logger } from '../utils/logger' +import { Readability } from '@omnivore/readability' const signToken = promisify(jwt.sign) @@ -69,7 +69,7 @@ interface FetchResult { title: string content?: string contentType?: string - readabilityResult?: unknown + readabilityResult?: Readability.ParseResult } const isFetchResult = (obj: unknown): obj is FetchResult => { @@ -238,60 +238,6 @@ const sendCreateArticleMutation = async (userId: string, input: unknown) => { } } -const sendSavePageMutation = async (userId: string, input: unknown) => { - const data = JSON.stringify({ - query: `mutation SavePage ($input: SavePageInput!){ - savePage(input:$input){ - ... on SaveSuccess{ - url - clientRequestId - } - ... on SaveError{ - errorCodes - } - } - }`, - variables: { - input, - }, - }) - - const auth = await signToken({ uid: userId }, JWT_SECRET) - try { - const response = await axios.post( - `${REST_BACKEND_ENDPOINT}/graphql`, - data, - { - headers: { - Cookie: `auth=${auth as string};`, - 'Content-Type': 'application/json', - }, - timeout: REQUEST_TIMEOUT, - } - ) - - if ( - response.data.data.savePage.errorCodes && - response.data.data.savePage.errorCodes.length > 0 - ) { - console.error( - 'error while saving page', - response.data.data.savePage.errorCodes[0] - ) - if (response.data.data.savePage.errorCodes[0] === 'UNAUTHORIZED') { - return { error: 'UNAUTHORIZED' } - } - - return null - } - - return response.data.data.savePage - } catch (error) { - console.error('error saving page', error) - return null - } -} - const sendImportStatusUpdate = async ( userId: string, taskId: string, @@ -409,26 +355,26 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { clientRequestId: articleSavingRequestId, title, originalContent: content, - parseResult: readabilityResult as ParseResult, - state: state as ArticleSavingRequestStatus, + parseResult: readabilityResult, + state: stringToRequestStatus(state), labels: labels?.map((name) => { return { name, } }), rssFeedUrl, - savedAt, - publishedAt, + savedAt: savedAt ? new Date(savedAt) : new Date(), + publishedAt: publishedAt ? new Date(publishedAt) : null, source, folder, }, user ) - if (result.__typename == 'SaveError') { - logger.error('Error saving page', { userId, url, result }) - throw new Error('Error saving page') - } + // if (result.__typename == 'SaveError') { + // logger.error('Error saving page', { userId, url, result }) + // throw new Error('Error saving page') + // } // if the readability result is not parsed, the import is failed isImported = !!readabilityResult diff --git a/packages/api/src/resolvers/article/index.ts b/packages/api/src/resolvers/article/index.ts index 08f0401a1..03cce990e 100644 --- a/packages/api/src/resolvers/article/index.ts +++ b/packages/api/src/resolvers/article/index.ts @@ -93,7 +93,6 @@ import { traceAs } from '../../tracing' import { analytics } from '../../utils/analytics' import { isSiteBlockedForParse } from '../../utils/blocked' import { - authorized, cleanUrl, errorHandler, generateSlug, @@ -103,6 +102,7 @@ import { titleForFilePath, userDataToUser, } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' import { contentConverter, getDistillerResult, diff --git a/packages/api/src/resolvers/article_saving_request/index.ts b/packages/api/src/resolvers/article_saving_request/index.ts index 731de1c72..9f4cf7c4a 100644 --- a/packages/api/src/resolvers/article_saving_request/index.ts +++ b/packages/api/src/resolvers/article_saving_request/index.ts @@ -19,11 +19,12 @@ import { } from '../../services/library_item' import { analytics } from '../../utils/analytics' import { - authorized, cleanUrl, isParsingTimeout, libraryItemToArticleSavingRequest, } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' + import { isErrorWithCode } from '../user' export const createArticleSavingRequestResolver = authorized< diff --git a/packages/api/src/resolvers/highlight/index.ts b/packages/api/src/resolvers/highlight/index.ts index 82618990c..9853ef7e4 100644 --- a/packages/api/src/resolvers/highlight/index.ts +++ b/packages/api/src/resolvers/highlight/index.ts @@ -34,7 +34,8 @@ import { updateHighlight, } from '../../services/highlights' import { analytics } from '../../utils/analytics' -import { authorized, highlightDataToHighlight } from '../../utils/helpers' +import { highlightDataToHighlight } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const createHighlightResolver = authorized< CreateHighlightSuccess, diff --git a/packages/api/src/resolvers/recommendations/index.ts b/packages/api/src/resolvers/recommendations/index.ts index 31f7da898..a23a573e2 100644 --- a/packages/api/src/resolvers/recommendations/index.ts +++ b/packages/api/src/resolvers/recommendations/index.ts @@ -40,7 +40,8 @@ import { import { findLibraryItemById } from '../../services/library_item' import { analytics } from '../../utils/analytics' import { enqueueRecommendation } from '../../utils/createTask' -import { authorized, userDataToUser } from '../../utils/helpers' +import { userDataToUser } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const createGroupResolver = authorized< CreateGroupSuccess, diff --git a/packages/api/src/resolvers/subscriptions/index.ts b/packages/api/src/resolvers/subscriptions/index.ts index 2b0e0b12f..b7e1507e4 100644 --- a/packages/api/src/resolvers/subscriptions/index.ts +++ b/packages/api/src/resolvers/subscriptions/index.ts @@ -44,11 +44,8 @@ import { unsubscribe } from '../../services/subscriptions' import { Merge } from '../../util' import { analytics } from '../../utils/analytics' import { enqueueRssFeedFetch } from '../../utils/createTask' -import { - authorized, - getAbsoluteUrl, - keysToCamelCase, -} from '../../utils/helpers' +import { getAbsoluteUrl, keysToCamelCase } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' import { parseFeed, parseOpml, RSS_PARSER_CONFIG } from '../../utils/parser' import { updateSubscription } from '../../services/update_subscription' diff --git a/packages/api/src/resolvers/update/index.ts b/packages/api/src/resolvers/update/index.ts index 255efd2b5..c38019d8b 100644 --- a/packages/api/src/resolvers/update/index.ts +++ b/packages/api/src/resolvers/update/index.ts @@ -5,7 +5,8 @@ import { UpdatePageSuccess, } from '../../generated/graphql' import { updateLibraryItem } from '../../services/library_item' -import { authorized, libraryItemToArticle } from '../../utils/helpers' +import { libraryItemToArticle } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const updatePageResolver = authorized< UpdatePageSuccess, diff --git a/packages/api/src/resolvers/upload_files/index.ts b/packages/api/src/resolvers/upload_files/index.ts index 9ea214584..b22395458 100644 --- a/packages/api/src/resolvers/upload_files/index.ts +++ b/packages/api/src/resolvers/upload_files/index.ts @@ -19,7 +19,9 @@ import { updateLibraryItem, } from '../../services/library_item' import { analytics } from '../../utils/analytics' -import { authorized, generateSlug } from '../../utils/helpers' +import { generateSlug } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' + import { contentReaderForLibraryItem, generateUploadFilePathName, diff --git a/packages/api/src/resolvers/user/index.ts b/packages/api/src/resolvers/user/index.ts index c15000a66..8303d41ea 100644 --- a/packages/api/src/resolvers/user/index.ts +++ b/packages/api/src/resolvers/user/index.ts @@ -43,9 +43,10 @@ import { userRepository } from '../../repository/user' import { createUser } from '../../services/create_user' import { sendVerificationEmail } from '../../services/send_emails' import { softDeleteUser } from '../../services/user' -import { authorized, userDataToUser } from '../../utils/helpers' +import { userDataToUser } from '../../utils/helpers' import { validateUsername } from '../../utils/usernamePolicy' import { WithDataSourcesContext } from '../types' +import { authorized } from '../../utils/gql-utils' export const updateUserResolver = authorized< UpdateUserSuccess, diff --git a/packages/api/src/services/save_page.ts b/packages/api/src/services/save_page.ts index 85f341df2..b03608dcd 100644 --- a/packages/api/src/services/save_page.ts +++ b/packages/api/src/services/save_page.ts @@ -5,14 +5,12 @@ import { Highlight } from '../entity/highlight' import { LibraryItem, LibraryItemState } from '../entity/library_item' import { User } from '../entity/user' import { homePageURL } from '../env' -import { - ArticleSavingRequestStatus, - Maybe, - PreparedDocumentInput, - SaveErrorCode, - SavePageInput, - SaveResult, -} from '../generated/graphql' +// import { +// ArticleSavingRequestStatus, +// PreparedDocumentInput, +// SaveErrorCode, +// SaveResult, +// } from '../generated/graphql' import { authTrx } from '../repository' import { enqueueThumbnailTask } from '../utils/createTask' import { @@ -30,6 +28,14 @@ import { createPageSaveRequest } from './create_page_save_request' import { createHighlight } from './highlights' import { createAndSaveLabelsInLibraryItem } from './labels' import { createLibraryItem, updateLibraryItem } from './library_item' +import { CreateLabelInput } from '../repository/label' +import { + ArticleSavingRequestStatus, + PreparedDocumentInput, + SaveErrorCode, + SavePageInput, + SaveResult, +} from '../generated/graphql' // where we can use APIs to fetch their underlying content. const FORCE_PUPPETEER_URLS = [ @@ -43,7 +49,7 @@ const ALREADY_PARSED_SOURCES = [ 'pocket', ] -const createSlug = (url: string, title?: Maybe | undefined) => { +const createSlug = (url: string, title?: string | null | undefined) => { const { pathname } = new URL(url) const croppedPathname = decodeURIComponent( pathname @@ -63,6 +69,15 @@ const shouldParseInBackend = (input: SavePageInput): boolean => { ) } +export const stringToRequestStatus = ( + str: string | undefined | null +): ArticleSavingRequestStatus | null => { + if (str && str in ArticleSavingRequestStatus) { + return str as ArticleSavingRequestStatus + } + return null +} + export const savePage = async ( input: SavePageInput, user: User @@ -95,7 +110,7 @@ export const savePage = async ( canonicalUrl: parseResult.canonicalUrl, savedAt: input.savedAt ? new Date(input.savedAt) : new Date(), publishedAt: input.publishedAt ? new Date(input.publishedAt) : undefined, - state: input.state || undefined, + state: stringToRequestStatus(input.state) || undefined, rssFeedUrl: input.rssFeedUrl, folder: input.folder, }) @@ -109,7 +124,7 @@ export const savePage = async ( userId: user.id, url: itemToSave.originalUrl, articleSavingRequestId: clientRequestId || undefined, - state: input.state || undefined, + state: stringToRequestStatus(input.state) || undefined, labels: input.labels || undefined, folder: input.folder || undefined, }) From db3d7dca3676d61dc1347d1543242351cb9f7349 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 07:43:50 +0800 Subject: [PATCH 16/29] Add missing file --- packages/api/src/utils/gql-utils.ts | 26 ++++++++++++++++++++++++++ 1 file changed, 26 insertions(+) create mode 100644 packages/api/src/utils/gql-utils.ts diff --git a/packages/api/src/utils/gql-utils.ts b/packages/api/src/utils/gql-utils.ts new file mode 100644 index 000000000..ef5986f28 --- /dev/null +++ b/packages/api/src/utils/gql-utils.ts @@ -0,0 +1,26 @@ +import { ResolverFn } from '../generated/graphql' +import { Claims, WithDataSourcesContext } from '../resolvers/types' + +export function authorized< + TSuccess, + TError extends { errorCodes: string[] }, + /* eslint-disable @typescript-eslint/no-explicit-any */ + TArgs = any, + TParent = any + /* eslint-enable @typescript-eslint/no-explicit-any */ +>( + resolver: ResolverFn< + TSuccess | TError, + TParent, + WithDataSourcesContext & { claims: Claims }, + TArgs + > +): ResolverFn { + return (parent, args, ctx, info) => { + const { claims } = ctx + if (claims?.uid) { + return resolver(parent, args, { ...ctx, claims, uid: claims.uid }, info) + } + return { errorCodes: ['UNAUTHORIZED'] } as TError + } +} From 9abd0a59daa8488d74d9a8fa1a34d73f58d968a1 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 13:16:35 +0800 Subject: [PATCH 17/29] Clean up imports --- packages/api/src/services/save_page.ts | 21 +++++++-------------- 1 file changed, 7 insertions(+), 14 deletions(-) diff --git a/packages/api/src/services/save_page.ts b/packages/api/src/services/save_page.ts index b03608dcd..d49fc8b57 100644 --- a/packages/api/src/services/save_page.ts +++ b/packages/api/src/services/save_page.ts @@ -5,12 +5,13 @@ import { Highlight } from '../entity/highlight' import { LibraryItem, LibraryItemState } from '../entity/library_item' import { User } from '../entity/user' import { homePageURL } from '../env' -// import { -// ArticleSavingRequestStatus, -// PreparedDocumentInput, -// SaveErrorCode, -// SaveResult, -// } from '../generated/graphql' +import { + ArticleSavingRequestStatus, + PreparedDocumentInput, + SaveErrorCode, + SavePageInput, + SaveResult, +} from '../generated/graphql' import { authTrx } from '../repository' import { enqueueThumbnailTask } from '../utils/createTask' import { @@ -28,14 +29,6 @@ import { createPageSaveRequest } from './create_page_save_request' import { createHighlight } from './highlights' import { createAndSaveLabelsInLibraryItem } from './labels' import { createLibraryItem, updateLibraryItem } from './library_item' -import { CreateLabelInput } from '../repository/label' -import { - ArticleSavingRequestStatus, - PreparedDocumentInput, - SaveErrorCode, - SavePageInput, - SaveResult, -} from '../generated/graphql' // where we can use APIs to fetch their underlying content. const FORCE_PUPPETEER_URLS = [ From eefceeb1ef3758ad3eaa61b41310ee282b4fcb0d Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 13:19:50 +0800 Subject: [PATCH 18/29] Move checks to have the PDF check --- packages/api/src/jobs/save_page.ts | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index cdeddc969..7eacbf3d5 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -309,12 +309,6 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { const { title, content, contentType, readabilityResult } = await getCachedFetchResult(url) - if (!title || !content || !contentType || !readabilityResult) { - throw new Error( - 'Invalid SavePage job, fetch result missing required data' - ) - } - // for pdf content, we need to upload the pdf if (contentType === 'application/pdf') { const uploadFileId = await uploadPdf(url, userId, articleSavingRequestId) @@ -339,6 +333,12 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { return true } + if (!title || !content || !contentType || !readabilityResult) { + throw new Error( + 'Invalid SavePage job, fetch result missing required data' + ) + } + const user = await userRepository.findById(userId) if (!user) { logger.error('Unable to save job, user can not be found.', { From 56a8ef23b6375f79f2550dae06fa366e5a90238e Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 13:22:14 +0800 Subject: [PATCH 19/29] Remove authorized function from imported helpers This creates some issues with init order when using the queue processor. --- packages/api/src/utils/helpers.ts | 24 ------------------------ 1 file changed, 24 deletions(-) diff --git a/packages/api/src/utils/helpers.ts b/packages/api/src/utils/helpers.ts index 72388667c..5b7ca9b83 100644 --- a/packages/api/src/utils/helpers.ts +++ b/packages/api/src/utils/helpers.ts @@ -78,30 +78,6 @@ export const stringToHash = (str: string, convertToUUID = false): string => { ).toLowerCase() } -export function authorized< - TSuccess, - TError extends { errorCodes: string[] }, - /* eslint-disable @typescript-eslint/no-explicit-any */ - TArgs = any, - TParent = any - /* eslint-enable @typescript-eslint/no-explicit-any */ ->( - resolver: ResolverFn< - TSuccess | TError, - TParent, - WithDataSourcesContext & { claims: Claims }, - TArgs - > -): ResolverFn { - return (parent, args, ctx, info) => { - const { claims } = ctx - if (claims?.uid) { - return resolver(parent, args, { ...ctx, claims, uid: claims.uid }, info) - } - return { errorCodes: ['UNAUTHORIZED'] } as TError - } -} - export const findDelimiter = ( text: string, delimiters = ['\t', ',', ':', ';'], From 5393743c9da1db2693f27ea5b665e7ecabc57fb8 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 13:33:05 +0800 Subject: [PATCH 20/29] Fix imports of authorized function --- packages/api/src/resolvers/api_key/index.ts | 2 +- packages/api/src/resolvers/features/index.ts | 2 +- packages/api/src/resolvers/filters/index.ts | 2 +- .../api/src/resolvers/importers/uploadImportFileResolver.ts | 2 +- packages/api/src/resolvers/integrations/index.ts | 2 +- packages/api/src/resolvers/labels/index.ts | 2 +- packages/api/src/resolvers/links/index.ts | 2 +- packages/api/src/resolvers/newsletters/index.ts | 2 +- packages/api/src/resolvers/popular_reads/index.ts | 2 +- packages/api/src/resolvers/recent_emails/index.ts | 2 +- packages/api/src/resolvers/recent_searches/index.ts | 2 +- packages/api/src/resolvers/rules/index.ts | 2 +- packages/api/src/resolvers/save/index.ts | 2 +- packages/api/src/resolvers/send_install_instructions/index.ts | 2 +- packages/api/src/resolvers/user_device_tokens/index.ts | 2 +- packages/api/src/resolvers/user_personalization/index.ts | 2 +- packages/api/src/resolvers/webhooks/index.ts | 2 +- 17 files changed, 17 insertions(+), 17 deletions(-) diff --git a/packages/api/src/resolvers/api_key/index.ts b/packages/api/src/resolvers/api_key/index.ts index 2188d1ea0..4ca0d4c3e 100644 --- a/packages/api/src/resolvers/api_key/index.ts +++ b/packages/api/src/resolvers/api_key/index.ts @@ -17,7 +17,7 @@ import { getRepository } from '../../repository' import { findApiKeys } from '../../services/api_key' import { analytics } from '../../utils/analytics' import { generateApiKey, hashApiKey } from '../../utils/auth' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const apiKeysResolver = authorized( async (_, __, { log, uid }) => { diff --git a/packages/api/src/resolvers/features/index.ts b/packages/api/src/resolvers/features/index.ts index 73f9036b0..6f92d922c 100644 --- a/packages/api/src/resolvers/features/index.ts +++ b/packages/api/src/resolvers/features/index.ts @@ -9,7 +9,7 @@ import { optInFeature, signFeatureToken, } from '../../services/features' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const optInFeatureResolver = authorized< OptInFeatureSuccess, diff --git a/packages/api/src/resolvers/filters/index.ts b/packages/api/src/resolvers/filters/index.ts index f1c743afd..b31760f7e 100644 --- a/packages/api/src/resolvers/filters/index.ts +++ b/packages/api/src/resolvers/filters/index.ts @@ -25,7 +25,7 @@ import { } from '../../generated/graphql' import { authTrx } from '../../repository' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const saveFilterResolver = authorized< SaveFilterSuccess, diff --git a/packages/api/src/resolvers/importers/uploadImportFileResolver.ts b/packages/api/src/resolvers/importers/uploadImportFileResolver.ts index f9dccfa37..fdb39e7d9 100644 --- a/packages/api/src/resolvers/importers/uploadImportFileResolver.ts +++ b/packages/api/src/resolvers/importers/uploadImportFileResolver.ts @@ -9,7 +9,7 @@ import { } from '../../generated/graphql' import { userRepository } from '../../repository/user' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' import { logger } from '../../utils/logger' import { countOfFilesWithPrefix, diff --git a/packages/api/src/resolvers/integrations/index.ts b/packages/api/src/resolvers/integrations/index.ts index dc2d0c4b5..0a849c5f0 100644 --- a/packages/api/src/resolvers/integrations/index.ts +++ b/packages/api/src/resolvers/integrations/index.ts @@ -37,7 +37,7 @@ import { enqueueExportToIntegration, enqueueImportFromIntegration, } from '../../utils/createTask' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const setIntegrationResolver = authorized< SetIntegrationSuccess, diff --git a/packages/api/src/resolvers/labels/index.ts b/packages/api/src/resolvers/labels/index.ts index 5e6e77c5a..6a6860940 100644 --- a/packages/api/src/resolvers/labels/index.ts +++ b/packages/api/src/resolvers/labels/index.ts @@ -37,7 +37,7 @@ import { updateLabel, } from '../../services/labels' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const labelsResolver = authorized( async (_obj, _params, { authTrx, log, uid }) => { diff --git a/packages/api/src/resolvers/links/index.ts b/packages/api/src/resolvers/links/index.ts index 8950222d1..657866d0c 100644 --- a/packages/api/src/resolvers/links/index.ts +++ b/packages/api/src/resolvers/links/index.ts @@ -8,7 +8,7 @@ import { } from '../../generated/graphql' import { updateLibraryItem } from '../../services/library_item' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' // export const updateLinkShareInfoResolver = authorized< // UpdateLinkShareInfoSuccess, diff --git a/packages/api/src/resolvers/newsletters/index.ts b/packages/api/src/resolvers/newsletters/index.ts index a7fc5a447..0415041eb 100644 --- a/packages/api/src/resolvers/newsletters/index.ts +++ b/packages/api/src/resolvers/newsletters/index.ts @@ -30,7 +30,7 @@ import { import { unsubscribeAll } from '../../services/subscriptions' import { Merge } from '../../util' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export type CreateNewsletterEmailSuccessPartial = Merge< CreateNewsletterEmailSuccess, diff --git a/packages/api/src/resolvers/popular_reads/index.ts b/packages/api/src/resolvers/popular_reads/index.ts index 30e1de08c..4e8773c55 100644 --- a/packages/api/src/resolvers/popular_reads/index.ts +++ b/packages/api/src/resolvers/popular_reads/index.ts @@ -5,7 +5,7 @@ import { MutationAddPopularReadArgs, } from '../../generated/graphql' import { addPopularRead } from '../../services/popular_reads' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const addPopularReadResolver = authorized< AddPopularReadSuccess, AddPopularReadError, diff --git a/packages/api/src/resolvers/recent_emails/index.ts b/packages/api/src/resolvers/recent_emails/index.ts index 8b91143d0..a4b86a929 100644 --- a/packages/api/src/resolvers/recent_emails/index.ts +++ b/packages/api/src/resolvers/recent_emails/index.ts @@ -13,7 +13,7 @@ import { } from '../../generated/graphql' import { updateReceivedEmail } from '../../services/received_emails' import { saveNewsletter } from '../../services/save_newsletter_email' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' import { generateUniqueUrl, parseEmailAddress } from '../../utils/parser' import { sendEmail } from '../../utils/sendEmail' diff --git a/packages/api/src/resolvers/recent_searches/index.ts b/packages/api/src/resolvers/recent_searches/index.ts index 952e97fd1..1f6faff6e 100644 --- a/packages/api/src/resolvers/recent_searches/index.ts +++ b/packages/api/src/resolvers/recent_searches/index.ts @@ -3,7 +3,7 @@ import { RecentSearchesSuccess, } from '../../generated/graphql' import { getRecentSearches } from '../../services/search_history' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const recentSearchesResolver = authorized< RecentSearchesSuccess, diff --git a/packages/api/src/resolvers/rules/index.ts b/packages/api/src/resolvers/rules/index.ts index 3ab1a9031..49f9c7537 100644 --- a/packages/api/src/resolvers/rules/index.ts +++ b/packages/api/src/resolvers/rules/index.ts @@ -14,7 +14,7 @@ import { SetRuleSuccess, } from '../../generated/graphql' import { deleteRule } from '../../services/rules' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const setRuleResolver = authorized< SetRuleSuccess, diff --git a/packages/api/src/resolvers/save/index.ts b/packages/api/src/resolvers/save/index.ts index f5508663d..5c22c86d1 100644 --- a/packages/api/src/resolvers/save/index.ts +++ b/packages/api/src/resolvers/save/index.ts @@ -12,7 +12,7 @@ import { saveFile } from '../../services/save_file' import { savePage } from '../../services/save_page' import { saveUrl } from '../../services/save_url' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const savePageResolver = authorized< SaveSuccess, diff --git a/packages/api/src/resolvers/send_install_instructions/index.ts b/packages/api/src/resolvers/send_install_instructions/index.ts index 62d134fbc..8bcc1a9c7 100644 --- a/packages/api/src/resolvers/send_install_instructions/index.ts +++ b/packages/api/src/resolvers/send_install_instructions/index.ts @@ -5,7 +5,7 @@ import { SendInstallInstructionsSuccess, } from '../../generated/graphql' import { userRepository } from '../../repository/user' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' import { sendEmail } from '../../utils/sendEmail' const INSTALL_INSTRUCTIONS_EMAIL_TEMPLATE_ID = diff --git a/packages/api/src/resolvers/user_device_tokens/index.ts b/packages/api/src/resolvers/user_device_tokens/index.ts index 661762f48..8fcc4d868 100644 --- a/packages/api/src/resolvers/user_device_tokens/index.ts +++ b/packages/api/src/resolvers/user_device_tokens/index.ts @@ -20,7 +20,7 @@ import { findDeviceTokensByUserId, } from '../../services/user_device_tokens' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' const PG_UNIQUE_CONSTRAINT_VIOLATION = '23505' diff --git a/packages/api/src/resolvers/user_personalization/index.ts b/packages/api/src/resolvers/user_personalization/index.ts index e9bb00a89..e696f7160 100644 --- a/packages/api/src/resolvers/user_personalization/index.ts +++ b/packages/api/src/resolvers/user_personalization/index.ts @@ -8,7 +8,7 @@ import { SetUserPersonalizationSuccess, SortOrder, } from '../../generated/graphql' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const setUserPersonalizationResolver = authorized< SetUserPersonalizationSuccess, diff --git a/packages/api/src/resolvers/webhooks/index.ts b/packages/api/src/resolvers/webhooks/index.ts index 606debe1e..a164797f3 100644 --- a/packages/api/src/resolvers/webhooks/index.ts +++ b/packages/api/src/resolvers/webhooks/index.ts @@ -22,7 +22,7 @@ import { import { authTrx } from '../../repository' import { deleteWebhook } from '../../services/webhook' import { analytics } from '../../utils/analytics' -import { authorized } from '../../utils/helpers' +import { authorized } from '../../utils/gql-utils' export const webhooksResolver = authorized( async (_obj, _params, { uid, log }) => { From 202d19c9d3ffa1a59455e9e9ef9f879048a67fc1 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 14:11:53 +0800 Subject: [PATCH 21/29] Fix api status mapping --- packages/api/src/jobs/save_page.ts | 5 +- packages/api/src/services/save_page.ts | 13 +-- packages/api/test/resolvers/article.test.ts | 89 +++++++++++++-------- 3 files changed, 61 insertions(+), 46 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index 7eacbf3d5..957f54944 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -3,10 +3,11 @@ import jwt from 'jsonwebtoken' import { promisify } from 'util' import { env } from '../env' import { redisDataSource } from '../redis_data_source' -import { savePage, stringToRequestStatus } from '../services/save_page' +import { savePage } from '../services/save_page' import { userRepository } from '../repository/user' import { logger } from '../utils/logger' import { Readability } from '@omnivore/readability' +import { ArticleSavingRequestStatus } from '../generated/graphql' const signToken = promisify(jwt.sign) @@ -356,7 +357,7 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { title, originalContent: content, parseResult: readabilityResult, - state: stringToRequestStatus(state), + state: state ? (state as ArticleSavingRequestStatus) : undefined, labels: labels?.map((name) => { return { name, diff --git a/packages/api/src/services/save_page.ts b/packages/api/src/services/save_page.ts index d49fc8b57..9959c8a40 100644 --- a/packages/api/src/services/save_page.ts +++ b/packages/api/src/services/save_page.ts @@ -62,15 +62,6 @@ const shouldParseInBackend = (input: SavePageInput): boolean => { ) } -export const stringToRequestStatus = ( - str: string | undefined | null -): ArticleSavingRequestStatus | null => { - if (str && str in ArticleSavingRequestStatus) { - return str as ArticleSavingRequestStatus - } - return null -} - export const savePage = async ( input: SavePageInput, user: User @@ -103,7 +94,7 @@ export const savePage = async ( canonicalUrl: parseResult.canonicalUrl, savedAt: input.savedAt ? new Date(input.savedAt) : new Date(), publishedAt: input.publishedAt ? new Date(input.publishedAt) : undefined, - state: stringToRequestStatus(input.state) || undefined, + state: input.state || undefined, rssFeedUrl: input.rssFeedUrl, folder: input.folder, }) @@ -117,7 +108,7 @@ export const savePage = async ( userId: user.id, url: itemToSave.originalUrl, articleSavingRequestId: clientRequestId || undefined, - state: stringToRequestStatus(input.state) || undefined, + state: input.state || undefined, labels: input.labels || undefined, folder: input.folder || undefined, }) diff --git a/packages/api/test/resolvers/article.test.ts b/packages/api/test/resolvers/article.test.ts index 76398e0d7..601b2a498 100644 --- a/packages/api/test/resolvers/article.test.ts +++ b/packages/api/test/resolvers/article.test.ts @@ -17,7 +17,7 @@ import { PageType, SyncUpdatedItemEdge, UpdateReason, - UploadFileStatus + UploadFileStatus, } from '../../src/generated/graphql' import { getRepository } from '../../src/repository' import { createGroup, deleteGroup } from '../../src/services/groups' @@ -25,7 +25,7 @@ import { createHighlight } from '../../src/services/highlights' import { createLabel, deleteLabels, - saveLabelsInLibraryItem + saveLabelsInLibraryItem, } from '../../src/services/labels' import { createLibraryItem, @@ -36,7 +36,7 @@ import { deleteLibraryItemsByUserId, findLibraryItemById, findLibraryItemByUrl, - updateLibraryItem + updateLibraryItem, } from '../../src/services/library_item' import { deleteUser } from '../../src/services/user' import * as createTask from '../../src/utils/createTask' @@ -570,23 +570,37 @@ describe('Article API', () => { ).expect(200) // Save a link, then archive it - let allLinks = await graphqlRequest(searchQuery('in:inbox'), authToken).expect( - 200 - ) + let allLinks = await graphqlRequest( + searchQuery('in:inbox'), + authToken + ).expect(200) const justSavedId = allLinks.body.data.search.edges[0].node.id await archiveLink(authToken, justSavedId) // test the negative case, ensuring the archive link wasn't returned - allLinks = await graphqlRequest(searchQuery('in:inbox'), authToken).expect(200) + allLinks = await graphqlRequest( + searchQuery('in:inbox'), + authToken + ).expect(200) expect(allLinks.body.data.search.edges[0]?.node?.url).to.not.eq(url) // Now save the link again, and ensure it is returned await graphqlRequest( - savePageQuery(url, title, originalContent, null, null, generateFakeUuid()), + savePageQuery( + url, + title, + originalContent, + null, + null, + generateFakeUuid() + ), authToken ).expect(200) - allLinks = await graphqlRequest(searchQuery('in:inbox'), authToken).expect(200) + allLinks = await graphqlRequest( + searchQuery('in:inbox'), + authToken + ).expect(200) expect(allLinks.body.data.search.edges[0].node.id).to.eq(justSavedId) expect(allLinks.body.data.search.edges[0].node.url).to.eq(url) }) @@ -610,6 +624,7 @@ describe('Article API', () => { ).expect(200) const savedItem = await findLibraryItemByUrl(url, user.id) + console.log('savedItem: ', savedItem) expect(savedItem?.archivedAt).to.not.be.null expect(savedItem?.labels?.map((l) => l.name)).to.eql(labels) }) @@ -778,15 +793,20 @@ describe('Article API', () => { context('when force is true', () => { before(async () => { - itemId = (await createLibraryItem({ - user: { id: user.id }, - originalUrl: 'https://blog.omnivore.app/setBookmarkArticle', - slug: 'test-with-omnivore', - readableContent: '

test

', - title: 'test title', - readingProgressBottomPercent: 100, - readingProgressTopPercent: 80, - }, user.id)).id + itemId = ( + await createLibraryItem( + { + user: { id: user.id }, + originalUrl: 'https://blog.omnivore.app/setBookmarkArticle', + slug: 'test-with-omnivore', + readableContent: '

test

', + title: 'test title', + readingProgressBottomPercent: 100, + readingProgressTopPercent: 80, + }, + user.id + ) + ).id }) after(async () => { @@ -2052,20 +2072,23 @@ describe('Article API', () => { ) }) - context('when since is -1000000000-01-01T00:00:00Z from android app', () => { - before(() => { - since = '-1000000000-01-01T00:00:00Z' - }) + context( + 'when since is -1000000000-01-01T00:00:00Z from android app', + () => { + before(() => { + since = '-1000000000-01-01T00:00:00Z' + }) - it('returns all', async () => { - const res = await graphqlRequest( - updatesSinceQuery(since), - authToken - ).expect(200) + it('returns all', async () => { + const res = await graphqlRequest( + updatesSinceQuery(since), + authToken + ).expect(200) - expect(res.body.data.updatesSince.edges.length).to.eql(5) - }) - }) + expect(res.body.data.updatesSince.edges.length).to.eql(5) + }) + } + ) context('returns highlights', () => { let highlight: Highlight @@ -2092,9 +2115,9 @@ describe('Article API', () => { expect(res.body.data.updatesSince.edges[0].node.highlights[0].id).to.eq( highlight.id ) - expect(res.body.data.updatesSince.edges[0].node.highlights[0].type).to.eq( - HighlightType.Highlight - ) + expect( + res.body.data.updatesSince.edges[0].node.highlights[0].type + ).to.eq(HighlightType.Highlight) }) }) }) From 6cedffc249a37792eae7c51b4d588b0f81a9999a Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 15:30:06 +0800 Subject: [PATCH 22/29] Only require original content in savePageJob --- packages/api/src/jobs/save_page.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index 957f54944..d62a0c2db 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -334,7 +334,7 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { return true } - if (!title || !content || !contentType || !readabilityResult) { + if (!content) { throw new Error( 'Invalid SavePage job, fetch result missing required data' ) From 6ab737432b3899f0dc7ba96c40bfb25ef643e1be Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 15:31:16 +0800 Subject: [PATCH 23/29] More logging --- packages/content-fetch/src/request_handler.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/content-fetch/src/request_handler.ts b/packages/content-fetch/src/request_handler.ts index 9672344d3..8f2db904f 100644 --- a/packages/content-fetch/src/request_handler.ts +++ b/packages/content-fetch/src/request_handler.ts @@ -138,7 +138,7 @@ export const contentFetchRequestHandler: RequestHandler = async (req, res) => { console.log('cacheFetchResult result', cacheResult) const jobs = await queueSavePageJob(savePageJobs) - console.log('save-page jobs queued', jobs.length) + console.log('save-page jobs queued', { jobs }) } catch (error) { if (error instanceof Error) { logRecord.error = error.message From 9d3a183c45adf912216d2378b790c63a045805e6 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Sun, 21 Jan 2024 15:44:41 +0800 Subject: [PATCH 24/29] Dont map label input --- packages/api/src/jobs/save_page.ts | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index d62a0c2db..9018f2e55 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -7,7 +7,10 @@ import { savePage } from '../services/save_page' import { userRepository } from '../repository/user' import { logger } from '../utils/logger' import { Readability } from '@omnivore/readability' -import { ArticleSavingRequestStatus } from '../generated/graphql' +import { + ArticleSavingRequestStatus, + CreateLabelInput, +} from '../generated/graphql' const signToken = promisify(jwt.sign) @@ -23,7 +26,7 @@ interface Data { url: string articleSavingRequestId: string state?: string - labels?: string[] + labels?: CreateLabelInput[] source: string folder: string rssFeedUrl?: string @@ -358,11 +361,7 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { originalContent: content, parseResult: readabilityResult, state: state ? (state as ArticleSavingRequestStatus) : undefined, - labels: labels?.map((name) => { - return { - name, - } - }), + labels: labels, rssFeedUrl, savedAt: savedAt ? new Date(savedAt) : new Date(), publishedAt: publishedAt ? new Date(publishedAt) : null, From 33c994ad5cf7719b8ee58a42df1fed0777622af8 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Mon, 22 Jan 2024 11:00:55 +0800 Subject: [PATCH 25/29] Remove RSS feed task handler url as this handled with jobs now --- packages/api/src/util.ts | 3 --- 1 file changed, 3 deletions(-) diff --git a/packages/api/src/util.ts b/packages/api/src/util.ts index 8b0064372..691b58063 100755 --- a/packages/api/src/util.ts +++ b/packages/api/src/util.ts @@ -68,7 +68,6 @@ export interface BackendEnv { textToSpeechTaskHandlerUrl: string recommendationTaskHandlerUrl: string thumbnailTaskHandlerUrl: string - rssFeedTaskHandlerUrl: string integrationExporterUrl: string integrationImporterUrl: string importerMetricsUrl: string @@ -149,7 +148,6 @@ const nullableEnvVars = [ 'RECOMMENDATION_TASK_HANDLER_URL', 'POCKET_CONSUMER_KEY', 'THUMBNAIL_TASK_HANDLER_URL', - 'RSS_FEED_TASK_HANDLER_URL', 'SENDGRID_VERIFICATION_TEMPLATE_ID', 'REMINDER_TASK_HANDLER_URL', 'TRUST_PROXY', @@ -247,7 +245,6 @@ export function getEnv(): BackendEnv { textToSpeechTaskHandlerUrl: parse('TEXT_TO_SPEECH_TASK_HANDLER_URL'), recommendationTaskHandlerUrl: parse('RECOMMENDATION_TASK_HANDLER_URL'), thumbnailTaskHandlerUrl: parse('THUMBNAIL_TASK_HANDLER_URL'), - rssFeedTaskHandlerUrl: parse('RSS_FEED_TASK_HANDLER_URL'), integrationExporterUrl: parse('INTEGRATION_EXPORTER_URL'), integrationImporterUrl: parse('INTEGRATION_IMPORTER_URL'), importerMetricsUrl: parse('IMPORTER_METRICS_COLLECTOR_URL'), From 0f78419e1a68beba3c53d09d41ba076dde273cc0 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Mon, 22 Jan 2024 11:26:59 +0800 Subject: [PATCH 26/29] Remove some debug --- packages/api/src/redis_data_source.ts | 2 -- 1 file changed, 2 deletions(-) diff --git a/packages/api/src/redis_data_source.ts b/packages/api/src/redis_data_source.ts index f6bedf6c0..ef7cf55e3 100644 --- a/packages/api/src/redis_data_source.ts +++ b/packages/api/src/redis_data_source.ts @@ -22,8 +22,6 @@ export class RedisDataSource { async initialize(): Promise { if (this.isInitialized) throw 'Error already initialized' - console.trace('initializing with options: ', { options: this.options }) - this.redisClient = createIORedisClient('app', this.options) this.workerRedisClient = createIORedisClient('worker', this.options) this.isInitialized = true From f7b17cb93f375d5f04a9d7f9001ec85738849319 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Mon, 22 Jan 2024 11:27:13 +0800 Subject: [PATCH 27/29] Only create queue object once in queue processor --- packages/api/src/jobs/rss/refreshAllFeeds.ts | 15 +++------------ packages/api/src/queue-processor.ts | 14 ++++++++++++++ 2 files changed, 17 insertions(+), 12 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshAllFeeds.ts b/packages/api/src/jobs/rss/refreshAllFeeds.ts index 7a87f4990..d166b1afc 100644 --- a/packages/api/src/jobs/rss/refreshAllFeeds.ts +++ b/packages/api/src/jobs/rss/refreshAllFeeds.ts @@ -1,6 +1,6 @@ import { Job, Queue } from 'bullmq' import { DataSource } from 'typeorm' -import { QUEUE_NAME } from '../../queue-processor' +import { QUEUE_NAME, getBackendQueue } from '../../queue-processor' import { redisDataSource } from '../../redis_data_source' import { RssSubscriptionGroup } from '../../utils/createTask' import { stringToHash } from '../../utils/helpers' @@ -105,17 +105,8 @@ const updateSubscriptionGroup = async ( await queueRSSRefreshFeedJob(jobid, payload) } -const createBackendQueue = (): Queue | undefined => { - if (!redisDataSource.workerRedisClient) { - throw new Error('Can not create queues, redis is not initialized') - } - return new Queue(QUEUE_NAME, { - connection: redisDataSource.workerRedisClient, - }) -} - export const queueRSSRefreshAllFeedsJob = async () => { - const queue = createBackendQueue() + const queue = getBackendQueue() if (!queue) { return false } @@ -135,7 +126,7 @@ export const queueRSSRefreshFeedJob = async ( payload: any, options = { priority: 'high' as QueuePriority } ): Promise => { - const queue = createBackendQueue() + const queue = getBackendQueue() if (!queue) { return undefined } diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 283b3b9c0..2f6b0daa9 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -15,6 +15,20 @@ import { CustomTypeOrmLogger } from './utils/logger' export const QUEUE_NAME = 'omnivore-backend-queue' +let backendQueue: Queue | undefined +export const getBackendQueue = (): Queue | undefined => { + if (backendQueue) { + return backendQueue + } + if (!redisDataSource.workerRedisClient) { + throw new Error('Can not create queues, redis is not initialized') + } + backendQueue = new Queue(QUEUE_NAME, { + connection: redisDataSource.workerRedisClient, + }) + return backendQueue +} + const main = async () => { console.log('[queue-processor]: starting queue processor') From db8284dc0b52f4f7b77fd4b407d30dfa51d8225b Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Mon, 22 Jan 2024 11:31:46 +0800 Subject: [PATCH 28/29] Wait until backend queue is ready before putting jobs into it --- packages/api/src/jobs/rss/refreshAllFeeds.ts | 4 ++-- packages/api/src/queue-processor.ts | 4 +++- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/packages/api/src/jobs/rss/refreshAllFeeds.ts b/packages/api/src/jobs/rss/refreshAllFeeds.ts index d166b1afc..6a78d8596 100644 --- a/packages/api/src/jobs/rss/refreshAllFeeds.ts +++ b/packages/api/src/jobs/rss/refreshAllFeeds.ts @@ -106,7 +106,7 @@ const updateSubscriptionGroup = async ( } export const queueRSSRefreshAllFeedsJob = async () => { - const queue = getBackendQueue() + const queue = await getBackendQueue() if (!queue) { return false } @@ -126,7 +126,7 @@ export const queueRSSRefreshFeedJob = async ( payload: any, options = { priority: 'high' as QueuePriority } ): Promise => { - const queue = getBackendQueue() + const queue = await getBackendQueue() if (!queue) { return undefined } diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 2f6b0daa9..dc6b27673 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -16,8 +16,9 @@ import { CustomTypeOrmLogger } from './utils/logger' export const QUEUE_NAME = 'omnivore-backend-queue' let backendQueue: Queue | undefined -export const getBackendQueue = (): Queue | undefined => { +export const getBackendQueue = async (): Promise => { if (backendQueue) { + await backendQueue.waitUntilReady() return backendQueue } if (!redisDataSource.workerRedisClient) { @@ -26,6 +27,7 @@ export const getBackendQueue = (): Queue | undefined => { backendQueue = new Queue(QUEUE_NAME, { connection: redisDataSource.workerRedisClient, }) + await backendQueue.waitUntilReady() return backendQueue } From eac93e6cb690a430641c2231a064558c5ea7a8f4 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Mon, 22 Jan 2024 12:02:35 +0800 Subject: [PATCH 29/29] Reduce logging --- packages/content-fetch/src/request_handler.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/content-fetch/src/request_handler.ts b/packages/content-fetch/src/request_handler.ts index 8f2db904f..9672344d3 100644 --- a/packages/content-fetch/src/request_handler.ts +++ b/packages/content-fetch/src/request_handler.ts @@ -138,7 +138,7 @@ export const contentFetchRequestHandler: RequestHandler = async (req, res) => { console.log('cacheFetchResult result', cacheResult) const jobs = await queueSavePageJob(savePageJobs) - console.log('save-page jobs queued', { jobs }) + console.log('save-page jobs queued', jobs.length) } catch (error) { if (error instanceof Error) { logRecord.error = error.message