From e73417d249ef4e028e072b3d69e8c0e80d62f3df Mon Sep 17 00:00:00 2001 From: Johannes Rieken Date: Fri, 22 Jun 2018 12:21:13 +0200 Subject: [PATCH] add AsyncEmitter #43768 --- src/vs/base/common/event.ts | 58 +++++++++++++-- src/vs/base/test/common/event.test.ts | 100 +++++++++++++++++++++++++- 2 files changed, 151 insertions(+), 7 deletions(-) diff --git a/src/vs/base/common/event.ts b/src/vs/base/common/event.ts index 060b728bcc6..198f69c641b 100644 --- a/src/vs/base/common/event.ts +++ b/src/vs/base/common/event.ts @@ -4,11 +4,11 @@ *--------------------------------------------------------------------------------------------*/ 'use strict'; -import { IDisposable, toDisposable, combinedDisposable, empty as EmptyDisposable } from 'vs/base/common/lifecycle'; -import { TPromise } from 'vs/base/common/winjs.base'; -import { once as onceFn } from 'vs/base/common/functional'; import { onUnexpectedError } from 'vs/base/common/errors'; +import { once as onceFn } from 'vs/base/common/functional'; +import { combinedDisposable, empty as EmptyDisposable, IDisposable, toDisposable } from 'vs/base/common/lifecycle'; import { LinkedList } from 'vs/base/common/linkedList'; +import { TPromise } from 'vs/base/common/winjs.base'; /** * To an event a function with one or zero parameters @@ -58,9 +58,9 @@ export class Emitter { private static readonly _noop = function () { }; private _event: Event; - private _listeners: LinkedList; - private _deliveryQueue: [Listener, T][]; private _disposed: boolean; + private _deliveryQueue: [Listener, T][]; + protected _listeners: LinkedList; constructor(private _options?: EmitterOptions) { @@ -159,6 +159,52 @@ export class Emitter { } } +export interface IWaitUntil { + waitUntil(thenable: Thenable): void; +} + +export class AsyncEmitter extends Emitter { + + private _asyncDeliveryQueue: [Listener, T, Thenable[]][]; + + async fireAsync(eventFn: (thenables: Thenable[], listener: Function) => T): TPromise { + if (!this._listeners) { + return; + } + + // put all [listener,event]-pairs into delivery queue + // then emit all event. an inner/nested event might be + // the driver of this + if (!this._asyncDeliveryQueue) { + this._asyncDeliveryQueue = []; + } + + for (let iter = this._listeners.iterator(), e = iter.next(); !e.done; e = iter.next()) { + let thenables: Thenable[] = []; + this._asyncDeliveryQueue.push([e.value, eventFn(thenables, typeof e.value === 'function' ? e.value : e.value[0]), thenables]); + } + + while (this._asyncDeliveryQueue.length > 0) { + const [listener, event, thenables] = this._asyncDeliveryQueue.shift(); + try { + if (typeof listener === 'function') { + listener.call(undefined, event); + } else { + listener[0].call(listener[1], event); + } + } catch (e) { + onUnexpectedError(e); + continue; + } + + // freeze thenables-collection to enforce sync-calls to + // wait until and then wait for all thenables to resolve + Object.freeze(thenables); + await TPromise.join(thenables); + } + } +} + export class EventMultiplexer implements IDisposable { private readonly emitter: Emitter; @@ -550,4 +596,4 @@ export function latch(event: Event): Event { cache = value; return shouldEmit; }); -} \ No newline at end of file +} diff --git a/src/vs/base/test/common/event.test.ts b/src/vs/base/test/common/event.test.ts index 2ef6786f9ae..1b582337edf 100644 --- a/src/vs/base/test/common/event.test.ts +++ b/src/vs/base/test/common/event.test.ts @@ -5,10 +5,11 @@ 'use strict'; import * as assert from 'assert'; -import { Event, Emitter, debounceEvent, EventBufferer, once, fromPromise, stopwatch, buffer, echo, EventMultiplexer, latch } from 'vs/base/common/event'; +import { Event, Emitter, debounceEvent, EventBufferer, once, fromPromise, stopwatch, buffer, echo, EventMultiplexer, latch, AsyncEmitter, IWaitUntil } from 'vs/base/common/event'; import { IDisposable } from 'vs/base/common/lifecycle'; import * as Errors from 'vs/base/common/errors'; import { TPromise } from 'vs/base/common/winjs.base'; +import { timeout } from 'vs/base/common/async'; namespace Samples { @@ -238,6 +239,103 @@ suite('Event', function () { }); }); +suite('AsyncEmitter', function () { + + test('event has waitUntil-function', async function () { + + interface E extends IWaitUntil { + foo: boolean; + bar: number; + } + + let emitter = new AsyncEmitter(); + + emitter.event(e => { + assert.equal(e.foo, true); + assert.equal(e.bar, 1); + assert.equal(typeof e.waitUntil, 'function'); + }); + + emitter.fireAsync(thenables => ({ + foo: true, + bar: 1, + waitUntil(t: Thenable) { thenables.push(t); } + })); + emitter.dispose(); + }); + + test('sequential delivery', async function () { + + interface E extends IWaitUntil { + foo: boolean; + } + + let globalState = 0; + let emitter = new AsyncEmitter(); + + emitter.event(e => { + e.waitUntil(timeout(10).then(_ => { + assert.equal(globalState, 0); + globalState += 1; + })); + }); + + emitter.event(e => { + e.waitUntil(timeout(1).then(_ => { + assert.equal(globalState, 1); + globalState += 1; + })); + }); + + await emitter.fireAsync(thenables => ({ + foo: true, + waitUntil(t) { + thenables.push(t); + } + })); + assert.equal(globalState, 2); + }); + + test('sequential, in-order delivery', async function () { + interface E extends IWaitUntil { + foo: number; + } + let events: number[] = []; + let done = false; + let emitter = new AsyncEmitter(); + + // e1 + emitter.event(e => { + e.waitUntil(timeout(10).then(async _ => { + if (e.foo === 1) { + await emitter.fireAsync(thenables => ({ + foo: 2, + waitUntil(t) { + thenables.push(t); + } + })); + assert.deepEqual(events, [1, 2]); + done = true; + } + })); + }); + + // e2 + emitter.event(e => { + events.push(e.foo); + e.waitUntil(timeout(7)); + }); + + await emitter.fireAsync(thenables => ({ + foo: 1, + waitUntil(t) { + thenables.push(t); + } + })); + assert.ok(done); + }); +}); + suite('Event utils', () => { suite('EventBufferer', () => {