From 01ebcbb16b4c95ca0b9904bbe6b3e0ac91cb97d3 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Fri, 10 May 2024 14:37:05 +0800 Subject: [PATCH 1/6] add bulk upload original content job --- packages/api/src/jobs/find_thumbnail.ts | 4 +- packages/api/src/jobs/save_page.ts | 48 +++++----------- packages/api/src/jobs/upload_content.ts | 56 +++++++++++++++++++ packages/api/src/queue-processor.ts | 3 + packages/api/src/resolvers/article/index.ts | 4 ++ .../resolvers/article_saving_request/index.ts | 16 +++++- .../src/resolvers/recommendations/index.ts | 13 ++++- packages/api/src/routers/article_router.ts | 4 +- packages/api/src/routers/page_router.ts | 6 +- packages/api/src/services/library_item.ts | 26 ++++++--- packages/api/src/utils/createTask.ts | 25 +++++++++ packages/api/src/utils/uploads.ts | 32 +++++++++++ 12 files changed, 189 insertions(+), 48 deletions(-) create mode 100644 packages/api/src/jobs/upload_content.ts diff --git a/packages/api/src/jobs/find_thumbnail.ts b/packages/api/src/jobs/find_thumbnail.ts index b7d65dd59..3a81a2d6b 100644 --- a/packages/api/src/jobs/find_thumbnail.ts +++ b/packages/api/src/jobs/find_thumbnail.ts @@ -127,7 +127,9 @@ export const _findThumbnail = (imagesSizes: (ImageSize | null)[]) => { export const findThumbnail = async (data: Data) => { const { libraryItemId, userId } = data - const item = await findLibraryItemById(libraryItemId, userId) + const item = await findLibraryItemById(libraryItemId, userId, { + select: ['thumbnail', 'readableContent'], + }) if (!item) { logger.info('page not found') return false diff --git a/packages/api/src/jobs/save_page.ts b/packages/api/src/jobs/save_page.ts index 9099e7468..89e554b9e 100644 --- a/packages/api/src/jobs/save_page.ts +++ b/packages/api/src/jobs/save_page.ts @@ -12,6 +12,7 @@ import { saveFile } from '../services/save_file' import { savePage } from '../services/save_page' import { uploadFile } from '../services/upload_file' import { logError, logger } from '../utils/logger' +import { downloadFromUrl, uploadToSignedUrl } from '../utils/uploads' const signToken = promisify(jwt.sign) @@ -47,39 +48,6 @@ const isFetchResult = (obj: unknown): obj is FetchResult => { return typeof obj === 'object' && obj !== null && 'finalUrl' in obj } -const uploadToSignedUrl = async ( - uploadSignedUrl: string, - contentType: string, - contentObjUrl: string -) => { - const maxContentLength = 10 * 1024 * 1024 // 10MB - - 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, - }) - - logger.info('uploading to signed url', { - uploadSignedUrl, - contentType, - }) - - // 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 ( url: string, userId: string, @@ -98,7 +66,19 @@ const uploadPdf = async ( throw new Error('error while getting upload id and signed url') } - await uploadToSignedUrl(result.uploadSignedUrl, 'application/pdf', url) + logger.info('downloading content', { + url, + }) + + const data = await downloadFromUrl(url, REQUEST_TIMEOUT) + + const uploadSignedUrl = result.uploadSignedUrl + const contentType = 'application/pdf' + logger.info('uploading to signed url', { + uploadSignedUrl, + contentType, + }) + await uploadToSignedUrl(uploadSignedUrl, data, contentType, REQUEST_TIMEOUT) logger.info('pdf uploaded successfully', { url, diff --git a/packages/api/src/jobs/upload_content.ts b/packages/api/src/jobs/upload_content.ts new file mode 100644 index 000000000..572f335e8 --- /dev/null +++ b/packages/api/src/jobs/upload_content.ts @@ -0,0 +1,56 @@ +import { findLibraryItemById } from '../services/library_item' +import { htmlToHighlightedMarkdown, htmlToMarkdown } from '../utils/parser' +import { uploadToSignedUrl } from '../utils/uploads' + +export const UPLOAD_CONTENT_JOB = 'UPLOAD_CONTENT_JOB' + +type ContentFormat = 'markdown' | 'highlightedMarkdown' | 'original' + +export interface UploadContentJobData { + libraryItemId: string + userId: string + format: ContentFormat + uploadUrl: string +} + +const convertContent = (content: string, format: ContentFormat): string => { + switch (format) { + case 'markdown': + return htmlToMarkdown(content) + case 'highlightedMarkdown': + return htmlToHighlightedMarkdown(content) + case 'original': + return content + default: + throw new Error('Unsupported format') + } +} + +const CONTENT_TYPES = { + markdown: 'text/markdown', + highlightedMarkdown: 'text/markdown', + original: 'text/html', +} + +export const uploadContentJob = async (data: UploadContentJobData) => { + const { libraryItemId, userId, format, uploadUrl } = data + const libraryItem = await findLibraryItemById(libraryItemId, userId, { + select: ['originalContent'], + }) + if (!libraryItem) { + throw new Error('Library item not found') + } + + if (!libraryItem.originalContent) { + throw new Error('Original content not found') + } + + const content = convertContent(libraryItem.originalContent, format) + + // 1 minute timeout + const timeout = 60000 + + const contentType = CONTENT_TYPES[format] + + await uploadToSignedUrl(uploadUrl, Buffer.from(content), contentType, timeout) +} diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 2d0bf5881..359bdfcf1 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -60,6 +60,7 @@ import { UPDATE_LABELS_JOB, } from './jobs/update_db' import { updatePDFContentJob } from './jobs/update_pdf_content' +import { uploadContentJob, UPLOAD_CONTENT_JOB } from './jobs/upload_content' import { redisDataSource } from './redis_data_source' import { CACHED_READING_POSITION_PREFIX } from './services/cached_reading_position' import { getJobPriority } from './utils/createTask' @@ -182,6 +183,8 @@ export const createWorker = (connection: ConnectionOptions) => return forwardEmailJob(job.data) case CREATE_DIGEST_JOB: return createDigest(job.data) + case UPLOAD_CONTENT_JOB: + return uploadContentJob(job.data) default: logger.warning(`[queue-processor] unhandled job: ${job.name}`) } diff --git a/packages/api/src/resolvers/article/index.ts b/packages/api/src/resolvers/article/index.ts index 660474c7b..026cb7988 100644 --- a/packages/api/src/resolvers/article/index.ts +++ b/packages/api/src/resolvers/article/index.ts @@ -399,6 +399,10 @@ export const getArticleResolver = authorized< 'recommendations.recommender', 'recommendations_recommender' ) + .leftJoinAndSelect( + 'recommendations_recommender.profile', + 'recommendations_recommender_profile' + ) .where('libraryItem.user_id = :uid', { uid }) // We allow the backend to use the ID instead of a slug to fetch the article diff --git a/packages/api/src/resolvers/article_saving_request/index.ts b/packages/api/src/resolvers/article_saving_request/index.ts index 467c00cc5..c99fa94aa 100644 --- a/packages/api/src/resolvers/article_saving_request/index.ts +++ b/packages/api/src/resolvers/article_saving_request/index.ts @@ -82,7 +82,21 @@ export const articleSavingRequestResolver = authorized< let libraryItem: LibraryItem | null = null if (id) { - libraryItem = await findLibraryItemById(id, uid) + libraryItem = await findLibraryItemById(id, uid, { + select: [ + 'id', + 'state', + 'originalUrl', + 'slug', + 'title', + 'author', + 'createdAt', + 'updatedAt', + ], + relations: { + user: true, + }, + }) } else if (url) { libraryItem = await findLibraryItemByUrl(cleanUrl(url), uid) } diff --git a/packages/api/src/resolvers/recommendations/index.ts b/packages/api/src/resolvers/recommendations/index.ts index e22401542..ef8b1b430 100644 --- a/packages/api/src/resolvers/recommendations/index.ts +++ b/packages/api/src/resolvers/recommendations/index.ts @@ -141,7 +141,14 @@ export const recommendResolver = authorized< MutationRecommendArgs >(async (_, { input }, { uid, log, signToken }) => { try { - const item = await findLibraryItemById(input.pageId, uid) + const item = await findLibraryItemById(input.pageId, uid, { + select: ['id'], + relations: { + highlights: { + user: true, + }, + }, + }) if (!item) { return { errorCodes: [RecommendErrorCode.NotFound], @@ -259,7 +266,9 @@ export const recommendHighlightsResolver = authorized< } } - const item = await findLibraryItemById(input.pageId, uid) + const item = await findLibraryItemById(input.pageId, uid, { + select: ['id'], + }) if (!item) { return { errorCodes: [RecommendHighlightsErrorCode.NotFound], diff --git a/packages/api/src/routers/article_router.ts b/packages/api/src/routers/article_router.ts index 6b3bab37f..fbe9ea4d5 100644 --- a/packages/api/src/routers/article_router.ts +++ b/packages/api/src/routers/article_router.ts @@ -94,7 +94,9 @@ export function articleRouter() { }) try { - const item = await findLibraryItemById(articleId, uid) + const item = await findLibraryItemById(articleId, uid, { + select: ['title', 'readableContent', 'itemLanguage'], + }) if (!item) { return res.status(404).send('Page not found') } diff --git a/packages/api/src/routers/page_router.ts b/packages/api/src/routers/page_router.ts index 459456292..bcf2e5b46 100644 --- a/packages/api/src/routers/page_router.ts +++ b/packages/api/src/routers/page_router.ts @@ -146,7 +146,11 @@ export function pageRouter() { return res.status(400).send({ errorCode: 'BAD_DATA' }) } - const item = await findLibraryItemById(itemId, claims.uid) + const item = await findLibraryItemById(itemId, claims.uid, { + relations: { + highlights: true, + }, + }) if (!item) { return res.status(404).send({ errorCode: 'NOT_FOUND' }) } diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index d91cfa10d..5ace2a4a0 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -782,17 +782,27 @@ export const findLibraryItemsByIds = async (ids: string[], userId: string) => { export const findLibraryItemById = async ( id: string, - userId: string + userId: string, + options?: { + select?: (keyof LibraryItem)[] + relations?: { + user?: boolean + labels?: boolean + highlights?: + | { + user?: boolean + } + | boolean + } + } ): Promise => { return authTrx( async (tx) => - tx - .createQueryBuilder(LibraryItem, 'library_item') - .leftJoinAndSelect('library_item.labels', 'labels') - .leftJoinAndSelect('library_item.highlights', 'highlights') - .leftJoinAndSelect('highlights.user', 'user') - .where('library_item.id = :id', { id }) - .getOne(), + tx.withRepository(libraryItemRepository).findOne({ + select: options?.select, + where: { id }, + relations: options?.relations, + }), undefined, userId ) diff --git a/packages/api/src/utils/createTask.ts b/packages/api/src/utils/createTask.ts index 47d1a7386..8740ba7f6 100644 --- a/packages/api/src/utils/createTask.ts +++ b/packages/api/src/utils/createTask.ts @@ -53,6 +53,10 @@ import { UPDATE_HIGHLIGHT_JOB, UPDATE_LABELS_JOB, } from '../jobs/update_db' +import { + UploadContentJobData, + UPLOAD_CONTENT_JOB, +} from '../jobs/upload_content' import { getBackendQueue, JOB_VERSION } from '../queue-processor' import { redisDataSource } from '../redis_data_source' import { writeDigest } from '../services/digest' @@ -94,6 +98,7 @@ export const getJobPriority = (jobName: string): number => { case `${REFRESH_FEED_JOB_NAME}_low`: case EXPORT_ITEM_JOB_NAME: case CREATE_DIGEST_JOB: + case UPLOAD_CONTENT_JOB: return 50 case EXPORT_ALL_ITEMS_JOB_NAME: case REFRESH_ALL_FEEDS_JOB_NAME: @@ -953,4 +958,24 @@ export const enqueueCreateDigest = async ( } } +export const enqueueBulkUploadContentJob = async ( + data: UploadContentJobData[] +) => { + const queue = await getBackendQueue() + if (!queue) { + return '' + } + + const jobs = data.map((d) => ({ + name: UPLOAD_CONTENT_JOB, + data: d, + opts: { + attempts: 3, + priority: getJobPriority(UPLOAD_CONTENT_JOB), + }, + })) + + return queue.addBulk(jobs) +} + export default createHttpTaskWithToken diff --git a/packages/api/src/utils/uploads.ts b/packages/api/src/utils/uploads.ts index 2d31b525f..64f94abe9 100644 --- a/packages/api/src/utils/uploads.ts +++ b/packages/api/src/utils/uploads.ts @@ -1,6 +1,7 @@ /* eslint-disable @typescript-eslint/no-unsafe-member-access */ /* eslint-disable @typescript-eslint/no-unsafe-assignment */ import { File, GetSignedUrlConfig, Storage } from '@google-cloud/storage' +import axios from 'axios' import { ContentReaderType } from '../entity/library_item' import { env } from '../env' import { PageType } from '../generated/graphql' @@ -33,6 +34,7 @@ const storage = env.fileUpload?.gcsUploadSAKeyFilePath ? new Storage({ keyFilename: env.fileUpload.gcsUploadSAKeyFilePath }) : new Storage() const bucketName = env.fileUpload.gcsUploadBucket +const maxContentLength = 10 * 1024 * 1024 // 10MB export const countOfFilesWithPrefix = async (prefix: string) => { const [files] = await storage.bucket(bucketName).getFiles({ prefix }) @@ -112,3 +114,33 @@ export const uploadToBucket = async ( export const createGCSFile = (filename: string): File => { return storage.bucket(bucketName).file(filename) } + +export const downloadFromUrl = async ( + contentObjUrl: string, + timeout?: number +) => { + // download the content as stream and max 10MB + const response = await axios.get(contentObjUrl, { + responseType: 'stream', + maxContentLength, + timeout, + }) + + return response.data +} + +export const uploadToSignedUrl = async ( + uploadSignedUrl: string, + data: Buffer, + contentType: string, + timeout?: number +) => { + // upload the stream to the signed url + await axios.put(uploadSignedUrl, data, { + headers: { + 'Content-Type': contentType, + }, + maxBodyLength: maxContentLength, + timeout, + }) +} From 0f184c4c21dd40c872fb9f9ed25b357843115bb4 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Fri, 10 May 2024 15:43:15 +0800 Subject: [PATCH 2/6] add get content api --- packages/api/src/jobs/upload_content.ts | 18 ++-- packages/api/src/routers/content_router.ts | 109 +++++++++++++++++++++ packages/api/src/server.ts | 3 + packages/api/src/services/library_item.ts | 16 ++- packages/api/src/utils/uploads.ts | 10 +- 5 files changed, 139 insertions(+), 17 deletions(-) create mode 100644 packages/api/src/routers/content_router.ts diff --git a/packages/api/src/jobs/upload_content.ts b/packages/api/src/jobs/upload_content.ts index 572f335e8..3c80f2bd6 100644 --- a/packages/api/src/jobs/upload_content.ts +++ b/packages/api/src/jobs/upload_content.ts @@ -1,16 +1,16 @@ import { findLibraryItemById } from '../services/library_item' import { htmlToHighlightedMarkdown, htmlToMarkdown } from '../utils/parser' -import { uploadToSignedUrl } from '../utils/uploads' +import { uploadToBucket } from '../utils/uploads' export const UPLOAD_CONTENT_JOB = 'UPLOAD_CONTENT_JOB' -type ContentFormat = 'markdown' | 'highlightedMarkdown' | 'original' +export type ContentFormat = 'markdown' | 'highlightedMarkdown' | 'original' export interface UploadContentJobData { libraryItemId: string userId: string format: ContentFormat - uploadUrl: string + filePath: string } const convertContent = (content: string, format: ContentFormat): string => { @@ -33,7 +33,7 @@ const CONTENT_TYPES = { } export const uploadContentJob = async (data: UploadContentJobData) => { - const { libraryItemId, userId, format, uploadUrl } = data + const { libraryItemId, userId, format, filePath } = data const libraryItem = await findLibraryItemById(libraryItemId, userId, { select: ['originalContent'], }) @@ -47,10 +47,8 @@ export const uploadContentJob = async (data: UploadContentJobData) => { const content = convertContent(libraryItem.originalContent, format) - // 1 minute timeout - const timeout = 60000 - - const contentType = CONTENT_TYPES[format] - - await uploadToSignedUrl(uploadUrl, Buffer.from(content), contentType, timeout) + await uploadToBucket(filePath, Buffer.from(content), { + contentType: CONTENT_TYPES[format], + timeout: 60000, // 1 minute + }) } diff --git a/packages/api/src/routers/content_router.ts b/packages/api/src/routers/content_router.ts new file mode 100644 index 000000000..63bd566c6 --- /dev/null +++ b/packages/api/src/routers/content_router.ts @@ -0,0 +1,109 @@ +import { Router } from 'express' +import { ContentFormat, UploadContentJobData } from '../jobs/upload_content' +import { findLibraryItemsByIds } from '../services/library_item' +import { getClaimsByToken, getTokenByRequest } from '../utils/auth' +import { enqueueBulkUploadContentJob } from '../utils/createTask' +import { logger } from '../utils/logger' +import { generateDownloadSignedUrl } from '../utils/uploads' + +export function contentRouter() { + const router = Router() + + interface GetContentRequest { + libraryItemIds: string[] + format: ContentFormat + } + + const isContentRequest = (data: any): data is GetContentRequest => { + return ( + typeof data === 'object' && + data !== null && + 'libraryItemIds' in data && + 'format' in data + ) + } + + // eslint-disable-next-line @typescript-eslint/no-misused-promises + router.get('/', async (req, res) => { + if (!isContentRequest(req.query)) { + logger.error('Bad request') + return res.status(400).send({ errorCode: 'BAD_REQUEST' }) + } + + const { libraryItemIds, format } = req.query + if ( + !Array.isArray(libraryItemIds) || + libraryItemIds.length === 0 || + libraryItemIds.length > 50 + ) { + logger.error('Library item ids are invalid') + return res.status(400).send({ errorCode: 'BAD_REQUEST' }) + } + + const token = getTokenByRequest(req) + // get claims from token + const claims = await getClaimsByToken(token) + if (!claims) { + logger.error('Token not found') + return res.status(401).send({ + error: 'UNAUTHORIZED', + }) + } + + // get user by uid from claims + const userId = claims.uid + + const libraryItems = await findLibraryItemsByIds(libraryItemIds, userId, { + select: ['id', 'updatedAt'], + }) + if (libraryItems.length === 0) { + logger.error('Library items not found') + return res.status(404).send({ errorCode: 'NOT_FOUND' }) + } + + // generate signed url for each library item + const data = await Promise.all( + libraryItems.map(async (libraryItem) => { + const filePath = `${userId}/${ + libraryItem.id + }.${libraryItem.updatedAt.getTime()}.${format}` + + try { + const downloadUrl = await generateDownloadSignedUrl(filePath, { + expires: Date.now() + 60 * 60 * 1000, // 1 hour + }) + + return { + libraryItemId: libraryItem.id, + userId, + filePath, + downloadUrl, + format, + } + } catch (error) { + logger.error('Error while generating signed url', error) + return { + libraryItemId: libraryItem.id, + error: 'Failed to generate download url', + } + } + }) + ) + + const validData = data.filter( + (d) => d.downloadUrl !== undefined && !('error' in d) + ) as UploadContentJobData[] + + await enqueueBulkUploadContentJob(validData) + + res.send({ + data: data.map((d) => ({ + libraryItemId: d.libraryItemId, + downloadUrl: d.downloadUrl, + error: d.error, + })), + }) + }) + + return router +} diff --git a/packages/api/src/server.ts b/packages/api/src/server.ts index 06d2fbdcb..88cb3b416 100755 --- a/packages/api/src/server.ts +++ b/packages/api/src/server.ts @@ -20,6 +20,7 @@ import { aiSummariesRouter } from './routers/ai_summary_router' import { articleRouter } from './routers/article_router' import { authRouter } from './routers/auth/auth_router' import { mobileAuthRouter } from './routers/auth/mobile/mobile_auth_router' +import { contentRouter } from './routers/content_router' import { digestRouter } from './routers/digest_router' import { explainRouter } from './routers/explain_router' import { integrationRouter } from './routers/integration_router' @@ -101,6 +102,8 @@ export const createApp = (): Express => { app.use('/api/integration', integrationRouter()) app.use('/api/tasks', taskRouter()) app.use('/api/digest', digestRouter()) + app.use('/api/content', contentRouter()) + app.use('/svc/pubsub/content', contentServiceRouter()) app.use('/svc/pubsub/links', linkServiceRouter()) app.use('/svc/pubsub/newsletters', newsletterServiceRouter()) diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index 5ace2a4a0..af6cf7e3a 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -764,10 +764,18 @@ export const findRecentLibraryItems = async ( ) } -export const findLibraryItemsByIds = async (ids: string[], userId: string) => { - const selectColumns = getColumns(libraryItemRepository) - .filter((column) => column !== 'originalContent') - .map((column) => `library_item.${column}`) +export const findLibraryItemsByIds = async ( + ids: string[], + userId: string, + options?: { + select?: (keyof LibraryItem)[] + } +) => { + const selectColumns = + options?.select || + getColumns(libraryItemRepository) + .filter((column) => column !== 'originalContent') + .map((column) => `library_item.${column}`) return authTrx( async (tx) => tx diff --git a/packages/api/src/utils/uploads.ts b/packages/api/src/utils/uploads.ts index 64f94abe9..49ad28aac 100644 --- a/packages/api/src/utils/uploads.ts +++ b/packages/api/src/utils/uploads.ts @@ -64,12 +64,16 @@ export const generateUploadSignedUrl = async ( } export const generateDownloadSignedUrl = async ( - filePathName: string + filePathName: string, + config?: { + expires?: number + } ): Promise => { const options: GetSignedUrlConfig = { version: 'v4', action: 'read', expires: Date.now() + 240 * 60 * 1000, // four hours + ...config, } const [url] = await storage .bucket(bucketName) @@ -102,13 +106,13 @@ export const generateUploadFilePathName = ( export const uploadToBucket = async ( filePath: string, data: Buffer, - options?: { contentType?: string; public?: boolean }, + options?: { contentType?: string; public?: boolean; timeout?: number }, selectedBucket?: string ): Promise => { await storage .bucket(selectedBucket || bucketName) .file(filePath) - .save(data, { ...options, timeout: 30000 }) + .save(data, { timeout: 30000, ...options }) // default timeout 30s } export const createGCSFile = (filename: string): File => { From 6e7a436ceac1b5648f362477fd33b61cb2f38944 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Fri, 10 May 2024 16:26:42 +0800 Subject: [PATCH 3/6] POST content API --- packages/api/src/jobs/upload_content.ts | 11 +++++++++++ packages/api/src/routers/content_router.ts | 8 +++++--- packages/api/src/services/library_item.ts | 2 +- packages/api/src/utils/createTask.ts | 4 ++-- 4 files changed, 19 insertions(+), 6 deletions(-) diff --git a/packages/api/src/jobs/upload_content.ts b/packages/api/src/jobs/upload_content.ts index 3c80f2bd6..78d9339b9 100644 --- a/packages/api/src/jobs/upload_content.ts +++ b/packages/api/src/jobs/upload_content.ts @@ -1,4 +1,5 @@ import { findLibraryItemById } from '../services/library_item' +import { logger } from '../utils/logger' import { htmlToHighlightedMarkdown, htmlToMarkdown } from '../utils/parser' import { uploadToBucket } from '../utils/uploads' @@ -33,22 +34,32 @@ const CONTENT_TYPES = { } export const uploadContentJob = async (data: UploadContentJobData) => { + logger.info('Uploading content to bucket', data) + const { libraryItemId, userId, format, filePath } = data const libraryItem = await findLibraryItemById(libraryItemId, userId, { select: ['originalContent'], }) if (!libraryItem) { + logger.error('Library item not found', data) throw new Error('Library item not found') } if (!libraryItem.originalContent) { + logger.error('Original content not found', data) throw new Error('Original content not found') } + logger.info('Converting content', data) const content = convertContent(libraryItem.originalContent, format) + console.time('uploadToBucket') + logger.info('Uploading content', data) await uploadToBucket(filePath, Buffer.from(content), { contentType: CONTENT_TYPES[format], timeout: 60000, // 1 minute }) + console.timeEnd('uploadToBucket') + + logger.info('Content uploaded', data) } diff --git a/packages/api/src/routers/content_router.ts b/packages/api/src/routers/content_router.ts index 63bd566c6..4d1709d06 100644 --- a/packages/api/src/routers/content_router.ts +++ b/packages/api/src/routers/content_router.ts @@ -24,13 +24,13 @@ export function contentRouter() { } // eslint-disable-next-line @typescript-eslint/no-misused-promises - router.get('/', async (req, res) => { - if (!isContentRequest(req.query)) { + router.post('/', async (req, res) => { + if (!isContentRequest(req.body)) { logger.error('Bad request') return res.status(400).send({ errorCode: 'BAD_REQUEST' }) } - const { libraryItemIds, format } = req.query + const { libraryItemIds, format } = req.body if ( !Array.isArray(libraryItemIds) || libraryItemIds.length === 0 || @@ -89,12 +89,14 @@ export function contentRouter() { } }) ) + logger.info('Signed urls generated', data) const validData = data.filter( (d) => d.downloadUrl !== undefined && !('error' in d) ) as UploadContentJobData[] await enqueueBulkUploadContentJob(validData) + logger.info('Bulk upload content job enqueued', validData) res.send({ data: data.map((d) => ({ diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index af6cf7e3a..a7376da30 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -772,7 +772,7 @@ export const findLibraryItemsByIds = async ( } ) => { const selectColumns = - options?.select || + options?.select?.map((column) => `library_item.${column}`) || getColumns(libraryItemRepository) .filter((column) => column !== 'originalContent') .map((column) => `library_item.${column}`) diff --git a/packages/api/src/utils/createTask.ts b/packages/api/src/utils/createTask.ts index 8740ba7f6..546b69d11 100644 --- a/packages/api/src/utils/createTask.ts +++ b/packages/api/src/utils/createTask.ts @@ -93,12 +93,12 @@ export const getJobPriority = (jobName: string): number => { return 5 case BULK_ACTION_JOB_NAME: case `${REFRESH_FEED_JOB_NAME}_high`: - return 10 case PROCESS_YOUTUBE_TRANSCRIPT_JOB_NAME: + case UPLOAD_CONTENT_JOB: + return 10 case `${REFRESH_FEED_JOB_NAME}_low`: case EXPORT_ITEM_JOB_NAME: case CREATE_DIGEST_JOB: - case UPLOAD_CONTENT_JOB: return 50 case EXPORT_ALL_ITEMS_JOB_NAME: case REFRESH_ALL_FEEDS_JOB_NAME: From 8034e188255c2e7f44e23df167acc8ff14848a17 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Fri, 10 May 2024 16:40:55 +0800 Subject: [PATCH 4/6] fix tests --- .../src/resolvers/article_saving_request/index.ts | 1 + packages/api/src/routers/content_router.ts | 2 +- packages/api/src/services/reports.ts | 8 ++++++-- packages/api/test/resolvers/article.test.ts | 6 +++++- packages/api/test/resolvers/labels.test.ts | 12 ++++++++++-- 5 files changed, 23 insertions(+), 6 deletions(-) diff --git a/packages/api/src/resolvers/article_saving_request/index.ts b/packages/api/src/resolvers/article_saving_request/index.ts index c99fa94aa..07d40aea5 100644 --- a/packages/api/src/resolvers/article_saving_request/index.ts +++ b/packages/api/src/resolvers/article_saving_request/index.ts @@ -92,6 +92,7 @@ export const articleSavingRequestResolver = authorized< 'author', 'createdAt', 'updatedAt', + 'savedAt', ], relations: { user: true, diff --git a/packages/api/src/routers/content_router.ts b/packages/api/src/routers/content_router.ts index 4d1709d06..d7605eb3a 100644 --- a/packages/api/src/routers/content_router.ts +++ b/packages/api/src/routers/content_router.ts @@ -64,7 +64,7 @@ export function contentRouter() { // generate signed url for each library item const data = await Promise.all( libraryItems.map(async (libraryItem) => { - const filePath = `${userId}/${ + const filePath = `content/${userId}/${ libraryItem.id }.${libraryItem.updatedAt.getTime()}.${format}` diff --git a/packages/api/src/services/reports.ts b/packages/api/src/services/reports.ts index e2aa573d5..33928c4f3 100644 --- a/packages/api/src/services/reports.ts +++ b/packages/api/src/services/reports.ts @@ -11,7 +11,9 @@ export const saveContentDisplayReport = async ( uid: string, input: ReportItemInput ): Promise => { - const item = await findLibraryItemById(input.pageId, uid) + const item = await findLibraryItemById(input.pageId, uid, { + select: ['id', 'readableContent', 'originalContent', 'originalUrl'], + }) if (!item) { logger.info('unable to submit report, item not found', input) return false @@ -53,7 +55,9 @@ export const saveAbuseReport = async ( uid: string, input: ReportItemInput ): Promise => { - const item = await findLibraryItemById(input.pageId, uid) + const item = await findLibraryItemById(input.pageId, uid, { + select: ['id'], + }) if (!item) { logger.info('unable to submit report, item not found', input) return false diff --git a/packages/api/test/resolvers/article.test.ts b/packages/api/test/resolvers/article.test.ts index e5d05e6ae..a9be6ed71 100644 --- a/packages/api/test/resolvers/article.test.ts +++ b/packages/api/test/resolvers/article.test.ts @@ -2345,7 +2345,11 @@ describe('Article API', () => { authToken ).expect(200) - const item = await findLibraryItemById(articleId, user.id) + const item = await findLibraryItemById(articleId, user.id, { + relations: { + labels: true, + }, + }) expect(item?.labels?.map((l) => l.name)).to.eql(['Favorites']) }) }) diff --git a/packages/api/test/resolvers/labels.test.ts b/packages/api/test/resolvers/labels.test.ts index 0dccc54d1..4020e2eaa 100644 --- a/packages/api/test/resolvers/labels.test.ts +++ b/packages/api/test/resolvers/labels.test.ts @@ -293,7 +293,11 @@ describe('Labels API', () => { labelId, }).expect(200) - const updatedItem = await findLibraryItemById(item.id, user.id) + const updatedItem = await findLibraryItemById(item.id, user.id, { + relations: { + labels: true, + }, + }) expect(updatedItem?.labels).not.deep.include(toDeleteLabel) }) }) @@ -545,7 +549,11 @@ describe('Labels API', () => { it('should update the item with the label', async () => { await graphqlRequest(query, authToken).expect(200) - const updatedItem = await findLibraryItemById(item.id, user.id) + const updatedItem = await findLibraryItemById(item.id, user.id, { + relations: { + labels: true, + }, + }) const updatedLabel = updatedItem?.labels?.filter( (l) => l.id === labelId )?.[0] From 6f7c58289a0121a7c95d074b3405c83361a2f389 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Fri, 10 May 2024 17:48:47 +0800 Subject: [PATCH 5/6] fix cors --- packages/api/src/routers/content_router.ts | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/packages/api/src/routers/content_router.ts b/packages/api/src/routers/content_router.ts index d7605eb3a..72218dda2 100644 --- a/packages/api/src/routers/content_router.ts +++ b/packages/api/src/routers/content_router.ts @@ -1,7 +1,9 @@ -import { Router } from 'express' +import cors from 'cors' +import express, { Router } from 'express' import { ContentFormat, UploadContentJobData } from '../jobs/upload_content' import { findLibraryItemsByIds } from '../services/library_item' import { getClaimsByToken, getTokenByRequest } from '../utils/auth' +import { corsConfig } from '../utils/corsConfig' import { enqueueBulkUploadContentJob } from '../utils/createTask' import { logger } from '../utils/logger' import { generateDownloadSignedUrl } from '../utils/uploads' @@ -23,8 +25,10 @@ export function contentRouter() { ) } + router.options('/', cors({ ...corsConfig, maxAge: 600 })) + // eslint-disable-next-line @typescript-eslint/no-misused-promises - router.post('/', async (req, res) => { + router.post('/', cors(corsConfig), async (req, res) => { if (!isContentRequest(req.body)) { logger.error('Bad request') return res.status(400).send({ errorCode: 'BAD_REQUEST' }) From 6bb81dd5c3041d48c70f11c46cb0ceaa2144dfac Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Sat, 11 May 2024 11:30:41 +0800 Subject: [PATCH 6/6] skip uploading if file already exists --- packages/api/src/routers/content_router.ts | 20 +++++++++++++++----- packages/api/src/utils/uploads.ts | 5 +++++ 2 files changed, 20 insertions(+), 5 deletions(-) diff --git a/packages/api/src/routers/content_router.ts b/packages/api/src/routers/content_router.ts index 72218dda2..ee17f432b 100644 --- a/packages/api/src/routers/content_router.ts +++ b/packages/api/src/routers/content_router.ts @@ -6,7 +6,7 @@ import { getClaimsByToken, getTokenByRequest } from '../utils/auth' import { corsConfig } from '../utils/corsConfig' import { enqueueBulkUploadContentJob } from '../utils/createTask' import { logger } from '../utils/logger' -import { generateDownloadSignedUrl } from '../utils/uploads' +import { generateDownloadSignedUrl, isFileExists } from '../utils/uploads' export function contentRouter() { const router = Router() @@ -77,12 +77,19 @@ export function contentRouter() { expires: Date.now() + 60 * 60 * 1000, // 1 hour }) + // check if file is already uploaded + const exists = await isFileExists(filePath) + if (exists) { + logger.info('File already exists', filePath) + } + return { libraryItemId: libraryItem.id, userId, filePath, downloadUrl, format, + exists, } } catch (error) { logger.error('Error while generating signed url', error) @@ -95,12 +102,15 @@ export function contentRouter() { ) logger.info('Signed urls generated', data) - const validData = data.filter( - (d) => d.downloadUrl !== undefined && !('error' in d) + // skip uploading if there is an error or file already exists + const uploadData = data.filter( + (d) => !('error' in d) && d.downloadUrl !== undefined && !d.exists ) as UploadContentJobData[] - await enqueueBulkUploadContentJob(validData) - logger.info('Bulk upload content job enqueued', validData) + if (uploadData.length > 0) { + await enqueueBulkUploadContentJob(uploadData) + logger.info('Bulk upload content job enqueued', uploadData) + } res.send({ data: data.map((d) => ({ diff --git a/packages/api/src/utils/uploads.ts b/packages/api/src/utils/uploads.ts index 49ad28aac..a60240c44 100644 --- a/packages/api/src/utils/uploads.ts +++ b/packages/api/src/utils/uploads.ts @@ -148,3 +148,8 @@ export const uploadToSignedUrl = async ( timeout, }) } + +export const isFileExists = async (filePath: string): Promise => { + const [exists] = await storage.bucket(bucketName).file(filePath).exists() + return exists +}