| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242 |
- /**
- * @node database-change
- * @name Database Change
- * @category triggers
- * @version 1.0.0
- * @description Start the workflow when a collection changes. The database reports the change, so nothing polls and nothing is missed between ticks
- * @icon database
- * @trigger
- */
- const configSchema = {
- type: 'object',
- properties: {
- collection: {
- dynamicOptions: {
- source: 'storage.collections',
- labelField: 'name',
- valueField: 'name',
- filter: {}
- },
- type: 'string',
- title: 'Collection',
- description: 'Collection to watch'
- },
- timestampField: {
- type: 'string',
- title: 'Timestamp Field',
- 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',
- default: '_updated_at'
- },
- filter: {
- type: 'object',
- title: 'Filter',
- description: 'Optional query restricting which documents are watched',
- additionalProperties: true
- },
- maxDocuments: {
- type: 'number',
- title: 'Max Documents',
- description: 'Most documents to report in one run, so a large backlog does not arrive as one enormous payload',
- default: 100
- },
- emitOnFirstRun: {
- type: 'boolean',
- title: 'Report Everything On First Run',
- 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',
- default: false
- },
- eventTypes: {
- type: 'array',
- title: 'Event Types',
- description: 'Which changes should start the workflow. Empty means all of them',
- items: { type: 'string', enum: ['insert', 'update', 'delete'] },
- default: []
- },
- cursorCollection: {
- dynamicOptions: {
- source: 'storage.collections',
- labelField: 'name',
- valueField: 'name',
- filter: { access: 'read-write' }
- },
- type: 'string',
- title: 'Cursor Collection',
- description: 'Only used by a manual run, which has no event to work from and asks what changed since last time instead',
- default: 'watch_cursors'
- }
- },
- required: ['collection']
- };
- const inputSchema = {
- type: 'object',
- properties: {
- data: { type: 'any' }
- }
- };
- const outputSchema = {
- type: 'object',
- properties: {
- documents: { type: 'array', description: 'Documents new or changed since the last run' },
- count: { type: 'number' },
- isFirstRun: { type: 'boolean' },
- cursor: { type: 'number', description: 'High-water timestamp stored for the next run. Zero for an event, which needs no cursor' },
- eventType: { type: 'string', description: 'insert, update or delete - only present when started by a change' },
- documentId: { type: 'string', description: 'The document that changed - only present when started by a change' },
- collection: { type: 'string' }
- }
- };
- async function execute(config, input, context) {
- const collection = config.collection;
- if (!collection) {
- throw new Error('Database Change: a collection is required');
- }
- // The usual path: the database said what changed and the workflow was
- // started because of it. Nothing to look up - the document is right here.
- //
- // A trigger node is handed the trigger data as its input; context carries
- // only the execution, node and workflow ids. Reading context.triggerData
- // gets undefined every time, which looks like "no event" and quietly sends
- // every run down the polling path.
- const triggerData = (input && typeof input === 'object' ? input : {});
- if (triggerData.eventType && triggerData.documentId) {
- const changed = triggerData.document && Object.keys(triggerData.document).length
- ? [triggerData.document]
- : [];
- smartbotic.log.info('Database Change: ' + triggerData.eventType + ' on ' +
- (triggerData.collection || collection) + '/' + triggerData.documentId);
- return {
- documents: changed,
- count: changed.length,
- isFirstRun: false,
- cursor: 0,
- eventType: triggerData.eventType,
- documentId: triggerData.documentId,
- collection: triggerData.collection || collection
- };
- }
- // No event, so this is somebody running the workflow by hand to see what it
- // does. Falling back to the old query means a manual run still produces
- // something to work with instead of an empty result that looks broken.
- const timestampField = config.timestampField || '_updated_at';
- const cursorCollection = config.cursorCollection || 'watch_cursors';
- const maxDocuments = Number(config.maxDocuments) || 100;
- const workflowId = (context && context.workflowId) || 'unknown';
- const nodeId = (context && context.nodeId) || 'database-change';
- const cursorId = workflowId + ':' + nodeId;
- const stored = smartbotic.storage.get(cursorCollection, cursorId);
- // A missing document is the expected shape of a first run. Anything else
- // that keeps storage.get from returning the cursor - no read access to
- // the collection, a connection error - is a real failure and must not be
- // swallowed into a false "first run" that then quietly overwrites state.
- // This mirrors nodes/triggers/file-watch.js so the two nodes agree on
- // what "first run" means.
- const cursorMissing = stored && typeof stored.error === 'string' &&
- stored.error.indexOf('Document not found') === 0;
- if (stored && stored.found !== true && !cursorMissing) {
- throw new Error('Database Change: could not read the cursor for ' + collection + ': ' +
- (stored.error || 'unknown storage error'));
- }
- const isFirstRun = !stored || stored.found !== true;
- const since = (!isFirstRun && stored.document && Number(stored.document.since)) || 0;
- const query = config.filter && typeof config.filter === 'object' ? config.filter : {};
- const result = smartbotic.storage.query(collection, query);
- // storage.query signals failure with an "error" field, not a "success"
- // flag - there is no success flag in its response at all, so checking
- // one would throw on every call, successful or not.
- if (!result || typeof result.error === 'string') {
- throw new Error('Database Change: could not query ' + collection + ': ' +
- ((result && result.error) || 'unknown error'));
- }
- const documents = result.documents || [];
- // The database stores _updated_at in nanoseconds while everything a node
- // sees is milliseconds, so the value is normalised rather than compared
- // against a cursor in different units.
- function stampOf(document) {
- const raw = Number(smartbotic.utils.getFieldValue(document, timestampField));
- if (!raw || isNaN(raw)) {
- return 0;
- }
- return raw > 1e15 ? Math.floor(raw / 1000000) : raw;
- }
- let highWater = since;
- let noStamp = 0;
- for (const document of documents) {
- const stamp = stampOf(document);
- if (stamp === 0) {
- noStamp++;
- }
- if (stamp > highWater) {
- highWater = stamp;
- }
- }
- if (noStamp > 0) {
- // A document with no readable value in timestampField sorts as if it
- // predates everything, and once the mark advances past 0 it can never
- // satisfy the strict > filter below - it is never reported, silently,
- // unless this is logged.
- smartbotic.log.warn('Database Change: ' + noStamp + ' document(s) in ' + collection +
- ' had no readable "' + timestampField + '" value and will never be reported');
- }
- let changed = [];
- if (isFirstRun && config.emitOnFirstRun !== true) {
- smartbotic.log.info('Database Change: first run on ' + collection + ', recorded the mark at ' +
- highWater + ' without reporting ' + documents.length + ' documents');
- } else {
- changed = documents
- .filter(function (document) { return stampOf(document) > since; })
- .sort(function (left, right) { return stampOf(left) - stampOf(right); });
- if (changed.length > maxDocuments) {
- let cut = maxDocuments;
- const boundary = stampOf(changed[cut - 1]);
- // Cutting through a group that shares one timestamp would strand
- // the rest of that group: the mark advances to the shared stamp
- // and the next run's strict > excludes them forever. The slice
- // grows to take the whole tie, even though that overshoots the
- // limit. If every changed document shares one stamp, cut grows to
- // all of them, which is correct - nothing is stranded.
- while (cut < changed.length && stampOf(changed[cut]) === boundary) {
- cut++;
- }
- smartbotic.log.warn('Database Change: ' + changed.length + ' documents changed, reporting ' +
- cut + (cut > maxDocuments ? ' (over the configured ' + maxDocuments + ', to avoid splitting a tied timestamp)' : '') +
- '. The rest arrive on the next run');
- changed = changed.slice(0, cut);
- // The mark only advances as far as what was actually reported, or
- // the remainder would be skipped rather than deferred.
- highWater = stampOf(changed[changed.length - 1]);
- }
- }
- const cursor = { since: highWater, updatedAt: Date.now(), collection: collection };
- const write = isFirstRun
- ? smartbotic.storage.insert(cursorCollection, cursor, cursorId)
- : smartbotic.storage.update(cursorCollection, cursorId, cursor);
- if (!write || write.success === false) {
- throw new Error('Database Change: could not save the cursor for ' + collection + ': ' +
- ((write && write.error) || 'unknown storage error'));
- }
- smartbotic.log.info('Database Change: ' + changed.length + ' documents from ' + collection);
- return {
- documents: changed,
- count: changed.length,
- isFirstRun: isFirstRun,
- cursor: highWater
- };
- }
- module.exports = { configSchema, inputSchema, outputSchema, execute };
|