Răsfoiți Sursa

feat: a File Watch trigger, reporting what changed in a directory

Also fixes fs.stat's mtime, which used file_time_type's epoch directly
instead of clock_cast'ing to system_clock like fs.readdir already does,
producing nonsensical (negative) timestamps.
fszontagh 1 lună în urmă
părinte
comite
74f237062b

+ 146 - 0
nodes/triggers/file-watch.js

@@ -0,0 +1,146 @@
+/**
+ * @node file-watch
+ * @name File Watch
+ * @category triggers
+ * @version 1.0.0
+ * @description Report files added or changed in a directory since the last run, paired with a schedule trigger for its cadence
+ * @icon folder-search
+ */
+
+const configSchema = {
+    type: 'object',
+    properties: {
+        directory: {
+            type: 'string',
+            title: 'Directory',
+            description: 'Absolute path to watch. Not recursive'
+        },
+        pattern: {
+            type: 'string',
+            title: 'Name Pattern',
+            description: 'Optional regular expression a file name must match, such as \\.csv$'
+        },
+        includeDirectories: {
+            type: 'boolean',
+            title: 'Include Directories',
+            description: 'Report subdirectories as well as files',
+            default: false
+        },
+        emitOnFirstRun: {
+            type: 'boolean',
+            title: 'Report Everything On First Run',
+            description: 'On the very first run there is nothing to compare against. Off records what is there and reports nothing, so adding this to a live workflow does not fire for every file already present',
+            default: false
+        },
+        cursorCollection: {
+            type: 'string',
+            title: 'Cursor Collection',
+            description: 'Collection holding the last-seen state',
+            default: 'watch_cursors'
+        }
+    },
+    required: ['directory']
+};
+
+const inputSchema = {
+    type: 'object',
+    properties: {
+        data: { type: 'any' }
+    }
+};
+
+const outputSchema = {
+    type: 'object',
+    properties: {
+        files: { type: 'array', description: 'Files new or changed since the last run' },
+        count: { type: 'number' },
+        isFirstRun: { type: 'boolean', description: 'True when no cursor existed yet' }
+    }
+};
+
+async function execute(config, input, context) {
+    const directory = config.directory;
+    if (!directory) {
+        throw new Error('File Watch: a directory is required');
+    }
+
+    const collection = config.cursorCollection || 'watch_cursors';
+    const workflowId = (context && context.workflowId) || 'unknown';
+    const nodeId = (context && context.nodeId) || 'file-watch';
+    const cursorId = workflowId + ':' + nodeId;
+
+    const listing = smartbotic.fs.readdir(directory);
+    if (!listing || listing.success !== true) {
+        throw new Error('File Watch: could not read ' + directory + ': ' +
+            ((listing && listing.error) || 'unknown error'));
+    }
+
+    let matcher = null;
+    if (config.pattern) {
+        try {
+            matcher = new RegExp(config.pattern);
+        } catch (e) {
+            throw new Error('File Watch: "' + config.pattern + '" is not a valid pattern: ' + e.message);
+        }
+    }
+
+    const current = {};
+    const candidates = [];
+    for (const entry of listing.entries) {
+        if (entry.isDirectory && config.includeDirectories !== true) {
+            continue;
+        }
+        if (matcher && !matcher.test(entry.name)) {
+            continue;
+        }
+        // Size and modification time together, because a file rewritten within
+        // the same second at the same length is not a change worth waking a
+        // workflow for, and a timestamp alone misses a rewrite that preserves
+        // mtime granularity.
+        current[entry.name] = entry.modifiedAt + ':' + entry.size;
+        candidates.push(entry);
+    }
+
+    const stored = smartbotic.storage.get(collection, 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.
+    const cursorMissing = stored && typeof stored.error === 'string' &&
+        stored.error.indexOf('Document not found') === 0;
+    if (stored && stored.found !== true && !cursorMissing) {
+        throw new Error('File Watch: could not read the cursor for ' + directory + ': ' +
+            (stored.error || 'unknown storage error'));
+    }
+    const isFirstRun = !stored || stored.found !== true;
+    const previous = (!isFirstRun && stored.document && stored.document.seen) || {};
+
+    let changed = [];
+    if (isFirstRun && config.emitOnFirstRun !== true) {
+        smartbotic.log.info('File Watch: first run on ' + directory + ', recorded ' +
+            candidates.length + ' entries without reporting them');
+    } else {
+        changed = candidates.filter(function (entry) {
+            return previous[entry.name] !== current[entry.name];
+        });
+    }
+
+    const cursor = { seen: current, updatedAt: Date.now(), directory: directory };
+    const write = isFirstRun
+        ? smartbotic.storage.insert(collection, cursor, cursorId)
+        : smartbotic.storage.update(collection, cursorId, cursor);
+    if (!write || write.success === false) {
+        throw new Error('File Watch: could not save the cursor for ' + directory + ': ' +
+            ((write && write.error) || 'unknown storage error'));
+    }
+
+    smartbotic.log.info('File Watch: ' + changed.length + ' new or changed in ' + directory);
+
+    return {
+        files: changed,
+        count: changed.length,
+        isFirstRun: isFirstRun
+    };
+}
+
+module.exports = { configSchema, inputSchema, outputSchema, execute };

+ 6 - 1
src/runner/engine/script_engine.cpp

@@ -2539,8 +2539,13 @@ void ScriptEngine::setupBuiltinAPIs() {
             auto status = std::filesystem::status(path_str);
             auto size = std::filesystem::is_regular_file(path_str) ? std::filesystem::file_size(path_str) : 0;
             auto mtime = std::filesystem::last_write_time(path_str);
+            // file_time_type's epoch is not guaranteed to match system_clock's
+            // (it does not on this libstdc++), so this has to go through
+            // clock_cast rather than time_since_epoch() directly, or mtime
+            // comes out shifted by whatever offset the filesystem clock uses.
+            auto sys_time = std::chrono::clock_cast<std::chrono::system_clock>(mtime);
             auto mtime_ms = std::chrono::duration_cast<std::chrono::milliseconds>(
-                mtime.time_since_epoch()
+                sys_time.time_since_epoch()
             ).count();
 
             JSValue response = JS_NewObject(ctx);

+ 30 - 0
tests/nodes/file-watch.json

@@ -0,0 +1,30 @@
+{
+  "name": "verify-file-watch",
+  "settings": {
+    "storagePermissions": {
+      "collections": { "watch_cursors": "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": "const dir = '/tmp/sb-file-watch-test';\nconst listing = smartbotic.fs.readdir(dir);\nif (listing.success) {\n    for (const e of listing.entries) { smartbotic.fs.unlink(e.path); }\n}\nsmartbotic.fs.mkdir(dir);\nsmartbotic.fs.writeFile(dir + '/existing.txt', smartbotic.utils.base64Encode('old'));\nsmartbotic.storage.delete('watch_cursors', context.workflowId + ':first');\nsmartbotic.storage.delete('watch_cursors', context.workflowId + ':second');\nconst before = smartbotic.fs.readdir(dir);\nconst seen = {};\nfor (const e of before.entries) { seen[e.name] = e.modifiedAt + ':' + e.size; }\nsmartbotic.storage.insert('watch_cursors', { seen: seen, updatedAt: Date.now(), directory: dir }, context.workflowId + ':second');\nconst st = smartbotic.fs.stat(dir + '/existing.txt');\nconst statMtimeSane = st.mtime > 1600000000000 && st.mtime < Date.now() + 60000;\nreturn { dir: dir, statMtimeSane: statMtimeSane };"}},
+    {"id": "first", "name": "First Run", "type": "file-watch", "position": {"x": 0, "y": 200},
+     "config": {"directory": "/tmp/sb-file-watch-test", "emitOnFirstRun": false}},
+    {"id": "add", "name": "Add A File", "type": "code", "position": {"x": 0, "y": 300},
+     "config": {"code": "smartbotic.fs.writeFile('/tmp/sb-file-watch-test/fresh.txt', smartbotic.utils.base64Encode('new'));\nreturn { added: 'fresh.txt' };"}},
+    {"id": "second", "name": "Second Run", "type": "file-watch", "position": {"x": 0, "y": 400},
+     "config": {"directory": "/tmp/sb-file-watch-test", "emitOnFirstRun": false}}
+  ],
+  "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"}
+  ],
+  "expect": {
+    "setup": {"status": "completed", "output": {"result": {"statMtimeSane": true}}},
+    "first": {"status": "completed", "output": {"count": 0, "isFirstRun": true}},
+    "second": {"status": "completed", "output": {"count": 1, "isFirstRun": false, "files": [{"name": "fresh.txt"}]}}
+  }
+}