From f052ab55a0ca01017b5339d671a0d8e27dd81f64 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Wed, 12 Jun 2024 13:01:36 +0800 Subject: [PATCH 1/3] monitor home feed job latency in prometheus --- packages/api/src/jobs/update_home.ts | 56 +++++++++++++++------------- packages/api/src/prometheus.ts | 8 ++++ packages/api/src/queue-processor.ts | 19 ++++++++-- 3 files changed, 54 insertions(+), 29 deletions(-) create mode 100644 packages/api/src/prometheus.ts diff --git a/packages/api/src/jobs/update_home.ts b/packages/api/src/jobs/update_home.ts index 02f994770..6c8ddad54 100644 --- a/packages/api/src/jobs/update_home.ts +++ b/packages/api/src/jobs/update_home.ts @@ -1,7 +1,9 @@ +import client from 'prom-client' import { LibraryItem } from '../entity/library_item' import { PublicItem } from '../entity/public_item' import { Subscription } from '../entity/subscription' import { User } from '../entity/user' +import { registerMetric } from '../prometheus' import { redisDataSource } from '../redis_data_source' import { findUnseenPublicItems } from '../services/home' import { searchLibraryItems } from '../services/library_item' @@ -434,6 +436,18 @@ const mixHomeItems = ( return sections } +// use prometheus to monitor the latency of each step +const latency = new client.Histogram({ + name: 'update_home_latency', + help: 'Latency of update home job', + labelNames: ['step'], + buckets: [0.1, 0.5, 1, 2, 5, 10], +}) + +registerMetric(latency) + +latency.observe(10) + export const updateHome = async (data: UpdateHomeJobData) => { const { userId, cursor } = data logger.info('Updating home for user', data) @@ -447,22 +461,19 @@ export const updateHome = async (data: UpdateHomeJobData) => { logger.info(`Updating home for user ${userId}`) - logger.profile('justAdded') + let end = latency.startTimer({ step: 'justAdded' }) const justAddedCandidates = await getJustAddedCandidates(userId) - logger.profile('justAdded', { - level: 'info', - message: `Found ${justAddedCandidates.length} just added candidates`, - }) + end() - logger.profile('selecting') + logger.info(`Found ${justAddedCandidates.length} just added candidates`) + + end = latency.startTimer({ step: 'select' }) const candidates = await selectCandidates( user, justAddedCandidates.map((c) => c.id) ) - logger.profile('selecting', { - level: 'info', - message: `Found ${candidates.length} candidates`, - }) + end() + logger.info(`Found ${candidates.length} candidates`) if (!justAddedCandidates.length && !candidates.length) { logger.info('No candidates found') @@ -471,26 +482,21 @@ export const updateHome = async (data: UpdateHomeJobData) => { // TODO: integrity check on candidates - logger.profile('ranking') + end = latency.startTimer({ step: 'ranking' }) const rankedCandidates = await rankCandidates(userId, candidates) - logger.profile('ranking', { - level: 'info', - message: `Ranked ${rankedCandidates.length} candidates`, - }) + end() - logger.profile('mixing') + logger.info(`Ranked ${rankedCandidates.length} candidates`) + + end = latency.startTimer({ step: 'mixing' }) const sections = mixHomeItems(justAddedCandidates, rankedCandidates) - logger.profile('mixing', { - level: 'info', - message: `Created ${sections.length} sections`, - }) + end() - logger.profile('saving') + logger.info(`Mixed ${sections.length} sections`) + + end = latency.startTimer({ step: 'saving' }) await appendSectionsToHome(userId, sections, cursor) - logger.profile('saving', { - level: 'info', - message: 'Sections appended to home', - }) + end() logger.info('Home updated for user', { userId }) } catch (error) { diff --git a/packages/api/src/prometheus.ts b/packages/api/src/prometheus.ts new file mode 100644 index 000000000..ffaa0516c --- /dev/null +++ b/packages/api/src/prometheus.ts @@ -0,0 +1,8 @@ +import client, { Metric } from 'prom-client' + +const registry = new client.Registry() + +export const registerMetric = (metric: Metric) => + registry.registerMetric(metric) + +export const getMetrics = async () => registry.metrics() diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index f77f92a55..42b775845 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -70,6 +70,7 @@ import { import { updateHome, UPDATE_HOME_JOB } from './jobs/update_home' import { updatePDFContentJob } from './jobs/update_pdf_content' import { uploadContentJob, UPLOAD_CONTENT_JOB } from './jobs/upload_content' +import { getMetrics } from './prometheus' import { redisDataSource } from './redis_data_source' import { CACHED_READING_POSITION_PREFIX } from './services/cached_reading_position' import { getJobPriority } from './utils/createTask' @@ -258,10 +259,15 @@ const main = async () => { } let output = '' - const metrics: JobType[] = ['active', 'failed', 'completed', 'prioritized'] - const counts = await queue.getJobCounts(...metrics) + const jobsTypes: JobType[] = [ + 'active', + 'failed', + 'completed', + 'prioritized', + ] + const counts = await queue.getJobCounts(...jobsTypes) - metrics.forEach((metric, idx) => { + jobsTypes.forEach((metric, idx) => { output += `# TYPE omnivore_queue_messages_${metric} gauge\n` output += `omnivore_queue_messages_${metric}{queue="${QUEUE_NAME}"} ${counts[metric]}\n` }) @@ -298,7 +304,12 @@ const main = async () => { output += `omnivore_queue_messages_oldest_job_age_seconds{queue="${QUEUE_NAME}"} ${0}\n` } - res.status(200).setHeader('Content-Type', 'text/plain').send(output) + const metrics = await getMetrics() + + res + .status(200) + .setHeader('Content-Type', 'text/plain') + .send(output + metrics) }) const server = app.listen(port, () => { From 57695c370c9f9bc597e9ba7816a7730de14d4749 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Wed, 12 Jun 2024 13:35:04 +0800 Subject: [PATCH 2/3] monitor each job latency in prometheus --- packages/api/src/jobs/update_home.ts | 4 +- packages/api/src/queue-processor.ts | 144 +++++++++++++++------------ 2 files changed, 83 insertions(+), 65 deletions(-) diff --git a/packages/api/src/jobs/update_home.ts b/packages/api/src/jobs/update_home.ts index 6c8ddad54..c80e18656 100644 --- a/packages/api/src/jobs/update_home.ts +++ b/packages/api/src/jobs/update_home.ts @@ -444,10 +444,10 @@ const latency = new client.Histogram({ buckets: [0.1, 0.5, 1, 2, 5, 10], }) -registerMetric(latency) - latency.observe(10) +registerMetric(latency) + export const updateHome = async (data: UpdateHomeJobData) => { const { userId, cursor } = data logger.info('Updating home for user', data) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 42b775845..4c4643070 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -12,6 +12,7 @@ import { Worker, } from 'bullmq' import express, { Express } from 'express' +import client from 'prom-client' import { appDataSource } from './data_source' import { env } from './env' import { TaskState } from './generated/graphql' @@ -70,7 +71,7 @@ import { import { updateHome, UPDATE_HOME_JOB } from './jobs/update_home' import { updatePDFContentJob } from './jobs/update_pdf_content' import { uploadContentJob, UPLOAD_CONTENT_JOB } from './jobs/upload_content' -import { getMetrics } from './prometheus' +import { getMetrics, registerMetric } from './prometheus' import { redisDataSource } from './redis_data_source' import { CACHED_READING_POSITION_PREFIX } from './services/cached_reading_position' import { getJobPriority } from './utils/createTask' @@ -79,6 +80,17 @@ import { logger } from './utils/logger' export const QUEUE_NAME = 'omnivore-backend-queue' export const JOB_VERSION = 'v001' +const jobLatency = new client.Histogram({ + name: 'omnivore_job_latency', + help: 'Latency of jobs in the queue', + labelNames: ['job_name'], + buckets: [0, 1, 5, 10, 50, 100, 500], +}) + +jobLatency.observe(10) + +registerMetric(jobLatency) + export const getBackendQueue = async ( name = QUEUE_NAME ): Promise => { @@ -139,71 +151,77 @@ 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 (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 + } + 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 createDigest(job.data) + case UPLOAD_CONTENT_JOB: + return uploadContentJob(job.data) + case UPDATE_HOME_JOB: + return updateHome(job.data) + case SCORE_LIBRARY_ITEM_JOB: + return scoreLibraryItem(job.data) + case GENERATE_PREVIEW_CONTENT_JOB: + return generatePreviewContent(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 createDigest(job.data) - case UPLOAD_CONTENT_JOB: - return uploadContentJob(job.data) - case UPDATE_HOME_JOB: - return updateHome(job.data) - case SCORE_LIBRARY_ITEM_JOB: - return scoreLibraryItem(job.data) - case GENERATE_PREVIEW_CONTENT_JOB: - return generatePreviewContent(job.data) - default: - logger.warning(`[queue-processor] unhandled job: ${job.name}`) } + + const end = jobLatency.startTimer({ job_name: job.name }) + await executeJob(job) + end() }, { connection, From a5295feb91806fb3ae5c4ad6ee6831e0d5d5a1c3 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Wed, 12 Jun 2024 13:37:23 +0800 Subject: [PATCH 3/3] remove console timer logs --- packages/api/src/jobs/ai/create_digest.ts | 15 --------------- 1 file changed, 15 deletions(-) diff --git a/packages/api/src/jobs/ai/create_digest.ts b/packages/api/src/jobs/ai/create_digest.ts index d499dbf57..7bdbd8870 100644 --- a/packages/api/src/jobs/ai/create_digest.ts +++ b/packages/api/src/jobs/ai/create_digest.ts @@ -166,7 +166,6 @@ const getCandidatesList = async ( userId: string, selectedLibraryItemIds?: string[] ): Promise => { - console.time('getCandidatesList') // use the queries from the digest definitions to lookup preferences // There should be a list of multiple queries we use. For now we can // hardcode these queries: @@ -222,8 +221,6 @@ const getCandidatesList = async ( dedupedCandidates.map((item) => item.title) ) - console.timeEnd('getCandidatesList') - if (dedupedCandidates.length === 0) { logger.info('No new candidates found') @@ -461,8 +458,6 @@ const generateSpeechFiles = ( summariesInHtml: string[], options: SSMLOptions ): SpeechFile[] => { - console.time('generateSpeechFiles') - const speechFiles = summariesInHtml.map((summary) => { const html = `
@@ -476,8 +471,6 @@ const generateSpeechFiles = ( }) }) - console.timeEnd('generateSpeechFiles') - return speechFiles } @@ -538,7 +531,6 @@ const uploadSummary = async ( digest: Digest, summaries: RankedItem[] ) => { - console.time('uploadSummary') logger.info('uploading summaries to gcs') const filename = `digest/${userId}/${digest.id}.json` @@ -560,7 +552,6 @@ const uploadSummary = async ( ) logger.info('uploaded summaries to gcs') - console.timeEnd('uploadSummary') } const sendPushNotification = async (userId: string, digest: Digest) => { @@ -749,8 +740,6 @@ const sendToChannels = async ( } export const createDigest = async (jobData: CreateDigestData) => { - console.time('createDigestJob') - // generate a unique id for the digest if not provided for scheduled jobs const digestId = jobData.id ?? uuid() @@ -804,9 +793,7 @@ export const createDigest = async (jobData: CreateDigestData) => { libraryItem: item, summary: '', })) - console.time('summarizeItems') const summaries = await summarizeItems(model, selections) - console.timeEnd('summarizeItems') const filteredSummaries = filterSummaries(summaries) const summariesInHtml = filteredSummaries.map((item) => { @@ -862,8 +849,6 @@ export const createDigest = async (jobData: CreateDigestData) => { // send notifications when digest is created await sendToChannels(user, digest, config?.channels) - - console.timeEnd('createDigestJob') } catch (error) { logger.error('createDigestJob error', error)