file-watch.js 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152
  1. /**
  2. * @node file-watch
  3. * @name File Watch
  4. * @category triggers
  5. * @version 1.0.0
  6. * @description Report files added or changed in a directory since the last run, paired with a schedule trigger for its cadence
  7. * @icon folder-search
  8. */
  9. const configSchema = {
  10. type: 'object',
  11. properties: {
  12. directory: {
  13. type: 'string',
  14. title: 'Directory',
  15. description: 'Absolute path to watch. Not recursive'
  16. },
  17. pattern: {
  18. type: 'string',
  19. title: 'Name Pattern',
  20. description: 'Optional regular expression a file name must match, such as \\.csv$'
  21. },
  22. includeDirectories: {
  23. type: 'boolean',
  24. title: 'Include Directories',
  25. description: 'Report subdirectories as well as files',
  26. default: false
  27. },
  28. emitOnFirstRun: {
  29. type: 'boolean',
  30. title: 'Report Everything On First Run',
  31. 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',
  32. default: false
  33. },
  34. cursorCollection: {
  35. dynamicOptions: {
  36. source: 'storage.collections',
  37. labelField: 'name',
  38. valueField: 'name',
  39. filter: { access: 'read-write' }
  40. },
  41. type: 'string',
  42. title: 'Cursor Collection',
  43. description: 'Collection holding the last-seen state',
  44. default: 'watch_cursors'
  45. }
  46. },
  47. required: ['directory']
  48. };
  49. const inputSchema = {
  50. type: 'object',
  51. properties: {
  52. data: { type: 'any' }
  53. }
  54. };
  55. const outputSchema = {
  56. type: 'object',
  57. properties: {
  58. files: { type: 'array', description: 'Files new or changed since the last run' },
  59. count: { type: 'number' },
  60. isFirstRun: { type: 'boolean', description: 'True when no cursor existed yet' }
  61. }
  62. };
  63. async function execute(config, input, context) {
  64. const directory = config.directory;
  65. if (!directory) {
  66. throw new Error('File Watch: a directory is required');
  67. }
  68. const collection = config.cursorCollection || 'watch_cursors';
  69. const workflowId = (context && context.workflowId) || 'unknown';
  70. const nodeId = (context && context.nodeId) || 'file-watch';
  71. const cursorId = workflowId + ':' + nodeId;
  72. const listing = smartbotic.fs.readdir(directory);
  73. if (!listing || listing.success !== true) {
  74. throw new Error('File Watch: could not read ' + directory + ': ' +
  75. ((listing && listing.error) || 'unknown error'));
  76. }
  77. let matcher = null;
  78. if (config.pattern) {
  79. try {
  80. matcher = new RegExp(config.pattern);
  81. } catch (e) {
  82. throw new Error('File Watch: "' + config.pattern + '" is not a valid pattern: ' + e.message);
  83. }
  84. }
  85. const current = {};
  86. const candidates = [];
  87. for (const entry of listing.entries) {
  88. if (entry.isDirectory && config.includeDirectories !== true) {
  89. continue;
  90. }
  91. if (matcher && !matcher.test(entry.name)) {
  92. continue;
  93. }
  94. // Size and modification time together, because a file rewritten within
  95. // the same second at the same length is not a change worth waking a
  96. // workflow for, and a timestamp alone misses a rewrite that preserves
  97. // mtime granularity.
  98. current[entry.name] = entry.modifiedAt + ':' + entry.size;
  99. candidates.push(entry);
  100. }
  101. const stored = smartbotic.storage.get(collection, cursorId);
  102. // A missing document is the expected shape of a first run. Anything else
  103. // that keeps storage.get from returning the cursor - no read access to
  104. // the collection, a connection error - is a real failure and must not be
  105. // swallowed into a false "first run" that then quietly overwrites state.
  106. const cursorMissing = stored && typeof stored.error === 'string' &&
  107. stored.error.indexOf('Document not found') === 0;
  108. if (stored && stored.found !== true && !cursorMissing) {
  109. throw new Error('File Watch: could not read the cursor for ' + directory + ': ' +
  110. (stored.error || 'unknown storage error'));
  111. }
  112. const isFirstRun = !stored || stored.found !== true;
  113. const previous = (!isFirstRun && stored.document && stored.document.seen) || {};
  114. let changed = [];
  115. if (isFirstRun && config.emitOnFirstRun !== true) {
  116. smartbotic.log.info('File Watch: first run on ' + directory + ', recorded ' +
  117. candidates.length + ' entries without reporting them');
  118. } else {
  119. changed = candidates.filter(function (entry) {
  120. return previous[entry.name] !== current[entry.name];
  121. });
  122. }
  123. const cursor = { seen: current, updatedAt: Date.now(), directory: directory };
  124. const write = isFirstRun
  125. ? smartbotic.storage.insert(collection, cursor, cursorId)
  126. : smartbotic.storage.update(collection, cursorId, cursor);
  127. if (!write || write.success === false) {
  128. throw new Error('File Watch: could not save the cursor for ' + directory + ': ' +
  129. ((write && write.error) || 'unknown storage error'));
  130. }
  131. smartbotic.log.info('File Watch: ' + changed.length + ' new or changed in ' + directory);
  132. return {
  133. files: changed,
  134. count: changed.length,
  135. isFirstRun: isFirstRun
  136. };
  137. }
  138. module.exports = { configSchema, inputSchema, outputSchema, execute };