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