rss-reader.js 33 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899
  1. /**
  2. * @node rss-reader
  3. * @name RSS Reader
  4. * @category data
  5. * @version 1.2.0
  6. * @description Read and parse RSS/Atom feeds with optional filtering by keyword, date, or count. Supports stateful new item detection.
  7. * @icon rss
  8. */
  9. const configSchema = {
  10. type: 'object',
  11. properties: {
  12. url: {
  13. type: 'string',
  14. title: 'Feed URL',
  15. description: 'URL of the RSS or Atom feed'
  16. },
  17. filterKeyword: {
  18. type: 'string',
  19. title: 'Keyword Filter',
  20. description: 'Filter items by keyword (matches title and description)'
  21. },
  22. filterNewerThan: {
  23. type: 'string',
  24. title: 'Newer Than',
  25. description: 'Only include items newer than this date (ISO 8601 format, e.g., 2024-01-15T00:00:00Z)'
  26. },
  27. maxItems: {
  28. type: 'number',
  29. title: 'Max Items',
  30. description: 'Maximum number of items to return (0 = unlimited)',
  31. default: 0
  32. },
  33. userAgent: {
  34. type: 'string',
  35. title: 'User-Agent',
  36. description: 'Sent with the request. Several feeds refuse an unnamed client outright - Reddit answers 403 - and some rate-limit a browser string harder than a named one. Identify yourself here',
  37. default: 'smartbotic-rss/1.0 (+https://smartbotics.ai)'
  38. },
  39. timeout: {
  40. type: 'number',
  41. title: 'Timeout (ms)',
  42. default: 30000
  43. },
  44. retries: {
  45. type: 'number', title: 'Retries', default: 2,
  46. description: 'Extra attempts when the feed times out or answers 429, 500, 502, 503 or 504. A feed that is briefly unreachable should not end a scheduled run'
  47. },
  48. retryDelayMs: {
  49. type: 'number', title: 'Retry Delay (ms)', default: 2000,
  50. description: 'Wait before the first retry. Doubles on each further attempt, and a Retry-After from the server wins over both'
  51. },
  52. retryMaxDelayMs: {
  53. type: 'number', title: 'Max Retry Delay (ms)', default: 30000,
  54. description: 'Upper bound for the backoff'
  55. },
  56. skipOnError: {
  57. type: 'boolean', title: 'Skip On Error', default: false,
  58. description: 'Return success false and no items instead of failing the workflow. For a workflow reading several feeds, this is what lets the others carry on when one is unreachable - without it the first bad feed ends the run'
  59. },
  60. skipTlsVerify: {
  61. type: 'boolean', title: 'Skip TLS certificate verification', default: false,
  62. description: 'Accept the feed even if its certificate has expired or does not chain to a ' +
  63. 'trusted root. The hostname is still checked, so this cannot be tricked into ' +
  64. 'accepting a certificate for a different site - only a broken one for the right ' +
  65. 'site. An attacker who can intercept this connection can impersonate the feed ' +
  66. 'host while this is on. Use it only for a feed you control, or a certificate ' +
  67. 'failure you have personally checked - not as a standing setting.'
  68. },
  69. detectNewItems: {
  70. type: 'boolean',
  71. title: 'Detect New Items',
  72. description: 'When enabled, only output items that have not been seen in previous executions',
  73. default: false
  74. },
  75. stateCollection: {
  76. type: 'string',
  77. title: 'State Collection',
  78. description: 'Select a read-write collection to store seen item identifiers',
  79. dynamicOptions: {
  80. source: 'storage.collections',
  81. labelField: 'name',
  82. valueField: 'name',
  83. filter: { access: 'read-write' }
  84. },
  85. showWhen: { field: 'detectNewItems', value: true }
  86. }
  87. },
  88. required: ['url']
  89. };
  90. const inputSchema = {
  91. type: 'object',
  92. properties: {
  93. url: {
  94. type: 'string',
  95. description: 'Override feed URL from input'
  96. },
  97. filterKeyword: {
  98. type: 'string',
  99. description: 'Override keyword filter from input'
  100. },
  101. filterNewerThan: {
  102. type: 'string',
  103. description: 'Override date filter from input'
  104. },
  105. maxItems: {
  106. type: 'number',
  107. description: 'Override max items from input'
  108. }
  109. }
  110. };
  111. const outputSchema = {
  112. type: 'object',
  113. properties: {
  114. success: { type: 'boolean', description: 'False when the feed could not be read and Skip On Error let the run carry on' },
  115. error: { type: 'string', description: 'Why it could not be read, when it could not' },
  116. feedTitle: { type: 'string', description: 'Title of the feed' },
  117. feedLink: { type: 'string', description: 'Link to the feed homepage' },
  118. feedDescription: { type: 'string', description: 'Description of the feed' },
  119. itemCount: { type: 'number', description: 'Number of items returned' },
  120. newItemCount: { type: 'number', description: 'Items this run will process. With Detect New Items off that is every item returned, so a workflow can branch on this whichever mode it is in' },
  121. totalFeedItemCount: { type: 'number', description: 'Total items in feed before filtering (when detectNewItems is enabled)' },
  122. items: {
  123. type: 'array',
  124. description: 'Array of feed items',
  125. items: {
  126. type: 'object',
  127. properties: {
  128. title: { type: 'string', description: 'Item title' },
  129. link: { type: 'string', description: 'Item link/URL' },
  130. description: { type: 'string', description: 'Item summary/description' },
  131. content: { type: 'string', description: 'Full content (content:encoded)' },
  132. pubDate: { type: 'string', description: 'Publication date' },
  133. author: { type: 'string', description: 'Author name' },
  134. categories: { type: 'array', items: { type: 'string' }, description: 'Categories/tags' },
  135. guid: { type: 'string', description: 'Unique identifier' },
  136. comments: { type: 'string', description: 'Comments URL' },
  137. source: { type: 'string', description: 'Source attribution' },
  138. enclosure: {
  139. type: 'object',
  140. description: 'Media enclosure (attachment)',
  141. properties: {
  142. url: { type: 'string', description: 'Media URL' },
  143. type: { type: 'string', description: 'MIME type' },
  144. length: { type: 'number', description: 'File size in bytes' }
  145. }
  146. },
  147. media: {
  148. type: 'object',
  149. description: 'Media content (media:content)',
  150. properties: {
  151. url: { type: 'string' },
  152. type: { type: 'string' },
  153. width: { type: 'number' },
  154. height: { type: 'number' }
  155. }
  156. },
  157. thumbnail: {
  158. type: 'object',
  159. description: 'Thumbnail image (media:thumbnail)',
  160. properties: {
  161. url: { type: 'string' },
  162. width: { type: 'number' },
  163. height: { type: 'number' }
  164. }
  165. }
  166. }
  167. }
  168. }
  169. }
  170. };
  171. /**
  172. * Simple XML parser utilities for RSS/Atom feeds
  173. * Handles basic XML structure without external dependencies
  174. */
  175. function xmlGetTagContent(xml, tagName) {
  176. var pattern = new RegExp('<' + tagName + '[^>]*>([\\s\\S]*?)<\\/' + tagName + '>', 'i');
  177. var match = xml.match(pattern);
  178. if (match) {
  179. return match[1].trim();
  180. }
  181. return null;
  182. }
  183. function xmlGetAllTagContents(xml, tagName) {
  184. var results = [];
  185. var regex = new RegExp('<' + tagName + '[^>]*>([\\s\\S]*?)<\\/' + tagName + '>', 'gi');
  186. var match;
  187. while ((match = regex.exec(xml)) !== null) {
  188. results.push(match[1].trim());
  189. }
  190. return results;
  191. }
  192. function xmlGetTagAttributes(xml, tagName) {
  193. // Match self-closing tag or opening tag
  194. var regex = new RegExp('<' + tagName + '([^>]*?)\\/?>', 'i');
  195. var match = xml.match(regex);
  196. if (!match) return null;
  197. var attrString = match[1];
  198. var attrs = {};
  199. // Parse attributes: name="value" or name='value'
  200. var attrRegex = /(\w+)=["']([^"']*)["']/g;
  201. var attrMatch;
  202. while ((attrMatch = attrRegex.exec(attrString)) !== null) {
  203. attrs[attrMatch[1]] = attrMatch[2];
  204. }
  205. return attrs;
  206. }
  207. function xmlDecodeHtmlEntities(text) {
  208. if (!text) return '';
  209. return text
  210. .replace(/&lt;/g, '<')
  211. .replace(/&gt;/g, '>')
  212. .replace(/&amp;/g, '&')
  213. .replace(/&quot;/g, '"')
  214. .replace(/&apos;/g, "'")
  215. .replace(/&#(\d+);/g, function(_, dec) { return String.fromCharCode(dec); })
  216. .replace(/&#x([0-9a-f]+);/gi, function(_, hex) { return String.fromCharCode(parseInt(hex, 16)); });
  217. }
  218. function xmlStripCdata(text) {
  219. if (!text) return '';
  220. return text.replace(/<!\[CDATA\[([\s\S]*?)\]\]>/g, '$1');
  221. }
  222. function xmlCleanText(text) {
  223. if (!text) return '';
  224. return xmlDecodeHtmlEntities(xmlStripCdata(text)).trim();
  225. }
  226. /**
  227. * Parse RSS 2.0 feed
  228. */
  229. function parseRss(xml) {
  230. var channel = xmlGetTagContent(xml, 'channel');
  231. if (!channel) {
  232. throw new Error('Invalid RSS feed: no channel element found');
  233. }
  234. var feedTitle = xmlGetTagContent(channel, 'title');
  235. var feedLink = xmlGetTagContent(channel, 'link');
  236. var feedDesc = xmlGetTagContent(channel, 'description');
  237. var feedLanguage = xmlGetTagContent(channel, 'language');
  238. var feedPubDate = xmlGetTagContent(channel, 'pubDate');
  239. var feedLastBuildDate = xmlGetTagContent(channel, 'lastBuildDate');
  240. var feedGenerator = xmlGetTagContent(channel, 'generator');
  241. var feedImage = xmlGetTagContent(channel, 'image');
  242. var feed = {
  243. title: xmlCleanText(feedTitle) || '',
  244. link: xmlCleanText(feedLink) || '',
  245. description: xmlCleanText(feedDesc) || '',
  246. language: xmlCleanText(feedLanguage) || '',
  247. pubDate: xmlCleanText(feedPubDate) || '',
  248. lastBuildDate: xmlCleanText(feedLastBuildDate) || '',
  249. generator: xmlCleanText(feedGenerator) || '',
  250. items: []
  251. };
  252. // Parse feed image if present
  253. if (feedImage) {
  254. feed.image = {
  255. url: xmlCleanText(xmlGetTagContent(feedImage, 'url')) || '',
  256. title: xmlCleanText(xmlGetTagContent(feedImage, 'title')) || '',
  257. link: xmlCleanText(xmlGetTagContent(feedImage, 'link')) || ''
  258. };
  259. }
  260. var itemRegex = /<item[^>]*>([\s\S]*?)<\/item>/gi;
  261. var match;
  262. while ((match = itemRegex.exec(channel)) !== null) {
  263. var itemXml = match[1];
  264. var catArray = xmlGetAllTagContents(itemXml, 'category');
  265. var categories = [];
  266. for (var ci = 0; ci < catArray.length; ci++) {
  267. categories.push(xmlCleanText(catArray[ci]));
  268. }
  269. var item = {
  270. title: xmlCleanText(xmlGetTagContent(itemXml, 'title')) || '',
  271. link: xmlCleanText(xmlGetTagContent(itemXml, 'link')) || '',
  272. description: xmlCleanText(xmlGetTagContent(itemXml, 'description')) || '',
  273. content: xmlCleanText(xmlGetTagContent(itemXml, 'content:encoded')) || '',
  274. pubDate: xmlCleanText(xmlGetTagContent(itemXml, 'pubDate')) || '',
  275. author: xmlCleanText(xmlGetTagContent(itemXml, 'author') ||
  276. xmlGetTagContent(itemXml, 'dc:creator')) || '',
  277. categories: categories,
  278. guid: xmlCleanText(xmlGetTagContent(itemXml, 'guid')) || '',
  279. comments: xmlCleanText(xmlGetTagContent(itemXml, 'comments')) || '',
  280. source: xmlCleanText(xmlGetTagContent(itemXml, 'source')) || ''
  281. };
  282. // Parse enclosure (media attachment)
  283. var enclosureAttrs = xmlGetTagAttributes(itemXml, 'enclosure');
  284. if (enclosureAttrs) {
  285. item.enclosure = {
  286. url: enclosureAttrs.url || '',
  287. type: enclosureAttrs.type || '',
  288. length: enclosureAttrs.length ? parseInt(enclosureAttrs.length, 10) : 0
  289. };
  290. }
  291. // Parse media:content (common in media RSS)
  292. var mediaAttrs = xmlGetTagAttributes(itemXml, 'media:content');
  293. if (mediaAttrs) {
  294. item.media = {
  295. url: mediaAttrs.url || '',
  296. type: mediaAttrs.type || mediaAttrs.medium || '',
  297. width: mediaAttrs.width ? parseInt(mediaAttrs.width, 10) : 0,
  298. height: mediaAttrs.height ? parseInt(mediaAttrs.height, 10) : 0
  299. };
  300. }
  301. // Parse media:thumbnail
  302. var thumbAttrs = xmlGetTagAttributes(itemXml, 'media:thumbnail');
  303. if (thumbAttrs) {
  304. item.thumbnail = {
  305. url: thumbAttrs.url || '',
  306. width: thumbAttrs.width ? parseInt(thumbAttrs.width, 10) : 0,
  307. height: thumbAttrs.height ? parseInt(thumbAttrs.height, 10) : 0
  308. };
  309. }
  310. feed.items.push(item);
  311. }
  312. return feed;
  313. }
  314. /**
  315. * Parse Atom feed
  316. */
  317. function parseAtom(xml) {
  318. var feedContent = xmlGetTagContent(xml, 'feed');
  319. if (!feedContent) {
  320. throw new Error('Invalid Atom feed: no feed element found');
  321. }
  322. function getAtomLink(content, rel) {
  323. var linkRegex = /<link([^>]*)>/gi;
  324. var match;
  325. while ((match = linkRegex.exec(content)) !== null) {
  326. var attrs = match[1];
  327. var relMatch = attrs.match(/rel=["']([^"']*)["']/);
  328. var hrefMatch = attrs.match(/href=["']([^"']*)["']/);
  329. if (hrefMatch) {
  330. var linkRel = relMatch ? relMatch[1] : 'alternate';
  331. if (linkRel === rel || (!rel && linkRel === 'alternate')) {
  332. return hrefMatch[1];
  333. }
  334. }
  335. }
  336. return '';
  337. }
  338. function getAtomLinkByType(content, type) {
  339. var linkRegex = /<link([^>]*)>/gi;
  340. var match;
  341. while ((match = linkRegex.exec(content)) !== null) {
  342. var attrs = match[1];
  343. var typeMatch = attrs.match(/type=["']([^"']*)["']/);
  344. var hrefMatch = attrs.match(/href=["']([^"']*)["']/);
  345. if (hrefMatch && typeMatch && typeMatch[1].indexOf(type) !== -1) {
  346. return {
  347. url: hrefMatch[1],
  348. type: typeMatch[1]
  349. };
  350. }
  351. }
  352. return null;
  353. }
  354. var feed = {
  355. title: xmlCleanText(xmlGetTagContent(feedContent, 'title')) || '',
  356. link: getAtomLink(feedContent, 'alternate') || getAtomLink(feedContent, null) || '',
  357. description: xmlCleanText(xmlGetTagContent(feedContent, 'subtitle')) || '',
  358. language: '',
  359. pubDate: xmlCleanText(xmlGetTagContent(feedContent, 'updated')) || '',
  360. lastBuildDate: '',
  361. generator: xmlCleanText(xmlGetTagContent(feedContent, 'generator')) || '',
  362. items: []
  363. };
  364. // Parse feed icon/logo
  365. var feedIcon = xmlGetTagContent(feedContent, 'icon');
  366. var feedLogo = xmlGetTagContent(feedContent, 'logo');
  367. if (feedIcon || feedLogo) {
  368. feed.image = {
  369. url: xmlCleanText(feedLogo || feedIcon) || '',
  370. title: feed.title,
  371. link: feed.link
  372. };
  373. }
  374. var entryRegex = /<entry[^>]*>([\s\S]*?)<\/entry>/gi;
  375. var match;
  376. while ((match = entryRegex.exec(feedContent)) !== null) {
  377. var entryXml = match[1];
  378. var categories = [];
  379. var catRegex = /<category([^>]*)>/gi;
  380. var catMatch;
  381. while ((catMatch = catRegex.exec(entryXml)) !== null) {
  382. var termMatch = catMatch[1].match(/term=["']([^"']*)["']/);
  383. if (termMatch) {
  384. categories.push(xmlCleanText(termMatch[1]));
  385. }
  386. }
  387. var author = '';
  388. var authorContent = xmlGetTagContent(entryXml, 'author');
  389. if (authorContent) {
  390. author = xmlCleanText(xmlGetTagContent(authorContent, 'name')) || '';
  391. }
  392. var pubDate = xmlCleanText(
  393. xmlGetTagContent(entryXml, 'published') ||
  394. xmlGetTagContent(entryXml, 'updated')
  395. ) || '';
  396. var summary = xmlCleanText(xmlGetTagContent(entryXml, 'summary')) || '';
  397. var content = xmlCleanText(xmlGetTagContent(entryXml, 'content')) || '';
  398. var item = {
  399. title: xmlCleanText(xmlGetTagContent(entryXml, 'title')) || '',
  400. link: getAtomLink(entryXml, 'alternate') || getAtomLink(entryXml, null) || '',
  401. description: summary || content,
  402. content: content,
  403. pubDate: pubDate,
  404. author: author,
  405. categories: categories,
  406. guid: xmlCleanText(xmlGetTagContent(entryXml, 'id')) || '',
  407. comments: '',
  408. source: ''
  409. };
  410. // Check for enclosure link
  411. var enclosureLink = getAtomLinkByType(entryXml, 'image');
  412. if (!enclosureLink) {
  413. enclosureLink = getAtomLinkByType(entryXml, 'audio');
  414. }
  415. if (!enclosureLink) {
  416. enclosureLink = getAtomLinkByType(entryXml, 'video');
  417. }
  418. if (enclosureLink) {
  419. item.enclosure = {
  420. url: enclosureLink.url,
  421. type: enclosureLink.type,
  422. length: 0
  423. };
  424. }
  425. // Parse media:content
  426. var mediaAttrs = xmlGetTagAttributes(entryXml, 'media:content');
  427. if (mediaAttrs) {
  428. item.media = {
  429. url: mediaAttrs.url || '',
  430. type: mediaAttrs.type || mediaAttrs.medium || '',
  431. width: mediaAttrs.width ? parseInt(mediaAttrs.width, 10) : 0,
  432. height: mediaAttrs.height ? parseInt(mediaAttrs.height, 10) : 0
  433. };
  434. }
  435. // Parse media:thumbnail
  436. var thumbAttrs = xmlGetTagAttributes(entryXml, 'media:thumbnail');
  437. if (thumbAttrs) {
  438. item.thumbnail = {
  439. url: thumbAttrs.url || '',
  440. width: thumbAttrs.width ? parseInt(thumbAttrs.width, 10) : 0,
  441. height: thumbAttrs.height ? parseInt(thumbAttrs.height, 10) : 0
  442. };
  443. }
  444. feed.items.push(item);
  445. }
  446. return feed;
  447. }
  448. /**
  449. * Parse date string to timestamp
  450. */
  451. function parseDate(dateStr) {
  452. if (!dateStr) return 0;
  453. var timestamp = Date.parse(dateStr);
  454. if (!isNaN(timestamp)) {
  455. return timestamp;
  456. }
  457. var rfc822Regex = /(\d{1,2})\s+(\w{3})\s+(\d{4})\s+(\d{2}):(\d{2}):(\d{2})/;
  458. var match = dateStr.match(rfc822Regex);
  459. if (match) {
  460. var months = {
  461. Jan: 0, Feb: 1, Mar: 2, Apr: 3, May: 4, Jun: 5,
  462. Jul: 6, Aug: 7, Sep: 8, Oct: 9, Nov: 10, Dec: 11
  463. };
  464. var day = parseInt(match[1], 10);
  465. var month = months[match[2]];
  466. var year = parseInt(match[3], 10);
  467. var hour = parseInt(match[4], 10);
  468. var minute = parseInt(match[5], 10);
  469. var second = parseInt(match[6], 10);
  470. if (month !== undefined) {
  471. return new Date(year, month, day, hour, minute, second).getTime();
  472. }
  473. }
  474. return 0;
  475. }
  476. /**
  477. * Filter items by keyword
  478. */
  479. function filterByKeyword(items, keyword) {
  480. if (!keyword) return items;
  481. var lowerKeyword = keyword.toLowerCase();
  482. var result = [];
  483. for (var i = 0; i < items.length; i++) {
  484. var item = items[i];
  485. var title = (item.title || '').toLowerCase();
  486. var description = (item.description || '').toLowerCase();
  487. if (title.indexOf(lowerKeyword) !== -1 || description.indexOf(lowerKeyword) !== -1) {
  488. result.push(item);
  489. }
  490. }
  491. return result;
  492. }
  493. /**
  494. * Filter items by date
  495. */
  496. function filterByDate(items, newerThanStr) {
  497. if (!newerThanStr) return items;
  498. var threshold = parseDate(newerThanStr);
  499. if (threshold === 0) {
  500. smartbotic.log.warn('Could not parse date filter: ' + newerThanStr);
  501. return items;
  502. }
  503. var result = [];
  504. for (var i = 0; i < items.length; i++) {
  505. var item = items[i];
  506. var itemDate = parseDate(item.pubDate);
  507. if (itemDate > threshold) {
  508. result.push(item);
  509. }
  510. }
  511. return result;
  512. }
  513. /**
  514. * Limit number of items
  515. */
  516. function limitItems(items, maxItems) {
  517. if (!maxItems || maxItems <= 0) return items;
  518. return items.slice(0, maxItems);
  519. }
  520. /**
  521. * Get unique identifier for an RSS/Atom item
  522. * Prefers GUID, falls back to link
  523. */
  524. function getItemIdentifier(item) {
  525. return item.guid || item.link || '';
  526. }
  527. /**
  528. * Generate a stable document ID from feed URL for state storage
  529. */
  530. function generateStateDocId(feedUrl) {
  531. // Create a simple hash from the feed URL for a stable doc ID
  532. var hash = 0;
  533. for (var i = 0; i < feedUrl.length; i++) {
  534. var char = feedUrl.charCodeAt(i);
  535. hash = ((hash << 5) - hash) + char;
  536. hash = hash & hash; // Convert to 32-bit integer
  537. }
  538. return 'rss-state-' + Math.abs(hash).toString(36);
  539. }
  540. /**
  541. * Load seen item identifiers from storage
  542. */
  543. function loadSeenItems(collection, feedUrl) {
  544. var docId = generateStateDocId(feedUrl);
  545. try {
  546. var result = smartbotic.storage.get(collection, docId);
  547. if (result.found && result.document && result.document.seenIds) {
  548. // Use plain object as a set (for QuickJS compatibility)
  549. var seenObj = {};
  550. var storedIds = result.document.seenIds;
  551. // Handle both array format and object format (in case of legacy data)
  552. if (Array.isArray(storedIds)) {
  553. smartbotic.log.debug('Loaded ' + storedIds.length + ' seen item IDs from storage');
  554. for (var i = 0; i < storedIds.length; i++) {
  555. seenObj[storedIds[i]] = true;
  556. }
  557. } else if (typeof storedIds === 'object' && storedIds !== null) {
  558. // Object format - keys are the IDs
  559. var keys = Object.keys(storedIds);
  560. smartbotic.log.debug('Loaded ' + keys.length + ' seen item IDs from storage (object format)');
  561. for (var j = 0; j < keys.length; j++) {
  562. seenObj[keys[j]] = true;
  563. }
  564. }
  565. return {
  566. seenIds: seenObj,
  567. docId: docId,
  568. exists: true,
  569. version: result.document._version || 1
  570. };
  571. }
  572. } catch (error) {
  573. smartbotic.log.warn('Error loading seen items from storage: ' + error.message);
  574. }
  575. return {
  576. seenIds: {},
  577. docId: docId,
  578. exists: false,
  579. version: 0
  580. };
  581. }
  582. /**
  583. * Save seen item identifiers to storage
  584. */
  585. function saveSeenItems(collection, docId, seenIds, feedUrl, exists, version) {
  586. // Convert object keys to array
  587. var seenArray = Object.keys(seenIds);
  588. var documentData = {
  589. feedUrl: feedUrl,
  590. seenIds: seenArray,
  591. lastUpdated: new Date().toISOString(),
  592. itemCount: seenArray.length
  593. };
  594. try {
  595. if (exists) {
  596. // Update existing document
  597. var updateResult = smartbotic.storage.update(collection, docId, documentData, version, false);
  598. if (!updateResult.success) {
  599. smartbotic.log.error('Failed to update seen items: ' + updateResult.error);
  600. return false;
  601. }
  602. smartbotic.log.debug('Updated ' + seenArray.length + ' seen item IDs in storage');
  603. } else {
  604. // Insert new document
  605. var insertResult = smartbotic.storage.insert(collection, documentData, docId, 0);
  606. if (!insertResult.success) {
  607. // If insert fails because document exists, try update instead
  608. if (insertResult.error && insertResult.error.indexOf('already exists') !== -1) {
  609. smartbotic.log.debug('Document exists, trying update instead');
  610. var retryResult = smartbotic.storage.update(collection, docId, documentData, 0, false);
  611. if (!retryResult.success) {
  612. smartbotic.log.error('Failed to update seen items on retry: ' + retryResult.error);
  613. return false;
  614. }
  615. smartbotic.log.debug('Updated ' + seenArray.length + ' seen item IDs in storage (retry)');
  616. } else {
  617. smartbotic.log.error('Failed to save seen items: ' + insertResult.error);
  618. return false;
  619. }
  620. } else {
  621. smartbotic.log.debug('Saved ' + seenArray.length + ' seen item IDs to storage');
  622. }
  623. }
  624. return true;
  625. } catch (error) {
  626. smartbotic.log.error('Error saving seen items to storage: ' + error.message);
  627. return false;
  628. }
  629. }
  630. /**
  631. * Filter items to only include new (unseen) items
  632. */
  633. function filterNewItems(items, seenIds) {
  634. var newItems = [];
  635. for (var i = 0; i < items.length; i++) {
  636. var id = getItemIdentifier(items[i]);
  637. if (id && !seenIds[id]) {
  638. newItems.push(items[i]);
  639. }
  640. }
  641. return newItems;
  642. }
  643. async function readFeed(config, input) {
  644. var url = input.url || config.url;
  645. var filterKeyword = input.filterKeyword || config.filterKeyword;
  646. var filterNewerThan = input.filterNewerThan || config.filterNewerThan;
  647. var maxItems = input.maxItems !== undefined ? input.maxItems : config.maxItems;
  648. var timeout = config.timeout || 30000;
  649. var detectNewItems = config.detectNewItems || false;
  650. var stateCollection = config.stateCollection;
  651. if (!url) {
  652. throw new Error('Feed URL is required');
  653. }
  654. // Validate stateful mode configuration
  655. if (detectNewItems && !stateCollection) {
  656. throw new Error('State collection is required when "Detect New Items" is enabled. Please select a read-write collection in the workflow storage settings.');
  657. }
  658. // No User-Agent at all is what several feeds reject, Reddit among them, so
  659. // there is always one.
  660. var userAgent = String(config.userAgent || '').trim() ||
  661. 'smartbotic-rss/1.0 (+https://smartbotics.ai)';
  662. smartbotic.log.info('Fetching RSS/Atom feed from ' + url);
  663. var response = smartbotic.http.request({
  664. method: 'GET',
  665. url: url,
  666. timeout: timeout,
  667. // Reading a feed is safe to repeat, and a feed that times out once is
  668. // the usual reason a scheduled run dies - the request is retried rather
  669. // than the whole pass being thrown away. Handled by the HTTP helper, so
  670. // the waiting happens without this node holding a loop.
  671. retries: config.retries !== undefined && config.retries !== null && config.retries !== ''
  672. ? Number(config.retries) : 2,
  673. retryDelayMs: Number(config.retryDelayMs) || 2000,
  674. retryMaxDelayMs: Number(config.retryMaxDelayMs) || 30000,
  675. skipTlsVerify: config.skipTlsVerify === true,
  676. headers: {
  677. 'Accept': 'application/rss+xml, application/atom+xml, application/xml, text/xml, */*',
  678. 'User-Agent': userAgent
  679. }
  680. });
  681. if (response.status < 200 || response.status >= 300) {
  682. var hint = '';
  683. if (response.status === 403) {
  684. hint = '. The feed refused this client - most feeds that do this want a descriptive ' +
  685. 'User-Agent, which is the User-Agent setting on this node';
  686. } else if (response.status === 429) {
  687. hint = '. The feed is rate limiting this address - poll it less often, or from ' +
  688. 'somewhere else. Reddit throttles hard and counts every request, including ' +
  689. 'failed ones';
  690. } else if (response.status === 404) {
  691. hint = '. Check the address: ' + url;
  692. }
  693. throw new Error('Failed to fetch feed: HTTP ' + response.status + hint);
  694. }
  695. var xml = typeof response.data === 'string' ? response.data : JSON.stringify(response.data);
  696. if (!xml || xml.trim().length === 0) {
  697. throw new Error('Empty response from feed URL');
  698. }
  699. var feed;
  700. if (xml.indexOf('<feed') !== -1 && xml.indexOf('xmlns') !== -1 && xml.indexOf('atom') !== -1) {
  701. smartbotic.log.debug('Detected Atom feed format');
  702. feed = parseAtom(xml);
  703. } else if (xml.indexOf('<rss') !== -1 || xml.indexOf('<channel') !== -1) {
  704. smartbotic.log.debug('Detected RSS feed format');
  705. feed = parseRss(xml);
  706. } else {
  707. throw new Error('Unknown feed format: neither RSS nor Atom detected');
  708. }
  709. var items = feed.items;
  710. var totalFeedItemCount = items.length;
  711. // Apply standard filters first (keyword, date)
  712. if (filterKeyword) {
  713. smartbotic.log.debug('Filtering by keyword: ' + filterKeyword);
  714. items = filterByKeyword(items, filterKeyword);
  715. }
  716. if (filterNewerThan) {
  717. smartbotic.log.debug('Filtering items newer than: ' + filterNewerThan);
  718. items = filterByDate(items, filterNewerThan);
  719. }
  720. // Stateful new item detection
  721. var newItemCount = items.length;
  722. var stateInfo = null;
  723. if (detectNewItems) {
  724. smartbotic.log.info('Detecting new items using collection: ' + stateCollection);
  725. // Load previously seen items
  726. stateInfo = loadSeenItems(stateCollection, url);
  727. // Filter to only new items
  728. var newItems = filterNewItems(items, stateInfo.seenIds);
  729. newItemCount = newItems.length;
  730. smartbotic.log.info('Found ' + newItemCount + ' new items out of ' + items.length + ' filtered items');
  731. // Update seen IDs with all current items (both new and previously seen)
  732. // This ensures we track all items from the current feed
  733. for (var i = 0; i < items.length; i++) {
  734. var id = getItemIdentifier(items[i]);
  735. if (id) {
  736. stateInfo.seenIds[id] = true;
  737. }
  738. }
  739. // Use only new items for output
  740. items = newItems;
  741. }
  742. // Apply max items limit last
  743. if (maxItems && maxItems > 0) {
  744. smartbotic.log.debug('Limiting to ' + maxItems + ' items');
  745. items = limitItems(items, maxItems);
  746. }
  747. smartbotic.log.info('Returning ' + items.length + ' items from feed');
  748. // Persist seen items state after successful processing
  749. if (detectNewItems && stateInfo) {
  750. var saved = saveSeenItems(
  751. stateCollection,
  752. stateInfo.docId,
  753. stateInfo.seenIds,
  754. url,
  755. stateInfo.exists,
  756. stateInfo.version
  757. );
  758. if (!saved) {
  759. smartbotic.log.warn('Failed to persist seen items state, but continuing with output');
  760. }
  761. }
  762. var result = {
  763. feedTitle: feed.title,
  764. feedLink: feed.link,
  765. feedDescription: feed.description,
  766. itemCount: items.length,
  767. items: items
  768. };
  769. // Always reported, in both modes.
  770. //
  771. // These used to appear only when Detect New Items was on, so a workflow
  772. // branching on newItemCount > 0 quietly started taking the other branch the
  773. // moment someone turned detection off: the field was not zero, it was gone,
  774. // and a comparison against a missing field is false. With detection off
  775. // every item returned is one this run will process, so newItemCount is the
  776. // item count - which is what a workflow asking "is there anything to do?"
  777. // means either way.
  778. result.newItemCount = detectNewItems ? newItemCount : items.length;
  779. result.totalFeedItemCount = detectNewItems ? totalFeedItemCount : items.length;
  780. return result;
  781. }
  782. async function execute(config, input) {
  783. try {
  784. var result = await readFeed(config, input);
  785. result.success = true;
  786. result.error = '';
  787. return result;
  788. } catch (err) {
  789. if (config.skipOnError !== true) {
  790. throw err;
  791. }
  792. var message = (err && err.message) ? err.message : String(err);
  793. smartbotic.log.warn('RSS Reader: ' + (config.url || '') + ' could not be read, skipping: ' + message);
  794. // Shaped like a successful read with nothing in it, so a merge or a
  795. // loop downstream does not have to know the difference - and carrying
  796. // the reason means a run that produced nothing can still say why.
  797. return {
  798. success: false,
  799. error: message,
  800. feedTitle: '',
  801. feedLink: '',
  802. feedDescription: '',
  803. itemCount: 0,
  804. items: [],
  805. newItemCount: 0,
  806. totalFeedItemCount: 0
  807. };
  808. }
  809. }
  810. module.exports = { configSchema, inputSchema, outputSchema, execute };