database-change.js 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201
  1. /**
  2. * @node database-change
  3. * @name Database Change
  4. * @category triggers
  5. * @version 1.0.0
  6. * @description Report documents added or changed in a collection since the last run, paired with a schedule trigger for its cadence
  7. * @icon database
  8. */
  9. const configSchema = {
  10. type: 'object',
  11. properties: {
  12. collection: {
  13. dynamicOptions: {
  14. source: 'storage.collections',
  15. labelField: 'name',
  16. valueField: 'name',
  17. filter: {}
  18. },
  19. type: 'string',
  20. title: 'Collection',
  21. description: 'Collection to watch'
  22. },
  23. timestampField: {
  24. type: 'string',
  25. title: 'Timestamp Field',
  26. description: 'Field holding the last-modified time. Defaults to the _updated_at the database maintains itself, in nanoseconds. A document with no readable value in this field is never reported',
  27. default: '_updated_at'
  28. },
  29. filter: {
  30. type: 'object',
  31. title: 'Filter',
  32. description: 'Optional query restricting which documents are watched',
  33. additionalProperties: true
  34. },
  35. maxDocuments: {
  36. type: 'number',
  37. title: 'Max Documents',
  38. description: 'Most documents to report in one run, so a large backlog does not arrive as one enormous payload',
  39. default: 100
  40. },
  41. emitOnFirstRun: {
  42. type: 'boolean',
  43. title: 'Report Everything On First Run',
  44. description: 'Off records the current high-water mark and reports nothing, so adding this to a live workflow does not fire for every document already there',
  45. default: false
  46. },
  47. cursorCollection: {
  48. dynamicOptions: {
  49. source: 'storage.collections',
  50. labelField: 'name',
  51. valueField: 'name',
  52. filter: { access: 'read-write' }
  53. },
  54. type: 'string',
  55. title: 'Cursor Collection',
  56. default: 'watch_cursors'
  57. }
  58. },
  59. required: ['collection']
  60. };
  61. const inputSchema = {
  62. type: 'object',
  63. properties: {
  64. data: { type: 'any' }
  65. }
  66. };
  67. const outputSchema = {
  68. type: 'object',
  69. properties: {
  70. documents: { type: 'array', description: 'Documents new or changed since the last run' },
  71. count: { type: 'number' },
  72. isFirstRun: { type: 'boolean' },
  73. cursor: { type: 'number', description: 'High-water timestamp stored for the next run' }
  74. }
  75. };
  76. async function execute(config, input, context) {
  77. const collection = config.collection;
  78. if (!collection) {
  79. throw new Error('Database Change: a collection is required');
  80. }
  81. const timestampField = config.timestampField || '_updated_at';
  82. const cursorCollection = config.cursorCollection || 'watch_cursors';
  83. const maxDocuments = Number(config.maxDocuments) || 100;
  84. const workflowId = (context && context.workflowId) || 'unknown';
  85. const nodeId = (context && context.nodeId) || 'database-change';
  86. const cursorId = workflowId + ':' + nodeId;
  87. const stored = smartbotic.storage.get(cursorCollection, cursorId);
  88. // A missing document is the expected shape of a first run. Anything else
  89. // that keeps storage.get from returning the cursor - no read access to
  90. // the collection, a connection error - is a real failure and must not be
  91. // swallowed into a false "first run" that then quietly overwrites state.
  92. // This mirrors nodes/triggers/file-watch.js so the two nodes agree on
  93. // what "first run" means.
  94. const cursorMissing = stored && typeof stored.error === 'string' &&
  95. stored.error.indexOf('Document not found') === 0;
  96. if (stored && stored.found !== true && !cursorMissing) {
  97. throw new Error('Database Change: could not read the cursor for ' + collection + ': ' +
  98. (stored.error || 'unknown storage error'));
  99. }
  100. const isFirstRun = !stored || stored.found !== true;
  101. const since = (!isFirstRun && stored.document && Number(stored.document.since)) || 0;
  102. const query = config.filter && typeof config.filter === 'object' ? config.filter : {};
  103. const result = smartbotic.storage.query(collection, query);
  104. // storage.query signals failure with an "error" field, not a "success"
  105. // flag - there is no success flag in its response at all, so checking
  106. // one would throw on every call, successful or not.
  107. if (!result || typeof result.error === 'string') {
  108. throw new Error('Database Change: could not query ' + collection + ': ' +
  109. ((result && result.error) || 'unknown error'));
  110. }
  111. const documents = result.documents || [];
  112. // The database stores _updated_at in nanoseconds while everything a node
  113. // sees is milliseconds, so the value is normalised rather than compared
  114. // against a cursor in different units.
  115. function stampOf(document) {
  116. const raw = Number(smartbotic.utils.getFieldValue(document, timestampField));
  117. if (!raw || isNaN(raw)) {
  118. return 0;
  119. }
  120. return raw > 1e15 ? Math.floor(raw / 1000000) : raw;
  121. }
  122. let highWater = since;
  123. let noStamp = 0;
  124. for (const document of documents) {
  125. const stamp = stampOf(document);
  126. if (stamp === 0) {
  127. noStamp++;
  128. }
  129. if (stamp > highWater) {
  130. highWater = stamp;
  131. }
  132. }
  133. if (noStamp > 0) {
  134. // A document with no readable value in timestampField sorts as if it
  135. // predates everything, and once the mark advances past 0 it can never
  136. // satisfy the strict > filter below - it is never reported, silently,
  137. // unless this is logged.
  138. smartbotic.log.warn('Database Change: ' + noStamp + ' document(s) in ' + collection +
  139. ' had no readable "' + timestampField + '" value and will never be reported');
  140. }
  141. let changed = [];
  142. if (isFirstRun && config.emitOnFirstRun !== true) {
  143. smartbotic.log.info('Database Change: first run on ' + collection + ', recorded the mark at ' +
  144. highWater + ' without reporting ' + documents.length + ' documents');
  145. } else {
  146. changed = documents
  147. .filter(function (document) { return stampOf(document) > since; })
  148. .sort(function (left, right) { return stampOf(left) - stampOf(right); });
  149. if (changed.length > maxDocuments) {
  150. let cut = maxDocuments;
  151. const boundary = stampOf(changed[cut - 1]);
  152. // Cutting through a group that shares one timestamp would strand
  153. // the rest of that group: the mark advances to the shared stamp
  154. // and the next run's strict > excludes them forever. The slice
  155. // grows to take the whole tie, even though that overshoots the
  156. // limit. If every changed document shares one stamp, cut grows to
  157. // all of them, which is correct - nothing is stranded.
  158. while (cut < changed.length && stampOf(changed[cut]) === boundary) {
  159. cut++;
  160. }
  161. smartbotic.log.warn('Database Change: ' + changed.length + ' documents changed, reporting ' +
  162. cut + (cut > maxDocuments ? ' (over the configured ' + maxDocuments + ', to avoid splitting a tied timestamp)' : '') +
  163. '. The rest arrive on the next run');
  164. changed = changed.slice(0, cut);
  165. // The mark only advances as far as what was actually reported, or
  166. // the remainder would be skipped rather than deferred.
  167. highWater = stampOf(changed[changed.length - 1]);
  168. }
  169. }
  170. const cursor = { since: highWater, updatedAt: Date.now(), collection: collection };
  171. const write = isFirstRun
  172. ? smartbotic.storage.insert(cursorCollection, cursor, cursorId)
  173. : smartbotic.storage.update(cursorCollection, cursorId, cursor);
  174. if (!write || write.success === false) {
  175. throw new Error('Database Change: could not save the cursor for ' + collection + ': ' +
  176. ((write && write.error) || 'unknown storage error'));
  177. }
  178. smartbotic.log.info('Database Change: ' + changed.length + ' documents from ' + collection);
  179. return {
  180. documents: changed,
  181. count: changed.length,
  182. isFirstRun: isFirstRun,
  183. cursor: highWater
  184. };
  185. }
  186. module.exports = { configSchema, inputSchema, outputSchema, execute };