diff --git a/packages/api/src/elastic/pages.ts b/packages/api/src/elastic/pages.ts index f0a19434a..5434dcfdf 100644 --- a/packages/api/src/elastic/pages.ts +++ b/packages/api/src/elastic/pages.ts @@ -377,13 +377,7 @@ export const searchPages = async ( }, ], should: [], - must_not: [ - { - term: { - state: ArticleSavingRequestStatus.Failed, - }, - }, - ], + must_not: [], }, }, sort: [ diff --git a/packages/api/src/elastic/types.ts b/packages/api/src/elastic/types.ts index c6909eb97..cd829d54a 100644 --- a/packages/api/src/elastic/types.ts +++ b/packages/api/src/elastic/types.ts @@ -196,7 +196,7 @@ export interface Page { createdAt: Date updatedAt?: Date publishedAt?: Date - savedAt?: Date + savedAt: Date sharedAt?: Date archivedAt?: Date | null siteName?: string diff --git a/packages/api/src/resolvers/article/index.ts b/packages/api/src/resolvers/article/index.ts index 02312f8b6..b7865a5c6 100644 --- a/packages/api/src/resolvers/article/index.ts +++ b/packages/api/src/resolvers/article/index.ts @@ -46,6 +46,7 @@ import { ContentParseError } from '../../utils/errors' import { authorized, generateSlug, + isParsingTimeout, pageError, stringToHash, userDataToUser, @@ -102,7 +103,7 @@ const FORCE_PUPPETEER_URLS = [ /twitter\.com\/(?:#!\/)?(\w+)\/status(?:es)?\/(\d+)(?:\/.*)?/, /^((?:https?:)?\/\/)?((?:www|m)\.)?((?:youtube\.com|youtu.be))(\/(?:[\w-]+\?v=|embed\/|v\/)?)([\w-]+)(\S+)?$/, ] -const UNPARSEABLE_CONTENT = 'We were unable to parse this page.' +const UNPARSEABLE_CONTENT = '

We were unable to parse this page.

' export type CreateArticlesSuccessPartial = Merge< CreateArticleSuccess, @@ -246,8 +247,8 @@ export const createArticleResolver = authorized< const saveTime = new Date() const slug = generateSlug(parsedContent?.title || croppedPathname) - let articleToSave: Page = { - id: '', + const articleToSave: Page = { + id: pageId || '', userId: uid, originalHtml: domContent, content: parsedContent?.content || '', @@ -317,63 +318,47 @@ export const createArticleResolver = authorized< ) } - const existingPage = await getPageByParam({ - userId: uid, - url: articleToSave.url, - state: ArticleSavingRequestStatus.Succeeded, - }) - if (existingPage) { - // update existing page in elastic - existingPage.slug = slug - existingPage.savedAt = saveTime - existingPage.archivedAt = archive ? saveTime : undefined - existingPage.url = uploadFileUrlOverride || articleToSave.url - existingPage.hash = articleToSave.hash - - await updatePage(existingPage.id, existingPage, { ...ctx, uid }) - - log.info('page updated in elastic', existingPage.id) - articleToSave = existingPage - } else { - // create new page in elastic - if (!pageId) { - pageId = await createPage(articleToSave, { ...ctx, uid }) - if (!pageId) { - return pageError( - { - errorCodes: [CreateArticleErrorCode.ElasticError], - }, - ctx, - pageId - ) - } - } else { - const updated = await updatePage(pageId, articleToSave, { - ...ctx, - uid, - }) - - if (!updated) { - return pageError( - { - errorCodes: [CreateArticleErrorCode.ElasticError], - }, - ctx, - pageId - ) - } + // create new page in elastic + if (!pageId) { + const newPageId = await createPage(articleToSave, { ...ctx, uid }) + if (!newPageId) { + return pageError( + { + errorCodes: [CreateArticleErrorCode.ElasticError], + }, + ctx, + pageId + ) } + articleToSave.id = newPageId + } else { + // update existing page's state from processing to succeeded + articleToSave.archivedAt = archive ? saveTime : undefined + articleToSave.url = uploadFileUrlOverride || articleToSave.url + const updated = await updatePage(pageId, articleToSave, { + ...ctx, + uid, + }) - log.info( - 'page created in elastic', - pageId, - articleToSave.url, - articleToSave.slug, - articleToSave.title - ) - articleToSave.id = pageId + if (!updated) { + return pageError( + { + errorCodes: [CreateArticleErrorCode.ElasticError], + }, + ctx, + pageId + ) + } } + log.info( + 'page created in elastic', + articleToSave.id, + articleToSave.url, + articleToSave.slug, + articleToSave.title + ) + const createdArticle: PartialArticle = { ...articleToSave, isArchived: !!articleToSave.archivedAt, @@ -408,7 +393,7 @@ export const getArticleResolver: ResolverFn< Record, WithDataSourcesContext, QueryArticleArgs -> = async (_obj, { slug }, { claims }) => { +> = async (_obj, { slug }, { claims, pubsub }) => { try { if (!claims?.uid) { return { errorCodes: [ArticleErrorCode.Unauthorized] } @@ -432,12 +417,8 @@ export const getArticleResolver: ResolverFn< return { errorCodes: [ArticleErrorCode.NotFound] } } - if ( - page.state === ArticleSavingRequestStatus.Processing && - new Date(page.createdAt).getTime() < new Date().getTime() - 1000 * 60 - ) { + if (isParsingTimeout(page)) { page.content = UNPARSEABLE_CONTENT - page.description = UNPARSEABLE_CONTENT } return { @@ -874,6 +855,7 @@ export const searchResolver = authorized< savedDateFilter: searchQuery.savedDateFilter, publishedDateFilter: searchQuery.publishedDateFilter, subscriptionFilter: searchQuery.subscriptionFilter, + includePending: true, }, claims.uid )) || [[], 0] diff --git a/packages/api/src/resolvers/article_saving_request/index.ts b/packages/api/src/resolvers/article_saving_request/index.ts index d4a542fb5..1b8f2764e 100644 --- a/packages/api/src/resolvers/article_saving_request/index.ts +++ b/packages/api/src/resolvers/article_saving_request/index.ts @@ -2,6 +2,7 @@ import { ArticleSavingRequestError, ArticleSavingRequestErrorCode, + ArticleSavingRequestStatus, ArticleSavingRequestSuccess, CreateArticleSavingRequestError, CreateArticleSavingRequestErrorCode, @@ -9,7 +10,11 @@ import { MutationCreateArticleSavingRequestArgs, QueryArticleSavingRequestArgs, } from '../../generated/graphql' -import { authorized, pageToArticleSavingRequest } from '../../utils/helpers' +import { + authorized, + isParsingTimeout, + pageToArticleSavingRequest, +} from '../../utils/helpers' import { createPageSaveRequest } from '../../services/create_page_save_request' import { getPageById } from '../../elastic/pages' import { isErrorWithCode } from '../user' @@ -62,8 +67,12 @@ export const articleSavingRequestResolver = authorized< user = await models.user.get(page.userId) // eslint-disable-next-line no-empty } catch (error) {} - if (user && page) + if (user && page) { + if (isParsingTimeout(page)) { + page.state = ArticleSavingRequestStatus.Succeeded + } return { articleSavingRequest: pageToArticleSavingRequest(user, page) } + } return { errorCodes: [ArticleSavingRequestErrorCode.NotFound] } }) diff --git a/packages/api/src/routers/svc/pdf_attachments.ts b/packages/api/src/routers/svc/pdf_attachments.ts index 4fcba0176..c408d821b 100644 --- a/packages/api/src/routers/svc/pdf_attachments.ts +++ b/packages/api/src/routers/svc/pdf_attachments.ts @@ -155,6 +155,7 @@ export function pdfAttachmentsRouter() { slug: generateSlug(title), id: '', createdAt: new Date(), + savedAt: new Date(), readingProgressPercent: 0, readingProgressAnchorIndex: 0, state: ArticleSavingRequestStatus.Succeeded, diff --git a/packages/api/src/services/create_page_save_request.ts b/packages/api/src/services/create_page_save_request.ts index 4c0170012..ca8c0c357 100644 --- a/packages/api/src/services/create_page_save_request.ts +++ b/packages/api/src/services/create_page_save_request.ts @@ -10,11 +10,11 @@ import { import { generateSlug, pageToArticleSavingRequest } from '../utils/helpers' import * as privateIpLib from 'private-ip' import { countByCreatedAt, createPage, getPageByParam } from '../elastic/pages' -import { ArticleSavingRequestStatus, Page, PageType } from '../elastic/types' +import { ArticleSavingRequestStatus, PageType } from '../elastic/types' import { createPubSubClient, PubsubClient } from '../datalayer/pubsub' import normalizeUrl from 'normalize-url' -const SAVING_DESCRIPTION = 'Your link is being saved...' +const SAVING_CONTENT = 'Your link is being saved...' const isPrivateIP = privateIpLib.default @@ -82,52 +82,47 @@ export const createPageSaveRequest = async ( // get priority by checking rate limit if not specified priority = priority || (await getPriorityByRateLimit(userId)) + // look for existing page url = normalizeUrl(url, { stripHash: true, stripWWW: false, }) - const createdTaskName = await enqueueParseRequest( - url, - userId, - articleSavingRequestId, - priority - ) - - const existingPage = await getPageByParam({ + let page = await getPageByParam({ userId, url, - state: ArticleSavingRequestStatus.Succeeded, }) - if (existingPage) { - console.log('Page already exists', url) - existingPage.taskName = createdTaskName - return pageToArticleSavingRequest(user, existingPage) + if (page) { + console.log('Page already exists', page) + articleSavingRequestId = page.id + } else { + page = { + id: articleSavingRequestId, + userId, + content: SAVING_CONTENT, + hash: '', + pageType: PageType.Unknown, + readingProgressAnchorIndex: 0, + readingProgressPercent: 0, + slug: generateSlug(url), + title: url, + url, + state: ArticleSavingRequestStatus.Processing, + createdAt: new Date(), + savedAt: new Date(), + } + + // create processing page + const pageId = await createPage(page, { pubsub, uid: userId }) + if (!pageId) { + console.log('Failed to create page', page) + return Promise.reject({ + errorCode: CreateArticleSavingRequestErrorCode.BadData, + }) + } } - const page: Page = { - id: articleSavingRequestId, - userId, - content: SAVING_DESCRIPTION, - createdAt: new Date(), - hash: '', - pageType: PageType.Unknown, - readingProgressAnchorIndex: 0, - readingProgressPercent: 0, - slug: generateSlug(url), - title: url, - url, - taskName: createdTaskName, - state: ArticleSavingRequestStatus.Processing, - description: SAVING_DESCRIPTION, - } - - const pageId = await createPage(page, { pubsub, uid: userId }) - if (!pageId) { - console.log('Failed to create page', page) - return Promise.reject({ - errorCode: CreateArticleSavingRequestErrorCode.BadData, - }) - } + // enqueue task to parse page + await enqueueParseRequest(url, userId, articleSavingRequestId, priority) return pageToArticleSavingRequest(user, page) } diff --git a/packages/api/src/services/save_email.ts b/packages/api/src/services/save_email.ts index 9088fa0e8..e62175c47 100644 --- a/packages/api/src/services/save_email.ts +++ b/packages/api/src/services/save_email.ts @@ -66,6 +66,7 @@ export const saveEmail = async ( publishedAt: validatedDate(parseResult.parsedContent?.publishedDate), slug: slug, createdAt: new Date(), + savedAt: new Date(), readingProgressAnchorIndex: 0, readingProgressPercent: 0, subscription: input.author, diff --git a/packages/api/src/services/save_file.ts b/packages/api/src/services/save_file.ts index b8cf39307..874ca8e1d 100644 --- a/packages/api/src/services/save_file.ts +++ b/packages/api/src/services/save_file.ts @@ -94,6 +94,7 @@ export const saveFile = async ( userId: saver.id, id: input.clientRequestId, createdAt: new Date(), + savedAt: new Date(), readingProgressPercent: 0, readingProgressAnchorIndex: 0, state: ArticleSavingRequestStatus.Succeeded, diff --git a/packages/api/src/services/save_page.ts b/packages/api/src/services/save_page.ts index f54cc351e..e6fd79110 100644 --- a/packages/api/src/services/save_page.ts +++ b/packages/api/src/services/save_page.ts @@ -91,10 +91,11 @@ export const savePage = async ( hash: stringToHash(parseResult.parsedContent?.content || input.url), image: parseResult.parsedContent?.previewImage, publishedAt: validatedDate(parseResult.parsedContent?.publishedDate), - createdAt: new Date(), readingProgressPercent: 0, readingProgressAnchorIndex: 0, state: ArticleSavingRequestStatus.Succeeded, + createdAt: new Date(), + savedAt: new Date(), } const existingPage = await getPageByParam({ diff --git a/packages/api/src/utils/helpers.ts b/packages/api/src/utils/helpers.ts index f5499fd56..73c0abdc9 100644 --- a/packages/api/src/utils/helpers.ts +++ b/packages/api/src/utils/helpers.ts @@ -202,6 +202,14 @@ export const pageToArticleSavingRequest = ( updatedAt: page.updatedAt || new Date(), }) +export const isParsingTimeout = (page: Page): boolean => { + return ( + // page processed more than 30 seconds ago + page.state === ArticleSavingRequestStatus.Processing && + new Date(page.savedAt).getTime() < new Date().getTime() - 1000 * 30 + ) +} + export const validatedDate = ( date: Date | string | undefined ): Date | undefined => { diff --git a/packages/api/test/elastic/index.test.ts b/packages/api/test/elastic/index.test.ts index 7815db81b..c313ad556 100644 --- a/packages/api/test/elastic/index.test.ts +++ b/packages/api/test/elastic/index.test.ts @@ -46,6 +46,7 @@ describe('elastic api', () => { slug: 'test slug', createdAt: new Date(), updatedAt: new Date(), + savedAt: new Date(), readingProgressPercent: 100, readingProgressAnchorIndex: 0, url: 'https://blog.omnivore.app/p/getting-started-with-omnivore', @@ -98,6 +99,7 @@ describe('elastic api', () => { slug: 'test', createdAt: new Date(), updatedAt: new Date(), + savedAt: new Date(), readingProgressPercent: 0, readingProgressAnchorIndex: 0, url: 'https://blog.omnivore.app/testUrl', @@ -202,6 +204,7 @@ describe('elastic api', () => { content: 'test', slug: 'test', createdAt: new Date(createdAt), + savedAt: new Date(), readingProgressPercent: 0, readingProgressAnchorIndex: 0, url: 'https://blog.omnivore.app/testCount', diff --git a/packages/api/test/resolvers/article.test.ts b/packages/api/test/resolvers/article.test.ts index c6b4e10fe..9fdc33baf 100644 --- a/packages/api/test/resolvers/article.test.ts +++ b/packages/api/test/resolvers/article.test.ts @@ -434,7 +434,7 @@ describe('Article API', () => { pageId, { state: ArticleSavingRequestStatus.Processing, - createdAt: new Date(Date.now() - 1000 * 60), + savedAt: new Date(Date.now() - 1000 * 60), }, ctx ) @@ -444,7 +444,7 @@ describe('Article API', () => { const res = await graphqlRequest(query, authToken).expect(200) expect(res.body.data.article.article.content).to.eql( - 'We were unable to parse this page.' + '

We were unable to parse this page.

' ) }) }) @@ -474,7 +474,7 @@ describe('Article API', () => { before(async () => { // Create some test pages for (let i = 0; i < 15; i++) { - const page = { + const page: Page = { id: '', hash: 'test hash', userId: user.id, @@ -488,6 +488,7 @@ describe('Article API', () => { readingProgressAnchorIndex: 0, url: url, savedAt: new Date(), + state: ArticleSavingRequestStatus.Succeeded, } as Page const pageId = await createPage(page, ctx) if (!pageId) { @@ -583,7 +584,7 @@ describe('Article API', () => { ) expect( res.body.data.articles.pageInfo.startCursor, - 'startCursor' + 'st artCursor' ).to.eql('5') expect(res.body.data.articles.pageInfo.endCursor, 'endCursor').to.eql( '10' @@ -596,6 +597,31 @@ describe('Article API', () => { // expect(res.body.data.articles.pageInfo.hasPreviousPage).to.eql(true) }) }) + + context('when there are pages with failed state', () => { + before(async () => { + for (let i = 0; i < 5; i++) { + await updatePage( + pages[i].id, + { + state: ArticleSavingRequestStatus.Failed, + }, + ctx + ) + } + after = '10' + }) + it('should include state=failed pages', async () => { + const res = await graphqlRequest(query, authToken).expect(200) + + expect(res.body.data.articles.edges.length).to.eql(5) + expect(res.body.data.articles.edges[0].node.id).to.eql(pages[4].id) + expect(res.body.data.articles.edges[1].node.id).to.eql(pages[3].id) + expect(res.body.data.articles.edges[2].node.id).to.eql(pages[2].id) + expect(res.body.data.articles.edges[3].node.id).to.eql(pages[1].id) + expect(res.body.data.articles.edges[4].node.id).to.eql(pages[0].id) + }) + }) }) describe('SavePage', () => { @@ -731,6 +757,7 @@ describe('Article API', () => { title: 'test title', content: '

test

', createdAt: new Date(), + savedAt: new Date(), url: 'https://blog.omnivore.app/setBookmarkArticle', slug: 'test-with-omnivore', readingProgressPercent: 0, diff --git a/packages/api/test/resolvers/article_saving_request.test.ts b/packages/api/test/resolvers/article_saving_request.test.ts index 43cb5cfa8..cc8aabc1e 100644 --- a/packages/api/test/resolvers/article_saving_request.test.ts +++ b/packages/api/test/resolvers/article_saving_request.test.ts @@ -96,7 +96,7 @@ describe('ArticleSavingRequest API', () => { const page = await getPageById( res.body.data.createArticleSavingRequest.articleSavingRequest.id ) - expect(page?.description).to.eq('Your link is being saved...') + expect(page?.content).to.eq('Your link is being saved...') }) it('returns an error if the url is invalid', async () => { diff --git a/packages/api/test/util.ts b/packages/api/test/util.ts index d235a2c66..8f1bab161 100644 --- a/packages/api/test/util.ts +++ b/packages/api/test/util.ts @@ -51,6 +51,7 @@ export const createTestElasticPage = async ( title: 'test title', content: '

test content

', createdAt: new Date(), + savedAt: new Date(), url: 'https://example.com/test-url', slug: 'test-with-omnivore', labels: labels, diff --git a/packages/db/migrate.ts b/packages/db/migrate.ts index 04659182b..9d5b2bfd6 100755 --- a/packages/db/migrate.ts +++ b/packages/db/migrate.ts @@ -87,7 +87,7 @@ const logAppliedMigrations = ( } export const INDEX_ALIAS = 'pages_alias' -export const client = new Client({ +export const esClient = new Client({ node: process.env.ELASTIC_URL || 'http://localhost:9200', auth: { username: process.env.ELASTIC_USERNAME || '', @@ -103,7 +103,7 @@ const updateMappings = async (): Promise => { ) // update mappings - await client.indices.putMapping({ + await esClient.indices.putMapping({ index: INDEX_ALIAS, body: JSON.parse(indexSettings).mappings, }) @@ -123,8 +123,39 @@ postgrator log('Starting updating elasticsearch index mappings...') updateMappings() - .then(() => console.log('\nUpdating elastic completed.')) + .then(() => console.log('\nUpdating elastic mappings completed.')) .catch((error) => { log(`${chalk.red('Updating failed: ')}${error.message}`, chalk.red) process.exit(1) }) + +log('Starting adding default state to pages in elasticsearch...') +esClient + .update_by_query({ + index: INDEX_ALIAS, + body: { + script: { + source: 'ctx._source.state = params.state', + lang: 'painless', + params: { + state: 'SUCCEEDED', + }, + }, + query: { + bool: { + must_not: [ + { + exists: { + field: 'state', + }, + }, + ], + }, + }, + }, + }) + .then(() => console.log('\nAdding default state completed.')) + .catch((error) => { + log(`${chalk.red('Adding failed: ')}${error.message}`, chalk.red) + process.exit(1) + })