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..8a532eb07 100644 --- a/packages/api/src/jobs/trigger_rule.ts +++ b/packages/api/src/jobs/trigger_rule.ts @@ -1,27 +1,30 @@ import { ReadingProgressDataSource } from '../datasources/reading_progress_data_source' -import { LibraryItem, LibraryItemState } from '../entity/library_item' +import { LibraryItemState } from '../entity/library_item' import { Rule, RuleAction, RuleActionType, RuleEventType } from '../entity/rule' import { addLabelsToLibraryItem } from '../services/labels' import { - SearchArgs, - searchLibraryItems, + filterItemEvents, + ItemEvent, 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' +import { parseSearchQuery } from '../utils/search' export interface TriggerRuleJobData { libraryItemId: string userId: string ruleEventType: RuleEventType + data: ItemEvent } interface RuleActionObj { + libraryItemId: string userId: string action: RuleAction - libraryItem: LibraryItem + data: ItemEvent } type RuleActionFunc = (obj: RuleActionObj) => Promise @@ -33,19 +36,19 @@ const addLabels = async (obj: RuleActionObj) => { return addLabelsToLibraryItem( labelIds, - obj.libraryItem.id, + obj.libraryItemId, obj.userId, 'system' ) } const deleteLibraryItem = async (obj: RuleActionObj) => { - return softDeleteLibraryItem(obj.libraryItem.id, obj.userId) + return softDeleteLibraryItem(obj.libraryItemId, obj.userId) } const archivePage = async (obj: RuleActionObj) => { return updateLibraryItem( - obj.libraryItem.id, + obj.libraryItemId, { archivedAt: new Date(), state: LibraryItemState.Archived }, obj.userId, undefined, @@ -56,7 +59,7 @@ const archivePage = async (obj: RuleActionObj) => { const markPageAsRead = async (obj: RuleActionObj) => { return readingProgressDataSource.updateReadingProgress( obj.userId, - obj.libraryItem.id, + obj.libraryItemId, { readingProgressPercent: 100, readingProgressTopPercent: 100, @@ -66,15 +69,15 @@ const markPageAsRead = async (obj: RuleActionObj) => { } const sendNotification = async (obj: RuleActionObj) => { - const item = obj.libraryItem + const item = obj.data const message = { - title: item.author || item.siteName || 'Omnivore', - body: item.title, + title: item.author?.toString() || item.siteName?.toString() || 'Omnivore', + body: item.title?.toString(), image: item.thumbnail, } const data = { - folder: item.folder, - libraryItemId: item.id, + folder: item.folder?.toString() || 'inbox', + libraryItemId: obj.libraryItemId, } return sendPushNotifications(obj.userId, message, 'rule', data) @@ -96,36 +99,41 @@ const getRuleAction = (actionType: RuleActionType): RuleActionFunc => { } const triggerActions = async ( + libraryItemId: string, userId: string, rules: Rule[], - data: TriggerRuleJobData + data: ItemEvent ) => { const actionPromises: Promise[] = [] for (const rule of rules) { - const itemId = data.libraryItemId - const searchArgs: SearchArgs = { - includeContent: false, - includeDeleted: false, - includePending: false, - size: 1, - query: `(${rule.filter}) AND includes:${itemId}`, - } + let filteredData: ItemEvent + + try { + const ast = parseSearchQuery(rule.filter) + // filter library item by rule filter + const results = filterItemEvents(ast, [data]) + if (results.length === 0) { + logger.info(`No items found for rule ${rule.id}`) + continue + } + + filteredData = results[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) - const libraryItems = await searchLibraryItems(searchArgs, userId) - if (libraryItems.count === 0) { - logger.info(`No pages found for rule ${rule.id}`) continue } - const libraryItem = libraryItems.libraryItems[0] - for (const action of rule.actions) { const actionFunc = getRuleAction(action.type) const actionObj: RuleActionObj = { + libraryItemId, userId, action, - libraryItem, + data: filteredData, } actionPromises.push(actionFunc(actionObj)) @@ -139,8 +147,8 @@ const triggerActions = async ( } } -export const triggerRule = async (data: TriggerRuleJobData) => { - const { userId, ruleEventType } = data +export const triggerRule = async (jobData: TriggerRuleJobData) => { + const { userId, ruleEventType, data, libraryItemId } = jobData // get rules by calling api const rules = await findEnabledRules(userId, ruleEventType) @@ -149,7 +157,7 @@ export const triggerRule = async (data: TriggerRuleJobData) => { return false } - await triggerActions(userId, rules, data) + await triggerActions(libraryItemId, userId, rules, data) return true } diff --git a/packages/api/src/pubsub.ts b/packages/api/src/pubsub.ts index f13bfe5e4..c64b45900 100644 --- a/packages/api/src/pubsub.ts +++ b/packages/api/src/pubsub.ts @@ -3,32 +3,19 @@ import express from 'express' import { RuleEventType } from './entity/rule' import { env } from './env' import { ReportType } from './generated/graphql' -import { Merge } from './util' +import { FeatureName, findFeatureByName } from './services/features' import { - enqueueAISummarizeJob, enqueueExportItem, enqueueProcessYouTubeVideo, 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') const client = new PubSub() -type EntityData> = Merge< - T, - { libraryItemId: string } -> - const isYouTubeVideoURL = (url: string | undefined): boolean => { if (!url) { return false @@ -42,8 +29,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()}`) @@ -75,16 +60,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 @@ -93,11 +79,6 @@ export const createPubSubClient = (): PubsubClient => { libraryItemIds: [libraryItemId], }) - const cleanData = deepDelete( - data as EntityData & Record, - [...fieldsToDelete] - ) - await enqueueWebhookJob({ userId, type, @@ -112,34 +93,30 @@ export const createPubSubClient = (): PubsubClient => { // }) } - if ( - 'originalUrl' in data && - isYouTubeVideoURL(data['originalUrl'] as string | undefined) - ) { + const isYoutubeVideo = (data: any): data is { originalUrl: string } => { + return 'originalUrl' in data + } + + if (isYoutubeVideo(data) && isYouTubeVideoURL(data['originalUrl'])) { await enqueueProcessYouTubeVideo({ userId, libraryItemId, }) } - - return publish( - 'entityCreated', - Buffer.from(JSON.stringify({ type, userId, ...cleanData })) - ) }, 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 @@ -148,32 +125,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, @@ -207,13 +172,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/resolvers/discover_feeds/add.ts b/packages/api/src/resolvers/discover_feeds/add.ts index 7cc72d729..5420bc8b5 100644 --- a/packages/api/src/resolvers/discover_feeds/add.ts +++ b/packages/api/src/resolvers/discover_feeds/add.ts @@ -2,7 +2,11 @@ /* eslint-disable @typescript-eslint/no-unsafe-member-access */ /* eslint-disable @typescript-eslint/no-unsafe-assignment */ /* eslint-disable @typescript-eslint/require-await */ -import { authorized } from '../../utils/gql-utils' +import axios from 'axios' +import { XMLParser } from 'fast-xml-parser' +import { QueryRunner } from 'typeorm' +import { v4 } from 'uuid' +import { appDataSource } from '../../data_source' import { AddDiscoverFeedError, AddDiscoverFeedErrorCode, @@ -10,13 +14,9 @@ import { DiscoverFeed, MutationAddDiscoverFeedArgs, } from '../../generated/graphql' -import { appDataSource } from '../../data_source' -import { QueryRunner } from 'typeorm' -import axios from 'axios' -import { RSS_PARSER_CONFIG } from '../../utils/parser' -import { XMLParser } from 'fast-xml-parser' import { EntityType } from '../../pubsub' -import { v4 } from 'uuid' +import { authorized } from '../../utils/gql-utils' +import { RSS_PARSER_CONFIG } from '../../utils/parser' const parser = new XMLParser({ ignoreAttributes: false, @@ -183,13 +183,14 @@ export const addDiscoverFeedResolver = authorized< } const result = await addNewSubscription(queryRunner, url, uid) - if (result.__typename == 'AddDiscoverFeedSuccess') { - await pubsub.entityCreated( - EntityType.RSS_FEED, - { feed: result.feed, libraryItemId: 'NA' }, - uid - ) - } + // TODO: Add pubsub for new feed + // if (result.__typename == 'AddDiscoverFeedSuccess') { + // await pubsub.entityCreated( + // EntityType.RSS_FEED, + // { feed: result.feed, libraryItemId: 'NA' }, + // uid + // ) + // } return result } catch (error) { 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, 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/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/integrations/notion.ts b/packages/api/src/services/integrations/notion.ts index 4d823e031..7e891287d 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 = @@ -256,7 +256,7 @@ export class NotionClient implements IntegrationClient { text: { content: highlight.quote || '', link: { - url: highlightUrl(item.slug, highlight.id), + url: getHighlightUrl(item.slug, highlight.id), }, }, annotations: { 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/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