Merge pull request #4050 from omnivore-app/monitor-home-feed-generation

monitor home feed job latency in prometheus
This commit is contained in:
Hongbo Wu 2024-06-12 13:39:18 +08:00 committed by GitHub
commit 54dabe449d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 134 additions and 106 deletions

View file

@ -166,7 +166,6 @@ const getCandidatesList = async (
userId: string,
selectedLibraryItemIds?: string[]
): Promise<LibraryItem[]> => {
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 = `
<div id="readability-content">
@ -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)

View file

@ -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],
})
latency.observe(10)
registerMetric(latency)
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) {

View file

@ -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()

View file

@ -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,6 +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, registerMetric } from './prometheus'
import { redisDataSource } from './redis_data_source'
import { CACHED_READING_POSITION_PREFIX } from './services/cached_reading_position'
import { getJobPriority } from './utils/createTask'
@ -78,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<Queue | undefined> => {
@ -138,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,
@ -258,10 +277,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 +322,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, () => {