From 9a4e9b90502de7b9d75f4e23cefc61ddaf91ccbc Mon Sep 17 00:00:00 2001 From: Benjamin Pasero Date: Tue, 11 Jan 2022 12:07:58 +0100 Subject: [PATCH] watcher - throttle events at their source to also prevent IPC spam --- src/vs/base/common/async.ts | 36 +++++++++---- src/vs/base/test/common/async.test.ts | 24 +++++++-- src/vs/platform/files/common/fileService.ts | 51 ++----------------- .../node/watcher/nodejs/nodejsWatcher.ts | 44 ++++++++++++---- .../node/watcher/parcel/parcelWatcher.ts | 31 +++++++++-- .../files/test/browser/fileService.test.ts | 41 +-------------- 6 files changed, 112 insertions(+), 115 deletions(-) diff --git a/src/vs/base/common/async.ts b/src/vs/base/common/async.ts index 6741bb8d0cfa..e8432db5bbad 100644 --- a/src/vs/base/common/async.ts +++ b/src/vs/base/common/async.ts @@ -959,11 +959,29 @@ export class RunOnceWorker extends RunOnceScheduler { } } +export interface IThrottledWorkerOptions { + + /** + * maximum of units the worker will pass onto handler at once + */ + maxWorkChunkSize: number; + + /** + * maximum of units the worker will keep in memory for processing + */ + maxBufferedWork: number | undefined; + + /** + * delay before processing the next round of chunks when chunk size exceeds limits + */ + throttleDelay: number; +} + /** * The `ThrottledWorker` will accept units of work `T` * to handle. The contract is: * * there is a maximum of units the worker can handle at once (via `maxWorkChunkSize`) - * * there is a maximum of units the worker will keep in memory for processing (via `maxPendingWork`) + * * there is a maximum of units the worker will keep in memory for processing (via `maxBufferedWork`) * * after having handled `maxWorkChunkSize` units, the worker needs to rest (via `throttleDelay`) */ export class ThrottledWorker extends Disposable { @@ -974,10 +992,8 @@ export class ThrottledWorker extends Disposable { private disposed = false; constructor( - private readonly maxWorkChunkSize: number, - private readonly maxBufferedWork: number | undefined, - private readonly throttleDelay: number, - private readonly handler: (units: readonly T[]) => void + private options: IThrottledWorkerOptions, + private readonly handler: (units: T[]) => void ) { super(); } @@ -1003,11 +1019,11 @@ export class ThrottledWorker extends Disposable { } // Check for reaching maximum of pending work - if (typeof this.maxBufferedWork === 'number') { + if (typeof this.options.maxBufferedWork === 'number') { // Throttled: simple check if pending + units exceeds max pending if (this.throttler.value) { - if (this.pending + units.length > this.maxBufferedWork) { + if (this.pending + units.length > this.options.maxBufferedWork) { return false; // work not accepted: too much pending work } } @@ -1015,7 +1031,7 @@ export class ThrottledWorker extends Disposable { // Unthrottled: same as throttled, but account for max chunk getting // worked on directly without being pending else { - if (this.pending + units.length - this.maxWorkChunkSize > this.maxBufferedWork) { + if (this.pending + units.length - this.options.maxWorkChunkSize > this.options.maxBufferedWork) { return false; // work not accepted: too much pending work } } @@ -1037,7 +1053,7 @@ export class ThrottledWorker extends Disposable { private doWork(): void { // Extract chunk to handle and handle it - this.handler(this.pendingWork.splice(0, this.maxWorkChunkSize)); + this.handler(this.pendingWork.splice(0, this.options.maxWorkChunkSize)); // If we have remaining work, schedule it after a delay if (this.pendingWork.length > 0) { @@ -1045,7 +1061,7 @@ export class ThrottledWorker extends Disposable { this.throttler.clear(); this.doWork(); - }, this.throttleDelay); + }, this.options.throttleDelay); this.throttler.value.schedule(); } } diff --git a/src/vs/base/test/common/async.test.ts b/src/vs/base/test/common/async.test.ts index b130e8ef1d2a..246bb95f9802 100644 --- a/src/vs/base/test/common/async.test.ts +++ b/src/vs/base/test/common/async.test.ts @@ -1068,7 +1068,11 @@ suite('Async', () => { } }; - const worker = new async.ThrottledWorker(5, undefined, 1, handler); + const worker = new async.ThrottledWorker({ + maxWorkChunkSize: 5, + maxBufferedWork: undefined, + throttleDelay: 1 + }, handler); // Work less than chunk size @@ -1176,7 +1180,11 @@ suite('Async', () => { let handled: number[] = []; const handler = (units: readonly number[]) => handled.push(...units); - const worker = new async.ThrottledWorker(5, 5, 1, handler); + const worker = new async.ThrottledWorker({ + maxWorkChunkSize: 5, + maxBufferedWork: 5, + throttleDelay: 1 + }, handler); let worked = worker.work([1, 2, 3]); assert.strictEqual(worked, true); @@ -1198,7 +1206,11 @@ suite('Async', () => { let handled: number[] = []; const handler = (units: readonly number[]) => handled.push(...units); - const worker = new async.ThrottledWorker(5, 5, 1, handler); + const worker = new async.ThrottledWorker({ + maxWorkChunkSize: 5, + maxBufferedWork: 5, + throttleDelay: 1 + }, handler); let worked = worker.work([1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11]); assert.strictEqual(worked, false); @@ -1213,7 +1225,11 @@ suite('Async', () => { let handled: number[] = []; const handler = (units: readonly number[]) => handled.push(...units); - const worker = new async.ThrottledWorker(5, undefined, 1, handler); + const worker = new async.ThrottledWorker({ + maxWorkChunkSize: 5, + maxBufferedWork: undefined, + throttleDelay: 1 + }, handler); worker.dispose(); const worked = worker.work([1, 2, 3]); diff --git a/src/vs/platform/files/common/fileService.ts b/src/vs/platform/files/common/fileService.ts index 31cf0ffe0136..d2adec89ea31 100644 --- a/src/vs/platform/files/common/fileService.ts +++ b/src/vs/platform/files/common/fileService.ts @@ -4,7 +4,7 @@ *--------------------------------------------------------------------------------------------*/ import { coalesce } from 'vs/base/common/arrays'; -import { Promises, ResourceQueue, ThrottledWorker } from 'vs/base/common/async'; +import { Promises, ResourceQueue } from 'vs/base/common/async'; import { bufferedStreamToBuffer, bufferToReadable, newWriteableBufferStream, readableToBuffer, streamToBuffer, VSBuffer, VSBufferReadable, VSBufferReadableBufferedStream, VSBufferReadableStream } from 'vs/base/common/buffer'; import { CancellationToken, CancellationTokenSource } from 'vs/base/common/cancellation'; import { Emitter } from 'vs/base/common/event'; @@ -18,7 +18,7 @@ import { extUri, extUriIgnorePathCase, IExtUri, isAbsolutePath } from 'vs/base/c import { consumeStream, isReadableBufferedStream, isReadableStream, listenStream, newWriteableStream, peekReadable, peekStream, transform } from 'vs/base/common/stream'; import { URI } from 'vs/base/common/uri'; import { localize } from 'vs/nls'; -import { ensureFileSystemProviderError, etag, ETAG_DISABLED, FileChangesEvent, FileDeleteOptions, FileOperation, FileOperationError, FileOperationEvent, FileOperationResult, FilePermission, FileSystemProviderCapabilities, FileSystemProviderErrorCode, FileType, hasFileAtomicReadCapability, hasFileFolderCopyCapability, hasFileReadStreamCapability, hasOpenReadWriteCloseCapability, hasReadWriteCapability, ICreateFileOptions, IFileChange, IFileContent, IFileService, IFileStat, IFileStatWithMetadata, IFileStreamContent, IFileSystemProvider, IFileSystemProviderActivationEvent, IFileSystemProviderCapabilitiesChangeEvent, IFileSystemProviderRegistrationEvent, IFileSystemProviderWithFileAtomicReadCapability, IFileSystemProviderWithFileReadStreamCapability, IFileSystemProviderWithFileReadWriteCapability, IFileSystemProviderWithOpenReadWriteCloseCapability, IReadFileOptions, IReadFileStreamOptions, IResolveFileOptions, IResolveFileResult, IResolveFileResultWithMetadata, IResolveMetadataFileOptions, IStat, IWatchOptions, IWriteFileOptions, NotModifiedSinceFileOperationError, toFileOperationResult, toFileSystemProviderErrorCode } from 'vs/platform/files/common/files'; +import { ensureFileSystemProviderError, etag, ETAG_DISABLED, FileChangesEvent, FileDeleteOptions, FileOperation, FileOperationError, FileOperationEvent, FileOperationResult, FilePermission, FileSystemProviderCapabilities, FileSystemProviderErrorCode, FileType, hasFileAtomicReadCapability, hasFileFolderCopyCapability, hasFileReadStreamCapability, hasOpenReadWriteCloseCapability, hasReadWriteCapability, ICreateFileOptions, IFileContent, IFileService, IFileStat, IFileStatWithMetadata, IFileStreamContent, IFileSystemProvider, IFileSystemProviderActivationEvent, IFileSystemProviderCapabilitiesChangeEvent, IFileSystemProviderRegistrationEvent, IFileSystemProviderWithFileAtomicReadCapability, IFileSystemProviderWithFileReadStreamCapability, IFileSystemProviderWithFileReadWriteCapability, IFileSystemProviderWithOpenReadWriteCloseCapability, IReadFileOptions, IReadFileStreamOptions, IResolveFileOptions, IResolveFileResult, IResolveFileResultWithMetadata, IResolveMetadataFileOptions, IStat, IWatchOptions, IWriteFileOptions, NotModifiedSinceFileOperationError, toFileOperationResult, toFileSystemProviderErrorCode } from 'vs/platform/files/common/files'; import { readFileIntoStream } from 'vs/platform/files/common/io'; import { ILogService } from 'vs/platform/log/common/log'; @@ -56,14 +56,13 @@ export class FileService extends Disposable implements IFileService { mark(`code/registerFilesystem/${scheme}`); const providerDisposables = new DisposableStore(); - + // Add provider with event this.provider.set(scheme, provider); this._onDidChangeFileSystemProviderRegistrations.fire({ added: true, scheme, provider }); // Forward events from provider - const providerFileChangeEventsWorker = providerDisposables.add(this.createProviderFileEventsWorker(this.isPathCaseSensitive(provider))); - providerDisposables.add(provider.onDidChangeFile(changes => this.onProviderDidChangeFile(providerFileChangeEventsWorker, changes))); + providerDisposables.add(provider.onDidChangeFile(changes => this._onDidFilesChange.fire(new FileChangesEvent(changes, !this.isPathCaseSensitive(provider))))); if (typeof provider.onDidWatchError === 'function') { providerDisposables.add(provider.onDidWatchError(error => this._onDidWatchError.fire(new Error(error)))); } @@ -993,20 +992,6 @@ export class FileService extends Disposable implements IFileService { //#region File Watching - /** - * Providers can send unlimited amount of `IFileChange` events - * and we want to protect against this to reduce CPU pressure. - * The following settings limit the amount of file changes we - * process at once. - * (https://github.com/microsoft/vscode/issues/124723) - */ - private static readonly FILE_EVENTS_THROTTLING = { - maxChangesChunkSize: 500 as const, // number of changes we process at once before... - coolDownDelay: 200 as const, // ...resting for 200ms until we process events again - maxChangesBufferSize: 30000 as const, // total number of changes we are willing to buffer in memory - warningscounter: 0 // keep track how many warnings we showed to reduce log spam - }; - private readonly _onDidFilesChange = this._register(new Emitter()); readonly onDidFilesChange = this._onDidFilesChange.event; @@ -1074,34 +1059,6 @@ export class FileService extends Disposable implements IFileService { }); } - private createProviderFileEventsWorker(caseSensitive: boolean): ThrottledWorker { - return new ThrottledWorker( - FileService.FILE_EVENTS_THROTTLING.maxChangesChunkSize, - FileService.FILE_EVENTS_THROTTLING.maxChangesBufferSize, - FileService.FILE_EVENTS_THROTTLING.coolDownDelay, - chunks => this._onDidFilesChange.fire(new FileChangesEvent(chunks, !caseSensitive)) - ); - } - - private onProviderDidChangeFile(worker: ThrottledWorker, changes: readonly IFileChange[]): void { - - // File events can be pretty much unbounded, depending on - // how many paths are watched and how large the changes are - // in them. As such, we use a `ThrottledWorker` that caps - // the number of changes we process at once as well as in - // total - - const worked = worker.work(changes); - - if (!worked && FileService.FILE_EVENTS_THROTTLING.warningscounter++ < 10) { - this.logService.warn(`[File watcher]: started ignoring events due to too many file change events at once (incoming: ${changes.length}, most recent change: ${changes[0].resource.toString()}). Use 'files.watcherExclude' setting to exclude folders with lots of changing files (e.g. compilation output).`); - } - - if (worker.pending > 0) { - this.logService.trace(`[File watcher]: started throttling events due to large amount of file change events at once (pending: ${worker.pending}, most recent change: ${changes[0].resource.toString()}). Use 'files.watcherExclude' setting to exclude folders with lots of changing files (e.g. compilation output).`); - } - } - override dispose(): void { super.dispose(); diff --git a/src/vs/platform/files/node/watcher/nodejs/nodejsWatcher.ts b/src/vs/platform/files/node/watcher/nodejs/nodejsWatcher.ts index 5df57d70b0b2..1dd6b878f579 100644 --- a/src/vs/platform/files/node/watcher/nodejs/nodejsWatcher.ts +++ b/src/vs/platform/files/node/watcher/nodejs/nodejsWatcher.ts @@ -4,7 +4,7 @@ *--------------------------------------------------------------------------------------------*/ import { watch } from 'fs'; -import { ThrottledDelayer } from 'vs/base/common/async'; +import { ThrottledDelayer, ThrottledWorker } from 'vs/base/common/async'; import { CancellationToken, CancellationTokenSource } from 'vs/base/common/cancellation'; import { isEqualOrParent } from 'vs/base/common/extpath'; import { parse } from 'vs/base/common/glob'; @@ -30,6 +30,20 @@ export class NodeJSFileWatcher extends Disposable implements INonRecursiveWatche // (same delay as Parcel is using) private static readonly FILE_CHANGES_HANDLER_DELAY = 50; + // Reduce likelyhood of spam from file events via throttling. + // These numbers are a bit more aggressive compared to the + // recursive watcher because we can have many individual + // node.js watchers per request. + // (https://github.com/microsoft/vscode/issues/124723) + private readonly throttledFileChangesWorker = new ThrottledWorker( + { + maxWorkChunkSize: 100, // only process up to 100 changes at once before... + throttleDelay: 200, // ...resting for 200ms until we process events again... + maxBufferedWork: 10000 // ...but never buffering more than 10000 events in memory + }, + events => this.onDidFilesChange(events) + ); + private readonly fileChangesDelayer = this._register(new ThrottledDelayer(NodeJSFileWatcher.FILE_CHANGES_HANDLER_DELAY)); private fileChangesBuffer: IDiskFileChange[] = []; @@ -372,16 +386,26 @@ export class NodeJSFileWatcher extends Disposable implements INonRecursiveWatche // Coalesce events: merge events of same kind const coalescedFileChanges = coalesceEvents(fileChanges); - // Logging - if (this.verboseLogging) { - for (const event of coalescedFileChanges) { - this.trace(`>> normalized ${event.type === FileChangeType.ADDED ? '[ADDED]' : event.type === FileChangeType.DELETED ? '[DELETED]' : '[CHANGED]'} ${event.path}`); - } - } - - // Broadcast to clients if (coalescedFileChanges.length > 0) { - this.onDidFilesChange(coalescedFileChanges); + + // Logging + if (this.verboseLogging) { + for (const event of coalescedFileChanges) { + this.trace(`>> normalized ${event.type === FileChangeType.ADDED ? '[ADDED]' : event.type === FileChangeType.DELETED ? '[DELETED]' : '[CHANGED]'} ${event.path}`); + } + } + + // Broadcast to clients via throttler + const worked = this.throttledFileChangesWorker.work(coalescedFileChanges); + + // Logging + if (!worked) { + this.warn(`started ignoring events due to too many file change events at once (incoming: ${coalescedFileChanges.length}, most recent change: ${coalescedFileChanges[0].path}). Use 'files.watcherExclude' setting to exclude folders with lots of changing files (e.g. compilation output).`); + } else { + if (this.throttledFileChangesWorker.pending > 0) { + this.trace(`started throttling events due to large amount of file change events at once (pending: ${this.throttledFileChangesWorker.pending}, most recent change: ${coalescedFileChanges[0].path}). Use 'files.watcherExclude' setting to exclude folders with lots of changing files (e.g. compilation output).`); + } + } } }); } catch (error) { diff --git a/src/vs/platform/files/node/watcher/parcel/parcelWatcher.ts b/src/vs/platform/files/node/watcher/parcel/parcelWatcher.ts index d9932fd6a001..d355481dcd00 100644 --- a/src/vs/platform/files/node/watcher/parcel/parcelWatcher.ts +++ b/src/vs/platform/files/node/watcher/parcel/parcelWatcher.ts @@ -6,7 +6,7 @@ import * as parcelWatcher from '@parcel/watcher'; import { existsSync, unlinkSync } from 'fs'; import { tmpdir } from 'os'; -import { DeferredPromise, RunOnceScheduler } from 'vs/base/common/async'; +import { DeferredPromise, RunOnceScheduler, ThrottledWorker } from 'vs/base/common/async'; import { CancellationToken, CancellationTokenSource } from 'vs/base/common/cancellation'; import { toErrorMessage } from 'vs/base/common/errorMessage'; import { Emitter } from 'vs/base/common/event'; @@ -88,6 +88,17 @@ export class ParcelWatcher extends Disposable implements IRecursiveWatcher { protected readonly watchers = new Map(); + // Reduce likelyhood of spam from file events via throttling. + // (https://github.com/microsoft/vscode/issues/124723) + private readonly throttledFileChangesWorker = new ThrottledWorker( + { + maxWorkChunkSize: 500, // only process up to 500 changes at once before... + throttleDelay: 200, // ...resting for 200ms until we process events again... + maxBufferedWork: 30000 // ...but never buffering more than 30000 events in memory + }, + events => this._onDidChangeFile.fire(events) + ); + private verboseLogging = false; private enospcErrorLogged = false; @@ -410,9 +421,9 @@ export class ParcelWatcher extends Disposable implements IRecursiveWatcher { } private emitEvents(events: IDiskFileChange[]): void { - - // Send outside - this._onDidChangeFile.fire(events); + if (events.length === 0) { + return; + } // Logging if (this.verboseLogging) { @@ -420,6 +431,18 @@ export class ParcelWatcher extends Disposable implements IRecursiveWatcher { this.trace(` >> normalized ${event.type === FileChangeType.ADDED ? '[ADDED]' : event.type === FileChangeType.DELETED ? '[DELETED]' : '[CHANGED]'} ${event.path}`); } } + + // Broadcast to clients via throttler + const worked = this.throttledFileChangesWorker.work(events); + + // Logging + if (!worked) { + this.warn(`started ignoring events due to too many file change events at once (incoming: ${events.length}, most recent change: ${events[0].path}). Use 'files.watcherExclude' setting to exclude folders with lots of changing files (e.g. compilation output).`); + } else { + if (this.throttledFileChangesWorker.pending > 0) { + this.trace(`started throttling events due to large amount of file change events at once (pending: ${this.throttledFileChangesWorker.pending}, most recent change: ${events[0].path}). Use 'files.watcherExclude' setting to exclude folders with lots of changing files (e.g. compilation output).`); + } + } } private normalizePath(request: IRecursiveWatchRequest): { realPath: string, realPathDiffers: boolean, realPathLength: number } { diff --git a/src/vs/platform/files/test/browser/fileService.test.ts b/src/vs/platform/files/test/browser/fileService.test.ts index c4b5963d6943..1dfd5216566a 100644 --- a/src/vs/platform/files/test/browser/fileService.test.ts +++ b/src/vs/platform/files/test/browser/fileService.test.ts @@ -9,7 +9,7 @@ import { CancellationToken, CancellationTokenSource } from 'vs/base/common/cance import { IDisposable, toDisposable } from 'vs/base/common/lifecycle'; import { consumeStream, newWriteableStream, ReadableStreamEvents } from 'vs/base/common/stream'; import { URI } from 'vs/base/common/uri'; -import { FileChangeType, FileOpenOptions, FileReadStreamOptions, FileSystemProviderCapabilities, FileType, IFileChange, IFileSystemProviderCapabilitiesChangeEvent, IFileSystemProviderRegistrationEvent, IStat } from 'vs/platform/files/common/files'; +import { FileOpenOptions, FileReadStreamOptions, FileSystemProviderCapabilities, FileType, IFileSystemProviderCapabilitiesChangeEvent, IFileSystemProviderRegistrationEvent, IStat } from 'vs/platform/files/common/files'; import { FileService } from 'vs/platform/files/common/fileService'; import { NullFileSystemProvider } from 'vs/platform/files/test/common/nullFileSystemProvider'; import { NullLogService } from 'vs/platform/log/common/log'; @@ -83,45 +83,6 @@ suite('File Service', () => { service.dispose(); }); - test('provider change events are throttled', async () => { - const service = new FileService(new NullLogService()); - - const provider = new NullFileSystemProvider(); - service.registerProvider('test', provider); - - await service.activateProvider('test'); - - let onDidFilesChangeFired = false; - service.onDidFilesChange(e => { - if (e.contains(URI.file('marker'))) { - onDidFilesChangeFired = true; - } - }); - - const throttledEvents: IFileChange[] = []; - for (let i = 0; i < 1000; i++) { - throttledEvents.push({ resource: URI.file(String(i)), type: FileChangeType.ADDED }); - } - throttledEvents.push({ resource: URI.file('marker'), type: FileChangeType.ADDED }); - - const nonThrottledEvents: IFileChange[] = []; - for (let i = 0; i < 100; i++) { - nonThrottledEvents.push({ resource: URI.file(String(i)), type: FileChangeType.ADDED }); - } - nonThrottledEvents.push({ resource: URI.file('marker'), type: FileChangeType.ADDED }); - - // 100 events are not throttled - provider.emitFileChangeEvents(nonThrottledEvents); - assert.strictEqual(onDidFilesChangeFired, true); - onDidFilesChangeFired = false; - - // 1000 events are throttled - provider.emitFileChangeEvents(throttledEvents); - assert.strictEqual(onDidFilesChangeFired, false); - - service.dispose(); - }); - test('watch', async () => { const service = new FileService(new NullLogService());