watcher - throttle events at their source to also prevent IPC spam

This commit is contained in:
Benjamin Pasero
2022-01-11 12:07:58 +01:00
parent 90ab95c46f
commit 9a4e9b9050
6 changed files with 112 additions and 115 deletions
+26 -10
View File
@@ -959,11 +959,29 @@ export class RunOnceWorker<T> 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<T> extends Disposable {
@@ -974,10 +992,8 @@ export class ThrottledWorker<T> 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<T> 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<T> 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<T> 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<T> extends Disposable {
this.throttler.clear();
this.doWork();
}, this.throttleDelay);
}, this.options.throttleDelay);
this.throttler.value.schedule();
}
}
+20 -4
View File
@@ -1068,7 +1068,11 @@ suite('Async', () => {
}
};
const worker = new async.ThrottledWorker<number>(5, undefined, 1, handler);
const worker = new async.ThrottledWorker<number>({
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<number>(5, 5, 1, handler);
const worker = new async.ThrottledWorker<number>({
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<number>(5, 5, 1, handler);
const worker = new async.ThrottledWorker<number>({
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<number>(5, undefined, 1, handler);
const worker = new async.ThrottledWorker<number>({
maxWorkChunkSize: 5,
maxBufferedWork: undefined,
throttleDelay: 1
}, handler);
worker.dispose();
const worked = worker.work([1, 2, 3]);
+4 -47
View File
@@ -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<FileChangesEvent>());
readonly onDidFilesChange = this._onDidFilesChange.event;
@@ -1074,34 +1059,6 @@ export class FileService extends Disposable implements IFileService {
});
}
private createProviderFileEventsWorker(caseSensitive: boolean): ThrottledWorker<IFileChange> {
return new ThrottledWorker<IFileChange>(
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<IFileChange>, 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();
@@ -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<IDiskFileChange>(
{
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<void>(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) {
@@ -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<string, IParcelWatcherInstance>();
// Reduce likelyhood of spam from file events via throttling.
// (https://github.com/microsoft/vscode/issues/124723)
private readonly throttledFileChangesWorker = new ThrottledWorker<IDiskFileChange>(
{
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 } {
@@ -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());