From 824b256d205cc7b91997a0f4073d02b76cd9080f Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Wed, 24 Apr 2024 15:55:54 +0800 Subject: [PATCH 1/3] fix memory leak from axios error --- packages/api/src/jobs/save_page.ts | 42 ++++++++++++++++++++++-------- packages/content-fetch/src/job.ts | 2 +- 2 files changed, 32 insertions(+), 12 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index 7e67024e9..d3d06c422 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -52,28 +52,42 @@ const uploadToSignedUrl = async ( contentType: string, contentObjUrl: string ) => { - logger.info('uploading to signed url', { - uploadSignedUrl, - contentType, - contentObjUrl, - }) + const maxContentLength = 10 * 1024 * 1024 // 10MB try { - const stream = await axios.get(contentObjUrl, { + logger.info('downloading content', { + contentObjUrl, + }) + + // download the content as stream and max 10MB + const response = await axios.get(contentObjUrl, { responseType: 'stream', + maxContentLength, timeout: REQUEST_TIMEOUT, }) - return await axios.put(uploadSignedUrl, stream.data, { + + logger.info('uploading to signed url', { + uploadSignedUrl, + contentType, + }) + + // upload the stream to the signed url + await axios.post(uploadSignedUrl, response.data, { headers: { 'Content-Type': contentType, }, - maxBodyLength: 1000000000, - maxContentLength: 100000000, + maxBodyLength: maxContentLength, timeout: REQUEST_TIMEOUT, }) + + return true } catch (error) { - logger.error('error uploading to signed url', error) - return null + if (axios.isAxiosError(error)) { + logger.error(`error uploading to signed url: ${error.message}`) + } else { + logger.error('error uploading to signed url', error) + } + return false } } @@ -104,6 +118,12 @@ const uploadPdf = async ( throw new Error('error while uploading pdf') } + logger.info('pdf uploaded successfully', { + url, + uploadFileId: result.id, + itemId: result.createdPageId, + }) + return { uploadFileId: result.id, itemId: result.createdPageId, diff --git a/packages/content-fetch/src/job.ts b/packages/content-fetch/src/job.ts index 88c263765..de3c2fb58 100644 --- a/packages/content-fetch/src/job.ts +++ b/packages/content-fetch/src/job.ts @@ -59,7 +59,7 @@ const getAttempts = (job: SavePageJob): number => { const getOpts = (job: SavePageJob): BulkJobOptions => { return { - jobId: `save-page_${job.userId}_${job.data.finalUrl}`, // make sure we don't have duplicate jobs + jobId: `${JOB_NAME}_${job.userId}_${job.data.finalUrl}`, // make sure we don't have duplicate jobs removeOnComplete: true, removeOnFail: true, attempts: getAttempts(job), From e680202b7e5a3a22235f5335ea5670579b9ba84a Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Wed, 24 Apr 2024 16:17:05 +0800 Subject: [PATCH 2/3] better error message --- packages/api/src/jobs/save_page.ts | 80 ++++++++++------------------ packages/api/src/utils/createTask.ts | 10 +--- packages/api/src/utils/logger.ts | 14 +++++ 3 files changed, 43 insertions(+), 61 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index d3d06c422..8f186a200 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -11,14 +11,14 @@ import { userRepository } from '../repository/user' import { saveFile } from '../services/save_file' import { savePage } from '../services/save_page' import { uploadFile } from '../services/upload_file' -import { logger } from '../utils/logger' +import { logError, logger } from '../utils/logger' const signToken = promisify(jwt.sign) const IMPORTER_METRICS_COLLECTOR_URL = env.queue.importerMetricsUrl const JWT_SECRET = env.server.jwtSecret -const MAX_ATTEMPTS = 2 +const MAX_ATTEMPTS = 1 const REQUEST_TIMEOUT = 30000 // 30 seconds interface Data { @@ -54,41 +54,30 @@ const uploadToSignedUrl = async ( ) => { const maxContentLength = 10 * 1024 * 1024 // 10MB - try { - logger.info('downloading content', { - contentObjUrl, - }) + logger.info('downloading content', { + contentObjUrl, + }) - // download the content as stream and max 10MB - const response = await axios.get(contentObjUrl, { - responseType: 'stream', - maxContentLength, - timeout: REQUEST_TIMEOUT, - }) + // download the content as stream and max 10MB + const response = await axios.get(contentObjUrl, { + responseType: 'stream', + maxContentLength, + timeout: REQUEST_TIMEOUT, + }) - logger.info('uploading to signed url', { - uploadSignedUrl, - contentType, - }) + logger.info('uploading to signed url', { + uploadSignedUrl, + contentType, + }) - // upload the stream to the signed url - await axios.post(uploadSignedUrl, response.data, { - headers: { - 'Content-Type': contentType, - }, - maxBodyLength: maxContentLength, - timeout: REQUEST_TIMEOUT, - }) - - return true - } catch (error) { - if (axios.isAxiosError(error)) { - logger.error(`error uploading to signed url: ${error.message}`) - } else { - logger.error('error uploading to signed url', error) - } - return false - } + // upload the stream to the signed url + await axios.put(uploadSignedUrl, response.data, { + headers: { + 'Content-Type': contentType, + }, + maxBodyLength: maxContentLength, + timeout: REQUEST_TIMEOUT, + }) } const uploadPdf = async ( @@ -109,14 +98,7 @@ const uploadPdf = async ( throw new Error('error while getting upload id and signed url') } - const uploaded = await uploadToSignedUrl( - result.uploadSignedUrl, - 'application/pdf', - url - ) - if (!uploaded) { - throw new Error('error while uploading pdf') - } + await uploadToSignedUrl(result.uploadSignedUrl, 'application/pdf', url) logger.info('pdf uploaded successfully', { url, @@ -154,7 +136,7 @@ const sendImportStatusUpdate = async ( } ) } catch (e) { - logger.error('error while sending import status update', e) + logError(e) } } @@ -288,20 +270,14 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { isImported = true isSaved = true } catch (e) { - if (e instanceof Error) { - logger.error(`error while saving page: ${e.message}`) - } else { - logger.error('error while saving page: unknown error') - } + logError(e) throw e } finally { - const lastAttempt = attemptsMade === MAX_ATTEMPTS - 1 - if (lastAttempt) { - logger.info(`last attempt reached ${data.url}`) - } + const lastAttempt = attemptsMade + 1 === MAX_ATTEMPTS if (taskId && (isSaved || lastAttempt)) { + logger.info('sending import status update') // send import status to update the metrics for importer await sendImportStatusUpdate(userId, taskId, isImported) } diff --git a/packages/api/src/utils/createTask.ts b/packages/api/src/utils/createTask.ts index b59914dc4..a8c6e2fb1 100644 --- a/packages/api/src/utils/createTask.ts +++ b/packages/api/src/utils/createTask.ts @@ -60,7 +60,7 @@ import { signFeatureToken } from '../services/features' import { OmnivoreAuthorizationHeader } from './auth' import { CreateTaskError } from './errors' import { stringToHash } from './helpers' -import { logger } from './logger' +import { logError, logger } from './logger' import View = google.cloud.tasks.v2.Task.View // Instantiates a client. @@ -106,14 +106,6 @@ export const getJobPriority = (jobName: string): number => { } } -const logError = (error: any): void => { - if (axios.isAxiosError(error)) { - logger.error(error.response) - } else { - logger.error(error) - } -} - const createHttpTaskWithToken = async ({ project = process.env.GOOGLE_CLOUD_PROJECT, queue = env.queue.name, diff --git a/packages/api/src/utils/logger.ts b/packages/api/src/utils/logger.ts index f4019893d..364cfffa3 100644 --- a/packages/api/src/utils/logger.ts +++ b/packages/api/src/utils/logger.ts @@ -1,6 +1,7 @@ /* eslint-disable @typescript-eslint/no-explicit-any */ /* eslint-disable @typescript-eslint/restrict-template-expressions */ import { LoggingWinston } from '@google-cloud/logging-winston' +import axios from 'axios' import jsonStringify from 'fast-safe-stringify' import { cloneDeep, isArray, isObject, isString, truncate } from 'lodash' import { DateTime } from 'luxon' @@ -168,6 +169,19 @@ export interface LogRecord { [key: string]: any } +export const logError = (error: any): void => { + if (axios.isAxiosError(error)) { + logger.error(error.message, { + response: error.response?.data, + stack: error.stack, + }) + } else if (error instanceof Error) { + logger.error(error.message, { stack: error.stack }) + } else { + logger.error(error) + } +} + export const logger = buildLogger('app') export default {} From e788c71eff6bc4e8471e6275114151dbccf080bd Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Wed, 24 Apr 2024 16:49:46 +0800 Subject: [PATCH 3/3] rename variables --- packages/api/src/jobs/save_page.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index 8f186a200..9099e7468 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -18,7 +18,7 @@ const signToken = promisify(jwt.sign) const IMPORTER_METRICS_COLLECTOR_URL = env.queue.importerMetricsUrl const JWT_SECRET = env.server.jwtSecret -const MAX_ATTEMPTS = 1 +const MAX_IMPORT_ATTEMPTS = 1 const REQUEST_TIMEOUT = 30000 // 30 seconds interface Data { @@ -274,7 +274,7 @@ export const savePageJob = async (data: Data, attemptsMade: number) => { throw e } finally { - const lastAttempt = attemptsMade + 1 === MAX_ATTEMPTS + const lastAttempt = attemptsMade + 1 === MAX_IMPORT_ATTEMPTS if (taskId && (isSaved || lastAttempt)) { logger.info('sending import status update')