src/vs/base/common/event.ts

1964 LOC · 1778 covered · 186 uncovered · 366 ranges · 17510 concepts · 128 introducers · 10852 tests

File neighbourhood

The centred file is linked to every concept that introduces one of its ranges, every test that runs code from the file, and the gray connector concepts standing between those tests and the file's own introducer concepts. Undirected links join concepts to every file where they introduce source and concepts to the tests they introduce; arrows show specialization between the displayed concepts and bridge only concepts omitted from this view. Concept colors match the source ranges below; connector concepts have no source color and are shown in gray.

Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the file

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the related-file, concept, and source links on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the related-file, concept, and source links on this page while the interactive map is unavailable.

1 > /*--------------------------------------------------------------------------------------------- event.ts ×93
2 > * Copyright (c) Microsoft Corporation. All rights reserved.
3 > * Licensed under the MIT License. See License.txt in the project root for license information.
4 > *--------------------------------------------------------------------------------------------*/
5 >
6 > import { CancelablePromise } from './async.js';
7 > import { CancellationToken } from './cancellation.js';
8 > import { diffSets } from './collections.js';
9 > import { onUnexpectedError } from './errors.js';
10 > import { createSingleCallFunction } from './functional.js';
11 > import { combinedDisposable, Disposable, DisposableMap, DisposableStore, IDisposable, toDisposable } from './lifecycle.js';
12 > import { LinkedList } from './linkedList.js';
13 > import { IObservable, IObservableWithChange, IObserver } from './observable.js';
14 > import { env } from './process.js';
15 > import { StopWatch } from './stopwatch.js';
16 > import { MicrotaskDelay } from './symbols.js';
17 >
18 >
19 > // -----------------------------------------------------------------------------------------------------------------------
20 > // Uncomment the next line to print warnings whenever an emitter with listeners is disposed. That is a sign of code smell.
21 > // -----------------------------------------------------------------------------------------------------------------------
22 > const _enableDisposeWithListenerWarning = false
23 > // || Boolean("TRUE") // causes a linter warning so that it cannot be pushed
24 > ;
25 >
26 >
27 > // -----------------------------------------------------------------------------------------------------------------------
28 > // Uncomment the next line to print warnings whenever a snapshotted event is used repeatedly without cleanup.
29 > // See https://github.com/microsoft/vscode/issues/142851
30 > // -----------------------------------------------------------------------------------------------------------------------
31 > const _enableSnapshotPotentialLeakWarning = false
32 > // || Boolean("TRUE") // causes a linter warning so that it cannot be pushed
33 > ;
34 >
35 >
36 > const _bufferLeakWarnCountThreshold = 100;
37 > const _bufferLeakWarnTimeThreshold = 60_000; // 1 minute
38 >
39 > function _isBufferLeakWarningEnabled(): boolean { event.ts ×10
40 > return !!env['VSCODE_DEV'];
41 > }
43 > /**
44 > * An event with zero or one parameters that can be subscribed to. The event is a function itself.
45 > */
46 > export interface Event<T> {
47 > (listener: (e: T) => unknown, thisArgs?: any, disposables?: IDisposable[] | DisposableStore): IDisposable;
48 > }
49 >
50 > export namespace Event {
51 > export const None: Event<any> = () => Disposable.None;
52 >
53 > function _addLeakageTraceLogic(options: EmitterOptions) {
54 > if (_enableSnapshotPotentialLeakWarning) { event.ts ×2
55 const { onDidAddListener: origListenerDidAdd } = options;
56 const stack = Stacktrace.create();
57 let count = 0;
58 options.onDidAddListener = () => {
59 if (++count === 2) {
60 console.warn('snapshotted emitter LIKELY used public and SHOULD HAVE BEEN created with DisposableStore. snapshotted here');
61 stack.print();
62 }
63 origListenerDidAdd?.();
64 };
65 }
66 > } event.ts ×2
68 > /**
69 > * Given an event, returns another event which debounces calls and defers the listeners to a later task via a shared
70 > * `setTimeout`. The event is converted into a signal (`Event<void>`) to avoid additional object creation as a
71 > * result of merging events and to try prevent race conditions that could arise when using related deferred and
72 > * non-deferred events.
73 > *
74 > * This is useful for deferring non-critical work (eg. general UI updates) to ensure it does not block critical work
75 > * (eg. latency of keypress to text rendered).
76 > *
77 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
78 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
79 > * returned event causes this utility to leak a listener on the original event.
80 > *
81 > * @param event The event source for the new event.
82 > * @param flushOnListenerRemove Whether to fire all debounced events when a listener is removed. If this is not
83 > * specified, some events could go missing. Use this if it's important that all events are processed, even if the
84 > * listener gets disposed before the debounced event fires.
85 > * @param disposable A disposable store to add the new EventEmitter to.
86 > */
87 > export function defer(event: Event<unknown>, flushOnListenerRemove?: boolean, disposable?: DisposableStore): Event<void> {
88 return debounce<unknown, void>(event, () => void 0, 0, undefined, flushOnListenerRemove ?? true, undefined, disposable);
89 }
91 > /**
92 > * Given an event, returns another event which only fires once.
93 > *
94 > * @param event The event source for the new event.
95 > */
96 > export function once<T>(event: Event<T>): Event<T> {
97 > return (listener, thisArgs = null, disposables?) => { event.ts ×3
98 > // we need this, in case the event fires during the listener call
99 > let didFire = false;
100 > let result: IDisposable | undefined = undefined;
101 > result = event(e => {
102 > if (didFire) { event.ts ×4
103 > return; ipc.ts ×43
104 > } else if (result) { event.ts ×4
105 > result.dispose(); event.ts ×1
106 > } else { event.ts ×4
107 > didFire = true; event.ts ×2
108 > }
109 > event.ts ×4
110 > return listener.call(thisArgs, e);
111 > }, null, disposables); event.ts ×3
112 >
113 > if (didFire) {
114 > result.dispose(); event.ts ×2
115 > }
116 > event.ts ×3
117 > return result;
118 > };
119 > }
121 > /**
122 > * Given an event, returns another event which only fires once, and only when the condition is met.
123 > *
124 > * @param event The event source for the new event.
125 > */
126 > export function onceIf<T>(event: Event<T>, condition: (e: T) => boolean): Event<T> {
127 return Event.once(Event.filter(event, condition));
128 }
130 > /**
131 > * Maps an event of one type into an event of another type using a mapping function, similar to how
132 > * `Array.prototype.map` works.
133 > *
134 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
135 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
136 > * returned event causes this utility to leak a listener on the original event.
137 > *
138 > * @param event The event source for the new event.
139 > * @param map The mapping function.
140 > * @param disposable A disposable store to add the new EventEmitter to.
141 > */
142 > export function map<I, O>(event: Event<I>, map: (i: I) => O, disposable?: DisposableStore): Event<O> {
143 > return snapshot((listener, thisArgs = null, disposables?) => event(i => listener.call(thisArgs, map(i)), null, disposables), disposable); event.ts ×1
144 > }
146 > /**
147 > * Wraps an event in another event that performs some function on the event object before firing.
148 > *
149 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
150 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
151 > * returned event causes this utility to leak a listener on the original event.
152 > *
153 > * @param event The event source for the new event.
154 > * @param each The function to perform on the event object.
155 > * @param disposable A disposable store to add the new EventEmitter to.
156 > */
157 > export function forEach<I>(event: Event<I>, each: (i: I) => void, disposable?: DisposableStore): Event<I> {
158 return snapshot((listener, thisArgs = null, disposables?) => event(i => { each(i); listener.call(thisArgs, i); }, null, disposables), disposable);
159 }
161 > /**
162 > * Wraps an event in another event that fires only when some condition is met.
163 > *
164 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
165 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
166 > * returned event causes this utility to leak a listener on the original event.
167 > *
168 > * @param event The event source for the new event.
169 > * @param filter The filter function that defines the condition. The event will fire for the object if this function
170 > * returns true.
171 > * @param disposable A disposable store to add the new EventEmitter to.
172 > */
173 > export function filter<T, U>(event: Event<T | U>, filter: (e: T | U) => e is T, disposable?: DisposableStore): Event<T>;
174 > export function filter<T>(event: Event<T>, filter: (e: T) => boolean, disposable?: DisposableStore): Event<T>;
175 > export function filter<T, R>(event: Event<T | R>, filter: (e: T | R) => e is R, disposable?: DisposableStore): Event<R>;
176 > export function filter<T>(event: Event<T>, filter: (e: T) => boolean, disposable?: DisposableStore): Event<T> {
177 > return snapshot((listener, thisArgs = null, disposables?) => event(e => filter(e) && listener.call(thisArgs, e), null, disposables), disposable); event.ts ×5
178 > }
180 > /**
181 > * Given an event, returns the same event but typed as `Event<void>`.
182 > */
183 > export function signal<T>(event: Event<T>): Event<void> {
184 > return event as Event<any> as Event<void>; claudeSessionCustomizationDiscovery.ts ×2
185 > }
187 > /**
188 > * Given a collection of events, returns a single event which emits whenever any of the provided events emit.
189 > */
190 > export function any<T>(...events: Event<T>[]): Event<T>;
191 > export function any(...events: Event<any>[]): Event<void>;
192 > export function any<T>(...events: Event<T>[]): Event<T> {
193 > return (listener, thisArgs = null, disposables?) => { event.ts ×2
194 > const disposable = combinedDisposable(...events.map(event => event(e => listener.call(thisArgs, e)))); event.ts ×4
195 > return addAndReturnDisposable(disposable, disposables);
196 > };
197 > } event.ts ×2
199 > /**
200 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
201 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
202 > * returned event causes this utility to leak a listener on the original event.
203 > */
204 > export function reduce<I, O>(event: Event<I>, merge: (last: O | undefined, event: I) => O, initial?: O, disposable?: DisposableStore): Event<O> {
205 let output: O | undefined = initial;
206
207 return map<I, O>(event, e => {
208 output = merge(output, e);
209 return output;
210 }, disposable);
211 }
213 > function snapshot<T>(event: Event<T>, disposable: DisposableStore | undefined): Event<T> {
214 > let listener: IDisposable | undefined; event.ts ×5
215 >
216 > const options: EmitterOptions | undefined = {
217 > onWillAddFirstListener() {
218 > listener = event(emitter.fire, emitter); event.ts ×1
219 > },
220 > onDidRemoveLastListener() { event.ts ×5
221 > listener?.dispose(); event.ts ×1
222 > }
223 > }; event.ts ×5
224 >
225 > if (!disposable) {
226 > _addLeakageTraceLogic(options); event.ts ×1
227 > }
228 > event.ts ×5
229 > const emitter = new Emitter<T>(options);
230 >
231 > disposable?.add(emitter);
232 >
233 > return emitter.event;
234 > }
236 > /**
237 > * Adds the IDisposable to the store if it's set, and returns it. Useful to
238 > * Event function implementation.
239 > */
240 > function addAndReturnDisposable<T extends IDisposable>(d: T, store: DisposableStore | IDisposable[] | undefined): T {
241 > if (store instanceof Array) { event.ts ×4
242 store.push(d);
243 > } else if (store) { event.ts ×4
244 store.add(d);
245 }
246 > return d; event.ts ×4
247 > }
249 > /**
250 > * Given an event, creates a new emitter that event that will debounce events based on {@link delay} and give an
251 > * array event object of all events that fired.
252 > *
253 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
254 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
255 > * returned event causes this utility to leak a listener on the original event.
256 > *
257 > * @param event The original event to debounce.
258 > * @param merge A function that reduces all events into a single event.
259 > * @param delay The number of milliseconds to debounce.
260 > * @param leading Whether to fire a leading event without debouncing.
261 > * @param flushOnListenerRemove Whether to fire all debounced events when a listener is removed. If this is not
262 > * specified, some events could go missing. Use this if it's important that all events are processed, even if the
263 > * listener gets disposed before the debounced event fires.
264 > * @param leakWarningThreshold See {@link EmitterOptions.leakWarningThreshold}.
265 > * @param disposable A disposable store to register the debounce emitter to.
266 > */
267 > export function debounce<T>(event: Event<T>, merge: (last: T | undefined, event: T) => T, delay?: number | typeof MicrotaskDelay, leading?: boolean, flushOnListenerRemove?: boolean, leakWarningThreshold?: number, disposable?: DisposableStore): Event<T>;
268 > export function debounce<I, O>(event: Event<I>, merge: (last: O | undefined, event: I) => O, delay?: number | typeof MicrotaskDelay, leading?: boolean, flushOnListenerRemove?: boolean, leakWarningThreshold?: number, disposable?: DisposableStore): Event<O>;
269 > export function debounce<I, O>(event: Event<I>, merge: (last: O | undefined, event: I) => O, delay: number | typeof MicrotaskDelay = 100, leading = false, flushOnListenerRemove = false, leakWarningThreshold?: number, disposable?: DisposableStore): Event<O> {
270 > let subscription: IDisposable; event.ts ×4
271 > let output: O | undefined = undefined;
272 > let handle: Timeout | undefined | null = undefined;
273 > let numDebouncedCalls = 0;
274 > let doFire: (() => void) | undefined;
275 >
276 > const options: EmitterOptions | undefined = {
277 > leakWarningThreshold,
278 > onWillAddFirstListener() {
279 > subscription = event(cur => {
280 > numDebouncedCalls++; event.ts ×4
281 > output = merge(output, cur);
282 >
283 > if (leading && !handle) {
284 > emitter.fire(output); event.ts ×1
285 > output = undefined;
286 > }
287 > event.ts ×4
288 > doFire = () => {
289 > const _output = output; event.ts ×2
290 > output = undefined;
291 > handle = undefined;
292 > if (!leading || numDebouncedCalls > 1) {
293 > emitter.fire(_output!); event.ts ×1
294 > }
295 > numDebouncedCalls = 0; event.ts ×2
296 > };
297 > event.ts ×4
298 > if (typeof delay === 'number') {
299 > if (handle) { event.ts ×2
300 > clearTimeout(handle); event.ts ×1
301 > }
302 > handle = setTimeout(doFire, delay); event.ts ×2
303 > } else { event.ts ×4
304 > if (handle === undefined) { event.ts ×1
305 > handle = null;
306 > queueMicrotask(doFire);
307 > }
308 > }
309 > }); event.ts ×4
310 > },
311 > onWillRemoveListener() {
312 > if (flushOnListenerRemove && numDebouncedCalls > 0) { event.ts ×2
313 > doFire?.(); event.ts ×1
314 > }
315 > }, event.ts ×2
316 > onDidRemoveLastListener() { event.ts ×4
317 > doFire = undefined;
318 > subscription.dispose();
319 > }
320 > };
321 >
322 > if (!disposable) {
323 > _addLeakageTraceLogic(options); event.ts ×1
324 > }
325 > event.ts ×4
326 > const emitter = new Emitter<O>(options);
327 >
328 > disposable?.add(emitter);
329 >
330 > return emitter.event;
331 > }
333 > /**
334 > * Debounces an event, firing after some delay (default=0) with an array of all event original objects.
335 > *
336 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
337 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
338 > * returned event causes this utility to leak a listener on the original event.
339 > *
340 > * @param event The event source for the new event.
341 > * @param delay The number of milliseconds to debounce.
342 > * @param flushOnListenerRemove Whether to fire all debounced events when a listener is removed. If this is not
343 > * specified, some events could go missing. Use this if it's important that all events are processed, even if the
344 > * listener gets disposed before the debounced event fires.
345 > * @param disposable A disposable store to add the new EventEmitter to.
346 > */
347 > export function accumulate<T>(event: Event<T>, delay: number | typeof MicrotaskDelay = 0, flushOnListenerRemove?: boolean, disposable?: DisposableStore): Event<T[]> {
348 > return Event.debounce<T, T[]>(event, (last, e) => { event.ts ×2
349 > if (!last) {
350 > return [e];
351 > }
352 > last.push(e); event.ts ×1
353 > return last;
354 > }, delay, undefined, flushOnListenerRemove ?? true, undefined, disposable); event.ts ×2
355 > }
357 > /**
358 > * Throttles an event, ensuring the event is fired at most once during the specified delay period.
359 > * Unlike debounce, throttle will fire immediately on the leading edge and/or after the delay on the trailing edge.
360 > *
361 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
362 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
363 > * returned event causes this utility to leak a listener on the original event.
364 > *
365 > * @param event The event source for the new event.
366 > * @param merge An accumulator function that merges events if multiple occur during the throttle period.
367 > * @param delay The number of milliseconds to throttle.
368 > * @param leading Whether to fire on the leading edge (immediately on first event).
369 > * @param trailing Whether to fire on the trailing edge (after delay with the last value).
370 > * @param leakWarningThreshold See {@link EmitterOptions.leakWarningThreshold}.
371 > * @param disposable A disposable store to register the throttle emitter to.
372 > */
373 > export function throttle<T>(event: Event<T>, merge: (last: T | undefined, event: T) => T, delay?: number | typeof MicrotaskDelay, leading?: boolean, trailing?: boolean, leakWarningThreshold?: number, disposable?: DisposableStore): Event<T>;
374 > export function throttle<I, O>(event: Event<I>, merge: (last: O | undefined, event: I) => O, delay?: number | typeof MicrotaskDelay, leading?: boolean, trailing?: boolean, leakWarningThreshold?: number, disposable?: DisposableStore): Event<O>;
375 > export function throttle<I, O>(event: Event<I>, merge: (last: O | undefined, event: I) => O, delay: number | typeof MicrotaskDelay = 100, leading = true, trailing = true, leakWarningThreshold?: number, disposable?: DisposableStore): Event<O> {
376 > let subscription: IDisposable; event.ts ×4
377 > let output: O | undefined = undefined;
378 > let handle: Timeout | undefined = undefined;
379 > let numThrottledCalls = 0;
380 >
381 > const options: EmitterOptions | undefined = {
382 > leakWarningThreshold,
383 > onWillAddFirstListener() {
384 > subscription = event(cur => {
385 > numThrottledCalls++;
386 > output = merge(output, cur);
387 >
388 > // If not currently throttling, fire immediately if leading is enabled
389 > if (handle === undefined) {
390 > if (leading) {
391 > emitter.fire(output); event.ts ×1
392 > output = undefined;
393 > numThrottledCalls = 0;
394 > }
395 > event.ts ×4
396 > // Set up the throttle period
397 > if (typeof delay === 'number') {
398 > handle = setTimeout(() => { event.ts ×2
399 > // Fire on trailing edge if there were calls during throttle period
400 > if (trailing && numThrottledCalls > 0) {
401 > emitter.fire(output!); event.ts ×1
402 > }
403 > output = undefined; event.ts ×2
404 > handle = undefined;
405 > numThrottledCalls = 0;
406 > }, delay);
407 > } else { event.ts ×4
408 > // Use a special marker to indicate microtask is pending event.ts ×1
409 > handle = 0 as unknown as Timeout;
410 > queueMicrotask(() => {
411 > // Fire on trailing edge if there were calls during throttle period
412 > if (trailing && numThrottledCalls > 0) {
413 > emitter.fire(output!);
414 > }
415 > output = undefined;
416 > handle = undefined;
417 > numThrottledCalls = 0;
418 > });
419 > }
420 > } event.ts ×4
421 > // If already throttling, just accumulate the value for trailing edge
422 > });
423 > },
424 > onDidRemoveLastListener() {
425 > subscription.dispose();
426 > }
427 > };
428 >
429 > if (!disposable) {
430 > _addLeakageTraceLogic(options);
431 > }
432 >
433 > const emitter = new Emitter<O>(options);
434 >
435 > disposable?.add(emitter);
436 >
437 > return emitter.event;
438 > }
440 > /**
441 > * Filters an event such that some condition is _not_ met more than once in a row, effectively ensuring duplicate
442 > * event objects from different sources do not fire the same event object.
443 > *
444 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
445 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
446 > * returned event causes this utility to leak a listener on the original event.
447 > *
448 > * @param event The event source for the new event.
449 > * @param equals The equality condition.
450 > * @param disposable A disposable store to add the new EventEmitter to.
451 > *
452 > * @example
453 > * ```
454 > * // Fire only one time when a single window is opened or focused
455 > * Event.latch(Event.any(onDidOpenWindow, onDidFocusWindow))
456 > * ```
457 > */
458 > export function latch<T>(event: Event<T>, equals: (a: T, b: T) => boolean = (a, b) => a === b, disposable?: DisposableStore): Event<T> {
459 > let firstCall = true; event.ts ×1
460 > let cache: T;
461 >
462 > return filter(event, value => {
463 > const shouldEmit = firstCall || !equals(value, cache);
464 > firstCall = false;
465 > cache = value;
466 > return shouldEmit;
467 > }, disposable);
468 > }
470 > /**
471 > * Splits an event whose parameter is a union type into 2 separate events for each type in the union.
472 > *
473 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
474 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
475 > * returned event causes this utility to leak a listener on the original event.
476 > *
477 > * @example
478 > * ```
479 > * const event = new EventEmitter<number | undefined>().event;
480 > * const [numberEvent, undefinedEvent] = Event.split(event, isUndefined);
481 > * ```
482 > *
483 > * @param event The event source for the new event.
484 > * @param isT A function that determines what event is of the first type.
485 > * @param disposable A disposable store to add the new EventEmitter to.
486 > */
487 > export function split<T, U>(event: Event<T | U>, isT: (e: T | U) => e is T, disposable?: DisposableStore): [Event<T>, Event<U>] {
488 return [
489 Event.filter(event, isT, disposable),
490 Event.filter(event, e => !isT(e), disposable) as Event<U>,
491 ];
492 }
494 > /**
495 > * Buffers an event until it has a listener attached.
496 > *
497 > * *NOTE* that this function returns an `Event` and it MUST be called with a `DisposableStore` whenever the returned
498 > * event is accessible to "third parties", e.g the event is a public property. Otherwise a leaked listener on the
499 > * returned event causes this utility to leak a listener on the original event.
500 > *
501 > * @param event The event source for the new event.
502 > * @param debugName A name for this buffer, used in leak detection warnings.
503 > * @param flushAfterTimeout Determines whether to flush the buffer after a timeout immediately or after a
504 > * `setTimeout` when the first event listener is added.
505 > * @param _buffer Internal: A source event array used for tests.
506 > *
507 > * @example
508 > * ```
509 > * // Start accumulating events, when the first listener is attached, flush
510 > * // the event after a timeout such that multiple listeners attached before
511 > * // the timeout would receive the event
512 > * this.onInstallExtension = Event.buffer(service.onInstallExtension, 'onInstallExtension', true);
513 > * ```
514 > */
515 > export function buffer<T>(event: Event<T>, debugName: string, flushAfterTimeout = false, _buffer: T[] = [], disposable?: DisposableStore): Event<T> {
516 > let buffer: T[] | null = _buffer.slice(); event.ts ×10
517 >
518 > // Dev-only leak detection: track when buffer was created and warn
519 > // if events accumulate without ever being consumed.
520 > let bufferLeakWarningData: { stack: Stacktrace; timerId: ReturnType<typeof setTimeout>; warned: boolean } | undefined;
521 > if (_isBufferLeakWarningEnabled()) {
522 bufferLeakWarningData = {
523 stack: Stacktrace.create(),
524 timerId: setTimeout(() => {
525 if (buffer && buffer.length > 0 && bufferLeakWarningData && !bufferLeakWarningData.warned) {
526 bufferLeakWarningData.warned = true;
527 console.warn(`[Event.buffer][${debugName}] potential LEAK detected: ${buffer.length} events buffered for ${_bufferLeakWarnTimeThreshold / 1000}s without being consumed. Buffered here:`);
528 bufferLeakWarningData.stack.print();
529 }
530 }, _bufferLeakWarnTimeThreshold),
531 warned: false
532 };
533 if (disposable) {
534 disposable.add(toDisposable(() => clearTimeout(bufferLeakWarningData!.timerId)));
535 }
536 }
538 > const clearLeakWarningTimer = () => {
539 > if (bufferLeakWarningData) {
540 clearTimeout(bufferLeakWarningData.timerId);
541 }
542 > }; event.ts ×10
543 >
544 > let listener: IDisposable | null = event(e => {
545 > if (buffer) { event.ts ×8
546 > buffer.push(e); event.ts ×1
547 > if (_isBufferLeakWarningEnabled() && bufferLeakWarningData && !bufferLeakWarningData.warned && buffer.length >= _bufferLeakWarnCountThreshold) {
548 bufferLeakWarningData.warned = true;
549 console.warn(`[Event.buffer][${debugName}] potential LEAK detected: ${buffer.length} events buffered without being consumed. Buffered here:`);
550 bufferLeakWarningData.stack.print();
551 }
552 > } else { event.ts ×8
553 > emitter.fire(e); event.ts ×1
554 > }
555 > }); event.ts ×10
556 >
557 > if (disposable) {
558 > disposable.add(listener); ipc.ts ×3
559 > }
561 > const flush = () => {
562 > buffer?.forEach(e => emitter.fire(e)); event.ts ×8
563 > buffer = null;
564 > clearLeakWarningTimer();
565 > };
567 > const emitter = new Emitter<T>({
568 > onWillAddFirstListener() {
569 > if (!listener) { event.ts ×8
570 > listener = event(e => emitter.fire(e)); event.ts ×1
571 > if (disposable) {
572 > disposable.add(listener);
573 > }
574 > }
575 > }, event.ts ×8
577 > onDidAddFirstListener() {
578 > if (buffer) { event.ts ×8
579 > if (flushAfterTimeout) {
580 > setTimeout(flush); event.ts ×1
581 > } else { event.ts ×8
582 > flush(); event.ts ×1
583 > }
584 > } event.ts ×8
585 > },
587 > onDidRemoveLastListener() {
588 > if (listener) {
589 > listener.dispose();
590 > }
591 > listener = null;
592 > clearLeakWarningTimer();
593 > }
594 > });
595 >
596 > if (disposable) {
597 > disposable.add(emitter); ipc.ts ×3
598 > }
600 > return emitter.event;
601 > }
602 > /** event.ts ×93
603 > * Wraps the event in an {@link IChainableEvent}, allowing a more functional programming style.
604 > *
605 > * @example
606 > * ```
607 > * // Normal
608 > * const onEnterPressNormal = Event.filter(
609 > * Event.map(onKeyPress.event, e => new StandardKeyboardEvent(e)),
610 > * e.keyCode === KeyCode.Enter
611 > * ).event;
612 > *
613 > * // Using chain
614 > * const onEnterPressChain = Event.chain(onKeyPress.event, $ => $
615 > * .map(e => new StandardKeyboardEvent(e))
616 > * .filter(e => e.keyCode === KeyCode.Enter)
617 > * );
618 > * ```
619 > */
620 > export function chain<T, R>(event: Event<T>, sythensize: ($: IChainableSythensis<T>) => IChainableSythensis<R>): Event<R> {
621 > const fn: Event<R> = (listener, thisArgs, disposables) => { event.ts ×4
622 > const cs = sythensize(new ChainableSynthesis()) as ChainableSynthesis;
623 > return event(function (value) {
624 > const result = cs.evaluate(value);
625 > if (result !== HaltChainable) {
626 > listener.call(thisArgs, result);
627 > }
628 > }, undefined, disposables);
629 > };
630 >
631 > return fn;
632 > }
634 > const HaltChainable = Symbol('HaltChainable');
635 >
636 > class ChainableSynthesis implements IChainableSythensis<any> {
637 > private readonly steps: ((input: any) => unknown)[] = []; event.ts ×4
639 > map<O>(fn: (i: any) => O): this {
640 > this.steps.push(fn); event.ts ×1
641 > return this;
642 > }
644 > forEach(fn: (i: any) => void): this {
645 this.steps.push(v => {
646 fn(v);
647 return v;
648 });
649 return this;
650 }
652 > filter(fn: (e: any) => boolean): this {
653 > this.steps.push(v => fn(v) ? v : HaltChainable); event.ts ×1
654 > return this;
655 > }
657 > reduce<R>(merge: (last: R | undefined, event: any) => R, initial?: R | undefined): this {
658 > let last = initial; event.ts ×1
659 > this.steps.push(v => {
660 > last = merge(last, v);
661 > return last;
662 > });
663 > return this;
664 > }
666 > latch(equals: (a: any, b: any) => boolean = (a, b) => a === b): ChainableSynthesis {
667 > let firstCall = true; event.ts ×1
668 > let cache: any;
669 > this.steps.push(value => {
670 > const shouldEmit = firstCall || !equals(value, cache);
671 > firstCall = false;
672 > cache = value;
673 > return shouldEmit ? value : HaltChainable;
674 > });
675 >
676 > return this;
677 > }
679 > public evaluate(value: any) {
680 > for (const step of this.steps) { event.ts ×4
681 > value = step(value);
682 > if (value === HaltChainable) {
683 > break; event.ts ×1
684 > }
685 > } event.ts ×4
686 >
687 > return value;
688 > }
689 > } event.ts ×93
690 >
691 > export interface IChainableSythensis<T> {
692 > map<O>(fn: (i: T) => O): IChainableSythensis<O>;
693 > forEach(fn: (i: T) => void): IChainableSythensis<T>;
694 > filter<R extends T>(fn: (e: T) => e is R): IChainableSythensis<R>;
695 > filter(fn: (e: T) => boolean): IChainableSythensis<T>;
696 > reduce<R>(merge: (last: R, event: T) => R, initial: R): IChainableSythensis<R>;
697 > reduce<R>(merge: (last: R | undefined, event: T) => R): IChainableSythensis<R>;
698 > latch(equals?: (a: T, b: T) => boolean): IChainableSythensis<T>;
699 > }
700 >
701 > export interface NodeEventEmitter {
702 > on(event: string | symbol, listener: Function): unknown;
703 > removeListener(event: string | symbol, listener: Function): unknown;
704 > }
705 >
706 > /**
707 > * Creates an {@link Event} from a node event emitter.
708 > */
709 > export function fromNodeEventEmitter<T>(emitter: NodeEventEmitter, eventName: string, map: (...args: any[]) => T = id => id): Event<T> {
710 const fn = (...args: unknown[]) => result.fire(map(...args));
711 const onFirstListenerAdd = () => emitter.on(eventName, fn);
712 const onLastListenerRemove = () => emitter.removeListener(eventName, fn);
713 const result = new Emitter<T>({ onWillAddFirstListener: onFirstListenerAdd, onDidRemoveLastListener: onLastListenerRemove });
714
715 return result.event;
716 }
718 > export interface DOMEventEmitter {
719 > addEventListener(event: string | symbol, listener: Function): void;
720 > removeEventListener(event: string | symbol, listener: Function): void;
721 > }
722 >
723 > /**
724 > * Creates an {@link Event} from a DOM event emitter.
725 > */
726 > export function fromDOMEventEmitter<T>(emitter: DOMEventEmitter, eventName: string, map: (...args: any[]) => T = id => id): Event<T> {
727 const fn = (...args: unknown[]) => result.fire(map(...args));
728 const onFirstListenerAdd = () => emitter.addEventListener(eventName, fn);
729 const onLastListenerRemove = () => emitter.removeEventListener(eventName, fn);
730 const result = new Emitter<T>({ onWillAddFirstListener: onFirstListenerAdd, onDidRemoveLastListener: onLastListenerRemove });
731
732 return result.event;
733 }
735 > /**
736 > * Creates a promise out of an event, using the {@link Event.once} helper.
737 > */
738 > export function toPromise<T>(event: Event<T>, disposables?: IDisposable[] | DisposableStore): CancelablePromise<T> {
739 > let cancelRef: () => void; event.ts ×3
740 > let listener: IDisposable;
741 > const promise = new Promise((resolve) => {
742 > listener = once(event)(resolve);
743 > addToDisposables(listener, disposables);
744 >
745 > // not resolved, matching the behavior of a normal disposal
746 > cancelRef = () => {
747 > disposeAndRemove(listener, disposables); event.ts ×1
748 > };
749 > }) as CancelablePromise<T>; event.ts ×3
750 > promise.cancel = cancelRef!;
751 >
752 > if (disposables) {
753 > promise.finally(() => disposeAndRemove(listener, disposables)); event.ts ×1
754 > }
755 > event.ts ×3
756 > return promise;
757 > }
759 > /**
760 > * A convenience function for forwarding an event to another emitter which
761 > * improves readability.
762 > *
763 > * This is similar to {@link Relay} but allows instantiating and forwarding
764 > * on a single line and also allows for multiple source events.
765 > * @param from The event to forward.
766 > * @param to The emitter to forward the event to.
767 > * @example
768 > * Event.forward(event, emitter);
769 > * // equivalent to
770 > * event(e => emitter.fire(e));
771 > * // equivalent to
772 > * event(emitter.fire, emitter);
773 > */
774 > export function forward<T>(from: Event<T>, to: Emitter<T>): IDisposable {
775 return from(e => to.fire(e));
776 }
778 > /**
779 > * Adds a listener to an event and calls the listener immediately with undefined as the event object.
780 > *
781 > * @example
782 > * ```
783 > * // Initialize the UI and update it when dataChangeEvent fires
784 > * runAndSubscribe(dataChangeEvent, () => this._updateUI());
785 > * ```
786 > */
787 > export function runAndSubscribe<T>(event: Event<T>, handler: (e: T) => unknown, initial: T): IDisposable;
788 > export function runAndSubscribe<T>(event: Event<T>, handler: (e: T | undefined) => unknown): IDisposable;
789 > export function runAndSubscribe<T>(event: Event<T>, handler: (e: T | undefined) => unknown, initial?: T): IDisposable {
790 > handler(initial); terminalSandboxEngine.ts ×2
791 > return event(e => handler(e));
792 > }
794 > class EmitterObserver<T> implements IObserver {
795 >
796 > readonly emitter: Emitter<T>;
797 >
798 > private _counter = 0;
799 > private _hasChanged = false;
800 >
801 > constructor(readonly _observable: IObservable<T>, store: DisposableStore | undefined) {
802 > const options: EmitterOptions = { event.ts ×3
803 > onWillAddFirstListener: () => {
804 > _observable.addObserver(this);
805 >
806 > // Communicate to the observable that we received its current value and would like to be notified about future changes.
807 > this._observable.reportChanges();
808 > },
809 > onDidRemoveLastListener: () => {
810 > _observable.removeObserver(this);
811 > }
812 > };
813 > if (!store) {
814 > _addLeakageTraceLogic(options);
815 > }
816 > this.emitter = new Emitter<T>(options);
817 > if (store) {
818 store.add(this.emitter);
819 }
820 > } event.ts ×3
822 > beginUpdate<T>(_observable: IObservable<T>): void {
823 > // assert(_observable === this.obs); event.ts ×3
824 > this._counter++;
825 > }
827 > handlePossibleChange<T>(_observable: IObservable<T>): void {
828 // assert(_observable === this.obs);
829 }
831 > handleChange<T, TChange>(_observable: IObservableWithChange<T, TChange>, _change: TChange): void {
832 > // assert(_observable === this.obs); event.ts ×3
833 > this._hasChanged = true;
834 > }
836 > endUpdate<T>(_observable: IObservable<T>): void {
837 > // assert(_observable === this.obs); event.ts ×3
838 > this._counter--;
839 > if (this._counter === 0) {
840 > this._observable.reportChanges();
841 > if (this._hasChanged) {
842 > this._hasChanged = false;
843 > this.emitter.fire(this._observable.get());
844 > }
845 > }
846 > }
847 > } event.ts ×93
848 >
849 > /**
850 > * Creates an event emitter that is fired when the observable changes.
851 > * Each listeners subscribes to the emitter.
852 > */
853 > export function fromObservable<T>(obs: IObservable<T>, store?: DisposableStore): Event<T> {
854 > const observer = new EmitterObserver(obs, store); event.ts ×3
855 > return observer.emitter.event;
856 > }
858 > /**
859 > * Each listener is attached to the observable directly.
860 > */
861 > export function fromObservableLight(observable: IObservable<unknown>): Event<void> {
862 > return (listener, thisArgs, disposables) => { event.ts ×2
863 > let count = 0; event.ts ×6
864 > let didChange = false;
865 > const observer: IObserver = {
866 > beginUpdate() {
867 > count++; event.ts ×3
868 > },
869 > endUpdate() { event.ts ×6
870 > count--; event.ts ×3
871 > if (count === 0) {
872 > observable.reportChanges();
873 > if (didChange) {
874 > didChange = false;
875 > listener.call(thisArgs);
876 > }
877 > }
878 > },
879 > handlePossibleChange() { event.ts ×6
880 // noop
881 },
882 > handleChange() { event.ts ×6
883 > didChange = true; event.ts ×3
884 > }
885 > }; event.ts ×6
886 > observable.addObserver(observer);
887 > observable.reportChanges();
888 > const disposable = {
889 > dispose() {
890 > observable.removeObserver(observer); event.ts ×1
891 > }
892 > }; event.ts ×6
893 >
894 > addToDisposables(disposable, disposables);
895 >
896 > return disposable;
897 > };
898 > } event.ts ×2
899 > } event.ts ×93
900 >
901 > export interface EmitterOptions {
902 > /**
903 > * Optional function that's called *before* the very first listener is added
904 > */
905 > onWillAddFirstListener?: Function;
906 > /**
907 > * Optional function that's called *after* the very first listener is added
908 > */
909 > onDidAddFirstListener?: Function;
910 > /**
911 > * Optional function that's called after a listener is added
912 > */
913 > onDidAddListener?: Function;
914 > /**
915 > * Optional function that's called *after* remove the very last listener
916 > */
917 > onDidRemoveLastListener?: Function;
918 > /**
919 > * Optional function that's called *before* a listener is removed
920 > */
921 > onWillRemoveListener?: Function;
922 > /**
923 > * Optional function that's called when a listener throws an error. Defaults to
924 > * {@link onUnexpectedError}
925 > */
926 > onListenerError?: (e: any) => void;
927 > /**
928 > * Number of listeners that are allowed before assuming a leak. Default to
929 > * a globally configured value
930 > *
931 > * @see setGlobalLeakWarningThreshold
932 > */
933 > leakWarningThreshold?: number;
934 > /**
935 > * Human-readable name for the emitter, included in leak warning error
936 > * messages to help identify which emitter is leaking in telemetry.
937 > */
938 > leakWarningName?: string;
939 > /**
940 > * Pass in a delivery queue, which is useful for ensuring
941 > * in order event delivery across multiple emitters.
942 > */
943 > deliveryQueue?: EventDeliveryQueue;
944 >
945 > /** ONLY enable this during development */
946 > _profName?: string;
947 > }
948 >
949 >
950 > export class EventProfiling {
951 >
952 > static readonly all = new Set<EventProfiling>();
953 >
954 > private static _idPool = 0;
955 >
956 > readonly name: string;
957 > public listenerCount: number = 0;
958 > public invocationCount = 0;
959 > public elapsedOverall = 0;
960 > public durations: number[] = [];
961 >
962 > private _stopWatch?: StopWatch;
963 >
964 > constructor(name: string) {
965 this.name = `${name}_${EventProfiling._idPool++}`;
966 EventProfiling.all.add(this);
967 }
969 > start(listenerCount: number): void {
970 this._stopWatch = new StopWatch();
971 this.listenerCount = listenerCount;
972 }
974 > stop(): void {
975 if (this._stopWatch) {
976 const elapsed = this._stopWatch.elapsed();
977 this.durations.push(elapsed);
978 this.elapsedOverall += elapsed;
979 this.invocationCount += 1;
980 this._stopWatch = undefined;
981 }
982 }
983 > } event.ts ×93
984 >
985 > let _globalLeakWarningThreshold = -1;
986 > export function setGlobalLeakWarningThreshold(n: number): IDisposable {
987 const oldValue = _globalLeakWarningThreshold;
988 _globalLeakWarningThreshold = n;
989 return {
990 dispose() {
991 _globalLeakWarningThreshold = oldValue;
992 }
993 };
994 }
996 > class LeakageMonitor {
997 >
998 > private static _idPool = 1;
999 >
1000 > private _stacks: Map<string, number> | undefined;
1001 > private _warnCountdown: number = 0;
1002 >
1003 > constructor(
1004 > private readonly _errorHandler: (err: Error) => void, event.ts ×3
1005 > readonly threshold: number,
1006 > readonly name: string = (LeakageMonitor._idPool++).toString(16).padStart(3, '0')
1007 > ) { }
1008 > event.ts ×93
1009 > dispose(): void {
1010 > this._stacks?.clear(); event.ts ×3
1011 > }
1012 > event.ts ×93
1013 > check(stack: Stacktrace, listenerCount: number): undefined | (() => void) {
1014 > event.ts ×8
1015 > const threshold = this.threshold;
1016 > if (threshold <= 0 || listenerCount < threshold) {
1017 > return undefined;
1018 > }
1019 >
1020 > if (!this._stacks) {
1021 > this._stacks = new Map();
1022 > }
1023 > const count = (this._stacks.get(stack.value) || 0);
1024 > this._stacks.set(stack.value, count + 1);
1025 > this._warnCountdown -= 1;
1026 >
1027 > if (this._warnCountdown <= 0) {
1028 > // only warn on first exceed and then every time the limit
1029 > // is exceeded by 50% again
1030 > this._warnCountdown = threshold * 0.5;
1031 >
1032 > const [topStack, topCount] = this.getMostFrequentStack()!;
1033 > const emitterName = /^[0-9a-f]+$/i.test(this.name) ? undefined : this.name;
1034 > const message = `[${this.name}] potential listener LEAK detected, having ${listenerCount} listeners already. MOST frequent listener (${topCount}):`;
1035 > console.warn(message);
1036 > console.warn(topStack);
1037 >
1038 > const kind = topCount / listenerCount > 0.3 ? 'dominated' : 'popular';
1039 > const error = new ListenerLeakError(kind, message, topStack, listenerCount, emitterName);
1040 > this._errorHandler(error);
1041 > }
1042 >
1043 > return () => {
1044 > const count = (this._stacks!.get(stack.value) || 0);
1045 > this._stacks!.set(stack.value, count - 1);
1046 > };
1047 > }
1048 > event.ts ×93
1049 > getMostFrequentStack(): [string, number] | undefined {
1050 > if (!this._stacks) { event.ts ×8
1051 return undefined;
1052 }
1053 > let topStack: [string, number] | undefined; event.ts ×8
1054 > let topCount: number = 0;
1055 > for (const [stack, count] of this._stacks) {
1056 > if (!topStack || topCount < count) {
1057 > topStack = [stack, count];
1058 > topCount = count;
1059 > }
1060 > }
1061 > return topStack;
1062 > }
1063 > } event.ts ×93
1064 >
1065 > class Stacktrace {
1066 >
1067 > static create() {
1068 > const err = new Error();
1069 > return new Stacktrace(err.stack ?? '');
1070 > }
1071 >
1072 > private constructor(readonly value: string) { }
1073 >
1074 > print() {
1075 console.warn(this.value.split('\n').slice(2).join('\n'));
1076 }
1077 > } event.ts ×93
1078 >
1079 > // error that is logged when going over the configured listener threshold
1080 > export class ListenerLeakError extends Error {
1081 > readonly kind: string;
1082 > readonly listenerCount: number;
1083 > /**
1084 > * The detailed message including listener count and most frequent stack.
1085 > * Available locally for debugging but intentionally not used as the error
1086 > * `message`. When `emitterName` is provided, errors group by emitter name
1087 > * and kind in telemetry; otherwise they group by kind alone.
1088 > */
1089 > readonly details: string;
1090 > constructor(kind: 'dominated' | 'popular', details: string, stack: string, listenerCount: number, emitterName?: string) {
1091 > super(emitterName event.ts ×8
1092 ? `[${emitterName}] potential listener LEAK detected, ${kind}`
1093 > : `potential listener LEAK detected, ${kind}`); event.ts ×8
1094 > this.name = 'ListenerLeakError';
1095 > this.kind = kind;
1096 > this.listenerCount = listenerCount;
1097 > this.details = details;
1098 > this.stack = stack;
1099 > }
1100 > event.ts ×93
1101 > static is(err: unknown): err is ListenerLeakError {
1102 return err instanceof ListenerLeakError
1103 || (err instanceof Error && typeof (err as Error & { kind: unknown; listenerCount: unknown }).kind === 'string' && typeof (err as Error & { kind: unknown; listenerCount: unknown }).listenerCount === 'number');
1104 }
1105 > } event.ts ×93
1106 >
1107 > // SEVERE error that is logged when having gone way over the configured listener
1108 > // threshold so that the emitter refuses to accept more listeners
1109 > export class ListenerRefusalError extends ListenerLeakError {
1110 > constructor(kind: 'dominated' | 'popular', details: string, stack: string, listenerCount: number, emitterName?: string) {
1111 > super(kind, details, stack, listenerCount, emitterName); event.ts ×8
1112 > this.name = 'ListenerRefusalError';
1113 > }
1114 > } event.ts ×93
1115 >
1116 > let id = 0;
1117 > class UniqueContainer<T> {
1118 > stack?: Stacktrace;
1119 > public id = id++;
1120 > constructor(public readonly value: T) { }
1121 > }
1122 > const compactionThreshold = 2;
1123 >
1124 > type ListenerContainer<T> = UniqueContainer<(data: T) => void>;
1125 > type ListenerOrListeners<T> = (ListenerContainer<T> | undefined)[] | ListenerContainer<T>;
1126 >
1127 > const forEachListener = <T>(listeners: ListenerOrListeners<T>, fn: (c: ListenerContainer<T>) => void) => {
1128 > if (listeners instanceof UniqueContainer) { event.ts ×8
1129 > fn(listeners); event.ts ×1
1130 > } else { event.ts ×8
1131 > for (let i = 0; i < listeners.length; i++) { event.ts ×6
1132 > const l = listeners[i];
1133 > if (l) {
1134 > fn(l);
1135 > }
1136 > }
1137 > }
1138 > }; event.ts ×8
1139 > event.ts ×93
1140 > /**
1141 > * The Emitter can be used to expose an Event to the public
1142 > * to fire it from the insides.
1143 > * Sample:
1144 > class Document {
1145 >
1146 > private readonly _onDidChange = new Emitter<(value:string)=>any>();
1147 >
1148 > public onDidChange = this._onDidChange.event;
1149 >
1150 > // getter-style
1151 > // get onDidChange(): Event<(value:string)=>any> {
1152 > // return this._onDidChange.event;
1153 > // }
1154 >
1155 > private _doIt() {
1156 > //...
1157 > this._onDidChange.fire(value);
1158 > }
1159 > }
1160 > */
1161 > export class Emitter<T> {
1162 >
1163 > private readonly _options?: EmitterOptions;
1164 > private readonly _leakageMon?: LeakageMonitor;
1165 > private readonly _perfMon?: EventProfiling;
1166 > private _disposed?: true;
1167 > private _event?: Event<T>;
1168 >
1169 > /**
1170 > * A listener, or list of listeners. A single listener is the most common
1171 > * for event emitters (#185789), so we optimize that special case to avoid
1172 > * wrapping it in an array (just like Node.js itself.)
1173 > *
1174 > * A list of listeners never 'downgrades' back to a plain function if
1175 > * listeners are removed, for two reasons:
1176 > *
1177 > * 1. That's complicated (especially with the deliveryQueue)
1178 > * 2. A listener with >1 listener is likely to have >1 listener again at
1179 > * some point, and swapping between arrays and functions may[citation needed]
1180 > * introduce unnecessary work and garbage.
1181 > *
1182 > * The array listeners can be 'sparse', to avoid reallocating the array
1183 > * whenever any listener is added or removed. If more than `1 / compactionThreshold`
1184 > * of the array is empty, only then is it resized.
1185 > */
1186 > protected _listeners?: ListenerOrListeners<T>;
1187 >
1188 > /**
1189 > * Always to be defined if _listeners is an array. It's no longer a true
1190 > * queue, but holds the dispatching 'state'. If `fire()` is called on an
1191 > * emitter, any work left in the _deliveryQueue is finished first.
1192 > */
1193 > private _deliveryQueue?: EventDeliveryQueuePrivate;
1194 > protected _size = 0;
1195 >
1196 > constructor(options?: EmitterOptions) {
1197 > this._options = options; event.ts ×2
1198 > this._leakageMon = (_globalLeakWarningThreshold > 0 || this._options?.leakWarningThreshold)
1199 > ? new LeakageMonitor(options?.onListenerError ?? onUnexpectedError, this._options?.leakWarningThreshold ?? _globalLeakWarningThreshold, this._options?.leakWarningName) : event.ts ×3
1200 > undefined; event.ts ×2
1201 > this._perfMon = this._options?._profName ? new EventProfiling(this._options._profName) : undefined;
1202 > this._deliveryQueue = this._options?.deliveryQueue as EventDeliveryQueuePrivate | undefined;
1203 > }
1204 > event.ts ×93
1205 > dispose() {
1206 > if (!this._disposed) { event.ts ×3
1207 > this._disposed = true;
1208 >
1209 > // It is bad to have listeners at the time of disposing an emitter, it is worst to have listeners keep the emitter
1210 > // alive via the reference that's embedded in their disposables. Therefore we loop over all remaining listeners and
1211 > // unset their subscriptions/disposables. Looping and blaming remaining listeners is done on next tick because the
1212 > // the following programming pattern is very popular:
1213 > //
1214 > // const someModel = this._disposables.add(new ModelObject()); // (1) create and register model
1215 > // this._disposables.add(someModel.onDidChange(() => { ... }); // (2) subscribe and register model-event listener
1216 > // ...later...
1217 > // this._disposables.dispose(); disposes (1) then (2): don't warn after (1) but after the "overall dispose" is done
1218 >
1219 > if (this._deliveryQueue?.current === this) {
1220 > this._deliveryQueue.reset(); event.ts ×1
1221 > }
1222 > if (this._listeners) { event.ts ×3
1223 > if (_enableDisposeWithListenerWarning) { event.ts ×2
1224 const listeners = this._listeners;
1225 queueMicrotask(() => {
1226 forEachListener(listeners, l => l.stack?.print());
1227 });
1228 }
1229 > event.ts ×2
1230 > this._listeners = undefined;
1231 > this._size = 0;
1232 > }
1233 > this._options?.onDidRemoveLastListener?.(); event.ts ×3
1234 > this._leakageMon?.dispose();
1235 > }
1236 > }
1237 > event.ts ×93
1238 > /**
1239 > * For the public to allow to subscribe
1240 > * to events from this Emitter
1241 > */
1242 > get event(): Event<T> {
1243 > this._event ??= (callback: (e: T) => unknown, thisArgs?: any, disposables?: IDisposable[] | DisposableStore) => { event.ts ×2
1244 > if (this._leakageMon && this._size > this._leakageMon.threshold ** 2) { event.ts ×8
1245 > const message = `[${this._leakageMon.name}] REFUSES to accept new listeners because it exceeded its threshold by far (${this._size} vs ${this._leakageMon.threshold})`; event.ts ×8
1246 > console.warn(message);
1247 >
1248 > const tuple = this._leakageMon.getMostFrequentStack() ?? ['UNKNOWN stack', -1];
1249 > const kind = tuple[1] / this._size > 0.3 ? 'dominated' : 'popular';
1250 > const error = new ListenerRefusalError(kind, `${message}. HINT: Stack shows most frequent listener (${tuple[1]}-times)`, tuple[0], this._size, this._options?.leakWarningName);
1251 > const errorHandler = this._options?.onListenerError || onUnexpectedError;
1252 > errorHandler(error);
1253 >
1254 > return Disposable.None;
1255 > }
1256 > event.ts ×8
1257 > if (this._disposed) {
1258 > // todo: should we warn if a listener is added to a disposed emitter? This happens often event.ts ×1
1259 > return Disposable.None;
1260 > }
1261 > event.ts ×8
1262 > if (thisArgs) {
1263 > callback = callback.bind(thisArgs); event.ts ×1
1264 > }
1265 > event.ts ×8
1266 > const contained = new UniqueContainer(callback);
1267 >
1268 > let removeMonitor: Function | undefined;
1269 > let stack: Stacktrace | undefined;
1270 > if (this._leakageMon && this._size >= Math.ceil(this._leakageMon.threshold * 0.2)) {
1271 > // check and record this emitter for potential leakage event.ts ×8
1272 > contained.stack = Stacktrace.create();
1273 > removeMonitor = this._leakageMon.check(contained.stack, this._size + 1);
1274 > }
1275 > event.ts ×8
1276 > if (_enableDisposeWithListenerWarning) {
1277 contained.stack = stack ?? Stacktrace.create();
1278 }
1279 > event.ts ×8
1280 > if (!this._listeners) {
1281 > this._options?.onWillAddFirstListener?.(this);
1282 > this._listeners = contained;
1283 > this._options?.onDidAddFirstListener?.(this);
1284 > } else if (this._listeners instanceof UniqueContainer) {
1285 > this._deliveryQueue ??= new EventDeliveryQueuePrivate(); event.ts ×2
1286 > this._listeners = [this._listeners, contained];
1287 > } else {
1288 > this._listeners.push(contained); event.ts ×1
1289 > }
1290 > this._options?.onDidAddListener?.(this); event.ts ×8
1291 >
1292 > this._size++;
1293 >
1294 >
1295 > const result = toDisposable(() => {
1296 > removeMonitor?.(); event.ts ×3
1297 > this._removeListener(contained);
1298 > }); event.ts ×8
1299 > addToDisposables(result, disposables);
1300 >
1301 > return result;
1302 > };
1303 > event.ts ×2
1304 > return this._event;
1305 > }
1306 > event.ts ×93
1307 > private _removeListener(listener: ListenerContainer<T>) {
1308 > this._options?.onWillRemoveListener?.(this); event.ts ×3
1309 >
1310 > if (!this._listeners) {
1311 > return; // expected if a listener gets disposed event.ts ×1
1312 > }
1313 > event.ts ×1
1314 > if (this._size === 1) {
1315 > this._listeners = undefined; event.ts ×1
1316 > this._options?.onDidRemoveLastListener?.(this);
1317 > this._size = 0;
1318 > return;
1319 > }
1320 > event.ts ×2
1321 > // size > 1 which requires that listeners be a list:
1322 > const listeners = this._listeners as (ListenerContainer<T> | undefined)[];
1323 >
1324 > const index = listeners.indexOf(listener);
1325 > if (index === -1) {
1326 console.log('disposed?', this._disposed);
1327 console.log('size?', this._size);
1328 console.log('arr?', JSON.stringify(this._listeners));
1329 throw new Error('Attempted to dispose unknown listener');
1330 }
1331 > event.ts ×2
1332 > this._size--;
1333 > listeners[index] = undefined;
1334 >
1335 > const adjustDeliveryQueue = this._deliveryQueue!.current === this;
1336 > if (this._size * compactionThreshold <= listeners.length) {
1337 > let n = 0; event.ts ×2
1338 > for (let i = 0; i < listeners.length; i++) {
1339 > if (listeners[i]) {
1340 > listeners[n++] = listeners[i];
1341 > } else if (adjustDeliveryQueue && n < this._deliveryQueue!.end) {
1342 > this._deliveryQueue!.end--; event.ts ×1
1343 > if (n < this._deliveryQueue!.i) {
1344 > this._deliveryQueue!.i--;
1345 > }
1346 > }
1347 > } event.ts ×2
1348 > listeners.length = n;
1349 > }
1350 > } event.ts ×3
1351 > event.ts ×93
1352 > private _deliver(listener: undefined | UniqueContainer<(value: T) => void>, value: T) {
1353 > if (!listener) { event.ts ×5
1354 > return; event.ts ×1
1355 > }
1356 > event.ts ×5
1357 > const errorHandler = this._options?.onListenerError || onUnexpectedError;
1358 > if (!errorHandler) {
1359 listener.value(value);
1360 return;
1361 }
1362 > event.ts ×5
1363 > try {
1364 > listener.value(value);
1365 > } catch (e) {
1366 > errorHandler(e); event.ts ×1
1367 > }
1368 > } event.ts ×5
1369 > event.ts ×93
1370 > /** Delivers items in the queue. Assumes the queue is ready to go. */
1371 > private _deliverQueue(dq: EventDeliveryQueuePrivate) {
1372 > const listeners = dq.current!._listeners! as (ListenerContainer<T> | undefined)[]; event.ts ×4
1373 > while (dq.i < dq.end) {
1374 > // important: dq.i is incremented before calling deliver() because it might reenter deliverQueue()
1375 > this._deliver(listeners[dq.i++], dq.value as T);
1376 > }
1377 > dq.reset();
1378 > }
1379 > event.ts ×93
1380 > /**
1381 > * To be kept private to fire an event to
1382 > * subscribers
1383 > */
1384 > fire(event: T): void {
1385 > if (this._deliveryQueue?.current) { event.ts ×4
1386 > this._deliverQueue(this._deliveryQueue); event.ts ×1
1387 > this._perfMon?.stop(); // last fire() will have starting perfmon, stop it before starting the next dispatch
1388 > }
1389 > event.ts ×4
1390 > this._perfMon?.start(this._size);
1391 >
1392 > if (!this._listeners) {
1393 > // no-op event.ts ×1
1394 > } else if (this._listeners instanceof UniqueContainer) { event.ts ×4
1395 > this._deliver(this._listeners, event); event.ts ×1
1396 > } else { event.ts ×5
1397 > const dq = this._deliveryQueue!; event.ts ×4
1398 > dq.enqueue(this, event, this._listeners.length);
1399 > this._deliverQueue(dq);
1400 > }
1401 > event.ts ×4
1402 > this._perfMon?.stop();
1403 > }
1404 > event.ts ×93
1405 > hasListeners(): boolean {
1406 > return this._size > 0; event.ts ×1
1407 > }
1408 > } event.ts ×93
1409 >
1410 > export interface EventDeliveryQueue {
1411 > _isEventDeliveryQueue: true;
1412 > }
1413 >
1414 > export const createEventDeliveryQueue = (): EventDeliveryQueue => new EventDeliveryQueuePrivate();
1415 >
1416 > class EventDeliveryQueuePrivate implements EventDeliveryQueue { event.ts ×2
1417 > declare _isEventDeliveryQueue: true;
1418 >
1419 > /**
1420 > * Index in current's listener list.
1421 > */
1422 > public i = -1;
1423 >
1424 > /**
1425 > * The last index in the listener's list to deliver.
1426 > */
1427 > public end = 0;
1428 > event.ts ×93
1429 > /**
1430 > * Emitter currently being dispatched on. Emitter._listeners is always an array.
1431 > */
1432 > public current?: Emitter<any>;
1433 > /**
1434 > * Currently emitting value. Defined whenever `current` is.
1435 > */
1436 > public value?: unknown;
1437 >
1438 > public enqueue<T>(emitter: Emitter<T>, value: T, end: number) {
1439 > this.i = 0; event.ts ×4
1440 > this.end = end;
1441 > this.current = emitter;
1442 > this.value = value;
1443 > }
1444 > event.ts ×93
1445 > public reset() {
1446 > this.i = this.end; // force any current emission loop to stop, mainly for during dispose event.ts ×4
1447 > this.current = undefined;
1448 > this.value = undefined;
1449 > }
1450 > } event.ts ×93
1451 >
1452 > export interface IWaitUntil {
1453 > token: CancellationToken;
1454 > waitUntil(thenable: Promise<unknown>): void;
1455 > }
1456 >
1457 > export type IWaitUntilData<T> = Omit<Omit<T, 'waitUntil'>, 'token'>;
1458 >
1459 > export class AsyncEmitter<T extends IWaitUntil> extends Emitter<T> {
1460 >
1461 > private _asyncDeliveryQueue?: LinkedList<[(ev: T) => void, IWaitUntilData<T>]>;
1462 >
1463 > async fireAsync(data: IWaitUntilData<T>, token: CancellationToken, promiseJoin?: (p: Promise<unknown>, listener: Function) => Promise<unknown>): Promise<void> {
1464 > if (!this._listeners) { event.ts ×8
1465 return;
1466 }
1467 > event.ts ×8
1468 > if (!this._asyncDeliveryQueue) {
1469 > this._asyncDeliveryQueue = new LinkedList();
1470 > }
1471 >
1472 > forEachListener(this._listeners, listener => this._asyncDeliveryQueue!.push([listener.value, data]));
1473 >
1474 > while (this._asyncDeliveryQueue.size > 0 && !token.isCancellationRequested) {
1475 >
1476 > const [listener, data] = this._asyncDeliveryQueue.shift()!;
1477 > const thenables: Promise<unknown>[] = [];
1478 >
1479 > // eslint-disable-next-line local/code-no-dangerous-type-assertions
1480 > const event = <T>{
1481 > ...data,
1482 > token,
1483 > waitUntil: (p: Promise<unknown>): void => {
1484 > if (Object.isFrozen(thenables)) { event.ts ×6
1485 throw new Error('waitUntil can NOT be called asynchronous');
1486 }
1487 > if (promiseJoin) { event.ts ×6
1488 p = promiseJoin(p, listener);
1489 }
1490 > thenables.push(p); event.ts ×6
1491 > }
1492 > }; event.ts ×8
1493 >
1494 > try {
1495 > listener(event);
1496 > } catch (e) {
1497 onUnexpectedError(e);
1498 continue;
1499 }
1500 > event.ts ×8
1501 > // freeze thenables-collection to enforce sync-calls to
1502 > // wait until and then wait for all thenables to resolve
1503 > Object.freeze(thenables);
1504 >
1505 > await Promise.allSettled(thenables).then(values => {
1506 > for (const value of values) {
1507 > if (value.status === 'rejected') { event.ts ×6
1508 > onUnexpectedError(value.reason); event.ts ×1
1509 > }
1510 > } event.ts ×6
1511 > }); event.ts ×8
1512 > }
1513 > }
1514 > } event.ts ×93
1515 >
1516 >
1517 > export class PauseableEmitter<T> extends Emitter<T> {
1518 >
1519 > private _isPaused = 0;
1520 > protected _eventQueue = new LinkedList<T>();
1521 > private _mergeFn?: (input: T[]) => T;
1522 >
1523 > public get isPaused(): boolean {
1524 > return this._isPaused !== 0;
1525 > }
1526 >
1527 > constructor(options?: EmitterOptions & { merge?: (input: T[]) => T }) {
1528 > super(options); event.ts ×1
1529 > this._mergeFn = options?.merge;
1530 > }
1531 > event.ts ×93
1532 > pause(): void {
1533 > this._isPaused++; event.ts ×4
1534 > }
1535 > event.ts ×93
1536 > resume(): void {
1537 > if (this._isPaused !== 0 && --this._isPaused === 0) { event.ts ×2
1538 > if (this._mergeFn) { event.ts ×4
1539 > // use the merge function to create a single composite event.ts ×2
1540 > // event. make a copy in case firing pauses this emitter
1541 > if (this._eventQueue.size > 0) {
1542 > const events = Array.from(this._eventQueue); event.ts ×1
1543 > this._eventQueue.clear();
1544 > super.fire(this._mergeFn(events));
1545 > }
1546 > event.ts ×2
1547 > } else { event.ts ×4
1548 > // no merging, fire each event individually and test event.ts ×2
1549 > // that this emitter isn't paused halfway through
1550 > while (!this._isPaused && this._eventQueue.size !== 0) {
1551 > super.fire(this._eventQueue.shift()!); event.ts ×1
1552 > }
1553 > } event.ts ×2
1554 > } event.ts ×4
1555 > } event.ts ×2
1556 > event.ts ×93
1557 > override fire(event: T): void {
1558 > if (this._size) { event.ts ×2
1559 > if (this._isPaused !== 0) { event.ts ×3
1560 > this._eventQueue.push(event); event.ts ×1
1561 > } else { event.ts ×3
1562 > super.fire(event); event.ts ×1
1563 > }
1564 > } event.ts ×3
1565 > } event.ts ×2
1566 > } event.ts ×93
1567 >
1568 > export class DebounceEmitter<T> extends PauseableEmitter<T> {
1569 >
1570 > private readonly _delay: number;
1571 > private _handle: Timeout | undefined;
1572 >
1573 > constructor(options: EmitterOptions & { merge: (input: T[]) => T; delay?: number }) {
1574 > super(options); event.ts ×1
1575 > this._delay = options.delay ?? 100;
1576 > }
1577 > event.ts ×93
1578 > override fire(event: T): void {
1579 > if (!this._handle) { event.ts ×1
1580 > this.pause();
1581 > this._handle = setTimeout(() => {
1582 > this._handle = undefined;
1583 > this.resume();
1584 > }, this._delay);
1585 > }
1586 > super.fire(event);
1587 > }
1588 > } event.ts ×93
1589 >
1590 > /**
1591 > * An emitter which queue all events and then process them at the
1592 > * end of the event loop.
1593 > */
1594 > export class MicrotaskEmitter<T> extends Emitter<T> {
1595 > private _queuedEvents: T[] = [];
1596 > private _mergeFn?: (input: T[]) => T;
1597 >
1598 > constructor(options?: EmitterOptions & { merge?: (input: T[]) => T }) {
1599 > super(options); event.ts ×1
1600 > this._mergeFn = options?.merge;
1601 > }
1602 > override fire(event: T): void { event.ts ×93
1603 > event.ts ×2
1604 > if (!this.hasListeners()) {
1605 > return; actions.ts ×7
1606 > }
1607 > event.ts ×3
1608 > this._queuedEvents.push(event);
1609 > if (this._queuedEvents.length === 1) {
1610 > queueMicrotask(() => {
1611 > if (this._mergeFn) {
1612 > super.fire(this._mergeFn(this._queuedEvents)); event.ts ×1
1613 > } else { event.ts ×3
1614 > this._queuedEvents.forEach(e => super.fire(e)); event.ts ×1
1615 > }
1616 > this._queuedEvents = []; event.ts ×3
1617 > });
1618 > }
1619 > } event.ts ×2
1620 > } event.ts ×93
1621 >
1622 > /**
1623 > * An event emitter that multiplexes many events into a single event.
1624 > *
1625 > * @example Listen to the `onData` event of all `Thing`s, dynamically adding and removing `Thing`s
1626 > * to the multiplexer as needed.
1627 > *
1628 > * ```typescript
1629 > * const anythingDataMultiplexer = new EventMultiplexer<{ data: string }>();
1630 > *
1631 > * const thingListeners = DisposableMap<Thing, IDisposable>();
1632 > *
1633 > * thingService.onDidAddThing(thing => {
1634 > * thingListeners.set(thing, anythingDataMultiplexer.add(thing.onData);
1635 > * });
1636 > * thingService.onDidRemoveThing(thing => {
1637 > * thingListeners.deleteAndDispose(thing);
1638 > * });
1639 > *
1640 > * anythingDataMultiplexer.event(e => {
1641 > * console.log('Something fired data ' + e.data)
1642 > * });
1643 > * ```
1644 > */
1645 > export class EventMultiplexer<T> implements IDisposable {
1646 >
1647 > private readonly emitter: Emitter<T>;
1648 > private hasListeners = false;
1649 > private events: { event: Event<T>; listener: IDisposable | null }[] = [];
1650 >
1651 > constructor() {
1652 > this.emitter = new Emitter<T>({ event.ts ×9
1653 > onWillAddFirstListener: () => this.onFirstListenerAdd(),
1654 > onDidRemoveLastListener: () => this.onLastListenerRemove()
1655 > });
1656 > }
1657 > event.ts ×93
1658 > get event(): Event<T> {
1659 > return this.emitter.event; event.ts ×9
1660 > }
1661 > event.ts ×93
1662 > add(event: Event<T>): IDisposable {
1663 > const e = { event: event, listener: null }; event.ts ×9
1664 > this.events.push(e);
1665 >
1666 > if (this.hasListeners) {
1667 > this.hook(e); event.ts ×1
1668 > }
1669 > event.ts ×9
1670 > const dispose = () => {
1671 > if (this.hasListeners) {
1672 > this.unhook(e); event.ts ×1
1673 > }
1674 > event.ts ×9
1675 > const idx = this.events.indexOf(e);
1676 > this.events.splice(idx, 1);
1677 > };
1678 >
1679 > return toDisposable(createSingleCallFunction(dispose));
1680 > }
1681 > event.ts ×93
1682 > private onFirstListenerAdd(): void {
1683 > this.hasListeners = true; event.ts ×9
1684 > this.events.forEach(e => this.hook(e));
1685 > }
1686 > event.ts ×93
1687 > private onLastListenerRemove(): void {
1688 > this.hasListeners = false; event.ts ×9
1689 > this.events.forEach(e => this.unhook(e));
1690 > }
1691 > event.ts ×93
1692 > private hook(e: { event: Event<T>; listener: IDisposable | null }): void {
1693 > e.listener = e.event(r => this.emitter.fire(r)); event.ts ×9
1694 > }
1695 > event.ts ×93
1696 > private unhook(e: { event: Event<T>; listener: IDisposable | null }): void {
1697 > e.listener?.dispose(); event.ts ×9
1698 > e.listener = null;
1699 > }
1700 > event.ts ×93
1701 > dispose(): void {
1702 > this.emitter.dispose(); event.ts ×2
1703 >
1704 > for (const e of this.events) {
1705 > e.listener?.dispose(); event.ts ×1
1706 > }
1707 > this.events = []; event.ts ×2
1708 > }
1709 > } event.ts ×93
1710 >
1711 > export interface IDynamicListEventMultiplexer<TEventType> extends IDisposable {
1712 > readonly event: Event<TEventType>;
1713 > }
1714 > export class DynamicListEventMultiplexer<TItem, TEventType> implements IDynamicListEventMultiplexer<TEventType> {
1715 > private readonly _store = new DisposableStore();
1716 >
1717 > readonly event: Event<TEventType>;
1718 >
1719 > constructor(
1720 > items: TItem[], event.ts ×4
1721 > onAddItem: Event<TItem>,
1722 > onRemoveItem: Event<TItem>,
1723 > getEvent: (item: TItem) => Event<TEventType>
1724 > ) {
1725 > const multiplexer = this._store.add(new EventMultiplexer<TEventType>());
1726 > const itemListeners = this._store.add(new DisposableMap<TItem, IDisposable>());
1727 >
1728 > function addItem(instance: TItem) {
1729 > itemListeners.set(instance, multiplexer.add(getEvent(instance)));
1730 > }
1731 >
1732 > // Existing items
1733 > for (const instance of items) {
1734 > addItem(instance);
1735 > }
1736 >
1737 > // Added items
1738 > this._store.add(onAddItem(instance => {
1739 > addItem(instance); event.ts ×1
1740 > })); event.ts ×4
1741 >
1742 > // Removed items
1743 > this._store.add(onRemoveItem(instance => {
1744 > itemListeners.deleteAndDispose(instance); event.ts ×1
1745 > })); event.ts ×4
1746 >
1747 > this.event = multiplexer.event;
1748 > }
1749 > event.ts ×93
1750 > dispose() {
1751 > this._store.dispose(); event.ts ×4
1752 > }
1753 > } event.ts ×93
1754 >
1755 > /**
1756 > * The EventBufferer is useful in situations in which you want
1757 > * to delay firing your events during some code.
1758 > * You can wrap that code and be sure that the event will not
1759 > * be fired during that wrap.
1760 > *
1761 > * ```
1762 > * const emitter: Emitter;
1763 > * const delayer = new EventDelayer();
1764 > * const delayedEvent = delayer.wrapEvent(emitter.event);
1765 > *
1766 > * delayedEvent(console.log);
1767 > *
1768 > * delayer.bufferEvents(() => {
1769 > * emitter.fire(); // event will not be fired yet
1770 > * });
1771 > *
1772 > * // event will only be fired at this point
1773 > * ```
1774 > */
1775 > export class EventBufferer {
1776 > event.ts ×4
1777 > private data: { buffers: Function[] }[] = [];
1778 > event.ts ×93
1779 > wrapEvent<T>(event: Event<T>): Event<T>;
1780 > wrapEvent<T>(event: Event<T>, reduce: (last: T | undefined, event: T) => T): Event<T>;
1781 > wrapEvent<T, O>(event: Event<T>, reduce: (last: O | undefined, event: T) => O, initial: O): Event<O>;
1782 > wrapEvent<T, O>(event: Event<T>, reduce?: (last: T | O | undefined, event: T) => T | O, initial?: O): Event<O | T> {
1783 > return (listener, thisArgs?, disposables?) => { event.ts ×4
1784 > return event(i => {
1785 > const data = this.data[this.data.length - 1];
1786 >
1787 > // Non-reduce scenario
1788 > if (!reduce) {
1789 > // Buffering case
1790 > if (data) {
1791 > data.buffers.push(() => listener.call(thisArgs, i)); event.ts ×2
1792 > } else { event.ts ×4
1793 > // Not buffering case
1794 > listener.call(thisArgs, i);
1795 > }
1796 > return;
1797 > }
1798
1799 // Reduce scenario
1800 const reduceData = data as typeof data & {
1801 /**
1802 * The accumulated items that will be reduced.
1803 */
1804 items?: T[];
1805 /**
1806 * The reduced result cached to be shared with other listeners.
1807 */
1808 reducedResult?: T | O;
1809 };
1810
1811 // Not buffering case
1812 if (!reduceData) {
1813 // TODO: Is there a way to cache this reduce call for all listeners?
1814 listener.call(thisArgs, reduce(initial, i));
1815 return;
1816 }
1817
1818 // Buffering case
1819 reduceData.items ??= [];
1820 reduceData.items.push(i);
1821 if (reduceData.buffers.length === 0) {
1822 // Include a single buffered function that will reduce all events when we're done buffering events
1823 data.buffers.push(() => {
1824 // cache the reduced result so that the value can be shared across all listeners
1825 reduceData.reducedResult ??= initial
1826 ? reduceData.items!.reduce(reduce as (last: O | undefined, event: T) => O, initial)
1827 : reduceData.items!.reduce(reduce as (last: T | undefined, event: T) => T);
1828 listener.call(thisArgs, reduceData.reducedResult);
1829 });
1830 }
1831 > }, undefined, disposables); event.ts ×4
1832 > };
1833 > }
1834 > event.ts ×93
1835 > bufferEvents<R = void>(fn: () => R): R {
1836 > const data = { buffers: new Array<Function>() }; event.ts ×2
1837 > this.data.push(data);
1838 > const r = fn();
1839 > this.data.pop();
1840 > data.buffers.forEach(flush => flush());
1841 > return r;
1842 > }
1843 > } event.ts ×93
1844 >
1845 > /**
1846 > * A Relay is an event forwarder which functions as a replugabble event pipe.
1847 > * Once created, you can connect an input event to it and it will simply forward
1848 > * events from that input event through its own `event` property. The `input`
1849 > * can be changed at any point in time.
1850 > */
1851 > export class Relay<T> implements IDisposable {
1852 > event.ts ×2
1853 > private listening = false;
1854 > private inputEvent: Event<T> = Event.None;
1855 > private inputEventListener: IDisposable = Disposable.None;
1856 >
1857 > private readonly emitter = new Emitter<T>({
1858 > onDidAddFirstListener: () => {
1859 > this.listening = true;
1860 > this.inputEventListener = this.inputEvent(this.emitter.fire, this.emitter);
1861 > },
1862 > onDidRemoveLastListener: () => {
1863 > this.listening = false;
1864 > this.inputEventListener.dispose();
1865 > }
1866 > });
1867 >
1868 > readonly event: Event<T> = this.emitter.event;
1869 > event.ts ×93
1870 > set input(event: Event<T>) {
1871 > this.inputEvent = event; event.ts ×2
1872 >
1873 > if (this.listening) {
1874 > this.inputEventListener.dispose();
1875 > this.inputEventListener = event(this.emitter.fire, this.emitter);
1876 > }
1877 > }
1878 > event.ts ×93
1879 > dispose() {
1880 > this.inputEventListener.dispose(); event.ts ×1
1881 > this.emitter.dispose();
1882 > }
1883 > } event.ts ×93
1884 >
1885 > export interface IValueWithChangeEvent<T> {
1886 > readonly onDidChange: Event<void>;
1887 > get value(): T;
1888 > }
1889 >
1890 > export class ValueWithChangeEvent<T> implements IValueWithChangeEvent<T> {
1891 > public static const<T>(value: T): IValueWithChangeEvent<T> {
1892 > return new ConstValueWithChangeEvent(value);
1893 > }
1894 >
1895 > private readonly _onDidChange = new Emitter<void>();
1896 > readonly onDidChange: Event<void> = this._onDidChange.event;
1897 >
1898 > constructor(private _value: T) { }
1899 >
1900 > get value(): T {
1901 return this._value;
1902 }
1903 > event.ts ×93
1904 > set value(value: T) {
1905 if (value !== this._value) {
1906 this._value = value;
1907 this._onDidChange.fire(undefined);
1908 }
1909 }
1910 > } event.ts ×93
1911 >
1912 > class ConstValueWithChangeEvent<T> implements IValueWithChangeEvent<T> {
1913 > public readonly onDidChange: Event<void> = Event.None;
1914 >
1915 > constructor(readonly value: T) { }
1916 > }
1917 >
1918 > /**
1919 > * @param handleItem Is called for each item in the set (but only the first time the item is seen in the set).
1920 > * The returned disposable is disposed if the item is no longer in the set.
1921 > */
1922 > export function trackSetChanges<T>(getData: () => ReadonlySet<T>, onDidChangeData: Event<unknown>, handleItem: (d: T) => IDisposable): IDisposable {
1923 > const map = new DisposableMap<T, IDisposable>(); debugStorage.ts ×17
1924 > let oldData = new Set(getData());
1925 > for (const d of oldData) {
1926 map.set(d, handleItem(d));
1927 }
1929 > const store = new DisposableStore();
1930 > store.add(onDidChangeData(() => {
1931 const newData = getData();
1932 const diff = diffSets(oldData, newData);
1933 for (const r of diff.removed) {
1934 map.deleteAndDispose(r);
1935 }
1936 for (const a of diff.added) {
1937 map.set(a, handleItem(a));
1938 }
1939 oldData = new Set(newData);
1940 > })); debugStorage.ts ×17
1941 > store.add(map);
1942 > return store;
1943 > }
1944 > event.ts ×93
1945 >
1946 > function addToDisposables(result: IDisposable, disposables: DisposableStore | IDisposable[] | undefined) { event.ts ×3
1947 > if (disposables instanceof DisposableStore) {
1948 > disposables.add(result); event.ts ×1
1949 > } else if (Array.isArray(disposables)) { event.ts ×3
1950 > disposables.push(result); event.ts ×1
1951 > }
1952 > } event.ts ×3
1953 > event.ts ×93
1954 > function disposeAndRemove(result: IDisposable, disposables: DisposableStore | IDisposable[] | undefined) { event.ts ×3
1955 > if (disposables instanceof DisposableStore) {
1956 > disposables.delete(result); event.ts ×1
1957 > } else if (Array.isArray(disposables)) { event.ts ×3
1958 > const index = disposables.indexOf(result); event.ts ×1
1959 > if (index !== -1) {
1960 > disposables.splice(index, 1);
1961 > }
1962 > }
1963 > result.dispose(); event.ts ×3
1964 > }