| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165 |
- /**
- * @node database-change
- * @name Database Change
- * @category triggers
- * @version 1.0.0
- * @description Report documents added or changed in a collection since the last run, paired with a schedule trigger for its cadence
- * @icon database
- */
- const configSchema = {
- type: 'object',
- properties: {
- collection: {
- 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',
- 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
- },
- cursorCollection: {
- type: 'string',
- title: 'Cursor Collection',
- 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' }
- }
- };
- async function execute(config, input, context) {
- const collection = config.collection;
- if (!collection) {
- throw new Error('Database Change: a collection is required');
- }
- 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;
- for (const document of documents) {
- const stamp = stampOf(document);
- if (stamp > highWater) {
- highWater = stamp;
- }
- }
- 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) {
- smartbotic.log.warn('Database Change: ' + changed.length + ' documents changed, reporting the oldest ' +
- maxDocuments + '. The rest arrive on the next run');
- changed = changed.slice(0, maxDocuments);
- // 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 };
|