From a77ade089049ba81daef06b5fb0bfbebed85a80b Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Mon, 22 Apr 2024 19:45:25 +0800 Subject: [PATCH 1/8] fix: export all items job stuck --- packages/api/src/jobs/integration/export_all_items.ts | 10 +++++----- packages/api/src/services/library_item.ts | 6 +++++- 2 files changed, 10 insertions(+), 6 deletions(-) diff --git a/packages/api/src/jobs/integration/export_all_items.ts b/packages/api/src/jobs/integration/export_all_items.ts index f3bf8950d..f642f8e39 100644 --- a/packages/api/src/jobs/integration/export_all_items.ts +++ b/packages/api/src/jobs/integration/export_all_items.ts @@ -50,9 +50,9 @@ export const exportAllItems = async (jobData: ExportAllItemsJobData) => { const maxItems = 100 const limit = 10 - let offset = 0 + let exported = 0 // get max 100 most recent items from the database - while (offset < maxItems) { + for (let offset = 0; offset < maxItems; offset += limit) { const libraryItems = await findRecentLibraryItems(userId, limit, offset) if (libraryItems.length === 0) { logger.info('no library items found', { @@ -92,17 +92,17 @@ export const exportAllItems = async (jobData: ExportAllItemsJobData) => { updated, }) - offset += libraryItems.length + exported += libraryItems.length logger.info('exported items', { ...jobData, - offset, + exported, }) } logger.info('exported all items', { ...jobData, - offset, + exported, }) // clear task name in integration diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index 832aea3fa..fb7f5957c 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -1216,7 +1216,11 @@ export const batchUpdateLibraryItems = async ( await authTrx( async (tx) => - getQueryBuilder(userId, tx).update(LibraryItem).set(values).execute(), + getQueryBuilder(userId, tx) + .take(searchArgs.size) + .update(LibraryItem) + .set(values) + .execute(), undefined, userId ) From de258ae0bf559e7b6bb18eeedd970ba59285db7f Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Mon, 22 Apr 2024 22:37:33 +0800 Subject: [PATCH 2/8] timeout if job takes more than 10 minutes --- packages/api/src/queue-processor.ts | 138 ++++++++++++---------- packages/api/src/services/library_item.ts | 45 +++---- packages/api/src/utils/helpers.ts | 6 + 3 files changed, 103 insertions(+), 86 deletions(-) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 5f0fc6a21..489fcabfe 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -19,6 +19,17 @@ import { aiSummarize, AI_SUMMARIZE_JOB_NAME } from './jobs/ai-summarize' import { createDigestJob, CREATE_DIGEST_JOB } from './jobs/ai/create_digest' import { bulkAction, BULK_ACTION_JOB_NAME } from './jobs/bulk_action' import { callWebhook, CALL_WEBHOOK_JOB_NAME } from './jobs/call_webhook' +import { + confirmEmailJob, + CONFIRM_EMAIL_JOB, + forwardEmailJob, + FORWARD_EMAIL_JOB, + saveAttachmentJob, + saveNewsletterJob, + SAVE_ATTACHMENT_JOB, + SAVE_NEWSLETTER_JOB, +} from './jobs/email/inbound_emails' +import { sendEmailJob, SEND_EMAIL_JOB } from './jobs/email/send_email' import { findThumbnail, THUMBNAIL_JOB } from './jobs/find_thumbnail' import { exportAllItems, @@ -37,7 +48,6 @@ import { import { refreshAllFeeds } from './jobs/rss/refreshAllFeeds' import { refreshFeed } from './jobs/rss/refreshFeed' import { savePageJob } from './jobs/save_page' -import { sendEmailJob, SEND_EMAIL_JOB } from './jobs/email/send_email' import { syncReadPositionsJob, SYNC_READ_POSITIONS_JOB_NAME, @@ -53,17 +63,8 @@ import { updatePDFContentJob } from './jobs/update_pdf_content' import { redisDataSource } from './redis_data_source' import { CACHED_READING_POSITION_PREFIX } from './services/cached_reading_position' import { getJobPriority } from './utils/createTask' +import { timeout } from './utils/helpers' import { logger } from './utils/logger' -import { - confirmEmailJob, - CONFIRM_EMAIL_JOB, - forwardEmailJob, - FORWARD_EMAIL_JOB, - saveAttachmentJob, - saveNewsletterJob, - SAVE_ATTACHMENT_JOB, - SAVE_NEWSLETTER_JOB, -} from './jobs/email/inbound_emails' export const QUEUE_NAME = 'omnivore-backend-queue' export const JOB_VERSION = 'v001' @@ -128,66 +129,73 @@ export const createWorker = (connection: ConnectionOptions) => new Worker( QUEUE_NAME, async (job: Job) => { - switch (job.name) { - case 'refresh-all-feeds': { - const queue = await getBackendQueue() - const counts = await queue?.getJobCounts('prioritized') - if (counts && counts.wait > 1000) { - return + const executeJob = async () => { + switch (job.name) { + case 'refresh-all-feeds': { + const queue = await getBackendQueue() + const counts = await queue?.getJobCounts('prioritized') + if (counts && counts.wait > 1000) { + return + } + return await refreshAllFeeds(appDataSource) } - return await refreshAllFeeds(appDataSource) + case 'refresh-feed': { + return await refreshFeed(job.data) + } + case 'save-page': { + return savePageJob(job.data, job.attemptsMade) + } + case 'update-pdf-content': { + return updatePDFContentJob(job.data) + } + case THUMBNAIL_JOB: + return findThumbnail(job.data) + case TRIGGER_RULE_JOB_NAME: + return triggerRule(job.data) + case UPDATE_LABELS_JOB: + return updateLabels(job.data) + case UPDATE_HIGHLIGHT_JOB: + return updateHighlight(job.data) + case SYNC_READ_POSITIONS_JOB_NAME: + return syncReadPositionsJob(job.data) + case BULK_ACTION_JOB_NAME: + return bulkAction(job.data) + case CALL_WEBHOOK_JOB_NAME: + return callWebhook(job.data) + case EXPORT_ITEM_JOB_NAME: + return exportItem(job.data) + case AI_SUMMARIZE_JOB_NAME: + return aiSummarize(job.data) + case PROCESS_YOUTUBE_VIDEO_JOB_NAME: + return processYouTubeVideo(job.data) + case PROCESS_YOUTUBE_TRANSCRIPT_JOB_NAME: + return processYouTubeTranscript(job.data) + case EXPORT_ALL_ITEMS_JOB_NAME: + return exportAllItems(job.data) + case SEND_EMAIL_JOB: + return sendEmailJob(job.data) + case CONFIRM_EMAIL_JOB: + return confirmEmailJob(job.data) + case SAVE_ATTACHMENT_JOB: + return saveAttachmentJob(job.data) + case SAVE_NEWSLETTER_JOB: + return saveNewsletterJob(job.data) + case FORWARD_EMAIL_JOB: + return forwardEmailJob(job.data) + case CREATE_DIGEST_JOB: + return createDigestJob(job.data) + default: + logger.warning(`[queue-processor] unhandled job: ${job.name}`) } - case 'refresh-feed': { - return await refreshFeed(job.data) - } - case 'save-page': { - return savePageJob(job.data, job.attemptsMade) - } - case 'update-pdf-content': { - return updatePDFContentJob(job.data) - } - case THUMBNAIL_JOB: - return findThumbnail(job.data) - case TRIGGER_RULE_JOB_NAME: - return triggerRule(job.data) - case UPDATE_LABELS_JOB: - return updateLabels(job.data) - case UPDATE_HIGHLIGHT_JOB: - return updateHighlight(job.data) - case SYNC_READ_POSITIONS_JOB_NAME: - return syncReadPositionsJob(job.data) - case BULK_ACTION_JOB_NAME: - return bulkAction(job.data) - case CALL_WEBHOOK_JOB_NAME: - return callWebhook(job.data) - case EXPORT_ITEM_JOB_NAME: - return exportItem(job.data) - case AI_SUMMARIZE_JOB_NAME: - return aiSummarize(job.data) - case PROCESS_YOUTUBE_VIDEO_JOB_NAME: - return processYouTubeVideo(job.data) - case PROCESS_YOUTUBE_TRANSCRIPT_JOB_NAME: - return processYouTubeTranscript(job.data) - case EXPORT_ALL_ITEMS_JOB_NAME: - return exportAllItems(job.data) - case SEND_EMAIL_JOB: - return sendEmailJob(job.data) - case CONFIRM_EMAIL_JOB: - return confirmEmailJob(job.data) - case SAVE_ATTACHMENT_JOB: - return saveAttachmentJob(job.data) - case SAVE_NEWSLETTER_JOB: - return saveNewsletterJob(job.data) - case FORWARD_EMAIL_JOB: - return forwardEmailJob(job.data) - case CREATE_DIGEST_JOB: - return createDigestJob(job.data) - default: - logger.warning(`[queue-processor] unhandled job: ${job.name}`) } + + // timeout if the job takes more than 10 minutes to execute + await Promise.race([timeout(1000 * 60 * 10), executeJob()]) }, { connection, + autorun: true, // start processing jobs immediately + lockDuration: 60_000, // 1 minute } ) diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index fb7f5957c..24b7a77f4 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -1117,6 +1117,13 @@ export const batchUpdateLibraryItems = async ( labelIds?: string[] | null, args?: unknown ) => { + if (!searchArgs.query) { + throw new Error('Search query is required') + } + + const searchQuery = parseSearchQuery(searchArgs.query) + const parameters: ObjectLiteral[] = [] + const queryString = buildQueryString(searchQuery, parameters) interface FolderArguments { folder: string } @@ -1140,19 +1147,17 @@ export const batchUpdateLibraryItems = async ( const getLibraryItemIds = async ( userId: string, em: EntityManager - ): Promise<{ id: string }[]> => { + ): Promise => { const queryBuilder = getQueryBuilder(userId, em) - return queryBuilder.select('library_item.id', 'id').getRawMany() - } + const libraryItems = await queryBuilder + .select('library_item.id', 'id') + .take(searchArgs.size) + .skip(searchArgs.from) + .getRawMany<{ id: string }>() - if (!searchArgs.query) { - throw new Error('Search query is required') + return libraryItems.map((item) => item.id) } - const searchQuery = parseSearchQuery(searchArgs.query) - const parameters: ObjectLiteral[] = [] - const queryString = buildQueryString(searchQuery, parameters) - const now = new Date().toISOString() // build the script let values: Record = {} @@ -1174,27 +1179,27 @@ export const batchUpdateLibraryItems = async ( throw new Error('Labels are required for this action') } - const libraryItems = await authTrx( + const libraryItemIds = await authTrx( async (tx) => getLibraryItemIds(userId, tx), undefined, userId ) // add labels to library items - for (const libraryItem of libraryItems) { - await addLabelsToLibraryItem(labelIds, libraryItem.id, userId) + for (const libraryItemId of libraryItemIds) { + await addLabelsToLibraryItem(labelIds, libraryItemId, userId) } return } case BulkActionType.MarkAsRead: { - const libraryItems = await authTrx( + const libraryItemIds = await authTrx( async (tx) => getLibraryItemIds(userId, tx), undefined, userId ) // update reading progress for library items - for (const libraryItem of libraryItems) { - await markItemAsRead(libraryItem.id, userId) + for (const libraryItemId of libraryItemIds) { + await markItemAsRead(libraryItemId, userId) } return @@ -1215,12 +1220,10 @@ export const batchUpdateLibraryItems = async ( } await authTrx( - async (tx) => - getQueryBuilder(userId, tx) - .take(searchArgs.size) - .update(LibraryItem) - .set(values) - .execute(), + async (tx) => { + const libraryItemIds = await getLibraryItemIds(userId, tx) + await tx.getRepository(LibraryItem).update(libraryItemIds, values) + }, undefined, userId ) diff --git a/packages/api/src/utils/helpers.ts b/packages/api/src/utils/helpers.ts index 785e5d5ef..2f27c0266 100644 --- a/packages/api/src/utils/helpers.ts +++ b/packages/api/src/utils/helpers.ts @@ -316,6 +316,12 @@ export const wait = (ms: number): Promise => { }) } +export const timeout = (ms: number): Promise => { + return new Promise((_resolve, reject) => { + setTimeout(() => reject(new Error('timeout')), ms) + }) +} + export const wordsCount = (text: string, isHtml = true): number => { try { return wordsCounter(text, { isHtml }).wordsCount From 1b6f4976d3c4696519fa9d6f549c6557a60ec961 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Tue, 23 Apr 2024 15:28:28 +0800 Subject: [PATCH 3/8] remove timeout --- packages/api/src/jobs/bulk_action.ts | 7 +- packages/api/src/queue-processor.ts | 118 +++++++++++----------- packages/api/src/services/library_item.ts | 10 +- packages/api/src/utils/helpers.ts | 6 -- 4 files changed, 68 insertions(+), 73 deletions(-) diff --git a/packages/api/src/jobs/bulk_action.ts b/packages/api/src/jobs/bulk_action.ts index 56a38f282..cea355d6e 100644 --- a/packages/api/src/jobs/bulk_action.ts +++ b/packages/api/src/jobs/bulk_action.ts @@ -23,9 +23,8 @@ export const bulkAction = async (data: BulkActionData) => { throw new Error('Queue not initialized') } const now = new Date().toISOString() - let offset = 0 - do { + for (let offset = 0; offset < count; offset += batchSize) { const searchArgs = { size: batchSize, query: `(${query}) AND updated:*..${now}`, // only process items that have not been updated @@ -36,9 +35,7 @@ export const bulkAction = async (data: BulkActionData) => { } catch (error) { logger.error('batch update error', error) } - - offset += batchSize - } while (offset < count) + } return true } diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 489fcabfe..471e03d0f 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -63,7 +63,6 @@ import { updatePDFContentJob } from './jobs/update_pdf_content' import { redisDataSource } from './redis_data_source' import { CACHED_READING_POSITION_PREFIX } from './services/cached_reading_position' import { getJobPriority } from './utils/createTask' -import { timeout } from './utils/helpers' import { logger } from './utils/logger' export const QUEUE_NAME = 'omnivore-backend-queue' @@ -129,68 +128,63 @@ export const createWorker = (connection: ConnectionOptions) => new Worker( QUEUE_NAME, async (job: Job) => { - const executeJob = async () => { - switch (job.name) { - case 'refresh-all-feeds': { - const queue = await getBackendQueue() - const counts = await queue?.getJobCounts('prioritized') - if (counts && counts.wait > 1000) { - return - } - return await refreshAllFeeds(appDataSource) + switch (job.name) { + case 'refresh-all-feeds': { + const queue = await getBackendQueue() + const counts = await queue?.getJobCounts('prioritized') + if (counts && counts.wait > 1000) { + return } - case 'refresh-feed': { - return await refreshFeed(job.data) - } - case 'save-page': { - return savePageJob(job.data, job.attemptsMade) - } - case 'update-pdf-content': { - return updatePDFContentJob(job.data) - } - case THUMBNAIL_JOB: - return findThumbnail(job.data) - case TRIGGER_RULE_JOB_NAME: - return triggerRule(job.data) - case UPDATE_LABELS_JOB: - return updateLabels(job.data) - case UPDATE_HIGHLIGHT_JOB: - return updateHighlight(job.data) - case SYNC_READ_POSITIONS_JOB_NAME: - return syncReadPositionsJob(job.data) - case BULK_ACTION_JOB_NAME: - return bulkAction(job.data) - case CALL_WEBHOOK_JOB_NAME: - return callWebhook(job.data) - case EXPORT_ITEM_JOB_NAME: - return exportItem(job.data) - case AI_SUMMARIZE_JOB_NAME: - return aiSummarize(job.data) - case PROCESS_YOUTUBE_VIDEO_JOB_NAME: - return processYouTubeVideo(job.data) - case PROCESS_YOUTUBE_TRANSCRIPT_JOB_NAME: - return processYouTubeTranscript(job.data) - case EXPORT_ALL_ITEMS_JOB_NAME: - return exportAllItems(job.data) - case SEND_EMAIL_JOB: - return sendEmailJob(job.data) - case CONFIRM_EMAIL_JOB: - return confirmEmailJob(job.data) - case SAVE_ATTACHMENT_JOB: - return saveAttachmentJob(job.data) - case SAVE_NEWSLETTER_JOB: - return saveNewsletterJob(job.data) - case FORWARD_EMAIL_JOB: - return forwardEmailJob(job.data) - case CREATE_DIGEST_JOB: - return createDigestJob(job.data) - default: - logger.warning(`[queue-processor] unhandled job: ${job.name}`) + return await refreshAllFeeds(appDataSource) } + case 'refresh-feed': { + return await refreshFeed(job.data) + } + case 'save-page': { + return savePageJob(job.data, job.attemptsMade) + } + case 'update-pdf-content': { + return updatePDFContentJob(job.data) + } + case THUMBNAIL_JOB: + return findThumbnail(job.data) + case TRIGGER_RULE_JOB_NAME: + return triggerRule(job.data) + case UPDATE_LABELS_JOB: + return updateLabels(job.data) + case UPDATE_HIGHLIGHT_JOB: + return updateHighlight(job.data) + case SYNC_READ_POSITIONS_JOB_NAME: + return syncReadPositionsJob(job.data) + case BULK_ACTION_JOB_NAME: + return bulkAction(job.data) + case CALL_WEBHOOK_JOB_NAME: + return callWebhook(job.data) + case EXPORT_ITEM_JOB_NAME: + return exportItem(job.data) + case AI_SUMMARIZE_JOB_NAME: + return aiSummarize(job.data) + case PROCESS_YOUTUBE_VIDEO_JOB_NAME: + return processYouTubeVideo(job.data) + case PROCESS_YOUTUBE_TRANSCRIPT_JOB_NAME: + return processYouTubeTranscript(job.data) + case EXPORT_ALL_ITEMS_JOB_NAME: + return exportAllItems(job.data) + case SEND_EMAIL_JOB: + return sendEmailJob(job.data) + case CONFIRM_EMAIL_JOB: + return confirmEmailJob(job.data) + case SAVE_ATTACHMENT_JOB: + return saveAttachmentJob(job.data) + case SAVE_NEWSLETTER_JOB: + return saveNewsletterJob(job.data) + case FORWARD_EMAIL_JOB: + return forwardEmailJob(job.data) + case CREATE_DIGEST_JOB: + return createDigestJob(job.data) + default: + logger.warning(`[queue-processor] unhandled job: ${job.name}`) } - - // timeout if the job takes more than 10 minutes to execute - await Promise.race([timeout(1000 * 60 * 10), executeJob()]) }, { connection, @@ -324,6 +318,10 @@ const main = async () => { console.log('completed job: ', job.jobId) }) + queueEvents.on('failed', async (job) => { + console.log('failed job: ', job.jobId) + }) + workerRedisClient.on('error', (error) => { console.trace('[queue-processor]: redis worker error', { error }) }) diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index 24b7a77f4..2c5e740a7 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -1146,9 +1146,15 @@ export const batchUpdateLibraryItems = async ( const getLibraryItemIds = async ( userId: string, - em: EntityManager + em: EntityManager, + forUpdate = false ): Promise => { const queryBuilder = getQueryBuilder(userId, em) + + if (forUpdate) { + queryBuilder.setLock('pessimistic_write') + } + const libraryItems = await queryBuilder .select('library_item.id', 'id') .take(searchArgs.size) @@ -1221,7 +1227,7 @@ export const batchUpdateLibraryItems = async ( await authTrx( async (tx) => { - const libraryItemIds = await getLibraryItemIds(userId, tx) + const libraryItemIds = await getLibraryItemIds(userId, tx, true) await tx.getRepository(LibraryItem).update(libraryItemIds, values) }, undefined, diff --git a/packages/api/src/utils/helpers.ts b/packages/api/src/utils/helpers.ts index 2f27c0266..785e5d5ef 100644 --- a/packages/api/src/utils/helpers.ts +++ b/packages/api/src/utils/helpers.ts @@ -316,12 +316,6 @@ export const wait = (ms: number): Promise => { }) } -export const timeout = (ms: number): Promise => { - return new Promise((_resolve, reject) => { - setTimeout(() => reject(new Error('timeout')), ms) - }) -} - export const wordsCount = (text: string, isHtml = true): number => { try { return wordsCounter(text, { isHtml }).wordsCount From 2f7c18b3639a779bec88b50dc76b12302b660839 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Tue, 23 Apr 2024 15:50:00 +0800 Subject: [PATCH 4/8] add more logs --- packages/api/src/queue-processor.ts | 6 ++++++ packages/api/src/server.ts | 3 +-- packages/api/src/services/library_item.ts | 6 +----- 3 files changed, 8 insertions(+), 7 deletions(-) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 471e03d0f..89448b96f 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -343,8 +343,14 @@ const main = async () => { }) }) await worker.close() + console.log('[queue-processor]: Worker closed') + await redisDataSource.shutdown() + console.log('[queue-processor]: Redis connection closed') + await appDataSource.destroy() + console.log('[queue-processor]: DB connection closed') + process.exit(0) } diff --git a/packages/api/src/server.ts b/packages/api/src/server.ts index 58cdafbda..87d12536a 100755 --- a/packages/api/src/server.ts +++ b/packages/api/src/server.ts @@ -178,9 +178,8 @@ const main = async (): Promise => { await apollo.stop() console.log('[api]: Express server stopped') - console.log('[posthog]: flushing events') await analytics.shutdownAsync() - console.log('[posthog]: events flushed') + console.log('[api]: Posthog events flushed') // Shutdown redis before DB because the quit sequence can // cause appDataSource to get reloaded in the callback diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index 2c5e740a7..5c4bd2d2b 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -650,11 +650,7 @@ export const buildQuery = ( queryBuilder.where('library_item.user_id = :userId', { userId }) // add select - selects.forEach((select, index) => { - if (index === 0) { - queryBuilder.select(select.column, select.alias) - } - + selects.forEach((select) => { queryBuilder.addSelect(select.column, select.alias) }) From 0faafa078693b9069f5ceae1ee45248ecf73b944 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Tue, 23 Apr 2024 16:04:59 +0800 Subject: [PATCH 5/8] fix rebase conflicts --- packages/api/src/services/library_item.ts | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index 5c4bd2d2b..a837ad066 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -647,17 +647,15 @@ export const buildQuery = ( args.useFolders ) } - queryBuilder.where('library_item.user_id = :userId', { userId }) - // add select - selects.forEach((select) => { - queryBuilder.addSelect(select.column, select.alias) - }) + queryBuilder.select(selects.map((select) => select.column)) if (args.includeContent) { queryBuilder.addSelect('library_item.readableContent') } + queryBuilder.where('library_item.user_id = :userId', { userId }) + if (!args.includePending) { queryBuilder.andWhere("library_item.state <> 'PROCESSING'") } From 002e455bbdf45dbc61c37d09aeb7c74f676d5d24 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Tue, 23 Apr 2024 17:54:52 +0800 Subject: [PATCH 6/8] fix tests --- packages/api/src/services/library_item.ts | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index a837ad066..47ec7bba8 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -627,7 +627,9 @@ export const buildQuery = ( // select all columns except content const selects: Select[] = getColumns(libraryItemRepository) .filter( - (select) => select !== 'readableContent' && select !== 'originalContent' + (select) => + select !== 'originalContent' && // exclude original content + (args.includeContent || select !== 'readableContent') // exclude content if not requested ) .map((column) => ({ column: `library_item.${column}` })) @@ -647,12 +649,14 @@ export const buildQuery = ( args.useFolders ) } - // add select - queryBuilder.select(selects.map((select) => select.column)) - if (args.includeContent) { - queryBuilder.addSelect('library_item.readableContent') - } + // add select + selects.forEach((select, index) => { + // select must be defined before adding additional selects + index === 0 + ? queryBuilder.select(select.column, select.alias) + : queryBuilder.addSelect(select.column, select.alias) + }) queryBuilder.where('library_item.user_id = :userId', { userId }) From 3fb7193fb56cd1a99a2b8fa7670b0dd5690fa609 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Tue, 23 Apr 2024 21:28:02 +0800 Subject: [PATCH 7/8] reduce getTextNodesBetween timeout to 60 seconds --- packages/api/src/utils/highlightGenerator.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/api/src/utils/highlightGenerator.ts b/packages/api/src/utils/highlightGenerator.ts index bfdd346c7..b3f4448ba 100644 --- a/packages/api/src/utils/highlightGenerator.ts +++ b/packages/api/src/utils/highlightGenerator.ts @@ -49,7 +49,7 @@ type FillNodeResponse = { } function getTextNodesBetween(rootNode: Node, startNode: Node, endNode: Node) { - const maxTime = 1000 * 60 * 10 // 10 minutes + const maxTime = 1000 * 60 // 60 seconds const start = Date.now() let textNodeStartingPoint = 0 let articleText = '' From 7f441b4ff37225a303ea5aefced98b9701591abe Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Tue, 23 Apr 2024 21:44:25 +0800 Subject: [PATCH 8/8] dedupe save-page job --- packages/content-fetch/src/job.ts | 33 ++++++++++++++++++++++--------- 1 file changed, 24 insertions(+), 9 deletions(-) diff --git a/packages/content-fetch/src/job.ts b/packages/content-fetch/src/job.ts index 8d88472f9..88c263765 100644 --- a/packages/content-fetch/src/job.ts +++ b/packages/content-fetch/src/job.ts @@ -4,9 +4,24 @@ import { redisDataSource } from './redis_data_source' const QUEUE_NAME = 'omnivore-backend-queue' const JOB_NAME = 'save-page' -interface savePageJob { +interface SavePageJobData { userId: string - data: unknown + url: string + finalUrl: string + articleSavingRequestId: string + state?: string + labels?: string[] + source: string + folder?: string + rssFeedUrl?: string + savedAt?: string + publishedAt?: string + taskId?: string +} + +interface SavePageJob { + userId: string + data: SavePageJobData isRss: boolean isImport: boolean priority: 'low' | 'high' @@ -16,7 +31,7 @@ const queue = new Queue(QUEUE_NAME, { connection: redisDataSource.queueRedisClient, }) -const getPriority = (job: savePageJob): number => { +const getPriority = (job: SavePageJob): number => { // we want to prioritized jobs by the expected time to complete // lower number means higher priority // priority 1: jobs that are expected to finish immediately @@ -33,7 +48,7 @@ const getPriority = (job: savePageJob): number => { return job.priority === 'low' ? 10 : 1 } -const getAttempts = (job: savePageJob): number => { +const getAttempts = (job: SavePageJob): number => { if (job.isRss || job.isImport) { // we don't want to retry rss or import jobs return 1 @@ -42,11 +57,11 @@ const getAttempts = (job: savePageJob): number => { return 3 } -const getOpts = (job: savePageJob): BulkJobOptions => { +const getOpts = (job: SavePageJob): BulkJobOptions => { return { - // jobId: `${job.userId}-${job.url}`, - // removeOnComplete: true, - // removeOnFail: true, + jobId: `save-page_${job.userId}_${job.data.finalUrl}`, // make sure we don't have duplicate jobs + removeOnComplete: true, + removeOnFail: true, attempts: getAttempts(job), priority: getPriority(job), backoff: { @@ -56,7 +71,7 @@ const getOpts = (job: savePageJob): BulkJobOptions => { } } -export const queueSavePageJob = async (savePageJobs: savePageJob[]) => { +export const queueSavePageJob = async (savePageJobs: SavePageJob[]) => { const jobs = savePageJobs.map((job) => ({ name: JOB_NAME, data: job.data,