/** * @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 };