database-change.js 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242
  1. /**
  2. * @node database-change
  3. * @name Database Change
  4. * @category triggers
  5. * @version 1.0.0
  6. * @description Start the workflow when a collection changes. The database reports the change, so nothing polls and nothing is missed between ticks
  7. * @icon database
  8. * @trigger
  9. */
  10. const configSchema = {
  11. type: 'object',
  12. properties: {
  13. collection: {
  14. dynamicOptions: {
  15. source: 'storage.collections',
  16. labelField: 'name',
  17. valueField: 'name',
  18. filter: {}
  19. },
  20. type: 'string',
  21. title: 'Collection',
  22. description: 'Collection to watch'
  23. },
  24. timestampField: {
  25. type: 'string',
  26. title: 'Timestamp Field',
  27. 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',
  28. default: '_updated_at'
  29. },
  30. filter: {
  31. type: 'object',
  32. title: 'Filter',
  33. description: 'Optional query restricting which documents are watched',
  34. additionalProperties: true
  35. },
  36. maxDocuments: {
  37. type: 'number',
  38. title: 'Max Documents',
  39. description: 'Most documents to report in one run, so a large backlog does not arrive as one enormous payload',
  40. default: 100
  41. },
  42. emitOnFirstRun: {
  43. type: 'boolean',
  44. title: 'Report Everything On First Run',
  45. 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',
  46. default: false
  47. },
  48. eventTypes: {
  49. type: 'array',
  50. title: 'Event Types',
  51. description: 'Which changes should start the workflow. Empty means all of them',
  52. items: { type: 'string', enum: ['insert', 'update', 'delete'] },
  53. default: []
  54. },
  55. cursorCollection: {
  56. dynamicOptions: {
  57. source: 'storage.collections',
  58. labelField: 'name',
  59. valueField: 'name',
  60. filter: { access: 'read-write' }
  61. },
  62. type: 'string',
  63. title: 'Cursor Collection',
  64. description: 'Only used by a manual run, which has no event to work from and asks what changed since last time instead',
  65. default: 'watch_cursors'
  66. }
  67. },
  68. required: ['collection']
  69. };
  70. const inputSchema = {
  71. type: 'object',
  72. properties: {
  73. data: { type: 'any' }
  74. }
  75. };
  76. const outputSchema = {
  77. type: 'object',
  78. properties: {
  79. documents: { type: 'array', description: 'Documents new or changed since the last run' },
  80. count: { type: 'number' },
  81. isFirstRun: { type: 'boolean' },
  82. cursor: { type: 'number', description: 'High-water timestamp stored for the next run. Zero for an event, which needs no cursor' },
  83. eventType: { type: 'string', description: 'insert, update or delete - only present when started by a change' },
  84. documentId: { type: 'string', description: 'The document that changed - only present when started by a change' },
  85. collection: { type: 'string' }
  86. }
  87. };
  88. async function execute(config, input, context) {
  89. const collection = config.collection;
  90. if (!collection) {
  91. throw new Error('Database Change: a collection is required');
  92. }
  93. // The usual path: the database said what changed and the workflow was
  94. // started because of it. Nothing to look up - the document is right here.
  95. //
  96. // A trigger node is handed the trigger data as its input; context carries
  97. // only the execution, node and workflow ids. Reading context.triggerData
  98. // gets undefined every time, which looks like "no event" and quietly sends
  99. // every run down the polling path.
  100. const triggerData = (input && typeof input === 'object' ? input : {});
  101. if (triggerData.eventType && triggerData.documentId) {
  102. const changed = triggerData.document && Object.keys(triggerData.document).length
  103. ? [triggerData.document]
  104. : [];
  105. smartbotic.log.info('Database Change: ' + triggerData.eventType + ' on ' +
  106. (triggerData.collection || collection) + '/' + triggerData.documentId);
  107. return {
  108. documents: changed,
  109. count: changed.length,
  110. isFirstRun: false,
  111. cursor: 0,
  112. eventType: triggerData.eventType,
  113. documentId: triggerData.documentId,
  114. collection: triggerData.collection || collection
  115. };
  116. }
  117. // No event, so this is somebody running the workflow by hand to see what it
  118. // does. Falling back to the old query means a manual run still produces
  119. // something to work with instead of an empty result that looks broken.
  120. const timestampField = config.timestampField || '_updated_at';
  121. const cursorCollection = config.cursorCollection || 'watch_cursors';
  122. const maxDocuments = Number(config.maxDocuments) || 100;
  123. const workflowId = (context && context.workflowId) || 'unknown';
  124. const nodeId = (context && context.nodeId) || 'database-change';
  125. const cursorId = workflowId + ':' + nodeId;
  126. const stored = smartbotic.storage.get(cursorCollection, cursorId);
  127. // A missing document is the expected shape of a first run. Anything else
  128. // that keeps storage.get from returning the cursor - no read access to
  129. // the collection, a connection error - is a real failure and must not be
  130. // swallowed into a false "first run" that then quietly overwrites state.
  131. // This mirrors nodes/triggers/file-watch.js so the two nodes agree on
  132. // what "first run" means.
  133. const cursorMissing = stored && typeof stored.error === 'string' &&
  134. stored.error.indexOf('Document not found') === 0;
  135. if (stored && stored.found !== true && !cursorMissing) {
  136. throw new Error('Database Change: could not read the cursor for ' + collection + ': ' +
  137. (stored.error || 'unknown storage error'));
  138. }
  139. const isFirstRun = !stored || stored.found !== true;
  140. const since = (!isFirstRun && stored.document && Number(stored.document.since)) || 0;
  141. const query = config.filter && typeof config.filter === 'object' ? config.filter : {};
  142. const result = smartbotic.storage.query(collection, query);
  143. // storage.query signals failure with an "error" field, not a "success"
  144. // flag - there is no success flag in its response at all, so checking
  145. // one would throw on every call, successful or not.
  146. if (!result || typeof result.error === 'string') {
  147. throw new Error('Database Change: could not query ' + collection + ': ' +
  148. ((result && result.error) || 'unknown error'));
  149. }
  150. const documents = result.documents || [];
  151. // The database stores _updated_at in nanoseconds while everything a node
  152. // sees is milliseconds, so the value is normalised rather than compared
  153. // against a cursor in different units.
  154. function stampOf(document) {
  155. const raw = Number(smartbotic.utils.getFieldValue(document, timestampField));
  156. if (!raw || isNaN(raw)) {
  157. return 0;
  158. }
  159. return raw > 1e15 ? Math.floor(raw / 1000000) : raw;
  160. }
  161. let highWater = since;
  162. let noStamp = 0;
  163. for (const document of documents) {
  164. const stamp = stampOf(document);
  165. if (stamp === 0) {
  166. noStamp++;
  167. }
  168. if (stamp > highWater) {
  169. highWater = stamp;
  170. }
  171. }
  172. if (noStamp > 0) {
  173. // A document with no readable value in timestampField sorts as if it
  174. // predates everything, and once the mark advances past 0 it can never
  175. // satisfy the strict > filter below - it is never reported, silently,
  176. // unless this is logged.
  177. smartbotic.log.warn('Database Change: ' + noStamp + ' document(s) in ' + collection +
  178. ' had no readable "' + timestampField + '" value and will never be reported');
  179. }
  180. let changed = [];
  181. if (isFirstRun && config.emitOnFirstRun !== true) {
  182. smartbotic.log.info('Database Change: first run on ' + collection + ', recorded the mark at ' +
  183. highWater + ' without reporting ' + documents.length + ' documents');
  184. } else {
  185. changed = documents
  186. .filter(function (document) { return stampOf(document) > since; })
  187. .sort(function (left, right) { return stampOf(left) - stampOf(right); });
  188. if (changed.length > maxDocuments) {
  189. let cut = maxDocuments;
  190. const boundary = stampOf(changed[cut - 1]);
  191. // Cutting through a group that shares one timestamp would strand
  192. // the rest of that group: the mark advances to the shared stamp
  193. // and the next run's strict > excludes them forever. The slice
  194. // grows to take the whole tie, even though that overshoots the
  195. // limit. If every changed document shares one stamp, cut grows to
  196. // all of them, which is correct - nothing is stranded.
  197. while (cut < changed.length && stampOf(changed[cut]) === boundary) {
  198. cut++;
  199. }
  200. smartbotic.log.warn('Database Change: ' + changed.length + ' documents changed, reporting ' +
  201. cut + (cut > maxDocuments ? ' (over the configured ' + maxDocuments + ', to avoid splitting a tied timestamp)' : '') +
  202. '. The rest arrive on the next run');
  203. changed = changed.slice(0, cut);
  204. // The mark only advances as far as what was actually reported, or
  205. // the remainder would be skipped rather than deferred.
  206. highWater = stampOf(changed[changed.length - 1]);
  207. }
  208. }
  209. const cursor = { since: highWater, updatedAt: Date.now(), collection: collection };
  210. const write = isFirstRun
  211. ? smartbotic.storage.insert(cursorCollection, cursor, cursorId)
  212. : smartbotic.storage.update(cursorCollection, cursorId, cursor);
  213. if (!write || write.success === false) {
  214. throw new Error('Database Change: could not save the cursor for ' + collection + ': ' +
  215. ((write && write.error) || 'unknown storage error'));
  216. }
  217. smartbotic.log.info('Database Change: ' + changed.length + ' documents from ' + collection);
  218. return {
  219. documents: changed,
  220. count: changed.length,
  221. isFirstRun: isFirstRun,
  222. cursor: highWater
  223. };
  224. }
  225. module.exports = { configSchema, inputSchema, outputSchema, execute };