From c05721c87e572ce314ce4763190ab4c431509937 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Sun, 17 Mar 2024 11:48:58 +0800 Subject: [PATCH 01/14] validate the filter before creating rules --- packages/api/src/resolvers/rules/index.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/packages/api/src/resolvers/rules/index.ts b/packages/api/src/resolvers/rules/index.ts index 49f9c7537..0efa67750 100644 --- a/packages/api/src/resolvers/rules/index.ts +++ b/packages/api/src/resolvers/rules/index.ts @@ -15,6 +15,7 @@ import { } from '../../generated/graphql' import { deleteRule } from '../../services/rules' import { authorized } from '../../utils/gql-utils' +import { parseSearchQuery } from '../../utils/search' export const setRuleResolver = authorized< SetRuleSuccess, @@ -22,6 +23,9 @@ export const setRuleResolver = authorized< MutationSetRuleArgs >(async (_, { input }, { authTrx, uid, log }) => { try { + // validate filter + parseSearchQuery(input.filter) + const rule = await authTrx((t) => t.getRepository(Rule).save({ ...input, From 2fc6fd6e9365f5523955f3a2c8229635750cd320 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Sun, 17 Mar 2024 11:56:06 +0800 Subject: [PATCH 02/14] Show error message if failed to creating rules --- packages/web/pages/settings/rules.tsx | 20 +++++++++++++------- 1 file changed, 13 insertions(+), 7 deletions(-) diff --git a/packages/web/pages/settings/rules.tsx b/packages/web/pages/settings/rules.tsx index dd0d52448..d38ee76f1 100644 --- a/packages/web/pages/settings/rules.tsx +++ b/packages/web/pages/settings/rules.tsx @@ -32,13 +32,19 @@ const CreateRuleModal = (props: CreateRuleModalProps): JSX.Element => { const name = form.getFieldValue('name') const filter = form.getFieldValue('filter') const eventTypes = form.getFieldValue('eventTypes') - await setRuleMutation({ - name, - filter, - actions: [], - enabled: true, - eventTypes, - }) + try { + await setRuleMutation({ + name, + filter, + actions: [], + enabled: true, + eventTypes, + }) + } catch (error) { + showErrorToast('Error creating rule') + return + } + form.resetFields() props.setIsModalOpen(false) props.revalidate() From 9af4235233c17d7cfd2e059330165fe2c811bb8f Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Sun, 17 Mar 2024 12:38:23 +0800 Subject: [PATCH 03/14] skip failed to parse rules and show the timestamp in the UI --- packages/api/src/entity/rule.ts | 3 +++ packages/api/src/generated/graphql.ts | 2 ++ packages/api/src/generated/schema.graphql | 1 + packages/api/src/jobs/trigger_rule.ts | 25 ++++++++++++++----- packages/api/src/schema.ts | 1 + packages/api/src/services/rules.ts | 14 ++++++++++- .../0168.do.add_failed_at_to_rule.sql | 9 +++++++ .../0168.undo.add_failed_at_to_rule.sql | 9 +++++++ .../networking/queries/useGetRulesQuery.tsx | 2 ++ packages/web/pages/settings/rules.tsx | 6 +++++ 10 files changed, 65 insertions(+), 7 deletions(-) create mode 100755 packages/db/migrations/0168.do.add_failed_at_to_rule.sql create mode 100755 packages/db/migrations/0168.undo.add_failed_at_to_rule.sql diff --git a/packages/api/src/entity/rule.ts b/packages/api/src/entity/rule.ts index 0147cbd19..51c63b972 100644 --- a/packages/api/src/entity/rule.ts +++ b/packages/api/src/entity/rule.ts @@ -59,4 +59,7 @@ export class Rule { @UpdateDateColumn({ default: () => 'CURRENT_TIMESTAMP' }) updatedAt!: Date + + @Column('timestamptz') + failedAt?: Date } diff --git a/packages/api/src/generated/graphql.ts b/packages/api/src/generated/graphql.ts index 4e9c0735d..a2a4d7ef8 100644 --- a/packages/api/src/generated/graphql.ts +++ b/packages/api/src/generated/graphql.ts @@ -2434,6 +2434,7 @@ export type Rule = { createdAt: Scalars['Date']; enabled: Scalars['Boolean']; eventTypes: Array; + failedAt?: Maybe; filter: Scalars['String']; id: Scalars['ID']; name: Scalars['String']; @@ -6322,6 +6323,7 @@ export type RuleResolvers; enabled?: Resolver; eventTypes?: Resolver, ParentType, ContextType>; + failedAt?: Resolver, ParentType, ContextType>; filter?: Resolver; id?: Resolver; name?: Resolver; diff --git a/packages/api/src/generated/schema.graphql b/packages/api/src/generated/schema.graphql index 60ad805cd..ab001939f 100644 --- a/packages/api/src/generated/schema.graphql +++ b/packages/api/src/generated/schema.graphql @@ -1825,6 +1825,7 @@ type Rule { createdAt: Date! enabled: Boolean! eventTypes: [RuleEventType!]! + failedAt: Date filter: String! id: ID! name: String! diff --git a/packages/api/src/jobs/trigger_rule.ts b/packages/api/src/jobs/trigger_rule.ts index 260b7868f..3247e3fa9 100644 --- a/packages/api/src/jobs/trigger_rule.ts +++ b/packages/api/src/jobs/trigger_rule.ts @@ -8,7 +8,7 @@ import { softDeleteLibraryItem, updateLibraryItem, } from '../services/library_item' -import { findEnabledRules } from '../services/rules' +import { findEnabledRules, markRuleAsFailed } from '../services/rules' import { sendPushNotifications } from '../services/user' import { logger } from '../utils/logger' @@ -112,14 +112,27 @@ const triggerActions = async ( query: `(${rule.filter}) AND includes:${itemId}`, } - const libraryItems = await searchLibraryItems(searchArgs, userId) - if (libraryItems.count === 0) { - logger.info(`No pages found for rule ${rule.id}`) + let libraryItem: LibraryItem + + try { + const { libraryItems, count } = await searchLibraryItems( + searchArgs, + userId + ) + if (count === 0) { + logger.info(`No pages found for rule ${rule.id}`) + continue + } + + libraryItem = libraryItems[0] + } catch (error) { + // failed to search for library items, mark rule as failed + logger.error('Error parsing filter in rules', error) + await markRuleAsFailed(rule.id, userId) + continue } - const libraryItem = libraryItems.libraryItems[0] - for (const action of rule.actions) { const actionFunc = getRuleAction(action.type) const actionObj: RuleActionObj = { diff --git a/packages/api/src/schema.ts b/packages/api/src/schema.ts index 6c20fb51f..e7eb46a62 100755 --- a/packages/api/src/schema.ts +++ b/packages/api/src/schema.ts @@ -2152,6 +2152,7 @@ const schema = gql` createdAt: Date! updatedAt: Date eventTypes: [RuleEventType!]! + failedAt: Date } type RuleAction { diff --git a/packages/api/src/services/rules.ts b/packages/api/src/services/rules.ts index bf3637be7..6393a8c40 100644 --- a/packages/api/src/services/rules.ts +++ b/packages/api/src/services/rules.ts @@ -1,4 +1,4 @@ -import { ArrayContains, ILike } from 'typeorm' +import { ArrayContains, ILike, IsNull, Not } from 'typeorm' import { Rule, RuleAction, RuleEventType } from '../entity/rule' import { authTrx, getRepository } from '../repository' @@ -62,5 +62,17 @@ export const findEnabledRules = async ( user: { id: userId }, enabled: true, eventTypes: ArrayContains([eventType]), + failedAt: IsNull(), // only rules that have not failed }) } + +export const markRuleAsFailed = async (id: string, userId: string) => { + return authTrx( + (t) => + t.getRepository(Rule).update(id, { + failedAt: new Date(), + }), + undefined, + userId + ) +} diff --git a/packages/db/migrations/0168.do.add_failed_at_to_rule.sql b/packages/db/migrations/0168.do.add_failed_at_to_rule.sql new file mode 100755 index 000000000..6f77f166c --- /dev/null +++ b/packages/db/migrations/0168.do.add_failed_at_to_rule.sql @@ -0,0 +1,9 @@ +-- Type: DO +-- Name: add_failed_at_to_rule +-- Description: Add failed_at column to rules table + +BEGIN; + +ALTER TABLE omnivore.rules ADD COLUMN failed_at timestamptz; + +COMMIT; diff --git a/packages/db/migrations/0168.undo.add_failed_at_to_rule.sql b/packages/db/migrations/0168.undo.add_failed_at_to_rule.sql new file mode 100755 index 000000000..de2fc769d --- /dev/null +++ b/packages/db/migrations/0168.undo.add_failed_at_to_rule.sql @@ -0,0 +1,9 @@ +-- Type: UNDO +-- Name: add_failed_at_to_rule +-- Description: Add failed_at column to rules table + +BEGIN; + +ALTER TABLE omnivore.rules DROP COLUMN failed_at; + +COMMIT; diff --git a/packages/web/lib/networking/queries/useGetRulesQuery.tsx b/packages/web/lib/networking/queries/useGetRulesQuery.tsx index f7619cdcf..2da718072 100644 --- a/packages/web/lib/networking/queries/useGetRulesQuery.tsx +++ b/packages/web/lib/networking/queries/useGetRulesQuery.tsx @@ -29,6 +29,7 @@ export interface Rule { createdAt: Date updatedAt: Date eventTypes: RuleEventType[] + failedAt?: Date } interface RulesQueryResponse { @@ -62,6 +63,7 @@ export function useGetRulesQuery(): RulesQueryResponse { createdAt updatedAt eventTypes + failedAt } } ... on RulesError { diff --git a/packages/web/pages/settings/rules.tsx b/packages/web/pages/settings/rules.tsx index d38ee76f1..5cf3e1636 100644 --- a/packages/web/pages/settings/rules.tsx +++ b/packages/web/pages/settings/rules.tsx @@ -235,6 +235,7 @@ export default function Rules(): JSX.Element { filter: rule.filter, actions: rule.actions, eventTypes: rule.eventTypes, + failedAt: rule.failedAt, } }) }, [rules]) @@ -318,6 +319,11 @@ export default function Rules(): JSX.Element { ), }, + { + title: 'Failed At', + dataIndex: 'failedAt', + key: 'failedAt', + }, { title: '', key: 'tools', From 208a5895ef165a45a6cfe1120343e99c15d313b4 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Mon, 18 Mar 2024 11:06:24 +0800 Subject: [PATCH 04/14] remove unnecessary fields from item events --- packages/api/src/pubsub.ts | 37 ++------------ .../api/src/services/integrations/notion.ts | 2 +- .../api/src/services/integrations/readwise.ts | 4 +- packages/api/src/services/library_item.ts | 49 +++++++++++++++---- packages/api/src/utils/helpers.ts | 4 -- 5 files changed, 47 insertions(+), 49 deletions(-) diff --git a/packages/api/src/pubsub.ts b/packages/api/src/pubsub.ts index f13bfe5e4..53f95878a 100644 --- a/packages/api/src/pubsub.ts +++ b/packages/api/src/pubsub.ts @@ -3,6 +3,7 @@ import express from 'express' import { RuleEventType } from './entity/rule' import { env } from './env' import { ReportType } from './generated/graphql' +import { FeatureName, findFeatureByName } from './services/features' import { Merge } from './util' import { enqueueAISummarizeJob, @@ -11,13 +12,7 @@ import { enqueueTriggerRuleJob, enqueueWebhookJob, } from './utils/createTask' -import { deepDelete } from './utils/helpers' import { buildLogger } from './utils/logger' -import { - FeatureName, - findFeatureByName, - getFeatureName, -} from './services/features' import { processYouTubeVideo } from './jobs/process-youtube-video' const logger = buildLogger('pubsub') @@ -42,8 +37,6 @@ const isYouTubeVideoURL = (url: string | undefined): boolean => { } export const createPubSubClient = (): PubsubClient => { - const fieldsToDelete = ['user'] as const - const publish = (topicName: string, msg: Buffer): Promise => { if (env.dev.isLocal) { logger.info(`Publishing ${topicName}: ${msg.toString()}`) @@ -93,11 +86,6 @@ export const createPubSubClient = (): PubsubClient => { libraryItemIds: [libraryItemId], }) - const cleanData = deepDelete( - data as EntityData & Record, - [...fieldsToDelete] - ) - await enqueueWebhookJob({ userId, type, @@ -121,11 +109,6 @@ export const createPubSubClient = (): PubsubClient => { libraryItemId, }) } - - return publish( - 'entityCreated', - Buffer.from(JSON.stringify({ type, userId, ...cleanData })) - ) }, entityUpdated: async >( type: EntityType, @@ -148,32 +131,20 @@ export const createPubSubClient = (): PubsubClient => { libraryItemIds: [libraryItemId], }) - const cleanData = deepDelete( - data as EntityData & Record, - [...fieldsToDelete] - ) - await enqueueWebhookJob({ userId, type, action: 'updated', data, }) - - return publish( - 'entityUpdated', - Buffer.from(JSON.stringify({ type, userId, ...cleanData })) - ) }, - entityDeleted: ( + entityDeleted: async ( type: EntityType, id: string, userId: string ): Promise => { - return publish( - 'entityDeleted', - Buffer.from(JSON.stringify({ type, id, userId })) - ) + logger.info(`entityDeleted: ${type} ${id} ${userId}`) + await Promise.resolve() }, reportSubmitted: ( submitterId: string, diff --git a/packages/api/src/services/integrations/notion.ts b/packages/api/src/services/integrations/notion.ts index 4d823e031..de517dca6 100644 --- a/packages/api/src/services/integrations/notion.ts +++ b/packages/api/src/services/integrations/notion.ts @@ -5,8 +5,8 @@ import { Integration } from '../../entity/integration' import { LibraryItem } from '../../entity/library_item' import { env } from '../../env' import { Merge } from '../../util' -import { highlightUrl } from '../../utils/helpers' import { logger } from '../../utils/logger' +import { getHighlightUrl } from '../highlights' import { IntegrationClient } from './integration' type AnnotationColor = diff --git a/packages/api/src/services/integrations/readwise.ts b/packages/api/src/services/integrations/readwise.ts index dfc5b43db..0c204b630 100644 --- a/packages/api/src/services/integrations/readwise.ts +++ b/packages/api/src/services/integrations/readwise.ts @@ -1,7 +1,7 @@ import axios from 'axios' import { LibraryItem } from '../../entity/library_item' -import { highlightUrl } from '../../utils/helpers' import { logger } from '../../utils/logger' +import { getHighlightUrl } from '../highlights' import { IntegrationClient } from './integration' interface ReadwiseHighlight { @@ -98,7 +98,7 @@ export class ReadwiseClient implements IntegrationClient { text: highlight.quote, title: item.title, author: item.author || undefined, - highlight_url: highlightUrl(item.slug, highlight.id), + highlight_url: getHighlightUrl(item.slug, highlight.id), highlighted_at: new Date(highlight.createdAt).toISOString(), category, image_url: item.thumbnail || undefined, diff --git a/packages/api/src/services/library_item.ts b/packages/api/src/services/library_item.ts index 52d44f7ff..7d13e2fa8 100644 --- a/packages/api/src/services/library_item.ts +++ b/packages/api/src/services/library_item.ts @@ -29,6 +29,26 @@ import { logger } from '../utils/logger' import { parseSearchQuery } from '../utils/search' import { addLabelsToLibraryItem } from './labels' +type ItemEvent = { libraryItemId: string; userId: string } +type IgnoredFields = + | 'user' + | 'uploadFile' + | 'labelNames' + | 'highlightAnnotations' + | 'previewContentType' + | 'links' + | 'recommenderNames' + | 'textContentHash' + +type CreateItemEvent = Merge< + Omit, IgnoredFields>, + ItemEvent +> +type UpdateItemEvent = Merge< + Omit, IgnoredFields>, + ItemEvent +> + enum ReadFilter { ALL = 'all', READ = 'read', @@ -833,19 +853,32 @@ export const updateLibraryItem = async ( userId ) - if (skipPubSub) { + if (skipPubSub || libraryItem.state === LibraryItemState.Processing) { return updatedLibraryItem } - await pubsub.entityUpdated>( + if (libraryItem.state === LibraryItemState.Succeeded) { + // send create event if the item was created + await pubsub.entityCreated( + EntityType.PAGE, + { + ...updatedLibraryItem, + libraryItemId: id, + userId, + }, + userId + ) + + return updatedLibraryItem + } + + await pubsub.entityUpdated( EntityType.PAGE, { ...libraryItem, id, libraryItemId: id, - // don't send original content and readable content - originalContent: undefined, - readableContent: undefined, + userId, }, userId ) @@ -999,14 +1032,12 @@ export const createOrUpdateLibraryItem = async ( return newLibraryItem } - await pubsub.entityCreated>( + await pubsub.entityCreated( EntityType.PAGE, { ...newLibraryItem, libraryItemId: newLibraryItem.id, - // don't send original content and readable content - originalContent: undefined, - readableContent: undefined, + userId, }, userId ) diff --git a/packages/api/src/utils/helpers.ts b/packages/api/src/utils/helpers.ts index 9e3e18cd3..d2fec6520 100644 --- a/packages/api/src/utils/helpers.ts +++ b/packages/api/src/utils/helpers.ts @@ -10,7 +10,6 @@ import { Highlight as HighlightData } from '../entity/highlight' import { LibraryItem, LibraryItemState } from '../entity/library_item' import { Recommendation as RecommendationData } from '../entity/recommendation' import { RegistrationType, User } from '../entity/user' -import { env } from '../env' import { Article, ArticleSavingRequest, @@ -404,6 +403,3 @@ export const setRecentlySavedItemInRedis = async ( }) } } - -export const highlightUrl = (slug: string, highlightId: string): string => - `${env.client.url}/me/${slug}#${highlightId}` From a49366b4cc249bec2f2eec55fcc47f53ab8891c8 Mon Sep 17 00:00:00 2001 From: Hongbo Wu Date: Mon, 18 Mar 2024 11:24:16 +0800 Subject: [PATCH 05/14] remove unnecessary fields from item events --- packages/api/src/jobs/trigger_rule.ts | 1 + packages/api/src/pubsub.ts | 31 ++++++++--------- packages/api/src/services/highlights.ts | 11 +++--- packages/api/src/services/labels.ts | 12 +++---- packages/api/src/services/library_item.ts | 41 +++++++---------------- 5 files changed, 41 insertions(+), 55 deletions(-) diff --git a/packages/api/src/jobs/trigger_rule.ts b/packages/api/src/jobs/trigger_rule.ts index 3247e3fa9..a75fb75f5 100644 --- a/packages/api/src/jobs/trigger_rule.ts +++ b/packages/api/src/jobs/trigger_rule.ts @@ -16,6 +16,7 @@ export interface TriggerRuleJobData { libraryItemId: string userId: string ruleEventType: RuleEventType + data: unknown } interface RuleActionObj { diff --git a/packages/api/src/pubsub.ts b/packages/api/src/pubsub.ts index 53f95878a..b210fbfa5 100644 --- a/packages/api/src/pubsub.ts +++ b/packages/api/src/pubsub.ts @@ -4,7 +4,6 @@ import { RuleEventType } from './entity/rule' import { env } from './env' import { ReportType } from './generated/graphql' import { FeatureName, findFeatureByName } from './services/features' -import { Merge } from './util' import { enqueueAISummarizeJob, enqueueExportItem, @@ -19,11 +18,6 @@ const logger = buildLogger('pubsub') const client = new PubSub() -type EntityData> = Merge< - T, - { libraryItemId: string } -> - const isYouTubeVideoURL = (url: string | undefined): boolean => { if (!url) { return false @@ -68,16 +62,17 @@ export const createPubSubClient = (): PubsubClient => { }, entityCreated: async >( type: EntityType, - data: EntityData, - userId: string + data: T, + userId: string, + libraryItemId: string ): Promise => { - const libraryItemId = data.libraryItemId // queue trigger rule job if (type === EntityType.PAGE) { await enqueueTriggerRuleJob({ userId, ruleEventType: RuleEventType.PageCreated, libraryItemId, + data, }) } // queue export item job @@ -112,17 +107,17 @@ export const createPubSubClient = (): PubsubClient => { }, entityUpdated: async >( type: EntityType, - data: EntityData, - userId: string + data: T, + userId: string, + libraryItemId: string ): Promise => { - const libraryItemId = data.libraryItemId - // queue trigger rule job if (type === EntityType.PAGE) { await enqueueTriggerRuleJob({ userId, ruleEventType: RuleEventType.PageUpdated, libraryItemId, + data, }) } // queue export item job @@ -178,13 +173,15 @@ export interface PubsubClient { ) => Promise entityCreated: >( type: EntityType, - data: EntityData, - userId: string + data: T, + userId: string, + libraryItemId: string ) => Promise entityUpdated: >( type: EntityType, - data: EntityData, - userId: string + data: T, + userId: string, + libraryItemId: string ) => Promise entityDeleted: (type: EntityType, id: string, userId: string) => Promise reportSubmitted( diff --git a/packages/api/src/services/highlights.ts b/packages/api/src/services/highlights.ts index 5883fb7cb..babf4d4a4 100644 --- a/packages/api/src/services/highlights.ts +++ b/packages/api/src/services/highlights.ts @@ -59,7 +59,8 @@ export const createHighlight = async ( await pubsub.entityCreated( EntityType.HIGHLIGHT, { ...newHighlight, pageId: libraryItemId }, - userId + userId, + libraryItemId ) await enqueueUpdateHighlight({ @@ -106,7 +107,8 @@ export const mergeHighlights = async ( await pubsub.entityCreated( EntityType.HIGHLIGHT, { ...newHighlight, pageId: libraryItemId }, - userId + userId, + libraryItemId ) await enqueueUpdateHighlight({ @@ -139,8 +141,9 @@ export const updateHighlight = async ( const libraryItemId = updatedHighlight.libraryItem.id await pubsub.entityUpdated( EntityType.HIGHLIGHT, - { ...highlight, id: highlightId, pageId: libraryItemId, libraryItemId }, - userId + { ...highlight, id: highlightId, pageId: libraryItemId }, + userId, + libraryItemId ) await enqueueUpdateHighlight({ diff --git a/packages/api/src/services/labels.ts b/packages/api/src/services/labels.ts index 285338745..f4bae13b7 100644 --- a/packages/api/src/services/labels.ts +++ b/packages/api/src/services/labels.ts @@ -11,13 +11,11 @@ import { findHighlightById } from './highlights' import { findLibraryItemIdsByLabelId } from './library_item' type AddLabelsToLibraryItemEvent = { - libraryItemId: string pageId: string labels: DeepPartial