database-change.js 7.9 KB

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