Browse Source

feat: a Database Change trigger, reporting documents changed since the last run

fszontagh 1 tháng trước cách đây
mục cha
commit
4a56fac9e6
2 tập tin đã thay đổi với 197 bổ sung và 0 xóa
  1. 165 0
      nodes/triggers/database-change.js
  2. 32 0
      tests/nodes/database-change.json

+ 165 - 0
nodes/triggers/database-change.js

@@ -0,0 +1,165 @@
+/**
+ * @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 };

+ 32 - 0
tests/nodes/database-change.json

@@ -0,0 +1,32 @@
+{
+  "name": "verify-database-change",
+  "settings": {
+    "storagePermissions": {
+      "collections": { "watch_cursors": "read-write", "dbchange_test": "read-write" }
+    }
+  },
+  "nodes": [
+    {"id": "n1", "name": "Trigger", "type": "click-trigger", "position": {"x": 0, "y": 0}, "config": {}},
+    {"id": "setup", "name": "Setup", "type": "code", "position": {"x": 0, "y": 100},
+     "config": {"code": "smartbotic.storage.delete('watch_cursors', context.workflowId + ':first');\nsmartbotic.storage.delete('watch_cursors', context.workflowId + ':second');\nsmartbotic.storage.delete('dbchange_test', 'before-1');\nsmartbotic.storage.delete('dbchange_test', 'after-1');\nsmartbotic.storage.insert('dbchange_test', { name: 'before', at: Date.now() }, 'before-1');\nconst beforeDoc = smartbotic.storage.get('dbchange_test', 'before-1');\nconst raw = Number(beforeDoc.document._updated_at);\nconst stamp = raw > 1e15 ? Math.floor(raw / 1000000) : raw;\nsmartbotic.storage.insert('watch_cursors', { since: stamp, updatedAt: Date.now(), collection: 'dbchange_test' }, context.workflowId + ':second');\nreturn { ready: true };"}},
+    {"id": "first", "name": "First Run", "type": "database-change", "position": {"x": 0, "y": 200},
+     "config": {"collection": "dbchange_test", "emitOnFirstRun": false}},
+    {"id": "add", "name": "Insert One", "type": "code", "position": {"x": 0, "y": 300},
+     "config": {"code": "smartbotic.utils.sleep(1100);\nsmartbotic.storage.insert('dbchange_test', { name: 'after', at: Date.now() }, 'after-1');\nreturn { inserted: 'after-1' };"}},
+    {"id": "second", "name": "Second Run", "type": "database-change", "position": {"x": 0, "y": 400},
+     "config": {"collection": "dbchange_test", "emitOnFirstRun": false}},
+    {"id": "cleanup", "name": "Cleanup", "type": "code", "position": {"x": 0, "y": 500},
+     "config": {"code": "smartbotic.storage.delete('dbchange_test', 'before-1');\nsmartbotic.storage.delete('dbchange_test', 'after-1');\nreturn { cleaned: true };"}}
+  ],
+  "connections": [
+    {"sourceNodeId": "n1", "sourceOutput": "main", "targetNodeId": "setup", "targetInput": "data"},
+    {"sourceNodeId": "setup", "sourceOutput": "main", "targetNodeId": "first", "targetInput": "data"},
+    {"sourceNodeId": "first", "sourceOutput": "main", "targetNodeId": "add", "targetInput": "data"},
+    {"sourceNodeId": "add", "sourceOutput": "main", "targetNodeId": "second", "targetInput": "data"},
+    {"sourceNodeId": "second", "sourceOutput": "main", "targetNodeId": "cleanup", "targetInput": "data"}
+  ],
+  "expect": {
+    "first": {"status": "completed", "output": {"count": 0, "isFirstRun": true}},
+    "second": {"status": "completed", "output": {"count": 1, "isFirstRun": false, "documents": [{"name": "after"}]}}
+  }
+}